From ab72789eef4e47bdee37d21d2b3a38d9eadcc629 Mon Sep 17 00:00:00 2001 From: p1neappleXpress Date: Sun, 27 Sep 2026 08:40:31 +0300 Subject: [PATCH] transport: remove native mailru/cupsonline/boards - scripted versions proven Per the project's own rule for this migration (script it, verify it live, then remove the native implementation): transport/mailru, transport/cupsonline and transport/yandex/boards.go are deleted. Each had already reached full live round-trip verification as a script transport (see the two preceding commits and CHANGELOG.md) before removal. yandex.go's shortStr helper was defined in boards.go but used by yandex.go too; moved rather than duplicated. yandex/vyandex stay native for now - their scripted happy path isn't proven yet. --transport=mailru/cupsonline/boards and their --*-url flags no longer exist; use --transport=script: (main.go's help text and transport_factory.go's factory were already updated to make these the only path). Co-Authored-By: Claude Sonnet 5 --- transport/cupsonline/cupsonline.go | 1494 ----------------------- transport/cupsonline/cupsonline_test.go | 115 -- transport/cupsonline/fakecups_test.go | 842 ------------- transport/cupsonline/live_test.go | 234 ---- transport/cupsonline/stop_test.go | 19 - transport/mailru/mailru.go | 591 --------- transport/yandex/boards.go | 1134 ----------------- transport/yandex/yandex.go | 10 + 8 files changed, 10 insertions(+), 4429 deletions(-) delete mode 100644 transport/cupsonline/cupsonline.go delete mode 100644 transport/cupsonline/cupsonline_test.go delete mode 100644 transport/cupsonline/fakecups_test.go delete mode 100644 transport/cupsonline/live_test.go delete mode 100644 transport/cupsonline/stop_test.go delete mode 100644 transport/mailru/mailru.go delete mode 100644 transport/yandex/boards.go diff --git a/transport/cupsonline/cupsonline.go b/transport/cupsonline/cupsonline.go deleted file mode 100644 index fc18e09..0000000 --- a/transport/cupsonline/cupsonline.go +++ /dev/null @@ -1,1494 +0,0 @@ -package cupsonline - -import ( - "bytes" - "context" - "encoding/base64" - "encoding/binary" - "encoding/json" - "errors" - "fmt" - "io" - "net/http" - "net/http/cookiejar" - "net/url" - "regexp" - "strings" - "sync" - "sync/atomic" - "time" - - "github.com/gorilla/websocket" - - "openflux/netbind" - "openflux/transport" - "openflux/utils" -) - -// baseRoomURL is where the rooms live. It's a variable only so tests can aim -// the transport at a local stand-in for cups.online. -var baseRoomURL = "https://interview.cups.online/live-coding/" - -const cupsUA = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/137.0.0.0 Safari/537.36" - -type CupsonlineConfig struct { - WSHandshakeTimeout time.Duration - WSReadTimeout time.Duration - WSWriteTimeout time.Duration - KeepAliveInterval time.Duration - ReconnectMinDelay time.Duration - ReconnectMaxDelay time.Duration - ReconnectMultiplier float64 - HTTPTimeout time.Duration - - // Комнату, похожую на закрытую, проверяем не чаще этого: от Min с - // удвоением до Max. Каждая такая проверка может заставить cups выдать - // новую комнату, так что долбить её нельзя. - RoomGoneRetryMin time.Duration - RoomGoneRetryMax time.Duration - - // Батчинг внутри одного WS. Данные едут в целых row/column курсоров - // (bytesPerNumber байт на число, два числа на курсор): cups.online больше - // не принимает строку в column, но массив целых курсоров ретранслирует - // как есть, дословно и по порядку. - // - // BatchMaxBytes — сколько байт пакетов копится за одну отправку. - // MaxMessageData — сколько байт едет в одном сообщении: на проводе число - // стоит ~4× своих байт, поэтому пачка режется на несколько сообщений, - // каждое заведомо под потолком сервера (проверено probe'ом: ~21 KB - // проходит, ~44 KB уже нет), а приёмная сторона склеивает их по порядку. - // MaxCursors — страховочный потолок массива в одном сообщении. - BatchMaxPackets int - BatchMaxBytes int - MaxMessageData int - MaxCursors int - BatchTimeout time.Duration - - // SendInterval is the least time between two messages on one channel. - // cups.online throttles inbound, and firing messages back to back stalls - // the write and drops the connection, so the transport paces itself below - // that. Four channels each send at this rate. - SendInterval time.Duration - - SendQueueSize int - - // Буферы сокета, а не лимит на сообщение: gorilla выделяет их целиком на - // каждое соединение. MaxMessageBytes — сколько может прийти одним фреймом. - ReadBufferSize int - WriteBufferSize int - MaxMessageBytes int - - MaxPayloadBytes int - - NumRooms int - RoomCreatePause time.Duration - - StatsInterval time.Duration -} - -func DefaultCupsonlineConfig() CupsonlineConfig { - return CupsonlineConfig{ - WSHandshakeTimeout: 15 * time.Second, - // Сервер пингует каждые ~25 с и отвечает на наш keepalive раз в 20 с, - // так что полторы минуты тишины — это мёртвое соединение. - WSReadTimeout: 90 * time.Second, - WSWriteTimeout: 15 * time.Second, - KeepAliveInterval: 20 * time.Second, - ReconnectMinDelay: 100 * time.Millisecond, - ReconnectMaxDelay: 10 * time.Second, - ReconnectMultiplier: 1.3, - HTTPTimeout: 30 * time.Second, - - RoomGoneRetryMin: time.Minute, - RoomGoneRetryMax: 30 * time.Minute, - - BatchMaxPackets: 256, - BatchMaxBytes: 11000, - // 4096 байт данных → 2 байта заголовка → ~342 курсора → провод ~17 KB, - // с запасом под проверенными ~21 KB и вчетверо меньше потолка. - MaxMessageData: 4096, - MaxCursors: 1000, - BatchTimeout: 2 * time.Millisecond, - // 18ms (~55 msg/s per channel) — calibrated against the live server - // (TestLiveThroughputCalibrate): the knee of throughput (~80 KB/s per - // channel; faster doesn't help, the server itself is the ceiling) with - // no reconnects and margin under the burst that stalled the old, - // unpaced code. - SendInterval: 18 * time.Millisecond, - - // В очереди кадры кодека (до ~8 KB), а не IP-пакеты: очередь длиннее - // — это уже не буфер, а секунды задержки. - SendQueueSize: 1024, - - ReadBufferSize: 64 << 10, - WriteBufferSize: 64 << 10, - MaxMessageBytes: 8 << 20, - - // Длина пакета во внутреннем кадре — uint16, так что один пакет не - // может быть больше 65535 байт. Кадр кодека (до ~8 KB) с запасом - // помещается и режется на сообщения при отправке. - MaxPayloadBytes: 65535, - - NumRooms: 4, - RoomCreatePause: 500 * time.Millisecond, - - StatsInterval: 5 * time.Second, - } -} - -var ( - reMetaConnToken = regexp.MustCompile(`]+name="centrifuge-connection-token"[^>]+content="([^"]+)"`) - reMetaConnURL = regexp.MustCompile(`]+name="centrifuge-connection-url"[^>]+content="([^"]+)"`) - reMetaSubURL = regexp.MustCompile(`]+name="centrifuge-subscription-token-url"[^>]+content="([^"]+)"`) - reDataRoomUUID = regexp.MustCompile(`data-room="\{"uuid":\s*"([0-9a-f-]{36})"`) - reDataUserUUID = regexp.MustCompile(`data-user="\{"uuid":\s*"([0-9a-f-]{36})"`) -) - -type cupsAuth struct { - roomUUID string - userUUID string - connToken string - connURL string - subURL string - subToken string - channel string - httpClient *http.Client - csrfToken string -} - -// authorize loads a room page and collects what it takes to join the room's -// channel. session carries the cookies of an earlier join over (nil starts a -// fresh one), so re-joining doesn't add a new participant every time. -func authorize(ctx context.Context, roomURL string, session *http.Client) (*cupsAuth, error) { - client := session - if client == nil { - jar, _ := cookiejar.New(nil) - client = &http.Client{ - Jar: jar, - Timeout: 30 * time.Second, - CheckRedirect: func(req *http.Request, via []*http.Request) error { - if len(via) >= 5 { - return fmt.Errorf("too many redirects") - } - return nil - }, - } - } - - req, _ := http.NewRequestWithContext(ctx, "GET", roomURL, nil) - req.Header.Set("User-Agent", cupsUA) - req.Header.Set("Accept", "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8") - req.Header.Set("Accept-Language", "ru-RU,ru;q=0.9") - - resp, err := client.Do(req) - if err != nil { - return nil, fmt.Errorf("GET room: %w", err) - } - body, _ := io.ReadAll(resp.Body) - resp.Body.Close() - if resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusGone { - return nil, fmt.Errorf("%w: GET room status %d", errRoomGone, resp.StatusCode) - } - if resp.StatusCode != 200 { - return nil, fmt.Errorf("GET room status %d", resp.StatusCode) - } - - html := string(body) - a := &cupsAuth{ - roomUUID: firstMatch(reDataRoomUUID, html), - userUUID: firstMatch(reDataUserUUID, html), - connToken: firstMatch(reMetaConnToken, html), - connURL: firstMatch(reMetaConnURL, html), - subURL: firstMatch(reMetaSubURL, html), - httpClient: client, - } - for _, c := range client.Jar.Cookies(mustParseURL(roomURL)) { - if c.Name == "csrftoken" { - a.csrfToken = c.Value - break - } - } - if a.roomUUID == "" || a.userUUID == "" { - // A closed room's page renders without the room/user data. - return nil, fmt.Errorf("%w: room/user uuid missing", errRoomGone) - } - if a.connToken == "" || a.connURL == "" || a.subURL == "" { - return nil, fmt.Errorf("centrifuge meta missing") - } - if a.csrfToken == "" { - return nil, fmt.Errorf("csrftoken missing") - } - a.channel = fmt.Sprintf("$shared_editor:room-%s", a.roomUUID) - - subBody, _ := json.Marshal(map[string]string{"channel": a.channel}) - req2, _ := http.NewRequestWithContext(ctx, "POST", a.subURL, bytes.NewReader(subBody)) - req2.Header.Set("User-Agent", cupsUA) - req2.Header.Set("Content-Type", "application/json") - req2.Header.Set("X-CSRFToken", a.csrfToken) - req2.Header.Set("Origin", originOf(roomURL)) - req2.Header.Set("Referer", roomURL) - - resp2, err := client.Do(req2) - if err != nil { - return nil, fmt.Errorf("POST sub token: %w", err) - } - body2, _ := io.ReadAll(resp2.Body) - resp2.Body.Close() - if resp2.StatusCode != 200 { - return nil, fmt.Errorf("sub token status %d", resp2.StatusCode) - } - var subResp struct { - Token string `json:"token"` - } - if err := json.Unmarshal(body2, &subResp); err != nil { - return nil, err - } - a.subToken = subResp.Token - if a.subToken == "" { - return nil, fmt.Errorf("empty sub token") - } - utils.Debugf("[CUPS] auth OK: room=%s user=%s", a.roomUUID, a.userUUID) - return a, nil -} - -var ( - // errRoomGone marks a room that no longer exists, as opposed to a network - // hiccup on the way to it. - errRoomGone = errors.New("комната не существует") - // errRefused marks Centrifugo turning our tokens down, with an error - // reply or by hanging up mid-handshake. Only a fresh join fixes that. - errRefused = errors.New("refused") - errNoRoom = errors.New("cups: нет подключённой комнаты") - errStopped = errors.New("cups: транспорт остановлен") -) - -func joinURL(roomUUID string) string { return baseRoomURL + "?room=" + roomUUID } - -// joinRoom enters an existing room. A room that's gone can come back as a -// fresh one (a different uuid) instead of an error, so that counts as gone. -func joinRoom(ctx context.Context, roomUUID string, session *http.Client) (*cupsAuth, error) { - a, err := authorize(ctx, joinURL(roomUUID), session) - if err != nil { - return nil, err - } - if a.roomUUID != roomUUID { - return nil, fmt.Errorf("%w: вместо неё выдана новая %s", errRoomGone, shortRoom(a.roomUUID)) - } - return a, nil -} - -// readReply waits for the reply to command id. Centrifugo reports refusals -// (expired token, unknown channel...) as {"error": ...} in the reply, which -// used to be ignored - the channel then just sat there never receiving -// anything. Centrifugo may also pack several messages into one frame, one per -// line; whatever else comes along (a ping, a push) goes to other. -func readReply(conn *websocket.Conn, timeout time.Duration, id int, what string, other func([]byte)) error { - deadline := time.Now().Add(timeout) - for { - conn.SetReadDeadline(deadline) - _, raw, err := conn.ReadMessage() - if err != nil { - var closeErr *websocket.CloseError - if errors.As(err, &closeErr) { - return fmt.Errorf("%w: %s: %v", errRefused, what, err) - } - return err - } - found := false - var refusal error - for _, line := range splitFrame(raw) { - if !found { - if ok, err := matchReply(line, id); ok { - found, refusal = true, err - continue - } - } - if other != nil { - other(line) - } - } - if refusal != nil { - return fmt.Errorf("%w: %s: %v", errRefused, what, refusal) - } - if found { - return nil - } - } -} - -// matchReply reports whether line is the reply to command id, and the error -// it carries if it is one. -func matchReply(line []byte, id int) (bool, error) { - if isPing(line) { - return false, nil - } - var reply struct { - ID int `json:"id"` - Error *struct { - Code int `json:"code"` - Message string `json:"message"` - } `json:"error"` - } - if json.Unmarshal(line, &reply) != nil || reply.ID != id { - return false, nil - } - if reply.Error != nil { - return true, fmt.Errorf("%d %s", reply.Error.Code, reply.Error.Message) - } - return true, nil -} - -// splitFrame takes apart a frame Centrifugo may have packed several messages -// into, one per line. Parsing such a frame whole fails, which silently lost -// every packet in it. -func splitFrame(raw []byte) [][]byte { - var out [][]byte - for _, line := range bytes.Split(raw, []byte{'\n'}) { - if line = bytes.TrimSpace(line); len(line) > 0 { - out = append(out, line) - } - } - return out -} - -// isPing reports a Centrifugo server ping, which must be answered in kind. -func isPing(line []byte) bool { return string(line) == "{}" } - -func shortRoom(roomUUID string) string { - if len(roomUUID) > 8 { - return roomUUID[:8] - } - return roomUUID -} - -// fmtUptime renders a duration the way the log reads best: "47ч12м", "3м05с". -func fmtUptime(d time.Duration) string { - d = d.Round(time.Second) - h, m, s := int(d.Hours()), int(d.Minutes())%60, int(d.Seconds())%60 - if h > 0 { - return fmt.Sprintf("%dч%02dм", h, m) - } - return fmt.Sprintf("%dм%02dс", m, s) -} - -// parseRoomList pulls room uuids out of --url: the packed base64 list the -// exit node prints, bare or as ?rooms=, or a single ?room=. -func parseRoomList(rawURL string) []string { - rawURL = strings.TrimSpace(rawURL) - if rawURL == "" { - return nil - } - if u, err := url.Parse(rawURL); err == nil { - q := u.Query() - if packed := q.Get("rooms"); packed != "" { - if ids, err := unpackRooms(packed); err == nil { - return ids - } - } else if single := q.Get("room"); single != "" { - return []string{single} - } - } - if ids, err := unpackRooms(rawURL); err == nil { - return ids - } - return nil -} - -func createRooms(ctx context.Context, baseURL string, n int, pause time.Duration) ([]*cupsAuth, error) { - if baseURL == "" { - baseURL = baseRoomURL - } - out := make([]*cupsAuth, 0, n) - delay := pause - if delay <= 0 { - delay = 150 * time.Millisecond - } - for i := 0; i < n; i++ { - var a *cupsAuth - var lastErr error - for attempt := 0; attempt < 6; attempt++ { - a, lastErr = authorize(ctx, baseURL, nil) - if lastErr == nil { - break - } - msg := lastErr.Error() - var wait time.Duration - if strings.Contains(msg, "403") || strings.Contains(msg, "429") { - wait = delay * time.Duration(1< 30*time.Second { - wait = 30 * time.Second - } - } else { - wait = delay * time.Duration(attempt+1) - } - utils.Debugf("[CUPS] room %d attempt %d failed: %v (wait %v)", i+1, attempt+1, lastErr, wait) - if !sleepCtx(ctx, wait) { - return nil, errStopped - } - } - if a == nil { - utils.Debugf("[CUPS] room %d skipped: %v", i+1, lastErr) - continue - } - out = append(out, a) - utils.Debugf("[CUPS] created room %d/%d: %s", i+1, n, a.roomUUID) - if i < n-1 && !sleepCtx(ctx, delay) { - return nil, errStopped - } - } - if len(out) == 0 { - return nil, fmt.Errorf("could not create any room") - } - return out, nil -} - -func packRooms(ids []string) string { - raw, _ := json.Marshal(ids) - return base64.RawURLEncoding.EncodeToString(raw) -} - -func unpackRooms(s string) ([]string, error) { - raw, err := base64.RawURLEncoding.DecodeString(s) - if err != nil { - return nil, fmt.Errorf("base64 decode: %w", err) - } - var ids []string - if err := json.Unmarshal(raw, &ids); err != nil { - return nil, fmt.Errorf("json decode: %w", err) - } - return ids, nil -} - -// sleepCtx waits for d, or less if ctx ends first; it reports whether the -// whole wait went by. -func sleepCtx(ctx context.Context, d time.Duration) bool { - t := time.NewTimer(d) - defer t.Stop() - select { - case <-t.C: - return true - case <-ctx.Done(): - return false - } -} - -// ---- per-channel stats ---- - -type channelStats struct { - idx int - roomUUID string - packetsSent atomic.Uint64 - packetsRecv atomic.Uint64 - bytesSent atomic.Uint64 - bytesRecv atomic.Uint64 - batchesSent atomic.Uint64 - batchesRecv atomic.Uint64 - reconnects atomic.Uint64 - dropped atomic.Uint64 -} - -// ---- WS ---- - -type cupsWS struct { - idx int - roomUUID string - joinedAt time.Time - - // authPtr is swapped on every re-join (fresh tokens), so it's atomic. It - // stays nil until the room has been entered at least once. - authPtr atomic.Pointer[cupsAuth] - // session keeps the cookies between joins (run's goroutine only). - session *http.Client - - // dead is set while the room looks closed; onRoomState reports the - // change so the operator can see when (and after how long) it happened. - dead atomic.Bool - goneDelay time.Duration // run's goroutine only - lastSend time.Time // sendLoop goroutine only; paces messages - onRoomState func() - - config CupsonlineConfig - conn *websocket.Conn - - writeMu sync.Mutex - rpcID atomic.Int64 - onData func([]byte) - - ctx context.Context - connected atomic.Bool - - // recvBuf reassembles the peer's byte stream across messages (one packet - // may be split over several). Touched only by the read goroutine, and - // reset on every reconnect - a packet split across a drop is lost, like - // on any link. - recvBuf []byte - - sendQueue chan []byte - - stats *channelStats -} - -func (w *cupsWS) nextID() int64 { return w.rpcID.Add(1) } - -func (w *cupsWS) auth() *cupsAuth { return w.authPtr.Load() } - -func (w *cupsWS) writeRaw(data []byte) error { - w.writeMu.Lock() - defer w.writeMu.Unlock() - if w.conn == nil { - return fmt.Errorf("ws not connected") - } - w.conn.SetWriteDeadline(time.Now().Add(w.config.WSWriteTimeout)) - if err := w.conn.WriteMessage(websocket.TextMessage, data); err != nil { - // A failed or timed-out write leaves the socket unusable. Closing it - // makes the reader notice and reconnect instead of sitting on a dead - // channel until the read timeout. - w.conn.Close() - return err - } - return nil -} - -func (w *cupsWS) writeJSON(v interface{}) error { - data, err := json.Marshal(v) - if err != nil { - return err - } - return w.writeRaw(data) -} - -// currentConn returns the currently active WebSocket, if any. -func (w *cupsWS) currentConn() *websocket.Conn { - w.writeMu.Lock() - defer w.writeMu.Unlock() - return w.conn -} - -// roomDeadAfterFails is how many joins in a row may fail to reach a -// subscribed channel (with fresh tokens each time) before the room is -// reported closed even though its page still loads. -const roomDeadAfterFails = 5 - -func (w *cupsWS) run() { - delay := w.config.ReconnectMinDelay - fails := 0 - // Start hands over freshly joined rooms; one it couldn't enter has no - // tokens yet and has to be joined first. - needJoin := w.auth() == nil - for w.ctx.Err() == nil { - if needJoin || w.dead.Load() { - if err := w.join(); err != nil { - if !sleepCtx(w.ctx, w.retryDelay(delay)) { - return - } - delay = w.backoff(delay) - continue - } - needJoin = false - } - - err := w.connectAndServe() - wasReady := w.connected.Swap(false) - if w.ctx.Err() != nil { - return - } - utils.Debugf("[CUPS] ws error (%s): %v", w.roomUUID, err) - if wasReady { - delay = w.config.ReconnectMinDelay - fails = 0 - } else if fails++; fails == roomDeadAfterFails { - w.markDead(fmt.Errorf("канал не поднимается %d раз подряд: %v", fails, err)) - } - w.stats.reconnects.Add(1) - // Tokens come from the room page once per join and can expire, so - // reconnecting with the old ones may never work again. Re-join when - // Centrifugo turned them down, and every other failed attempt in case - // it said so just by hanging up. A connection that was working and - // dropped just reconnects: a join is two page loads, and a reconnect - // storm of those is how an IP gets rate-limited. - needJoin = errors.Is(err, errRefused) || (!wasReady && fails%2 == 0) - if !sleepCtx(w.ctx, w.retryDelay(delay)) { - return - } - delay = w.backoff(delay) - } -} - -// retryDelay is the wait before the next attempt: the usual backoff, or the -// much slower probe interval while the room looks closed. -func (w *cupsWS) retryDelay(delay time.Duration) time.Duration { - if !w.dead.Load() { - return delay - } - d := w.goneDelay - w.goneDelay = min(2*d, w.config.RoomGoneRetryMax) - return d -} - -func (w *cupsWS) backoff(delay time.Duration) time.Duration { - return min(time.Duration(float64(delay)*w.config.ReconnectMultiplier), w.config.ReconnectMaxDelay) -} - -// join enters the room for fresh tokens, in the same session as before. This -// is also where a closed room shows itself. -func (w *cupsWS) join() error { - a, err := joinRoom(w.ctx, w.roomUUID, w.session) - if err != nil { - // Next time from a clean session, in case this one is what's broken. - w.session = nil - if errors.Is(err, errRoomGone) { - w.markDead(err) - } else if w.ctx.Err() == nil { - utils.Debugf("[CUPS] join %s failed: %v", w.roomUUID, err) - } - return err - } - w.session = a.httpClient - w.authPtr.Store(a) - return nil -} - -func (w *cupsWS) markDead(err error) { - if w.dead.Swap(true) { - return - } - utils.Infof("[ERROR] Cups: комната %s закрылась (в работе %s): %v", - shortRoom(w.roomUUID), fmtUptime(time.Since(w.joinedAt)), err) - if w.onRoomState != nil { - w.onRoomState() - } -} - -func (w *cupsWS) markAlive() { - w.goneDelay = w.config.RoomGoneRetryMin - if !w.dead.Swap(false) { - return - } - utils.Infof("[CUPS] Cups: комната %s снова доступна", shortRoom(w.roomUUID)) - if w.onRoomState != nil { - w.onRoomState() - } -} - -func (w *cupsWS) connectAndServe() error { - a := w.auth() - wsURL := strings.Replace(a.connURL, "https://", "wss://", 1) - wsURL = strings.Replace(wsURL, "http://", "ws://", 1) - wsURL = strings.TrimRight(wsURL, "/") + "/websocket" - - header := http.Header{} - header.Set("Origin", originOf(a.connURL)) - header.Set("User-Agent", cupsUA) - for _, c := range a.httpClient.Jar.Cookies(mustParseURL(a.connURL)) { - header.Add("Cookie", c.Name+"="+c.Value) - } - - dialer := websocket.Dialer{ - NetDialContext: netbind.DialContext, - HandshakeTimeout: w.config.WSHandshakeTimeout, - ReadBufferSize: w.config.ReadBufferSize, - WriteBufferSize: w.config.WriteBufferSize, - } - conn, _, err := dialer.DialContext(w.ctx, wsURL, header) - if err != nil { - return fmt.Errorf("dial: %w", err) - } - // Stop mustn't have to wait for the server's next message to take the - // reader down. - unhook := context.AfterFunc(w.ctx, func() { conn.Close() }) - defer unhook() - w.writeMu.Lock() - w.conn = conn - w.writeMu.Unlock() - defer func() { - w.writeMu.Lock() - w.conn = nil - w.writeMu.Unlock() - conn.Close() - }() - - conn.SetReadLimit(int64(w.config.MaxMessageBytes)) - // Fresh socket, fresh stream: drop any half-assembled packet from before. - w.recvBuf = w.recvBuf[:0] - - if err := w.writeJSON(map[string]interface{}{ - "id": 1, "connect": map[string]interface{}{"token": a.connToken, "name": "js"}, - }); err != nil { - return err - } - if err := readReply(conn, w.config.WSHandshakeTimeout, 1, "connect", w.handleReply); err != nil { - return err - } - if err := w.writeJSON(map[string]interface{}{ - "id": 2, "subscribe": map[string]interface{}{"channel": a.channel, "token": a.subToken}, - }); err != nil { - return err - } - if err := readReply(conn, w.config.WSHandshakeTimeout, 2, "subscribe", w.handleReply); err != nil { - return err - } - w.connected.Store(true) - w.markAlive() - utils.Debugf("[CUPS] WS ready: %s", a.roomUUID) - - kaStop := make(chan struct{}) - go w.keepAliveLoop(kaStop) - defer close(kaStop) - - for { - conn.SetReadDeadline(time.Now().Add(w.config.WSReadTimeout)) - _, raw, err := conn.ReadMessage() - if err != nil { - return err - } - w.handleMessage(raw) - } -} - -func (w *cupsWS) keepAliveLoop(stop chan struct{}) { - t := time.NewTicker(w.config.KeepAliveInterval) - defer t.Stop() - for { - select { - case <-stop: - return - case <-w.ctx.Done(): - return - case <-t.C: - _ = w.writeJSON(map[string]interface{}{ - "rpc": map[string]interface{}{ - "method": "shared_editor_ping", - "data": map[string]interface{}{"room": w.auth().roomUUID, "user": w.auth().userUUID}, - }, - "id": w.nextID(), - }) - } - } -} - -func (w *cupsWS) sendLoop() { - batch := make([][]byte, 0, w.config.BatchMaxPackets) - totalBytes := 0 - timer := time.NewTimer(w.config.BatchTimeout) - timer.Stop() - defer timer.Stop() - - flush := func() { - timer.Stop() - if len(batch) == 0 { - return - } - if !w.connected.Load() { - // The channel went down with this batch still here. New traffic - // has already moved to a room that's up, and by the time this one - // is back the batch is stale, so it's dropped rather than - // delivered late and out of order. - w.stats.dropped.Add(uint64(len(batch))) - } else if err := w.sendBatch(batch); err != nil { - w.stats.dropped.Add(uint64(len(batch))) - utils.Debugf("[CUPS] batch send (%s): %v", w.roomUUID, err) - } - clear(batch) - batch = batch[:0] - totalBytes = 0 - } - - for { - select { - case <-w.ctx.Done(): - return - case pkt := <-w.sendQueue: - // Flush before a packet that would overflow the batch, not after - // it: otherwise a batch ends up a whole packet over the limit. - if len(batch) > 0 && totalBytes+2+len(pkt) > w.config.BatchMaxBytes { - flush() - } - batch = append(batch, pkt) - totalBytes += 2 + len(pkt) - if len(batch) >= w.config.BatchMaxPackets || totalBytes >= w.config.BatchMaxBytes { - flush() - } else if len(batch) == 1 { - timer.Reset(w.config.BatchTimeout) - } - case <-timer.C: - flush() - } - } -} - -func (w *cupsWS) sendBatch(batch [][]byte) error { - var buf bytes.Buffer - var hdr [2]byte - totalRaw := 0 - for _, p := range batch { - binary.BigEndian.PutUint16(hdr[:], uint16(len(p))) - buf.Write(hdr[:]) - buf.Write(p) - totalRaw += len(p) - } - blob := buf.Bytes() - - // One message can only carry so much before the server drops it, so a - // batch is cut into chunks and each goes as its own message. The far side - // stitches them back together in order. - for off := 0; off < len(blob); off += w.config.MaxMessageData { - end := min(off+w.config.MaxMessageData, len(blob)) - if err := w.pace(); err != nil { - return err - } - if err := w.sendChunk(blob[off:end]); err != nil { - return err - } - } - w.stats.packetsSent.Add(uint64(len(batch))) - w.stats.bytesSent.Add(uint64(totalRaw)) - w.stats.batchesSent.Add(1) - return nil -} - -// pace waits out the rest of SendInterval since the last message, so the -// channel never floods the server. Returns errStopped if the transport is -// stopped while waiting. -func (w *cupsWS) pace() error { - if w.config.SendInterval <= 0 { - return nil - } - if d := w.config.SendInterval - time.Since(w.lastSend); d > 0 { - if !sleepCtx(w.ctx, d) { - return errStopped - } - } - w.lastSend = time.Now() - return nil -} - -// sendChunk sends one message: a 2-byte length of the chunk, then the chunk, -// packed into cursor integers. The length lets the peer strip the zero -// padding exactly before stitching the byte stream back. -func (w *cupsWS) sendChunk(chunk []byte) error { - payload := make([]byte, 2+len(chunk)) - binary.BigEndian.PutUint16(payload[:2], uint16(len(chunk))) - copy(payload[2:], chunk) - - msg := map[string]interface{}{ - "rpc": map[string]interface{}{ - "method": "shared_editor_change_cursors", - "data": map[string]interface{}{ - "cursors": cursorsFromBytes(payload), - "ranges": []interface{}{}, - "room": w.auth().roomUUID, - "user": w.auth().userUUID, - }, - }, - "id": w.nextID(), - } - return w.writeJSON(msg) -} - -// bytesPerNumber is how many bytes ride in one row/column integer. cups.online -// keeps a cursor integer exact up to 2^53-1; 6 bytes (48 bits) stays well -// inside that and is byte-aligned. -const bytesPerNumber = 6 - -// cursorsFromBytes packs a blob into {row, column} integer pairs. The data -// rides in the numbers because cups.online rejects a non-numeric column but -// relays a cursor array of integers verbatim and in order. Trailing padding is -// zero, which the length-prefixed framing on the far side reads as "no more -// packets" and stops. -func cursorsFromBytes(blob []byte) []map[string]interface{} { - nums := packNumbers(blob) - cursors := make([]map[string]interface{}, 0, (len(nums)+1)/2) - for i := 0; i < len(nums); i += 2 { - c := map[string]interface{}{"row": nums[i], "column": uint64(0)} - if i+1 < len(nums) { - c["column"] = nums[i+1] - } - cursors = append(cursors, c) - } - return cursors -} - -// packNumbers cuts a blob into big-endian bytesPerNumber-byte integers, the -// last one zero-padded. -func packNumbers(blob []byte) []uint64 { - n := (len(blob) + bytesPerNumber - 1) / bytesPerNumber - out := make([]uint64, n) - for i := range out { - var v uint64 - for j := 0; j < bytesPerNumber; j++ { - v <<= 8 - if idx := i*bytesPerNumber + j; idx < len(blob) { - v |= uint64(blob[idx]) - } - } - out[i] = v - } - return out -} - -// bytesFromNumbers reverses packNumbers. -func bytesFromNumbers(nums []uint64) []byte { - out := make([]byte, 0, len(nums)*bytesPerNumber) - for _, v := range nums { - var b [bytesPerNumber]byte - for j := bytesPerNumber - 1; j >= 0; j-- { - b[j] = byte(v) - v >>= 8 - } - out = append(out, b[:]...) - } - return out -} - -// asCursorInt reads a JSON number back as an integer. Values stay under 2^48, -// which float64 (how encoding/json hands back a number) holds exactly. -func asCursorInt(v interface{}) (uint64, bool) { - f, ok := v.(float64) - if !ok || f < 0 { - return 0, false - } - return uint64(f), true -} - -func (w *cupsWS) handleMessage(raw []byte) { - for _, line := range splitFrame(raw) { - w.handleReply(line) - } -} - -// handleReply handles one message: a server ping, or a push that may carry -// a batch from the peer. -func (w *cupsWS) handleReply(raw []byte) { - if isPing(raw) { - _ = w.writeRaw([]byte("{}")) - return - } - var obj map[string]interface{} - if err := json.Unmarshal(raw, &obj); err != nil { - return - } - push, ok := obj["push"].(map[string]interface{}) - if !ok { - return - } - pub, _ := push["pub"].(map[string]interface{}) - data, _ := pub["data"].(map[string]interface{}) - if data == nil { - return - } - if t, _ := data["type"].(string); t != "cursors_update" { - return - } - payload, _ := data["payload"].(map[string]interface{}) - if payload == nil { - return - } - if uuid, _ := payload["user_uuid"].(string); uuid == w.auth().userUUID { - return - } - cursors, _ := payload["cursors"].([]interface{}) - if len(cursors) == 0 { - return - } - // Rebuild the blob from the cursor integers, two per cursor in order, the - // reverse of cursorsFromBytes. - nums := make([]uint64, 0, len(cursors)*2) - for _, c := range cursors { - m, _ := c.(map[string]interface{}) - if m == nil { - return - } - row, ok := asCursorInt(m["row"]) - if !ok { - return - } - col, _ := asCursorInt(m["column"]) - nums = append(nums, row, col) - } - decoded := bytesFromNumbers(nums) - // Strip the padding: the first two bytes are how many real bytes follow. - if len(decoded) < 2 { - return - } - dataLen := int(binary.BigEndian.Uint16(decoded[:2])) - if 2+dataLen > len(decoded) { - return - } - chunk := decoded[2 : 2+dataLen] - - // Append to the running stream and pull out whole packets; a packet split - // across messages completes once the rest of it arrives. - w.recvBuf = append(w.recvBuf, chunk...) - buf := w.recvBuf - adv, count, total := 0, 0, 0 - for len(buf)-adv >= 2 { - ln := int(binary.BigEndian.Uint16(buf[adv : adv+2])) - if ln == 0 { - // Not a real length - the stream is out of sync; drop it. - adv = len(buf) - break - } - if len(buf)-adv < 2+ln { - break // rest of this packet hasn't arrived yet - } - pkt := buf[adv+2 : adv+2+ln] - if w.onData != nil { - w.onData(pkt) - } - adv += 2 + ln - count++ - total += ln - } - // Keep only the unconsumed tail. onData ran already, so moving the bytes - // now can't disturb a packet still in flight. - w.recvBuf = append(w.recvBuf[:0], buf[adv:]...) - // A stream that never yields a packet must not grow without bound. - if len(w.recvBuf) > w.config.MaxPayloadBytes+w.config.MaxMessageData { - utils.Debugf("[CUPS] recv stream out of sync (%s), resetting", w.roomUUID) - w.recvBuf = w.recvBuf[:0] - } - if count > 0 { - w.stats.packetsRecv.Add(uint64(count)) - w.stats.bytesRecv.Add(uint64(total)) - w.stats.batchesRecv.Add(1) - } -} - -func (w *cupsWS) Send(data []byte) error { - select { - case w.sendQueue <- data: - return nil - case <-w.ctx.Done(): - return errStopped - } -} - -// ---- Transport ---- - -type CupsonlineTransport struct { - *transport.BaseTransport - - baseURL string - // roomIDs are the rooms to join, from --url. Empty on an exit node - // means "create new ones"; a client can't start without them. - roomIDs []string - config CupsonlineConfig - isClient bool - clientErr error - - // ctx ends with Stop and takes every room's goroutines down with it. - ctx context.Context - cancel context.CancelFunc - - wss []*cupsWS - // rr round-robins outbound frames across the connected rooms; see Send. - rr atomic.Uint32 - - statsStart time.Time -} - -func NewCupsonlineTransport(rawURL string, cfg transport.TransportConfig, isClient bool) *CupsonlineTransport { - ctx, cancel := context.WithCancel(context.Background()) - t := &CupsonlineTransport{ - BaseTransport: transport.NewBaseTransport(cfg), - config: DefaultCupsonlineConfig(), - isClient: isClient, - ctx: ctx, - cancel: cancel, - statsStart: time.Now(), - } - - t.baseURL = baseRoomURL - // Both roles accept a room list: the client always needs one, and an - // exit node given one re-joins those rooms instead of creating new - // ones, so a restart doesn't hand the phone a new string every time. - t.roomIDs = parseRoomList(rawURL) - if isClient && len(t.roomIDs) == 0 { - t.clientErr = fmt.Errorf("client mode: --url must contain base64 room list") - } - return t -} - -// roomSlot is one room the transport will keep a channel to. -type roomSlot struct { - id string - auth *cupsAuth // nil if it couldn't be entered yet - err error // why it couldn't -} - -func (t *CupsonlineTransport) Start() error { - if t.clientErr != nil { - return t.clientErr - } - if err := t.BaseTransport.Start(); err != nil { - return err - } - - slots, err := t.enterRooms() - if err != nil { - return err - } - if t.ctx.Err() != nil { - return errStopped - } - - joinedAt := time.Now() - wss := make([]*cupsWS, len(slots)) - for i, s := range slots { - ws := &cupsWS{ - idx: i, - roomUUID: s.id, - joinedAt: joinedAt, - goneDelay: t.config.RoomGoneRetryMin, - onRoomState: t.reportRooms, - config: t.config, - ctx: t.ctx, - sendQueue: make(chan []byte, t.config.SendQueueSize), - stats: &channelStats{idx: i, roomUUID: s.id}, - } - // Command ids 1 and 2 are connect and subscribe. - ws.rpcID.Store(2) - if s.auth != nil { - ws.authPtr.Store(s.auth) - ws.session = s.auth.httpClient - } - // Already reported by enterRooms; run() only probes it now and then. - ws.dead.Store(errors.Is(s.err, errRoomGone)) - ws.onData = t.CallReceive - wss[i] = ws - } - t.wss = wss - for _, ws := range wss { - go ws.run() - go ws.sendLoop() - } - - utils.Debugf("[CUPS] transport started: %d channels", len(t.wss)) - t.SetConnected(true) - - go t.statsLoop() - return nil -} - -// enterRooms joins the rooms from --url, or, on an exit node without them, -// creates new ones and prints the string for the client. -func (t *CupsonlineTransport) enterRooms() ([]roomSlot, error) { - if len(t.roomIDs) > 0 { - slots := t.joinListed() - joined, allGone := 0, true - for _, s := range slots { - if s.auth != nil { - joined++ - } else if !errors.Is(s.err, errRoomGone) { - allGone = false - } - } - if t.ctx.Err() != nil { - return nil, errStopped - } - utils.Infof("[CUPS] Cups: зашли в %d из %d комнат", joined, len(slots)) - // The rooms that failed get a channel all the same and keep being - // retried. Dropping them left the peer's traffic in those rooms with - // nobody listening. - if joined > 0 { - return slots, nil - } - switch { - case t.isClient && allGone: - utils.Infof("[ERROR] Cups: все комнаты закрыты - пересоздайте комнаты на ноде (запуск без --url) и обновите строку в профиле") - return nil, fmt.Errorf("no rooms joined: all rooms are gone") - case t.isClient: - utils.Infof("[ERROR] Cups: ни одна комната не отвечает - проверьте сеть") - return nil, fmt.Errorf("no rooms joined") - case !allGone: - // New rooms because of a network hiccup would cost the phone its - // string for nothing. Fail, and let the restart try again. - utils.Infof("[ERROR] Cups: сохранённые комнаты не отвечают - новые не создаю, чтобы строка на телефоне осталась прежней") - return nil, fmt.Errorf("saved rooms unreachable") - } - utils.Infof("[ERROR] Cups: сохранённые комнаты закрыты, создаю новые - строку на телефоне нужно будет заменить") - } - - auths, err := createRooms(t.ctx, t.baseURL, t.config.NumRooms, t.config.RoomCreatePause) - if err != nil { - return nil, fmt.Errorf("create rooms: %w", err) - } - slots := make([]roomSlot, len(auths)) - ids := make([]string, len(auths)) - for i, a := range auths { - slots[i] = roomSlot{id: a.roomUUID, auth: a} - ids[i] = a.roomUUID - } - fmt.Printf("\n=== COPY THIS TO CLIENT ===\n") - fmt.Printf("%s\n", packRooms(ids)) - fmt.Printf("===========================\n") - fmt.Printf("Save it and pass it back as --url to reuse these rooms after a restart.\n\n") - return slots, nil -} - -// joinListed enters all the listed rooms at once rather than one by one: on -// a phone each join is two round trips, and Start blocks the app until -// they're done. -func (t *CupsonlineTransport) joinListed() []roomSlot { - slots := make([]roomSlot, len(t.roomIDs)) - var wg sync.WaitGroup - for i, id := range t.roomIDs { - slots[i].id = id - wg.Add(1) - go func() { - defer wg.Done() - slots[i].auth, slots[i].err = joinRoom(t.ctx, id, nil) - }() - } - wg.Wait() - for _, s := range slots { - if s.err != nil && t.ctx.Err() == nil { - utils.Infof("[ERROR] Cups: не удалось зайти в комнату %s: %v", shortRoom(s.id), s.err) - } - } - return slots -} - -// Stop can come at any time, even before Start is done, and more than once. -func (t *CupsonlineTransport) Stop() error { - t.cancel() - t.SetConnected(false) - return t.BaseTransport.Stop() -} - -// Send spreads frames across all the connected rooms, round-robin, so the -// four channels carry the load together instead of one. It can't keep a flow -// in one room - what arrives here is already a codec frame mixing many flows, -// with no IP header to hash - so frames from different rooms can arrive out of -// order; a whole frame always stays within one room (reassembly is per-room), -// and TCP above puts the rest back in order. A frame is never split across -// rooms. -func (t *CupsonlineTransport) Send(data []byte) error { - if len(data) > t.config.MaxPayloadBytes { - return fmt.Errorf("cups: посылка %d байт больше предела %d", len(data), t.config.MaxPayloadBytes) - } - ws := t.pickRoom() - if ws == nil { - return errNoRoom - } - return ws.Send(data) -} - -// pickRoom returns the next connected room in round-robin order, or nil when -// none is up. -func (t *CupsonlineTransport) pickRoom() *cupsWS { - n := len(t.wss) - if n == 0 { - return nil - } - start := int(t.rr.Add(1)) - for i := 0; i < n; i++ { - ws := t.wss[(start+i)%n] - if ws.connected.Load() { - return ws - } - } - return nil -} - -func (t *CupsonlineTransport) Receive(callback func([]byte)) { - t.BaseTransport.Receive(callback) -} - -// IsConnected is true while at least one room's channel is up, so the app -// stops claiming "connected" once every room has gone away. -func (t *CupsonlineTransport) IsConnected() bool { - if !t.BaseTransport.IsConnected() { - return false - } - for _, ws := range t.wss { - if ws.connected.Load() { - return true - } - } - return false -} - -// reportRooms logs how many rooms are still alive whenever one of them -// closes or comes back. -func (t *CupsonlineTransport) reportRooms() { - alive := 0 - for _, ws := range t.wss { - if !ws.dead.Load() { - alive++ - } - } - utils.Infof("[CUPS] Cups: живых комнат %d/%d", alive, len(t.wss)) - if alive == 0 { - utils.Infof("[ERROR] Cups: все комнаты недоступны - пересоздайте комнаты на ноде (запуск без --url) и обновите строку в профиле") - } -} - -func (t *CupsonlineTransport) Stats() transport.TransportStats { - var sent, recv, bytesSent, bytesRecv, reconnects uint64 - for _, ws := range t.wss { - sent += ws.stats.packetsSent.Load() - recv += ws.stats.packetsRecv.Load() - bytesSent += ws.stats.bytesSent.Load() - bytesRecv += ws.stats.bytesRecv.Load() - reconnects += ws.stats.reconnects.Load() - } - base := t.BaseTransport.Stats() - return transport.TransportStats{ - BytesSent: bytesSent, - BytesReceived: bytesRecv, - PacketsSent: sent, - PacketsRecv: recv, - Reconnects: reconnects, - Connected: t.IsConnected(), - Uptime: base.Uptime, - } -} - -// statsLoop — печатает per-channel статистику каждые StatsInterval секунд. -func (t *CupsonlineTransport) statsLoop() { - tick := time.NewTicker(t.config.StatsInterval) - defer tick.Stop() - - lastSent := make([]uint64, len(t.wss)) - lastRecv := make([]uint64, len(t.wss)) - lastBytesSent := make([]uint64, len(t.wss)) - lastBytesRecv := make([]uint64, len(t.wss)) - var lastTotalSent, lastTotalRecv, lastTotalBytesSent, lastTotalBytesRecv uint64 - - for { - select { - case <-t.ctx.Done(): - return - case <-tick.C: - var totalSent, totalRecv, totalBytesSent, totalBytesRecv uint64 - for i, ws := range t.wss { - s := ws.stats.packetsSent.Load() - r := ws.stats.packetsRecv.Load() - bs := ws.stats.bytesSent.Load() - br := ws.stats.bytesRecv.Load() - - dS := s - lastSent[i] - dR := r - lastRecv[i] - dBS := bs - lastBytesSent[i] - dBR := br - lastBytesRecv[i] - - lastSent[i] = s - lastRecv[i] = r - lastBytesSent[i] = bs - lastBytesRecv[i] = br - - totalSent += s - totalRecv += r - totalBytesSent += bs - totalBytesRecv += br - - utils.Debugf("[CH-%02d %s] tx=%d pkt/s (%.1f KB/s) rx=%d pkt/s (%.1f KB/s) drop=%d reconn=%d conn=%v", - i, shortRoom(ws.roomUUID), - dS/uint64(t.config.StatsInterval.Seconds()), - float64(dBS)/t.config.StatsInterval.Seconds()/1024, - dR/uint64(t.config.StatsInterval.Seconds()), - float64(dBR)/t.config.StatsInterval.Seconds()/1024, - ws.stats.dropped.Load(), - ws.stats.reconnects.Load(), - ws.connected.Load(), - ) - } - dtS := totalSent - lastTotalSent - dtR := totalRecv - lastTotalRecv - dtBS := totalBytesSent - lastTotalBytesSent - dtBR := totalBytesRecv - lastTotalBytesRecv - lastTotalSent = totalSent - lastTotalRecv = totalRecv - lastTotalBytesSent = totalBytesSent - lastTotalBytesRecv = totalBytesRecv - - utils.Debugf("[CUPS-TOTAL] tx=%d pkt/s (%.1f KB/s) rx=%d pkt/s (%.1f KB/s) uptime=%v", - dtS/uint64(t.config.StatsInterval.Seconds()), - float64(dtBS)/t.config.StatsInterval.Seconds()/1024, - dtR/uint64(t.config.StatsInterval.Seconds()), - float64(dtBR)/t.config.StatsInterval.Seconds()/1024, - time.Since(t.statsStart).Round(time.Second), - ) - } - } -} - -func (t *CupsonlineTransport) RoomUUIDs() []string { - out := make([]string, 0, len(t.wss)) - for _, ws := range t.wss { - out = append(out, ws.roomUUID) - } - return out -} - -// ---- helpers ---- - -func firstMatch(re *regexp.Regexp, s string) string { - m := re.FindStringSubmatch(s) - if len(m) < 2 { - return "" - } - return m[1] -} - -func originOf(rawURL string) string { - u, err := url.Parse(rawURL) - if err != nil { - return "" - } - return u.Scheme + "://" + u.Host -} - -func mustParseURL(rawURL string) *url.URL { - u, err := url.Parse(rawURL) - if err != nil { - panic(err) - } - return u -} - -// ---- CookieExchanger ---- - -// FetchCookies returns the cookies of the first room that has been entered, -// as name -> value. Every room keeps its own session (see cupsWS.session), -// so re-joining doesn't add a new participant; they all live on one site. -func (t *CupsonlineTransport) FetchCookies() (map[string]string, error) { - u, err := url.Parse(baseRoomURL) - if err != nil { - return nil, err - } - out := make(map[string]string) - for _, ws := range t.wss { - if a := ws.auth(); a != nil { - for _, c := range a.httpClient.Jar.Cookies(u) { - out[c.Name] = c.Value - } - break - } - } - return out, nil -} - -// ApplyCookies puts the provided values into every room's session and drops -// the connections, so they come back carrying them. -func (t *CupsonlineTransport) ApplyCookies(values map[string]string) error { - if len(values) == 0 { - return nil - } - u, _ := url.Parse(baseRoomURL) - cookies := make([]*http.Cookie, 0, len(values)) - for k, v := range values { - cookies = append(cookies, &http.Cookie{Name: k, Value: v, Path: "/"}) - } - utils.Debugf("[CUPS] applied %d cookies, forcing reconnects", len(cookies)) - for _, ws := range t.wss { - if a := ws.auth(); a != nil { - a.httpClient.Jar.SetCookies(u, cookies) - } - if conn := ws.currentConn(); conn != nil { - _ = conn.Close() - } - } - return nil -} diff --git a/transport/cupsonline/cupsonline_test.go b/transport/cupsonline/cupsonline_test.go deleted file mode 100644 index e45e702..0000000 --- a/transport/cupsonline/cupsonline_test.go +++ /dev/null @@ -1,115 +0,0 @@ -package cupsonline - -import ( - "context" - "errors" - "net/http" - "net/http/httptest" - "strings" - "testing" - "time" - - "github.com/gorilla/websocket" - - "openflux/transport" -) - -const ( - roomA = "11111111-1111-1111-1111-111111111111" - roomB = "22222222-2222-2222-2222-222222222222" -) - -func TestParseRoomListAcceptsEveryForm(t *testing.T) { - packed := packRooms([]string{roomA, roomB}) - for name, in := range map[string]string{ - "bare base64": packed, - "?rooms=": "https://example/?rooms=" + packed, - "surrounding ": " " + packed + "\n", - } { - got := parseRoomList(in) - if strings.Join(got, ",") != roomA+","+roomB { - t.Errorf("%s: got %q", name, got) - } - } - if got := parseRoomList("https://example/?room=" + roomA); len(got) != 1 || got[0] != roomA { - t.Errorf("?room=: got %q", got) - } - for _, in := range []string{"", "http://#", "not base64 at all"} { - if got := parseRoomList(in); len(got) != 0 { - t.Errorf("%q: want no rooms, got %q", in, got) - } - } -} - -// An exit node given the printed list must re-join those rooms rather than -// create new ones; without one it creates. A client can't start without one. -func TestExitNodeReusesRoomsFromURL(t *testing.T) { - packed := packRooms([]string{roomA, roomB}) - cfg := transport.DefaultConfig() - - exit := NewCupsonlineTransport(packed, cfg, false) - if len(exit.roomIDs) != 2 || exit.clientErr != nil { - t.Fatalf("exit with a list: roomIDs=%q err=%v", exit.roomIDs, exit.clientErr) - } - fresh := NewCupsonlineTransport("http://#", cfg, false) - if len(fresh.roomIDs) != 0 || fresh.clientErr != nil { - t.Fatalf("exit without a list must create rooms: roomIDs=%q err=%v", fresh.roomIDs, fresh.clientErr) - } - if client := NewCupsonlineTransport("", cfg, true); client.clientErr == nil { - t.Fatal("a client without a room list must refuse to start") - } -} - -func TestAuthorizeReportsGoneRooms(t *testing.T) { - for name, handler := range map[string]http.HandlerFunc{ - "404": func(w http.ResponseWriter, r *http.Request) { http.NotFound(w, r) }, - "no uuid": func(w http.ResponseWriter, r *http.Request) { w.Write([]byte("room closed")) }, - "server 502": func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusBadGateway) }, - } { - srv := httptest.NewServer(handler) - _, err := authorize(context.Background(), srv.URL, nil) - srv.Close() - gone := errors.Is(err, errRoomGone) - if wantGone := name != "server 502"; gone != wantGone { - t.Errorf("%s: gone=%v want %v (err %v)", name, gone, wantGone, err) - } - } -} - -func TestReadReplySurfacesRefusal(t *testing.T) { - srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - conn, err := (&websocket.Upgrader{}).Upgrade(w, r, nil) - if err != nil { - return - } - defer conn.Close() - conn.WriteMessage(websocket.TextMessage, []byte(`{"id":1,"connect":{"client":"x"}}`)) - conn.WriteMessage(websocket.TextMessage, []byte(`{"id":2,"error":{"code":109,"message":"token expired"}}`)) - conn.ReadMessage() - })) - defer srv.Close() - conn, _, err := websocket.DefaultDialer.Dial("ws"+strings.TrimPrefix(srv.URL, "http"), nil) - if err != nil { - t.Fatal(err) - } - defer conn.Close() - - if err := readReply(conn, time.Second, 1, "connect", nil); err != nil { - t.Fatalf("a normal reply must pass: %v", err) - } - err = readReply(conn, time.Second, 2, "subscribe", nil) - if err == nil || !strings.Contains(err.Error(), "token expired") || !errors.Is(err, errRefused) { - t.Fatalf("a refusal must come back as errRefused, got %v", err) - } -} - -func TestFmtUptime(t *testing.T) { - for d, want := range map[time.Duration]string{ - 47*time.Hour + 12*time.Minute: "47ч12м", - 3*time.Minute + 5*time.Second: "3м05с", - } { - if got := fmtUptime(d); got != want { - t.Errorf("%v: got %q want %q", d, got, want) - } - } -} diff --git a/transport/cupsonline/fakecups_test.go b/transport/cupsonline/fakecups_test.go deleted file mode 100644 index 41708bb..0000000 --- a/transport/cupsonline/fakecups_test.go +++ /dev/null @@ -1,842 +0,0 @@ -package cupsonline - -import ( - "bytes" - "crypto/rand" - "encoding/json" - "fmt" - "net/http" - "net/http/httptest" - "runtime" - "strings" - "sync" - "testing" - "time" - - "github.com/gorilla/websocket" - - "openflux/transport" -) - -// fakeCups stands in for cups.online: room pages, subscription tokens and a -// Centrifugo that relays cursor updates to everyone in a room. The tests run -// an exit node and a client through it end to end, and break it on purpose. -type fakeCups struct { - srv *httptest.Server - - mu sync.Mutex - users map[string]string // session cookie -> user uuid - rooms map[string]bool // room uuid -> still open - subs map[string]map[*fakeConn]bool - conns map[*fakeConn]bool - pageHits map[string]int - failPages map[string]int // room -> page loads still to fail with 502 - failJoins bool // every existing room's page fails with 502; creating still works - // redirectGone makes a closed room's page hand out a brand-new room, the - // way cups does, instead of a 404. - redirectGone bool - // pack sends a ping and the push in one frame, one per line. - pack bool - // stall makes page loads hang until it's closed, like a dead network. - stall chan struct{} - created int - maxCursors int -} - -type fakeConn struct { - ws *websocket.Conn - mu sync.Mutex -} - -func (c *fakeConn) send(msg string) { - c.mu.Lock() - defer c.mu.Unlock() - c.ws.WriteMessage(websocket.TextMessage, []byte(msg)) -} - -func newFakeCups(t *testing.T) *fakeCups { - f := &fakeCups{ - users: map[string]string{}, - rooms: map[string]bool{}, - subs: map[string]map[*fakeConn]bool{}, - conns: map[*fakeConn]bool{}, - pageHits: map[string]int{}, - failPages: map[string]int{}, - } - mux := http.NewServeMux() - mux.HandleFunc("/live-coding/", f.page) - mux.HandleFunc("/sub/", f.subToken) - mux.HandleFunc("/connection/websocket", f.centrifugo) - f.srv = httptest.NewServer(mux) - - old := baseRoomURL - baseRoomURL = f.srv.URL + "/live-coding/" - t.Cleanup(func() { - baseRoomURL = old - f.mu.Lock() - for c := range f.conns { - c.ws.Close() - } - f.mu.Unlock() - f.srv.Close() - }) - return f -} - -func newUUID() string { - var b [16]byte - rand.Read(b[:]) - return fmt.Sprintf("%x-%x-%x-%x-%x", b[0:4], b[4:6], b[6:8], b[8:10], b[10:16]) -} - -func (f *fakeCups) page(w http.ResponseWriter, r *http.Request) { - var stall chan struct{} - f.locked(func() { stall = f.stall }) - if stall != nil { - select { - case <-stall: - case <-r.Context().Done(): - } - return - } - - f.mu.Lock() - defer f.mu.Unlock() - - user := "" - if c, err := r.Cookie("sessionid"); err == nil { - user = f.users[c.Value] - } - if user == "" { - session := newUUID() - user = newUUID() - f.users[session] = user - http.SetCookie(w, &http.Cookie{Name: "sessionid", Value: session, Path: "/"}) - } - http.SetCookie(w, &http.Cookie{Name: "csrftoken", Value: "csrf-" + user, Path: "/"}) - - room := r.URL.Query().Get("room") - if room != "" { - f.pageHits[room]++ - } - open, known := f.rooms[room] - switch { - case room != "" && (f.failJoins || f.failPages[room] > 0): - f.failPages[room]-- - http.Error(w, "bad gateway", http.StatusBadGateway) - return - case room == "" || (known && !open && f.redirectGone): - room = newUUID() - f.rooms[room] = true - f.created++ - case !open: - http.NotFound(w, r) - return - } - fmt.Fprintf(w, ` - - - -
`, - user, f.srv.URL, f.srv.URL, room, user) -} - -func (f *fakeCups) subToken(w http.ResponseWriter, r *http.Request) { - c, err := r.Cookie("csrftoken") - if r.Method != http.MethodPost || err != nil || r.Header.Get("X-CSRFToken") != c.Value { - http.Error(w, "forbidden", http.StatusForbidden) - return - } - var body struct { - Channel string `json:"channel"` - } - json.NewDecoder(r.Body).Decode(&body) - json.NewEncoder(w).Encode(map[string]string{"token": "sub-" + body.Channel}) -} - -func (f *fakeCups) centrifugo(w http.ResponseWriter, r *http.Request) { - ws, err := (&websocket.Upgrader{}).Upgrade(w, r, nil) - if err != nil { - return - } - c := &fakeConn{ws: ws} - f.mu.Lock() - f.conns[c] = true - f.mu.Unlock() - defer func() { - f.mu.Lock() - delete(f.conns, c) - for _, s := range f.subs { - delete(s, c) - } - f.mu.Unlock() - ws.Close() - }() - - for { - _, raw, err := ws.ReadMessage() - if err != nil { - return - } - if string(raw) == "{}" { - continue - } - var cmd struct { - ID int `json:"id"` - Connect any `json:"connect"` - Subscribe *struct { - Channel string `json:"channel"` - } `json:"subscribe"` - RPC *struct { - Method string `json:"method"` - Data struct { - Cursors []struct { - Row json.Number `json:"row"` - Column json.Number `json:"column"` - } `json:"cursors"` - Room string `json:"room"` - User string `json:"user"` - } `json:"data"` - } `json:"rpc"` - } - if json.Unmarshal(raw, &cmd) != nil { - return - } - switch { - case cmd.Connect != nil: - c.send(fmt.Sprintf(`{"id":%d,"connect":{"client":"x","ping":25,"pong":true}}`, cmd.ID)) - case cmd.Subscribe != nil: - room := strings.TrimPrefix(cmd.Subscribe.Channel, "$shared_editor:room-") - f.mu.Lock() - open := f.rooms[room] - if open { - if f.subs[room] == nil { - f.subs[room] = map[*fakeConn]bool{} - } - f.subs[room][c] = true - } - f.mu.Unlock() - if !open { - c.send(fmt.Sprintf(`{"id":%d,"error":{"code":103,"message":"permission denied"}}`, cmd.ID)) - continue - } - c.send(fmt.Sprintf(`{"id":%d,"subscribe":{}}`, cmd.ID)) - case cmd.RPC != nil: - c.send(fmt.Sprintf(`{"id":%d,"rpc":{}}`, cmd.ID)) - if d := cmd.RPC.Data; cmd.RPC.Method == "shared_editor_change_cursors" && len(d.Cursors) > 0 { - // Model the real server: a non-numeric row/column is rejected, - // anything else is relayed verbatim and in order. - cursors := make([]any, len(d.Cursors)) - for i, cur := range d.Cursors { - row, errR := cur.Row.Int64() - col, errC := cur.Column.Int64() - if errR != nil || errC != nil { - c.send(fmt.Sprintf(`{"id":%d,"error":{"code":400,"message":"column must be a number"}}`, cmd.ID)) - cursors = nil - break - } - cursors[i] = map[string]any{"row": row, "column": col} - } - if cursors != nil { - f.relay(d.Room, d.User, cursors) - } - } - } - } -} - -// relay hands a cursor update to everyone in the room, the sender included, -// the way Centrifugo does. -func (f *fakeCups) relay(room, user string, cursors []any) { - push, _ := json.Marshal(map[string]any{"push": map[string]any{ - "channel": "$shared_editor:room-" + room, - "pub": map[string]any{"data": map[string]any{ - "type": "cursors_update", - "payload": map[string]any{ - "user_uuid": user, - "cursors": cursors, - }, - }}, - }}) - msg := string(push) - - f.mu.Lock() - f.maxCursors = max(f.maxCursors, len(cursors)) - if f.pack { - msg = "{}\n" + msg - } - var to []*fakeConn - for c := range f.subs[room] { - to = append(to, c) - } - f.mu.Unlock() - for _, c := range to { - c.send(msg) - } -} - -// closeRoom closes a room for good and hangs up on everyone in it. -func (f *fakeCups) closeRoom(room string) { - f.mu.Lock() - defer f.mu.Unlock() - f.rooms[room] = false - for c := range f.subs[room] { - c.ws.Close() - } - delete(f.subs, room) -} - -// dropConns hangs up on everyone in a room that stays open, like a network -// blip. -func (f *fakeCups) dropConns(room string) { - f.mu.Lock() - defer f.mu.Unlock() - for c := range f.subs[room] { - c.ws.Close() - } -} - -func (f *fakeCups) locked(fn func()) { - f.mu.Lock() - defer f.mu.Unlock() - fn() -} - -func (f *fakeCups) hits(room string) (n int) { - f.locked(func() { n = f.pageHits[room] }) - return -} - -func (f *fakeCups) roomsCreated() (n int) { - f.locked(func() { n = f.created }) - return -} - -func (f *fakeCups) openConns() (n int) { - f.locked(func() { n = len(f.conns) }) - return -} - -// ---- helpers ---- - -func fastConfig() CupsonlineConfig { - c := DefaultCupsonlineConfig() - c.NumRooms = 3 - c.RoomCreatePause = time.Millisecond - c.ReconnectMinDelay = 20 * time.Millisecond - c.ReconnectMaxDelay = 200 * time.Millisecond - c.RoomGoneRetryMin = 300 * time.Millisecond - c.RoomGoneRetryMax = 300 * time.Millisecond - c.WSHandshakeTimeout = 2 * time.Second - c.StatsInterval = time.Hour - c.SendInterval = 0 // no throttling against the in-process fake - return c -} - -type peer struct { - *CupsonlineTransport - got chan []byte -} - -func startPeer(t *testing.T, url string, isClient bool, cfg CupsonlineConfig) (*peer, error) { - tr := NewCupsonlineTransport(url, transport.DefaultConfig(), isClient) - tr.config = cfg - p := &peer{CupsonlineTransport: tr, got: make(chan []byte, 4096)} - tr.Receive(func(b []byte) { p.got <- append([]byte(nil), b...) }) - t.Cleanup(func() { tr.Stop() }) - if err := tr.Start(); err != nil { - return nil, err - } - return p, nil -} - -func mustStart(t *testing.T, url string, isClient bool, cfg CupsonlineConfig) *peer { - t.Helper() - p, err := startPeer(t, url, isClient, cfg) - if err != nil { - t.Fatalf("start (client=%v): %v", isClient, err) - } - return p -} - -// startPair brings up an exit node that creates the rooms and a client that -// joins them from the printed string, and waits until every channel is up. -func startPair(t *testing.T, cfg CupsonlineConfig) (exit, client *peer) { - t.Helper() - exit = mustStart(t, "", false, cfg) - client = mustStart(t, packRooms(exit.RoomUUIDs()), true, cfg) - waitFor(t, 5*time.Second, "all channels up", func() bool { - return allConnected(exit) && allConnected(client) - }) - return exit, client -} - -func allConnected(p *peer) bool { - for _, ws := range p.wss { - if !ws.connected.Load() { - return false - } - } - return true -} - -func roomDead(p *peer, room string) bool { - for _, ws := range p.wss { - if ws.roomUUID == room { - return ws.dead.Load() - } - } - return false -} - -func waitFor(t *testing.T, within time.Duration, what string, cond func() bool) { - t.Helper() - deadline := time.Now().Add(within) - for !cond() { - if time.Now().After(deadline) { - t.Fatalf("timed out waiting for %s", what) - } - time.Sleep(10 * time.Millisecond) - } -} - -// deliver keeps sending until a copy gets through: right after a room dies -// the first few can be lost, like packets on any link. -func deliver(t *testing.T, from, to *peer, payload []byte, within time.Duration) { - t.Helper() - deadline := time.Now().Add(within) - for time.Now().Before(deadline) { - from.Send(payload) - wait := time.After(100 * time.Millisecond) - drain: - for { - select { - case got := <-to.got: - if bytes.Equal(got, payload) { - return - } - case <-wait: - break drain - } - } - } - t.Fatalf("%q never got through", payload) -} - -func next(t *testing.T, p *peer) []byte { - t.Helper() - select { - case got := <-p.got: - return got - case <-time.After(5 * time.Second): - t.Fatal("nothing arrived") - return nil - } -} - -// drain discards anything already delivered to a peer, so one phase's -// leftovers don't bleed into the next. -func drain(p *peer) { - for { - select { - case <-p.got: - default: - return - } - } -} - -// expectSet reads len(want) payloads and checks they are exactly want as a -// multiset. Order isn't asserted: frames spread across rooms can arrive in a -// different order, and only that every one arrives intact and once matters. -func expectSet(t *testing.T, p *peer, want [][]byte) { - t.Helper() - remaining := make(map[string]int, len(want)) - for _, w := range want { - remaining[string(w)]++ - } - for range want { - got := next(t, p) - k := string(got) - if remaining[k] == 0 { - t.Fatalf("unexpected or duplicate payload of %d bytes", len(got)) - } - remaining[k]-- - } -} - -// roomsUsed counts how many of a peer's channels have sent at least one packet. -func roomsUsed(p *peer) (n int) { - for _, ws := range p.wss { - if ws.stats.packetsSent.Load() > 0 { - n++ - } - } - return -} - -// frame looks like what really reaches the transport: a codec frame (here -// the batched one, 0x02...), not an IP packet. -func frame(i int) []byte { - return append([]byte{0x02, 0x01, byte(i >> 8), byte(i)}, bytes.Repeat([]byte{0x5a}, 100)...) -} - -// ---- tests ---- - -func TestFramesCrossBothWays(t *testing.T) { - newFakeCups(t) - exit, client := startPair(t, fastConfig()) - - var want [][]byte - for i := 0; i < 100; i++ { - f := frame(i) - want = append(want, f) - if err := client.Send(f); err != nil { - t.Fatal(err) - } - } - expectSet(t, exit, want) - // Spreading is the point: 100 frames must not all funnel through one room. - if used := roomsUsed(client); used < 2 { - t.Fatalf("frames used only %d room(s), expected them spread", used) - } - - exit.Send(frame(1000)) - if got := next(t, client); !bytes.Equal(got, frame(1000)) { - t.Fatalf("reply: got %x", got[:4]) - } -} - -// A room closes: the traffic it would have carried has to keep flowing through -// the rooms still up, both ways, instead of the tunnel going dark. -func TestTrafficSurvivesLosingARoom(t *testing.T) { - f := newFakeCups(t) - exit, client := startPair(t, fastConfig()) - deliver(t, client, exit, []byte("before-up"), 2*time.Second) - deliver(t, exit, client, []byte("before-down"), 2*time.Second) - - gone := client.wss[0].roomUUID - f.closeRoom(gone) - waitFor(t, 5*time.Second, "the closed room reported on both sides", func() bool { - return roomDead(client, gone) && roomDead(exit, gone) - }) - - // The surviving rooms still carry traffic both ways. - deliver(t, client, exit, []byte("after-up"), 5*time.Second) - deliver(t, exit, client, []byte("after-down"), 5*time.Second) - if !client.IsConnected() || !exit.IsConnected() { - t.Fatal("other rooms are up, the transport must still say connected") - } -} - -func TestClosedRoomIsProbedSlowly(t *testing.T) { - f := newFakeCups(t) - f.redirectGone = true // every probe makes cups hand out a new room - cfg := fastConfig() - cfg.RoomGoneRetryMin, cfg.RoomGoneRetryMax = 500*time.Millisecond, 500*time.Millisecond - exit, client := startPair(t, cfg) - - room := client.wss[0].roomUUID - f.closeRoom(room) - waitFor(t, 5*time.Second, "room reported closed on both sides", func() bool { - return roomDead(client, room) && roomDead(exit, room) - }) - - hits, created := f.hits(room), f.roomsCreated() - time.Sleep(2 * time.Second) - // Both peers probe it: at most 2 s / 500 ms + 1 each. - if n := f.hits(room) - hits; n > 10 { - t.Fatalf("closed room probed %d times in 2 s", n) - } - if n := f.roomsCreated() - created; n > 10 { - t.Fatalf("%d rooms made by probing a closed one in 2 s", n) - } -} - -// A connection that drops comes back with the tokens it has. Loading the -// room page on every reconnect turns a flaky network into a flood of page -// loads, and that's how an IP gets rate-limited. -func TestDroppedConnectionReconnectsWithoutReloadingThePage(t *testing.T) { - f := newFakeCups(t) - exit, client := startPair(t, fastConfig()) - room := client.wss[0].roomUUID - hits := f.hits(room) - - for i := 0; i < 5; i++ { - f.dropConns(room) - waitFor(t, 5*time.Second, "reconnect", func() bool { - return client.wss[0].stats.reconnects.Load() > uint64(i) && allConnected(client) && allConnected(exit) - }) - } - if n := f.hits(room) - hits; n != 0 { - t.Fatalf("room page loaded %d times over 5 dropped connections", n) - } - deliver(t, client, exit, []byte("still works"), 2*time.Second) -} - -// A room the client couldn't enter at start still gets a channel and joins -// once it can. Dropping it left the exit node's traffic in it unheard. -func TestClientRetriesRoomItCouldNotEnterAtStart(t *testing.T) { - f := newFakeCups(t) - cfg := fastConfig() - exit := mustStart(t, "", false, cfg) - ids := exit.RoomUUIDs() - f.locked(func() { f.failPages[ids[1]] = 3 }) - - client := mustStart(t, packRooms(ids), true, cfg) - if len(client.wss) != len(ids) { - t.Fatalf("client keeps %d of %d rooms", len(client.wss), len(ids)) - } - waitFor(t, 5*time.Second, "late room joined", func() bool { return client.wss[1].connected.Load() }) - waitFor(t, 5*time.Second, "exit node up", func() bool { return allConnected(exit) }) - - // The late room must carry its share once it joins: send enough that the - // round-robin certainly uses it, and check it did and everything arrived. - before := exit.wss[1].stats.packetsSent.Load() - var want [][]byte - for i := 0; i < 12; i++ { - p := []byte{byte(i), 0xaa} - want = append(want, p) - if err := exit.Send(p); err != nil { - t.Fatal(err) - } - } - expectSet(t, client, want) - if exit.wss[1].stats.packetsSent.Load() == before { - t.Fatal("the late room carried no traffic") - } -} - -func TestExitKeepsSavedRoomsThroughNetworkTrouble(t *testing.T) { - f := newFakeCups(t) - cfg := fastConfig() - first := mustStart(t, "", false, cfg) - ids := first.RoomUUIDs() - packed := packRooms(ids) - first.Stop() - created := f.roomsCreated() - - f.locked(func() { f.failJoins = true }) - if _, err := startPeer(t, packed, false, cfg); err == nil { - t.Fatal("exit node must refuse to start while its rooms don't answer") - } - if f.roomsCreated() != created { - t.Fatal("new rooms made over a network error: the phone's string is lost") - } - - f.locked(func() { f.failJoins = false }) - again := mustStart(t, packed, false, cfg) - if strings.Join(again.RoomUUIDs(), ",") != strings.Join(ids, ",") || f.roomsCreated() != created { - t.Fatal("restart with the saved string must reuse the same rooms") - } - again.Stop() - - for _, id := range ids { - f.closeRoom(id) - } - fresh := mustStart(t, packed, false, cfg) - if f.roomsCreated() != created+cfg.NumRooms { - t.Fatalf("all saved rooms gone: want %d new rooms, got %d", cfg.NumRooms, f.roomsCreated()-created) - } - for _, id := range fresh.RoomUUIDs() { - if strings.Contains(packed, id) { - t.Fatal("reused a closed room") - } - } -} - -func TestClientWithoutReachableRoomsFails(t *testing.T) { - f := newFakeCups(t) - exit := mustStart(t, "", false, fastConfig()) - packed := packRooms(exit.RoomUUIDs()) - f.locked(func() { f.failJoins = true }) - if _, err := startPeer(t, packed, true, fastConfig()); err == nil { - t.Fatal("client with no room to enter must not report success") - } -} - -func TestBatchesStayUnderCursorLimit(t *testing.T) { - f := newFakeCups(t) - cfg := fastConfig() - exit, client := startPair(t, cfg) - - body := bytes.Repeat([]byte{0xab}, 7000) - var want [][]byte - for i := 0; i < 200; i++ { - p := append([]byte{byte(i)}, body...) - want = append(want, p) - if err := client.Send(p); err != nil { - t.Fatal(err) - } - } - expectSet(t, exit, want) - f.locked(func() { - if f.maxCursors > cfg.MaxCursors { - t.Fatalf("%d cursors in one message, limit %d", f.maxCursors, cfg.MaxCursors) - } - }) - if err := client.Send(make([]byte, cfg.MaxPayloadBytes+1)); err == nil { - t.Fatal("a payload that can't fit in one message must be refused, not sent") - } -} - -// A packet larger than one message must be split on send and stitched back -// on receive, in order and intact. This is the throughput path that broke -// against the real server when a whole batch went in one oversized message. -func TestLargePacketsSplitAndReassemble(t *testing.T) { - f := newFakeCups(t) - cfg := fastConfig() - exit, client := startPair(t, cfg) - - // Each packet spans several messages; distinct contents so a swap shows. - const n = 40 - sizes := []int{cfg.MaxMessageData - 5, cfg.MaxMessageData, cfg.MaxMessageData + 5, 3*cfg.MaxMessageData + 1, 8192} - want := make([][]byte, n) - for i := 0; i < n; i++ { - size := sizes[i%len(sizes)] - p := make([]byte, size) - for j := range p { - p[j] = byte(i*31 + j) - } - want[i] = p - if err := client.Send(p); err != nil { - t.Fatalf("packet %d (%d bytes): %v", i, size, err) - } - } - // Each packet reassembles within whichever room carried it; across rooms - // they may arrive in a different order, so check the set, not the order. - expectSet(t, exit, want) - f.locked(func() { - if f.maxCursors > cfg.MaxCursors { - t.Fatalf("%d cursors in one message, over the %d limit", f.maxCursors, cfg.MaxCursors) - } - }) -} - -func TestNumberPackingRoundTrips(t *testing.T) { - for _, n := range []int{0, 1, 2, 5, 6, 7, 11, 12, 13, 100, 4096} { - in := make([]byte, n) - for i := range in { - in[i] = byte(i*7 + 1) // never all-zero, so nothing looks like padding - } - out := bytesFromNumbers(packNumbers(in)) - if len(out) < len(in) || !bytes.Equal(out[:len(in)], in) { - t.Fatalf("len %d did not round-trip", n) - } - // padding is only trailing zeros - for _, b := range out[len(in):] { - if b != 0 { - t.Fatalf("len %d: non-zero padding", n) - } - } - } - // values must stay inside what cups keeps exact (< 2^53). - for _, v := range packNumbers(bytes.Repeat([]byte{0xff}, 4096)) { - if v >= 1<<(8*bytesPerNumber) { - t.Fatalf("number %d exceeds %d bytes", v, bytesPerNumber) - } - } -} - -// Centrifugo may pack several messages into one frame; parsing it whole -// used to lose every packet in it. -func TestPackedFramesAreTakenApart(t *testing.T) { - f := newFakeCups(t) - f.pack = true - exit, client := startPair(t, fastConfig()) - deliver(t, client, exit, []byte("up"), 2*time.Second) - deliver(t, exit, client, []byte("down"), 2*time.Second) -} - -func TestStopHangsUpRightAway(t *testing.T) { - f := newFakeCups(t) - _, client := startPair(t, fastConfig()) - before := f.openConns() - - client.Stop() - waitFor(t, time.Second, "client sockets closed", func() bool { - return f.openConns() == before-len(client.wss) - }) - if client.IsConnected() { - t.Fatal("stopped transport says connected") - } - if err := client.Send([]byte("x")); err == nil { - t.Fatal("send after stop must fail") - } - client.Stop() // twice must be harmless -} - -func TestStopDuringStartAbortsIt(t *testing.T) { - f := newFakeCups(t) - exit := mustStart(t, "", false, fastConfig()) - packed := packRooms(exit.RoomUUIDs()) - - stall := make(chan struct{}) - defer close(stall) - f.locked(func() { f.stall = stall }) - - tr := NewCupsonlineTransport(packed, transport.DefaultConfig(), true) - tr.config = fastConfig() - done := make(chan error, 1) - go func() { done <- tr.Start() }() - time.Sleep(100 * time.Millisecond) - tr.Stop() - select { - case err := <-done: - if err == nil { - t.Fatal("Start after Stop must not report success") - } - case <-time.After(2 * time.Second): - t.Fatal("Stop didn't cut Start short") - } -} - -// gorilla allocates the socket buffers whole for every connection: at the -// old 32 MB each, four rooms took 256 MB on the phone, again on every -// reconnect. -func TestChannelsDontHoardMemory(t *testing.T) { - newFakeCups(t) - var before, after runtime.MemStats - runtime.GC() - runtime.ReadMemStats(&before) - exit, client := startPair(t, fastConfig()) - runtime.ReadMemStats(&after) - conns := len(exit.wss) + len(client.wss) - if grew := int64(after.HeapAlloc) - int64(before.HeapAlloc); grew > 64<<20 { - t.Fatalf("%d channels took %d MB of heap", conns, grew>>20) - } -} - -func TestPickRoomSpreadsAndSkipsDown(t *testing.T) { - tr := &CupsonlineTransport{} - for i := 0; i < 3; i++ { - ws := &cupsWS{idx: i, roomUUID: fmt.Sprintf("room-%d-0000", i)} - ws.connected.Store(true) - tr.wss = append(tr.wss, ws) - } - - // Round-robin visits every connected room. - seen := map[*cupsWS]bool{} - for i := 0; i < 30; i++ { - seen[tr.pickRoom()] = true - } - if len(seen) != 3 { - t.Fatalf("round-robin used %d of 3 rooms", len(seen)) - } - - // A disconnected room is skipped. - tr.wss[1].connected.Store(false) - for i := 0; i < 30; i++ { - if tr.pickRoom() == tr.wss[1] { - t.Fatal("pickRoom returned a disconnected room") - } - } - - // None up -> nil. - for _, ws := range tr.wss { - ws.connected.Store(false) - } - if tr.pickRoom() != nil { - t.Fatal("no room up: pickRoom must return nil") - } -} diff --git a/transport/cupsonline/live_test.go b/transport/cupsonline/live_test.go deleted file mode 100644 index 780d64b..0000000 --- a/transport/cupsonline/live_test.go +++ /dev/null @@ -1,234 +0,0 @@ -package cupsonline - -import ( - "bytes" - "context" - "fmt" - "os" - "strings" - "testing" - "time" -) - -// TestLiveCups runs an exit node and a client through the real cups.online. -// It creates rooms there, so it only runs when asked to: -// -// OPENFLUX_CUPS_LIVE=1 go test -run TestLiveCups -v -timeout 15m ./transport/cupsonline/ -// -// OPENFLUX_CUPS_LIVE_IDLE sets how long the channels sit idle (default 2m, -// longer than WSReadTimeout on purpose). -func TestLiveCups(t *testing.T) { - if os.Getenv("OPENFLUX_CUPS_LIVE") == "" { - t.Skip("set OPENFLUX_CUPS_LIVE=1 to run against the real cups.online") - } - idle := 2 * time.Minute - if v := os.Getenv("OPENFLUX_CUPS_LIVE_IDLE"); v != "" { - d, err := time.ParseDuration(v) - if err != nil { - t.Fatal(err) - } - idle = d - } - cfg := DefaultCupsonlineConfig() - cfg.StatsInterval = time.Hour - - start := time.Now() - exit := mustStart(t, "", false, cfg) - packed := packRooms(exit.RoomUUIDs()) - client := mustStart(t, packed, true, cfg) - waitFor(t, 60*time.Second, "all channels up", func() bool { - return allConnected(exit) && allConnected(client) - }) - t.Logf("%d rooms created and joined on both sides in %v", len(exit.wss), time.Since(start).Round(time.Millisecond)) - - t.Run("frames both ways", func(t *testing.T) { - drain(exit) - // Frames spread across rooms and may arrive in a different order, so - // check they all arrive intact, not the order. - var want [][]byte - for i := 0; i < 50; i++ { - f := frame(i) - want = append(want, f) - client.Send(f) - } - expectSet(t, exit, want) - - drain(exit) - drain(client) - var rtts []time.Duration - for i := 0; i < 10; i++ { - sent := time.Now() - client.Send(frame(500 + i)) - next(t, exit) - exit.Send(frame(600 + i)) - next(t, client) - rtts = append(rtts, time.Since(sent)) - } - t.Logf("round trip client->exit->client: %v", rtts) - }) - - t.Run("throughput", func(t *testing.T) { - // Feed packets steadily rather than all at once; the transport paces - // itself under that so the server doesn't drop the connection. Frames - // spread across all rooms and can arrive in a different order, so this - // checks that every one arrives intact, not the order. - drain(exit) - const n, size = 200, 8 << 10 - body := bytes.Repeat([]byte{0x5a}, size) - seen := make([]bool, n) - done := make(chan int, 1) - go func() { - count := 0 - for count < n { - select { - case p := <-exit.got: - id := int(p[0])<<8 | int(p[1]) - if id < 0 || id >= n || seen[id] || len(p) != size+2 { - done <- -1 - return - } - seen[id] = true - count++ - case <-time.After(30 * time.Second): - done <- count - return - } - } - done <- n - }() - sent := time.Now() - for i := 0; i < n; i++ { - p := append([]byte{byte(i >> 8), byte(i)}, body...) - if err := client.Send(p); err != nil { - t.Fatal(err) - } - time.Sleep(15 * time.Millisecond) // offered load; the transport paces the rest - } - if got := <-done; got != n { - t.Fatalf("only %d of %d packets arrived intact", got, n) - } - took := time.Since(sent) - t.Logf("client->exit: %d KB across %d rooms in %v = %.0f KB/s", - n*size>>10, len(client.wss), took.Round(time.Millisecond), float64(n*size>>10)/took.Seconds()) - }) - - t.Run("idle longer than the read timeout", func(t *testing.T) { - reconnects := func() (n uint64) { - for _, p := range []*peer{exit, client} { - for _, ws := range p.wss { - n += ws.stats.reconnects.Load() - } - } - return - } - before := reconnects() - t.Logf("sitting idle for %v (read timeout %v)", idle, cfg.WSReadTimeout) - time.Sleep(idle) - if n := reconnects() - before; n != 0 { - t.Errorf("%d reconnects while idle: keepalive doesn't keep the read side busy", n) - } - deliver(t, client, exit, []byte("after idle"), 10*time.Second) - deliver(t, exit, client, []byte("after idle back"), 10*time.Second) - }) - - t.Run("dropped socket comes back without a page reload", func(t *testing.T) { - ws := client.wss[0] - was := ws.auth() - ws.writeMu.Lock() - if ws.conn != nil { - ws.conn.Close() - } - ws.writeMu.Unlock() - - deliver(t, client, exit, []byte("during reconnect"), 15*time.Second) - waitFor(t, 30*time.Second, "dropped channel back", func() bool { return ws.connected.Load() }) - if ws.auth() != was { - t.Error("reconnect after a plain drop re-joined the room (page reload)") - } - }) - - t.Run("re-join keeps the session", func(t *testing.T) { - a := client.wss[0].auth() - again, err := joinRoom(context.Background(), a.roomUUID, a.httpClient) - if err != nil { - t.Fatal(err) - } - if again.userUUID != a.userUUID { - t.Errorf("same session, new user: %s -> %s", a.userUUID, again.userUUID) - } - }) - - t.Run("exit restart reuses its rooms", func(t *testing.T) { - exit.Stop() - again := mustStart(t, packed, false, cfg) - if strings.Join(again.RoomUUIDs(), ",") != strings.Join(exit.RoomUUIDs(), ",") { - t.Fatalf("restarted exit node is in other rooms: %v vs %v", again.RoomUUIDs(), exit.RoomUUIDs()) - } - waitFor(t, 60*time.Second, "restarted exit node up", func() bool { return allConnected(again) }) - deliver(t, client, again, []byte("to restarted exit"), 15*time.Second) - deliver(t, again, client, []byte("from restarted exit"), 15*time.Second) - }) -} - -// TestLiveThroughputCalibrate finds how fast one channel may send before -// cups.online starts dropping the connection. It drives a real pair at a few -// SendInterval values and reports, per value, how much arrived and how many -// times the channel had to reconnect. Pick the fastest interval with no loss -// and no reconnects for DefaultCupsonlineConfig.SendInterval. -// -// OPENFLUX_CUPS_LIVE=1 go test -run TestLiveThroughputCalibrate -v -timeout 15m ./transport/cupsonline/ -func TestLiveThroughputCalibrate(t *testing.T) { - if os.Getenv("OPENFLUX_CUPS_LIVE") == "" { - t.Skip("set OPENFLUX_CUPS_LIVE=1 to run against the real cups.online") - } - cfg := DefaultCupsonlineConfig() - cfg.StatsInterval = time.Hour - - exit := mustStart(t, "", false, cfg) - // Measure one channel: the client joins only the first room, so the whole - // offered load rides that single channel at the interval under test. - oneRoom := packRooms(exit.RoomUUIDs()[:1]) - waitFor(t, 60*time.Second, "exit up", func() bool { return allConnected(exit) }) - - const m, size = 120, 8 << 10 - body := bytes.Repeat([]byte{0x5a}, size) - - t.Logf("%-10s %-12s %-12s %-10s", "interval", "arrived", "reconnects", "KB/s") - for _, iv := range []time.Duration{40, 25, 18, 12, 8, 5} { - interval := iv * time.Millisecond - pcfg := cfg - pcfg.SendInterval = interval - cl, err := startPeer(t, oneRoom, true, pcfg) - if err != nil { - t.Fatalf("interval %v: client start: %v", interval, err) - } - waitFor(t, 60*time.Second, "client up", func() bool { return allConnected(cl) }) - drain(exit) // clear the previous phase's leftovers - reconBefore := cl.wss[0].stats.reconnects.Load() - - start := time.Now() - for i := 0; i < m; i++ { - if err := cl.Send(append([]byte{byte(i >> 8), byte(i)}, body...)); err != nil { - break - } - } - arrived := 0 - deadline := time.After(30 * time.Second) - collect: - for arrived < m { - select { - case <-exit.got: - arrived++ - case <-deadline: - break collect - } - } - elapsed := time.Since(start) - recon := cl.wss[0].stats.reconnects.Load() - reconBefore - rate := float64(arrived*size>>10) / elapsed.Seconds() - t.Logf("%-10v %-12s %-12d %-10.0f", interval, fmt.Sprintf("%d/%d", arrived, m), recon, rate) - cl.Stop() - time.Sleep(2 * time.Second) // let the server settle between phases - } - exit.Stop() -} diff --git a/transport/cupsonline/stop_test.go b/transport/cupsonline/stop_test.go deleted file mode 100644 index 11dea1f..0000000 --- a/transport/cupsonline/stop_test.go +++ /dev/null @@ -1,19 +0,0 @@ -package cupsonline - -import ( - "testing" - - "openflux/transport" -) - -// A Session stops a carrier through every wrapper around it; the second -// Stop must not close the channels again. -func TestStopTwice(t *testing.T) { - c := NewCupsonlineTransport("", transport.TransportConfig{}, false) - if err := c.Stop(); err != nil { - t.Fatal(err) - } - if err := c.Stop(); err != nil { - t.Fatal(err) - } -} diff --git a/transport/mailru/mailru.go b/transport/mailru/mailru.go deleted file mode 100644 index 691bf25..0000000 --- a/transport/mailru/mailru.go +++ /dev/null @@ -1,591 +0,0 @@ -// Package mailru implements a transport that tunnels packets through -// Mail.ru's cloud document editor (docs.datacloudmail.ru), the same -// coauthoring backend family as Yandex.Docs. Two peers open the same -// public document and smuggle packets through the "cursor" field of the -// collaborative editing protocol. -package mailru - -import ( - "bytes" - "encoding/base64" - "encoding/json" - "fmt" - "io" - "math/rand" - "net" - "net/http" - "net/http/cookiejar" - "net/url" - "regexp" - "strings" - "sync" - "sync/atomic" - "time" - - "github.com/gorilla/websocket" - - "openflux/netbind" - "openflux/transport" - "openflux/utils" -) - -const mailruUserAgent = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/137.0.0.0 Safari/537.36" - -var cursorPayloadRe = regexp.MustCompile(`"cursor":"[^;]+;([^"]+)"`) - -type MailruDocsInfo struct { - Token string - DocKey string - WsURL string - FileType string - DocURL string - DocTitle string - Permissions map[string]interface{} - CallbackURL string - EditorUserID string -} - -type DocSession struct { - Info MailruDocsInfo - Conn *websocket.Conn - WriteQueue chan []byte - UserID string - writeMu sync.Mutex -} - -func (s *DocSession) safeWrite(messageType int, data []byte) error { - s.writeMu.Lock() - defer s.writeMu.Unlock() - return s.Conn.WriteMessage(messageType, data) -} - -type MailruDocsTransport struct { - *transport.BaseTransport - - weblink string - session *DocSession - - userCounter atomic.Int32 - baseUserID string - - cookieJar *cookiejar.Jar - jarMu sync.RWMutex -} - -// NewMailruDocsTransport accepts either a bare weblink ("AbCdEfGh1/IjKlMnOp2") -// or a full public URL ("https://cloud.mail.ru/public/AbCdEfGh1/IjKlMnOp2"), -// normalizing the latter to the former. -func NewMailruDocsTransport(weblink string, config transport.TransportConfig) *MailruDocsTransport { - t := &MailruDocsTransport{ - BaseTransport: transport.NewBaseTransport(config), - weblink: normalizeWeblink(weblink), - } - t.baseUserID = randUserID() - jar, _ := cookiejar.New(nil) - t.cookieJar = jar - return t -} - -func normalizeWeblink(weblink string) string { - weblink = strings.TrimSpace(weblink) - for _, prefix := range []string{ - "https://cloud.mail.ru/public/", - "http://cloud.mail.ru/public/", - "https://cloud.mail.ru/", - "http://cloud.mail.ru/", - } { - if strings.HasPrefix(weblink, prefix) { - return strings.Trim(strings.TrimPrefix(weblink, prefix), "/") - } - } - return weblink -} - -func (t *MailruDocsTransport) Start() error { - if err := t.BaseTransport.Start(); err != nil { - return err - } - - t.baseUserID = randUserID() - utils.SafeGo("mailru.keepAlive", t.keepAliveLoop) - t.connectToDoc(0) - - return nil -} - -func (t *MailruDocsTransport) Send(data []byte) error { - if !t.IsConnected() { - return fmt.Errorf("transport not connected") - } - - t.Mu.RLock() - session := t.session - t.Mu.RUnlock() - - if session == nil { - return fmt.Errorf("no active session") - } - - select { - case session.WriteQueue <- data: - t.RecordSend(len(data)) - return nil - default: - return fmt.Errorf("write queue full") - } -} - -func (t *MailruDocsTransport) connectToDoc(attempt int) { - if !t.IsRunning() { - return - } - - utils.Debugf("[M-DOCS] connectToDoc attempt %d", attempt) - - go func() { - defer func() { - if r := recover(); r != nil { - utils.Debugf("[PANIC] recovered in mailru.connect: %v", r) - } - }() - t.Mu.Lock() - existingSession := t.session - t.Mu.Unlock() - - var userID string - if existingSession != nil { - userID = existingSession.UserID - } else { - suffix := fmt.Sprintf("%03d", t.userCounter.Add(1)%1000) - userID = t.baseUserID + suffix - } - - info, err := t.fetchDocInfo(t.weblink) - if err != nil { - utils.Debugf("[M-DOCS] fetchDocInfo failed: %v", err) - t.scheduleReconnect(attempt) - return - } - - dialer := websocket.Dialer{ - HandshakeTimeout: 15 * time.Second, - NetDialContext: netbind.Wrap(&net.Dialer{ - Timeout: 10 * time.Second, - KeepAlive: 30 * time.Second, - }).DialContext, - } - headers := http.Header{} - headers.Set("User-Agent", mailruUserAgent) - headers.Set("Origin", "https://docs.datacloudmail.ru") - - utils.Debugf("[M-DOCS] WebSocket dial %s", info.WsURL) - conn, resp, err := dialer.Dial(info.WsURL, headers) - if err != nil { - status := 0 - if resp != nil { - status = resp.StatusCode - } - utils.Debugf("[M-DOCS] WebSocket dial failed (http %d): %v", status, err) - t.scheduleReconnect(attempt) - return - } - utils.Debugf("[M-DOCS] WebSocket connected") - - writeQueue := make(chan []byte, t.GetConfig().MaxQueueSize) - if existingSession != nil { - writeQueue = existingSession.WriteQueue - } - - session := &DocSession{ - Info: info, - Conn: conn, - WriteQueue: writeQueue, - UserID: userID, - } - - t.Mu.Lock() - t.session = session - t.SetConnected(true) - t.Mu.Unlock() - - if existingSession == nil { - utils.SafeGo("mailru.writer", t.writerLoop) - } - - // Auth - fired immediately, same as the Yandex.Docs transport. No - // need to wait for the server's own "0{"/"40" handshake frames - // first: Mail.ru's coauthoring server buffers and processes these - // once its own session state catches up, and waiting for explicit - // acks here only stretches the outage window on every reconnect - // (Mail.ru can delay a fresh joiner's auth confirmation by up to - // ~30s while it reconciles with the other participant). - auth1 := fmt.Sprintf(`40{"token":"%s"}`, info.Token) - session.safeWrite(websocket.TextMessage, []byte(auth1)) - - authMsg := map[string]interface{}{ - "type": "auth", - "docid": info.DocKey, - "documentCallbackUrl": info.CallbackURL, - "token": "fghhfgsjdgfjs", - "user": map[string]interface{}{ - "id": info.EditorUserID, - "username": userID, - "indexUser": -1, - }, - "editorType": 0, - "lastOtherSaveTime": -1, - "block": []interface{}{}, - "documentFormatSave": 65, - "view": false, - "isCloseCoAuthoring": false, - "openCmd": map[string]interface{}{ - "c": "open", - "id": info.DocKey, - "userid": info.EditorUserID, - "format": info.FileType, - "url": info.DocURL, - "title": info.DocTitle, - "lcid": 25, - "nobase64": true, - "convertToOrigin": ".pdf.xps.oxps.djvu", - }, - "lang": "ru", - "mode": "edit", - "permissions": info.Permissions, - "IsAnonymousUser": false, - "timezoneOffset": -180, - "coEditingMode": "fast", - "jwtOpen": info.Token, - "time": 1000, - "supportAuthChangesAck": true, - } - messagePart, _ := json.Marshal([]interface{}{"message", authMsg}) - session.safeWrite(websocket.TextMessage, []byte(fmt.Sprintf("42%s", string(messagePart)))) - - connectedAt := time.Now() - for t.IsRunning() { - _, message, err := conn.ReadMessage() - if err != nil { - utils.Debugf("[M-DOCS] Read error: %v", err) - t.SetConnected(false) - conn.Close() - - next := attempt - if time.Since(connectedAt) > 15*time.Second { - next = -1 - } - t.scheduleReconnect(next) - return - } - t.handleMessage(session, message) - } - }() -} - -func (t *MailruDocsTransport) writerLoop() { - // The write queue is created once and preserved across reconnects, so we - // capture it and block on it instead of polling with a sleep. - var queue chan []byte - for t.IsRunning() && queue == nil { - t.Mu.Lock() - if t.session != nil { - queue = t.session.WriteQueue - } - t.Mu.Unlock() - if queue == nil { - time.Sleep(5 * time.Millisecond) - } - } - if queue == nil { - return - } - - var pending []byte - for t.IsRunning() { - if pending == nil { - packet, ok := <-queue - if !ok { - return - } - pending = packet - } - - t.Mu.RLock() - session := t.session - t.Mu.RUnlock() - if session == nil || session.Conn == nil { - // Mid-reconnect: hold the packet and retry rather than drop it. - time.Sleep(15 * time.Millisecond) - continue - } - - payload := base64.StdEncoding.EncodeToString(pending) - msg := fmt.Sprintf(`42["message",{"type":"cursor","cursor":"18;%s"}]`, payload) - if err := session.safeWrite(websocket.TextMessage, []byte(msg)); err != nil { - utils.Debugf("[M-DOCS] Write error: %v", err) - time.Sleep(15 * time.Millisecond) - continue // keep pending; the reconnect will bring up a new conn - } - pending = nil - } -} - -func (t *MailruDocsTransport) keepAliveLoop() { - ticker := time.NewTicker(t.GetConfig().KeepAliveInterval) - defer ticker.Stop() - keepAliveMsg := `42["message",{"type":"cursor","cursor":"18;---KA---"}]` - - for t.IsRunning() { - <-ticker.C - t.Mu.Lock() - session := t.session - t.Mu.Unlock() - - if session != nil && session.Conn != nil { - if err := session.safeWrite(websocket.TextMessage, []byte(keepAliveMsg)); err != nil { - utils.Debugf("[M-DOCS] Keep-alive failed: %v", err) - t.SetConnected(false) - } - } - } -} - -func (t *MailruDocsTransport) handleMessage(session *DocSession, data []byte) { - text := string(data) - - if strings.Contains(text, "---KA---") { - return - } - - // Socket.IO ping - respond with pong - if text == "2" { - if session != nil && session.Conn != nil { - session.safeWrite(websocket.TextMessage, []byte("3")) - } - return - } - if text == "3" { - return - } - - if strings.Contains(text, `"type":"auth"`) && strings.Contains(text, `"result":1`) { - utils.Debugf("[M-DOCS] Auth OK for user %s", session.UserID) - return - } - - if strings.Contains(text, "cursor") { - base64Str := t.extractBase64String(text) - if base64Str == "" { - return - } - - decoded, err := base64.StdEncoding.DecodeString(base64Str) - if err != nil { - utils.Debugf("[M-DOCS] Base64 decode error: %v", err) - return - } - - t.RecordReceive(len(decoded)) - t.CallReceive(decoded) - } -} - -func (t *MailruDocsTransport) extractBase64String(response string) string { - matches := cursorPayloadRe.FindStringSubmatch(response) - if len(matches) > 1 { - return matches[1] - } - return "" -} - -func (t *MailruDocsTransport) scheduleReconnect(attempt int) { - next := attempt + 1 - if !t.IsRunning() || next >= t.GetConfig().MaxReconnectAttempts { - return - } - - d := reconnectBackoff(next) - utils.Debugf("[M-DOCS] reconnecting in %v (attempt %d)", d, next) - time.Sleep(d) - if !t.IsRunning() { - return - } - - t.RecordReconnect() - t.connectToDoc(next) -} - -// reconnectBackoff returns an exponential backoff with jitter, capped at 15s. -func reconnectBackoff(n int) time.Duration { - if n < 1 { - n = 1 - } - shift := n - 1 - if shift > 5 { - shift = 5 - } - d := 500 * time.Millisecond * time.Duration(1< 15*time.Second { - d = 15 * time.Second - } - // add up to +50% jitter - d += time.Duration(rand.Int63n(int64(d/2) + 1)) - return d -} - -// fetchDocInfo POSTs to Mail.ru's public-document editor API and parses the -// response into the fields needed to open the collaborative WebSocket. -func (t *MailruDocsTransport) fetchDocInfo(weblink string) (MailruDocsInfo, error) { - t.jarMu.RLock() - jar := t.cookieJar - t.jarMu.RUnlock() - if jar == nil { - var err error - jar, err = cookiejar.New(nil) - if err != nil { - return MailruDocsInfo{}, err - } - } - client := &http.Client{Jar: jar, Timeout: 15 * time.Second} - - reqBody := map[string]string{ - "x-email": "anonym", - "public": "/" + weblink, - "platform": "desktop_web", - } - jsonData, _ := json.Marshal(reqBody) - - apiURL := "https://cloud.mail.ru/api/v4/r7/edit" - utils.Debugf("[M-DOCS] fetchDocInfo POST %s", apiURL) - - req, _ := http.NewRequest("POST", apiURL, bytes.NewBuffer(jsonData)) - req.Header.Set("Content-Type", "application/json") - req.Header.Set("Accept", "application/json, text/plain, */*") - req.Header.Set("User-Agent", mailruUserAgent) - req.Header.Set("X-Api-Version", "4") - req.Header.Set("Referer", fmt.Sprintf("https://cloud.mail.ru/public/%s?weblink=%s", weblink, weblink)) - - resp, err := client.Do(req) - if err != nil { - return MailruDocsInfo{}, err - } - defer resp.Body.Close() - - if resp.StatusCode != http.StatusOK { - return MailruDocsInfo{}, fmt.Errorf("API returned status %d", resp.StatusCode) - } - - bodyBytes, _ := io.ReadAll(resp.Body) - - var res map[string]interface{} - if err := json.Unmarshal(bodyBytes, &res); err != nil { - return MailruDocsInfo{}, fmt.Errorf("failed to parse JSON: %w", err) - } - - apiBase, _ := res["api"].(string) - token, _ := res["token"].(string) - - document, ok := res["document"].(map[string]interface{}) - if !ok || document == nil { - return MailruDocsInfo{}, fmt.Errorf("document object missing") - } - - docKey, _ := document["key"].(string) - fileType, _ := document["fileType"].(string) - docURL, _ := document["url"].(string) - docTitle, _ := document["title"].(string) - // document.permissions is an object of booleans (comment/edit/download/…), - // not a number - sending it as anything else makes the editor server - // reject the auth message with "access deny". - permissions, _ := document["permissions"].(map[string]interface{}) - if permissions == nil { - permissions = make(map[string]interface{}) - } - - editorConfig, ok := res["editorConfig"].(map[string]interface{}) - if !ok || editorConfig == nil { - return MailruDocsInfo{}, fmt.Errorf("editorConfig object missing") - } - callbackURL, _ := editorConfig["callbackUrl"].(string) - - userObj, _ := editorConfig["user"].(map[string]interface{}) - var editorUserID string - if userObj != nil { - editorUserID, _ = userObj["id"].(string) - } - - wsBase := strings.Replace(apiBase, "https://", "wss://", 1) - wsURL := fmt.Sprintf("%s/doc/%s/c/?EIO=4&transport=websocket", wsBase, docKey) - - return MailruDocsInfo{ - Token: token, - DocKey: docKey, - WsURL: wsURL, - FileType: fileType, - DocURL: docURL, - DocTitle: docTitle, - Permissions: permissions, - CallbackURL: callbackURL, - EditorUserID: editorUserID, - }, nil -} - -func randUserID() string { - return fmt.Sprintf("%010d", rand.New(rand.NewSource(time.Now().UnixNano())).Intn(1000000000)) -} - -// ---- CookieExchanger ---- - -// FetchCookies returns a snapshot of the transport's current cookie jar as -// name -> value. Used by the exit node to answer a SubtypeCookiesRequest. -func (t *MailruDocsTransport) FetchCookies() (map[string]string, error) { - t.jarMu.RLock() - jar := t.cookieJar - t.jarMu.RUnlock() - if jar == nil { - return nil, fmt.Errorf("mailru: cookie jar is nil") - } - u, err := url.Parse("https://cloud.mail.ru/") - if err != nil { - return nil, err - } - out := make(map[string]string) - for _, c := range jar.Cookies(u) { - out[c.Name] = c.Value - } - return out, nil -} - -// ApplyCookies replaces the transport's cookie jar with the provided values -// and forces the current session to reconnect. -func (t *MailruDocsTransport) ApplyCookies(values map[string]string) error { - if len(values) == 0 { - return nil - } - u, _ := url.Parse("https://cloud.mail.ru/") - jar, _ := cookiejar.New(nil) - cookies := make([]*http.Cookie, 0, len(values)) - for k, v := range values { - cookies = append(cookies, &http.Cookie{Name: k, Value: v, Path: "/"}) - } - jar.SetCookies(u, cookies) - - t.jarMu.Lock() - t.cookieJar = jar - t.jarMu.Unlock() - - utils.Debugf("[M-DOCS] applied %d cookies, forcing reconnect", len(cookies)) - - t.Mu.Lock() - session := t.session - t.session = nil - t.SetConnected(false) - t.Mu.Unlock() - if session != nil && session.Conn != nil { - _ = session.Conn.Close() - } - if t.IsRunning() { - t.scheduleReconnect(0) - } - return nil -} diff --git a/transport/yandex/boards.go b/transport/yandex/boards.go deleted file mode 100644 index bd4af6f..0000000 --- a/transport/yandex/boards.go +++ /dev/null @@ -1,1134 +0,0 @@ -package yandex - -import ( - "bytes" - "crypto/rand" - "encoding/base64" - "encoding/hex" - "encoding/json" - "fmt" - "io" - mrand "math/rand" - "net" - "net/http" - "net/http/cookiejar" - "net/url" - "strings" - "sync" - "sync/atomic" - "time" - - "github.com/gorilla/websocket" - - "openflux/netbind" - "openflux/transport" - "openflux/utils" -) - -const ( - boardsBase = "boards.yandex.ru" - boardsUA = "Mozilla/5.0 (Linux; Android 15; Pixel 9) AppleWebKit/537.36 " + - "(KHTML, like Gecko) Chrome/153.0.0.0 Mobile Safari/537.36" - boardsSocketHostDefault = "socket33.boards.yandex.ru" - - // Engine.io heartbeat. Сервер шлёт pingInterval=25000, pingTimeout=30000. - // Мы пингуем сами каждые 20с, чтобы NAT не рвал idle-соединение. - boardsPingInterval = 20 * time.Second - - // Дедлайн чтения в основном цикле. С запасом над boardsPingInterval. - boardsReadDeadline = 90 * time.Second - - // Таймаут ожидания 431 в handshake. - boardsHandshakeWait = 15 * time.Second -) - -// ---- внутренние типы ---- - -type boardsInfo struct { - hash string - name string - userHash string // payload.u из JWT - jwt string - cookies []*http.Cookie - wsHost string - session string // пустой до participant-connected/431 - dashboard string - currentSlide string -} - -type boardsSession struct { - Info boardsInfo - Conn *websocket.Conn - Queue chan []byte - writeMu sync.Mutex - ack atomic.Int64 - - // participant — тот, что мы отправили в subscribe. Это userHash из JWT. - participant atomic.Pointer[string] - - // creatorHash — серверный хэш (из participant-connected/431). - creatorHash atomic.Pointer[string] -} - -func (s *boardsSession) safeWrite(msgType int, data []byte) error { - s.writeMu.Lock() - defer s.writeMu.Unlock() - return s.Conn.WriteMessage(msgType, data) -} - -func shortStr(s string, n int) string { - if len(s) <= n { - return s - } - return s[:n] -} - -func (s *boardsSession) writeEventObj(ns string, obj interface{}) error { - ack := s.ack.Add(1) - 1 - payload, err := json.Marshal([]interface{}{ns, obj}) - if err != nil { - return err - } - msg := fmt.Sprintf("42%d%s", ack, payload) - return s.safeWrite(websocket.TextMessage, []byte(msg)) -} - -func (s *boardsSession) writeRaw(raw string) error { - return s.safeWrite(websocket.TextMessage, []byte(raw)) -} - -// ---- транспорт ---- - -type BoardsTransport struct { - *transport.BaseTransport - - url string - - session atomic.Pointer[boardsSession] - - onDataMu sync.RWMutex - onData func([]byte) - - closeOnce sync.Once - done chan struct{} - - cookieJar *cookiejar.Jar - jarMu sync.RWMutex - - errNotifier func(err error, transportName, url, reason string) -} - -func NewBoardsTransport(rawURL string, config transport.TransportConfig) *BoardsTransport { - jar, _ := cookiejar.New(nil) - return &BoardsTransport{ - BaseTransport: transport.NewBaseTransport(config), - url: rawURL, - done: make(chan struct{}), - cookieJar: jar, - } -} - -// SetErrorNotifier installs a callback for out-of-band errors. Called once -// by the manager. Boards currently never returns sentinel errors from -// fetchDocInfo (its captcha path is the old showcaptchafast, which the -// internal PoW solver handles), but the hook is wired for parity with the -// other transports. -func (t *BoardsTransport) SetErrorNotifier(fn func(err error, transportName, url, reason string)) { - t.errNotifier = fn -} - -func (t *BoardsTransport) Start() error { - if err := t.BaseTransport.Start(); err != nil { - return err - } - hash := extractBoardsHash(t.url) - if hash == "" { - return fmt.Errorf("boards: no hash in URL %q", t.url) - } - - name := randomGuestName() - info, err := t.authorize(hash, name) - if err != nil { - return fmt.Errorf("boards auth: %w", err) - } - utils.Debugf("[BOARDS] auth OK: hash=%s name=%q userHash=%s dashboard=%q wsHost=%s", - info.hash, info.name, info.userHash, info.dashboard, info.wsHost) - - t.closeOnce = sync.Once{} - t.done = make(chan struct{}) - utils.SafeGo("boards.connect", func() { t.connectLoop(info) }) - - return nil -} - -func (t *BoardsTransport) Stop() error { - t.closeOnce.Do(func() { close(t.done) }) - if s := t.session.Load(); s != nil && s.Conn != nil { - s.Conn.Close() - } - t.SetConnected(false) - return t.BaseTransport.Stop() -} - -func (t *BoardsTransport) Send(data []byte) error { - if len(data) == 0 { - return nil - } - s := t.session.Load() - if s == nil { - return fmt.Errorf("boards: no session") - } - cp := make([]byte, len(data)) - copy(cp, data) - select { - case s.Queue <- cp: - return nil - case <-t.done: - return fmt.Errorf("boards: closed") - default: - return fmt.Errorf("boards: queue full") - } -} - -func (t *BoardsTransport) Receive(cb func([]byte)) { - t.onDataMu.Lock() - t.onData = cb - t.onDataMu.Unlock() -} - -func (t *BoardsTransport) IsConnected() bool { - return t.BaseTransport.IsConnected() -} - -// ---- авторизация ---- -// -// 1) GET /whiteboard/?hash= → начальные cookies -// (может редиректнуть на showcaptchafast — тогда проходим капчу) -// 2) POST /api request-guest-token → Set-Cookie token_= -// 3) POST /api get-whiteboard-info → dashboard/current_slide/ws_host - -// errCaptchaRequired — sentinel error: сервер требует капчу. -var errCaptchaRequired = fmt.Errorf("captcha required") - -// getAllowCaptcha делает GET и, если встречает редирект на showcaptchafast, -// возвращает errCaptchaRequired. -func (t *BoardsTransport) getAllowCaptcha(client *http.Client, u, hash string) error { - req, _ := http.NewRequest("GET", u, nil) - req.Header.Set("User-Agent", boardsUA) - req.Header.Set("Accept", "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8") - req.Header.Set("Accept-Language", "en-US,en;q=0.9") - req.Header.Set("Referer", "https://"+boardsBase+"/guest/?hash="+hash) - resp, err := client.Do(req) - if err != nil { - return err - } - defer resp.Body.Close() - - if resp.StatusCode >= 300 && resp.StatusCode < 400 { - loc := resp.Header.Get("Location") - if strings.Contains(loc, "showcaptchafast") { - utils.Debugf("[BOARDS] redirect to captcha: %s", shortStr(loc, 100)) - return errCaptchaRequired - } - } - io.Copy(io.Discard, resp.Body) - return nil -} - -func (t *BoardsTransport) authorize(hash, name string) (boardsInfo, error) { - jar := t.jar() - client := &http.Client{ - Jar: jar, - Timeout: 15 * time.Second, - CheckRedirect: func(req *http.Request, via []*http.Request) error { - return http.ErrUseLastResponse - }, - } - - docURL := "https://" + boardsBase + "/whiteboard/?hash=" + hash - - // GET whiteboard — может вернуть 302 на showcaptchafast - if err := t.getAllowCaptcha(client, docURL, hash); err != nil { - if err == errCaptchaRequired { - utils.Debugf("[BOARDS] captcha required, solving...") - if _, cerr := solveCaptcha(docURL, jar, boardsUA); cerr != nil { - return boardsInfo{}, fmt.Errorf("captcha solve: %w", cerr) - } - utils.Debugf("[BOARDS] captcha solved, re-fetching whiteboard") - - // После капчи повторяем GET /whiteboard — сервер выдаёт - // свежие cookies, нужные для последующих /api запросов. - if err := t.getAllowCaptcha(client, docURL, hash); err != nil && err != errCaptchaRequired { - return boardsInfo{}, fmt.Errorf("GET whiteboard (post-captcha): %w", err) - } - } else { - return boardsInfo{}, err - } - } - - if err := t.postAPI(client, hash, "request-guest-token", - map[string]string{"name": name, "hash": hash}); err != nil { - return boardsInfo{}, fmt.Errorf("request-guest-token: %w", err) - } - - u, _ := url.Parse("https://" + boardsBase) - var jwt string - for _, c := range jar.Cookies(u) { - if c.Name == "token_"+hash { - jwt = c.Value - } - } - if jwt == "" { - return boardsInfo{}, fmt.Errorf("token_%s not found", hash) - } - payload := jwtPayload(jwt) - userHash, _ := payload["u"].(string) - - state, err := t.getWhiteboardInfo(client, hash) - if err != nil { - utils.Debugf("[BOARDS] get-whiteboard-info failed: %v", err) - } - - var cookies []*http.Cookie - cookies = append(cookies, jar.Cookies(u)...) - - wsHost := state["ws_host"] - if wsHost == "" { - wsHost = boardsSocketHostDefault - } - - return boardsInfo{ - hash: hash, - name: name, - userHash: userHash, - jwt: jwt, - cookies: cookies, - wsHost: wsHost, - session: "", - dashboard: state["dashboard"], - currentSlide: state["current_slide"], - }, nil -} - -func (t *BoardsTransport) get(client *http.Client, u, accept, referer string) error { - req, _ := http.NewRequest("GET", u, nil) - req.Header.Set("User-Agent", boardsUA) - req.Header.Set("Accept", accept) - req.Header.Set("Accept-Language", "en-US,en;q=0.9") - req.Header.Set("Referer", referer) - resp, err := client.Do(req) - if err != nil { - return err - } - io.Copy(io.Discard, resp.Body) - resp.Body.Close() - return nil -} - -// apiRequest builds a POST /api call: the action and its content (JSON, -// base64-encoded) in a JSON body. The API used to take a form and now -// answers 415 "Request body must use a JSON Content-Type" to one. -func apiRequest(hash, action string, content interface{}) *http.Request { - raw, _ := json.Marshal(content) - body, _ := json.Marshal(map[string]string{ - "action": action, - "content": base64.StdEncoding.EncodeToString(raw), - }) - req, _ := http.NewRequest("POST", "https://"+boardsBase+"/api", bytes.NewReader(body)) - req.Header.Set("User-Agent", boardsUA) - req.Header.Set("Content-Type", "application/json") - req.Header.Set("X-Requested-With", "XMLHttpRequest") - req.Header.Set("Accept", "application/json, text/javascript, */*; q=0.01") - req.Header.Set("Referer", "https://"+boardsBase+"/guest/?hash="+hash) - req.Header.Set("Origin", "https://"+boardsBase) - return req -} - -func (t *BoardsTransport) postAPI(client *http.Client, hash, action string, content interface{}) error { - resp, err := client.Do(apiRequest(hash, action, content)) - if err != nil { - return err - } - defer resp.Body.Close() - if resp.StatusCode != 200 { - body, _ := io.ReadAll(resp.Body) - n := len(body) - if n > 200 { - n = 200 - } - return fmt.Errorf("status %d: %s", resp.StatusCode, string(body[:n])) - } - return nil -} - -func (t *BoardsTransport) getWhiteboardInfo(client *http.Client, hash string) (map[string]string, error) { - resp, err := client.Do(apiRequest(hash, "get-whiteboard-info", map[string]string{"hash": hash})) - if err != nil { - return nil, err - } - defer resp.Body.Close() - body, _ := io.ReadAll(resp.Body) - if resp.StatusCode != 200 { - return nil, fmt.Errorf("status %d", resp.StatusCode) - } - - var info struct { - Presentation struct { - Properties struct { - CurrentSlide string `json:"current_slide"` - } `json:"properties"` - Items string `json:"items"` - } `json:"presentation"` - SocketServers []struct { - IP string `json:"ip"` - } `json:"socket_servers"` - } - if err := json.Unmarshal(body, &info); err != nil { - return nil, err - } - out := map[string]string{} - if info.Presentation.Properties.CurrentSlide != "" { - out["current_slide"] = info.Presentation.Properties.CurrentSlide - out["dashboard"] = info.Presentation.Properties.CurrentSlide - } - if info.Presentation.Items != "" { - if rawItems, err := base64.StdEncoding.DecodeString(info.Presentation.Items); err == nil { - var arr []map[string]interface{} - if json.Unmarshal(rawItems, &arr) == nil && len(arr) > 0 { - if h, _ := arr[0]["hash"].(string); h != "" && out["dashboard"] == "" { - out["dashboard"] = h - out["current_slide"] = h - } - } - } - } - if len(info.SocketServers) > 0 { - out["ws_host"] = info.SocketServers[0].IP - } - return out, nil -} - -func jwtPayload(jwt string) map[string]interface{} { - parts := strings.Split(jwt, ".") - if len(parts) < 2 { - return nil - } - pad := (4 - len(parts[1])%4) % 4 - b := parts[1] + strings.Repeat("=", pad) - raw, err := base64.URLEncoding.DecodeString(b) - if err != nil { - return nil - } - var m map[string]interface{} - _ = json.Unmarshal(raw, &m) - return m -} - -func extractBoardsHash(rawURL string) string { - u, err := url.Parse(rawURL) - if err != nil { - return "" - } - return u.Query().Get("hash") -} - -func randomGuestName() string { - var b [3]byte - if _, err := rand.Read(b[:]); err != nil { - return fmt.Sprintf("guest_%06x", time.Now().UnixNano()&0xffffff) - } - return "guest_" + hex.EncodeToString(b[:]) -} - -func randomHex(n int) string { - buf := make([]byte, n/2) - if _, err := rand.Read(buf); err != nil { - return fmt.Sprintf("%0*x", n, mrand.Int63()) - } - return hex.EncodeToString(buf) -} - -// ---- WS lifecycle ---- - -func (t *BoardsTransport) connectLoop(info boardsInfo) { - attempt := 0 - for { - select { - case <-t.done: - return - default: - } - if err := t.connectAndServe(info); err != nil { - utils.Debugf("[BOARDS] ws error: %v", err) - } - t.SetConnected(false) - select { - case <-t.done: - return - case <-time.After(reconnectBackoffBoards(attempt)): - } - attempt++ - if attempt > 10 { - attempt = 10 - } - } -} - -func reconnectBackoffBoards(n int) time.Duration { - if n < 1 { - n = 1 - } - shift := n - 1 - if shift > 4 { - shift = 4 - } - d := 500 * time.Millisecond * time.Duration(1< 15*time.Second { - d = 15 * time.Second - } - d += time.Duration(mrand.Int63n(int64(d/2) + 1)) - return d -} - -func (t *BoardsTransport) connectAndServe(info boardsInfo) error { - wsURL := fmt.Sprintf("wss://%s/socket.io/?EIO=4&transport=websocket", info.wsHost) - - header := http.Header{} - header.Set("User-Agent", boardsUA) - header.Set("Origin", "https://"+boardsBase) - header.Set("Accept-Language", "en-US,en;q=0.9") - - var cookieParts []string - for _, c := range info.cookies { - cookieParts = append(cookieParts, c.Name+"="+c.Value) - } - if !strings.Contains(strings.Join(cookieParts, ";"), "token_"+info.hash) { - cookieParts = append(cookieParts, "token_"+info.hash+"="+info.jwt) - } - header.Set("Cookie", strings.Join(cookieParts, "; ")) - - dialer := websocket.Dialer{ - HandshakeTimeout: 15 * time.Second, - NetDialContext: netbind.Wrap(&net.Dialer{ - Timeout: 10 * time.Second, - KeepAlive: 30 * time.Second, - }).DialContext, - } - utils.Debugf("[BOARDS] dial %s", wsURL) - conn, resp, err := dialer.Dial(wsURL, header) - if err != nil { - status := 0 - if resp != nil { - status = resp.StatusCode - } - return fmt.Errorf("dial %s (http %d): %w", wsURL, status, err) - } - utils.Debugf("[BOARDS] WS connected: %s", info.wsHost) - - participant := info.userHash - creator := info.userHash - sess := &boardsSession{ - Info: info, - Conn: conn, - Queue: make(chan []byte, t.GetConfig().MaxQueueSize), - } - sess.participant.Store(&participant) - sess.creatorHash.Store(&creator) - t.session.Store(sess) - - if err := t.handshake(sess); err != nil { - conn.Close() - t.session.Store(nil) - return fmt.Errorf("handshake: %w", err) - } - - _ = conn.SetReadDeadline(time.Time{}) - - t.SetConnected(true) - utils.SafeGo("boards.writer", func() { t.writerLoop(sess) }) - - kaStop := make(chan struct{}) - utils.SafeGo("boards.keepalive", func() { t.keepAliveLoop(sess, kaStop) }) - utils.SafeGo("boards.ping", func() { t.pingLoop(sess, kaStop) }) - defer close(kaStop) - - for { - select { - case <-t.done: - return nil - default: - } - - _ = conn.SetReadDeadline(time.Now().Add(boardsReadDeadline)) - _, msg, err := conn.ReadMessage() - if err != nil { - return fmt.Errorf("read: %w", err) - } - _ = conn.SetReadDeadline(time.Now().Add(boardsReadDeadline)) - - t.handleMessage(sess, msg) - } -} - -func (t *BoardsTransport) handshake(sess *boardsSession) error { - conn := sess.Conn - _ = conn.SetReadDeadline(time.Now().Add(15 * time.Second)) - - if _, err := readRaw(conn); err != nil { - return fmt.Errorf("engine.io hello: %w", err) - } - - if err := sess.writeRaw("40"); err != nil { - return err - } - if _, err := readRaw(conn); err != nil { - return fmt.Errorf("socket.io connect ack: %w", err) - } - - if err := sess.writeEventObj("im", map[string]interface{}{ - "operation": "subscribe", "user": nil, - }); err != nil { - return err - } - - deadline := time.Now().Add(10 * time.Second) - for time.Now().Before(deadline) { - _ = conn.SetReadDeadline(time.Now().Add(5 * time.Second)) - m, err := readRaw(conn) - if err != nil { - return fmt.Errorf("wait subscribed: %w", err) - } - if bytes.Contains(m, []byte(`"subscribed"`)) { - break - } - } - - part := *sess.participant.Load() - utils.Debugf("[BOARDS] subscribe-slide-dashboard participant=%s", shortStr(part, 8)) - if err := t.sendSubscribe(sess, part); err != nil { - return err - } - - _ = conn.SetReadDeadline(time.Now().Add(boardsHandshakeWait)) - deadline = time.Now().Add(boardsHandshakeWait) - for time.Now().Before(deadline) { - m, err := readRaw(conn) - if err != nil { - return fmt.Errorf("wait 431: %w", err) - } - t.handleMessage(sess, m) - if bytes.HasPrefix(m, []byte("431[")) { - break - } - if bytes.Contains(m, []byte(`"subscribed":true`)) && - bytes.Contains(m, []byte(`"dashboard_link"`)) { - break - } - } - - _ = conn.SetReadDeadline(time.Time{}) - utils.Debugf("[BOARDS] handshake done") - return nil -} - -func readRaw(conn *websocket.Conn) ([]byte, error) { - _, msg, err := conn.ReadMessage() - if err != nil { - return nil, err - } - return msg, nil -} - -func (t *BoardsTransport) sendSubscribe(sess *boardsSession, participant string) error { - data := map[string]interface{}{ - "session": sess.Info.session, - "dashboard": sess.Info.dashboard, - "presentation": sess.Info.hash, - "properties": map[string]interface{}{ - "guest_mode": true, - "guest_role": 1, - "guest_password": nil, - "guest_password_expiration_date": nil, - "current_slide": sess.Info.currentSlide, - }, - "participant_team_role": -1, - "participant": participant, - "options": map[string]interface{}{ - "type": "landing", - "participant": map[string]interface{}{ - "hash": participant, - "partner": "yandex", - "userHash": sess.Info.userHash, - "name": sess.Info.name, - "additional": map[string]interface{}{"guest": true}, - "module": "yandex", - "presentation": sess.Info.hash, - "identityCandidates": map[string]interface{}{ - "uidHash": nil, - "legacyHash": participant, - "uid": nil, - "partner": "yandex", - }, - "module_type": "yandex", - "participantCaptionName": sess.Info.name, - }, - "intermediate": "", - "device": map[string]interface{}{ - "screen": "674 x 619", - "screen_width": 674, - "screen_height": 619, - "browser": "Chrome", - "browserVersion": "153.0.0.0", - "browserMajorVersion": 153, - "mobile": true, - "os": "Android", - "osVersion": "15", - "osMajorVersion": 15, - "cookies": true, - "flashVersion": "no check", - "agent": "Chrome", - "appVersion": boardsUA, - "userAgent": boardsUA, - "appName": "Netscape", - "platform": "MacIntel", - }, - }, - } - return sess.writeEventObj("dashboard", map[string]interface{}{ - "action": "subscribe-slide-dashboard", - "data": data, - "participant": participant, - }) -} - -// ---- отправка ---- -// -// Канал — notify-position. Формат (объект, как в оригинале): -// -// 42N["dashboard",{ -// "action":"notify-position", -// "data":{ -// "position":{"x":"","y":123.0}, -// "vpt":{"translate":{"x":0,"y":0},"scale":1,"whyrugay":1} -// }, -// "participant":"" -// }] -// -// payload — base64 от tunnel-пакета, кладётся в data.position.x. - -func (t *BoardsTransport) writerLoop(sess *boardsSession) { - queue := sess.Queue - for { - select { - case <-t.done: - return - case pkt := <-queue: - if err := t.sendNotifyPosition(sess, pkt); err != nil { - utils.Debugf("[BOARDS] notify-position: %v", err) - } else { - t.RecordSend(len(pkt)) - } - } - } -} - -// sendNotifyPosition — точная копия формата из Python-скрипта: -// -// cmd_pos(x, y) -> send_dashboard("notify-position", -// {"position": {"x": x, "y": y}, -// "vpt": {"translate": {"x": 0, "y": 0}, "scale": 1}}) -// -// send_dashboard добавляет participant в envelope. -// payload кладём в position.x как base64-строку, position.y = 123.0, -// vpt.whyrugay = 1 (как в оригинале). -func (t *BoardsTransport) sendNotifyPosition(sess *boardsSession, pkt []byte) error { - b64 := base64.StdEncoding.EncodeToString(pkt) - - data := map[string]interface{}{ - "position": map[string]interface{}{ - "x": b64, - "y": 123.0, - }, - "vpt": map[string]interface{}{ - "translate": map[string]interface{}{"x": 0, "y": 0}, - "scale": 1, - "whyrugay": 1, - }, - } - - obj := map[string]interface{}{ - "action": "notify-position", - "data": data, - "participant": *sess.participant.Load(), - } - return sess.writeEventObj("dashboard", obj) -} - -// ---- keepalive ---- - -func (t *BoardsTransport) keepAliveLoop(sess *boardsSession, stop chan struct{}) { - tick := time.NewTicker(boardsPingInterval) - defer tick.Stop() - for { - select { - case <-stop: - return - case <-t.done: - return - case <-tick.C: - obj := map[string]interface{}{ - "action": "heartbeat", - "data": map[string]interface{}{}, - "participant": *sess.participant.Load(), - } - if err := sess.writeEventObj("dashboard", obj); err != nil { - utils.Debugf("[BOARDS] heartbeat: %v", err) - return - } - } - } -} - -func (t *BoardsTransport) pingLoop(sess *boardsSession, stop chan struct{}) { - tick := time.NewTicker(boardsPingInterval) - defer tick.Stop() - for { - select { - case <-stop: - return - case <-t.done: - return - case <-tick.C: - if err := sess.writeRaw("2"); err != nil { - utils.Debugf("[BOARDS] engine.io ping: %v", err) - return - } - } - } -} - -// ---- приём ---- - -func (t *BoardsTransport) handleMessage(sess *boardsSession, raw []byte) { - if len(raw) == 1 && raw[0] == '2' { - utils.Debugf("[BOARDS] ping -> pong") - _ = sess.writeRaw("3") - return - } - if len(raw) == 1 && raw[0] == '3' { - return - } - - if bytes.HasPrefix(raw, []byte("43")) { - t.handle431(sess, raw) - return - } - if !bytes.HasPrefix(raw, []byte("42[")) { - return - } - idx := bytes.IndexByte(raw, '[') - if idx < 0 { - return - } - var arr []json.RawMessage - if err := json.Unmarshal(raw[idx:], &arr); err != nil { - return - } - if len(arr) < 2 { - return - } - - var envelope struct { - Action string `json:"action"` - Data json.RawMessage `json:"data"` - Participant string `json:"participant"` - } - if err := json.Unmarshal(arr[1], &envelope); err != nil { - return - } - - switch envelope.Action { - case "participant-connected": - t.handleParticipantConnected(sess, envelope.Data) - case "notify-position": - t.handleNotifyPosition(sess, envelope.Data, envelope.Participant) - case "server-modify-objects", "modify-objects": - t.handleServerModifyObjects(sess, envelope.Data, envelope.Action) - } -} - -func (t *BoardsTransport) handleParticipantConnected(sess *boardsSession, raw json.RawMessage) { - var d struct { - Participant struct { - Hash string `json:"hash"` - Session string `json:"session"` - Name string `json:"name"` - } `json:"participant"` - } - if err := json.Unmarshal(raw, &d); err != nil { - return - } - utils.Debugf("[BOARDS] participant-connected: name=%q hash=%s session=%s", - d.Participant.Name, shortStr(d.Participant.Hash, 8), shortStr(d.Participant.Session, 8)) - - if d.Participant.Session != "" && sess.Info.session == "" { - sess.Info.session = d.Participant.Session - } - - if d.Participant.Name == sess.Info.name && d.Participant.Hash != "" { - h := d.Participant.Hash - sess.creatorHash.Store(&h) - utils.Debugf("[BOARDS] creatorHash from participant-connected: %s", shortStr(h, 8)) - } -} - -// handleNotifyPosition — основной канал приёма. Пробуем сначала объект -// (position.x), затем массив (data[4]). -// -// Объект: -// -// data.position.x — base64 -// -// Массив: -// -// data[2] — имя отправителя, data[4] — base64 -// -// Своё эхо фильтруем по: -// - envelope.participant == наш participant (для объекта) -// - data[2] == наше имя (для массива) -func (t *BoardsTransport) handleNotifyPosition(sess *boardsSession, raw json.RawMessage, envelopePart string) { - // Сначала пробуем объект (position.x) - var objForm struct { - Position struct { - X json.RawMessage `json:"x"` - Y json.RawMessage `json:"y"` - } `json:"position"` - } - if err := json.Unmarshal(raw, &objForm); err == nil { - var xStr string - if err := json.Unmarshal(objForm.Position.X, &xStr); err == nil && xStr != "" { - // Своё эхо по participant в envelope - myPart := *sess.participant.Load() - if envelopePart != "" && envelopePart == myPart { - return - } - decoded, derr := base64.StdEncoding.DecodeString(xStr) - if derr != nil || len(decoded) == 0 { - return - } - utils.Debugf("[BOARDS<-] notify-position obj from=%s pktlen=%d", - shortStr(envelopePart, 8), len(decoded)) - t.RecordReceive(len(decoded)) - t.onDataMu.RLock() - cb := t.onData - t.onDataMu.RUnlock() - if cb != nil { - cb(decoded) - } - return - } - } - - // Массив (data[2]=name, data[4]=base64) - var arr []json.RawMessage - if err := json.Unmarshal(raw, &arr); err != nil { - return - } - if len(arr) < 5 { - return - } - var sender string - _ = json.Unmarshal(arr[2], &sender) - if sender == sess.Info.name { - return - } - var b64 string - if err := json.Unmarshal(arr[4], &b64); err != nil || b64 == "" { - return - } - decoded, err := base64.StdEncoding.DecodeString(b64) - if err != nil || len(decoded) == 0 { - return - } - utils.Debugf("[BOARDS<-] notify-position arr from=%q pktlen=%d", sender, len(decoded)) - t.RecordReceive(len(decoded)) - t.onDataMu.RLock() - cb := t.onData - t.onDataMu.RUnlock() - if cb != nil { - cb(decoded) - } -} - -func (t *BoardsTransport) handleServerModifyObjects(sess *boardsSession, raw json.RawMessage, action string) { - var d struct { - Dashboard string `json:"dashboard"` - Name string `json:"name"` - Session string `json:"session"` - Objects []struct { - Attributes map[string]interface{} `json:"_attributes_"` - } `json:"objects"` - } - if err := json.Unmarshal(raw, &d); err != nil { - return - } - if len(d.Objects) == 0 { - return - } - - myPart := *sess.participant.Load() - myUser := sess.Info.userHash - myName := sess.Info.name - - for _, o := range d.Objects { - val, _ := o.Attributes["value"].(string) - if val == "" { - continue - } - id, _ := o.Attributes["id"].(string) - creator, _ := o.Attributes["creatorHash"].(string) - typ, _ := o.Attributes["type"].(string) - - if creator != "" && (creator == myPart || creator == myUser) { - utils.Debugf("[BOARDS] %s: own echo (creator=%s), skip", - action, shortStr(creator, 8)) - continue - } - if d.Name != "" && d.Name == myName { - utils.Debugf("[BOARDS] %s: own echo (name=%q), skip", action, d.Name) - continue - } - - decoded, err := base64.StdEncoding.DecodeString(val) - if err != nil || len(decoded) == 0 { - utils.Debugf("[BOARDS] %s: non-base64 value id=%s type=%s, skip", - action, shortStr(id, 8), typ) - continue - } - - utils.Debugf("[BOARDS<-] %s from=%q creator=%s id=%s pktlen=%d", - action, d.Name, shortStr(creator, 8), shortStr(id, 8), len(decoded)) - - t.RecordReceive(len(decoded)) - t.onDataMu.RLock() - cb := t.onData - t.onDataMu.RUnlock() - if cb != nil { - cb(decoded) - } - } -} - -func (t *BoardsTransport) handle431(sess *boardsSession, raw []byte) { - idx := bytes.IndexByte(raw, '[') - if idx < 0 { - return - } - body := raw[idx:] - if !bytes.Contains(body, []byte(`"dashboard_link"`)) && - !bytes.Contains(body, []byte(`"participantHash"`)) && - !bytes.Contains(body, []byte(`"creatorHash"`)) { - return - } - var arr []json.RawMessage - if err := json.Unmarshal(body, &arr); err != nil { - return - } - if len(arr) == 0 { - return - } - var snap struct { - Participant struct { - Hash string `json:"hash"` - Session string `json:"session"` - ParticipantHash string `json:"participantHash"` - CreatorHash string `json:"creatorHash"` - DashboardLink struct { - Session string `json:"session"` - Dashboard string `json:"dashboard"` - } `json:"dashboard_link"` - } `json:"participant"` - } - if err := json.Unmarshal(arr[0], &snap); err != nil { - return - } - - if snap.Participant.Session != "" && sess.Info.session == "" { - sess.Info.session = snap.Participant.Session - } - if snap.Participant.DashboardLink.Session != "" && sess.Info.session == "" { - sess.Info.session = snap.Participant.DashboardLink.Session - } - if snap.Participant.DashboardLink.Dashboard != "" && sess.Info.dashboard == "" { - sess.Info.dashboard = snap.Participant.DashboardLink.Dashboard - sess.Info.currentSlide = sess.Info.dashboard - utils.Debugf("[BOARDS] dashboard from 431: %s", shortStr(sess.Info.dashboard, 8)) - } - - if snap.Participant.CreatorHash != "" { - ch := snap.Participant.CreatorHash - sess.creatorHash.Store(&ch) - utils.Debugf("[BOARDS] creatorHash from 431: %s", shortStr(ch, 8)) - } - if snap.Participant.ParticipantHash != "" { - ch := snap.Participant.ParticipantHash - sess.creatorHash.Store(&ch) - utils.Debugf("[BOARDS] creatorHash from 431/participantHash: %s", shortStr(ch, 8)) - } -} - -// jar returns the transport's shared cookie jar (never nil). -func (t *BoardsTransport) jar() *cookiejar.Jar { - t.jarMu.RLock() - jar := t.cookieJar - t.jarMu.RUnlock() - if jar == nil { - jar, _ = cookiejar.New(nil) - t.jarMu.Lock() - t.cookieJar = jar - t.jarMu.Unlock() - } - return jar -} - -// ---- CookieExchanger ---- - -// FetchCookies returns a snapshot of the transport's current cookie jar as -// name -> value for boards.yandex.ru. -func (t *BoardsTransport) FetchCookies() (map[string]string, error) { - u, err := url.Parse("https://" + boardsBase + "/") - if err != nil { - return nil, err - } - out := make(map[string]string) - for _, c := range t.jar().Cookies(u) { - out[c.Name] = c.Value - } - return out, nil -} - -// ApplyCookies replaces the transport's cookie jar and forces the current WS -// session to reconnect. -func (t *BoardsTransport) ApplyCookies(values map[string]string) error { - if len(values) == 0 { - return nil - } - u, _ := url.Parse("https://" + boardsBase + "/") - jar, _ := cookiejar.New(nil) - cookies := make([]*http.Cookie, 0, len(values)) - for k, v := range values { - cookies = append(cookies, &http.Cookie{Name: k, Value: v, Path: "/"}) - } - jar.SetCookies(u, cookies) - - t.jarMu.Lock() - t.cookieJar = jar - t.jarMu.Unlock() - - utils.Debugf("[BOARDS] applied %d cookies, forcing reconnect", len(cookies)) - - if s := t.session.Load(); s != nil && s.Conn != nil { - _ = s.Conn.Close() - } - return nil -} diff --git a/transport/yandex/yandex.go b/transport/yandex/yandex.go index 7ce1ac1..be7a9b1 100644 --- a/transport/yandex/yandex.go +++ b/transport/yandex/yandex.go @@ -472,6 +472,16 @@ func (t *YandexDocsTransport) scheduleReconnectNoCaptcha(attempt int) { // participant-list messages during a failure streak). The floor here (was // 500ms) is raised to slow that churn down; this doesn't change steady-state // throughput since successful connects never hit backoff at all. +// shortStr truncates s to n bytes for log lines - a redirect chain or +// Location header can be arbitrarily long, and the full value is rarely +// what a log reader needs. +func shortStr(s string, n int) string { + if len(s) <= n { + return s + } + return s[:n] +} + func reconnectBackoff(n int) time.Duration { if n < 1 { n = 1