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:<name> (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 <noreply@anthropic.com>
This commit is contained in:
p1neappleXpress
2026-09-27 08:40:31 +03:00
co-authored by Claude Sonnet 5
parent 70e7a2042c
commit ab72789eef
8 changed files with 10 additions and 4429 deletions
File diff suppressed because it is too large Load Diff
-115
View File
@@ -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("<html>room closed</html>")) },
"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)
}
}
}
-842
View File
@@ -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, `<html><head>
<meta name="centrifuge-connection-token" content="conn-%s">
<meta name="centrifuge-connection-url" content="%s/connection">
<meta name="centrifuge-subscription-token-url" content="%s/sub/">
</head><body><div data-room="{&quot;uuid&quot;: &quot;%s&quot;}" data-user="{&quot;uuid&quot;: &quot;%s&quot;}"></div></body></html>`,
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")
}
}
-234
View File
@@ -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()
}
-19
View File
@@ -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)
}
}
-591
View File
@@ -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<<uint(shift))
if d > 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
}
File diff suppressed because it is too large Load Diff
+10
View File
@@ -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