// The relay command's networked helpers, driven against loopback WebSocket // servers instead of real relays: syncOnce (two peers: a remote feed and the // local relay it forwards into), outboxDiscover (the two REQs that build the // write-relay map), and the crawler (crawlRelay, crawlPublishBatch, // crawlPass). The remaining main.mx blocks are the subcommands that loop // forever or spawn the binary, which stay integration-only. // // The server side speaks just enough of the protocol: the 101 upgrade with a // correct Sec-WebSocket-Accept, unmasked text frames, and masked client frames // read back. It validates what the client sends and withholds the reply when // the request is wrong, so a passing assertion means the client really sent it. package main import ( "os" "syscall" "testing" "git.smesh.lol/nostr/pkg/envelope" "git.smesh.lol/nostr/pkg/ws" ) // --- loopback WebSocket server --- func mwBind() (fd int32, port int32, err error) { fd, err = syscall.Socket(syscall.AF_INET, syscall.SOCK_STREAM, 0) if err != nil { return 0, 0, err } syscall.SetsockoptInt(fd, syscall.SOL_SOCKET, syscall.SO_REUSEADDR, 1) tv := syscall.Timeval{Sec: 5} syscall.SetsockoptTimeval(fd, syscall.SOL_SOCKET, syscall.SO_RCVTIMEO, &tv) sa := &syscall.SockaddrInet4{Port: 0, Addr: [4]byte{127, 0, 0, 1}} if err = syscall.Bind(fd, sa); err != nil { syscall.Close(fd) return 0, 0, err } if err = syscall.Listen(fd, 8); err != nil { syscall.Close(fd) return 0, 0, err } got, gerr := syscall.Getsockname(fd) if gerr != nil { syscall.Close(fd) return 0, 0, gerr } sa4, ok := got.(*syscall.SockaddrInet4) if !ok { syscall.Close(fd) return 0, 0, syscall.EINVAL } return fd, sa4.Port, nil } func mwPortStr(port int32) (s string) { if port == 0 { return "0" } buf := []byte{:6} n := int32(6) for port > 0 { n-- buf[n] = byte('0' + port%10) port = port / 10 } return string(buf[n:]) } func mwListen(t *testing.T) (fd int32, base string) { t.Helper() fd, port, err := mwBind() if err != nil { t.Fatalf("listen: %s", err.Error()) return 0, "" } return fd, "ws://127.0.0.1:" | mwPortStr(port) } // mwHeader returns the value of an HTTP header, tested case-insensitively. func mwHeader(req []byte, name string) (v string) { n := int32(len(name)) for i := int32(0); i+n <= int32(len(req)); i++ { if string(req[i:i+n]) == name { j := i + n for j < int32(len(req)) && req[j] != '\r' && req[j] != '\n' { j++ } return string(req[i+n : j]) } } return "" } func mwHasHeadEnd(b []byte) (ok bool) { for i := int32(0); i+3 < int32(len(b)); i++ { if b[i] == '\r' && b[i+1] == '\n' && b[i+2] == '\r' && b[i+3] == '\n' { return true } } return false } func mwReadHead(fd int32) (req []byte) { buf := []byte{:2048} for { n, err := syscall.Read(fd, buf) if n <= 0 || err != nil { return req } req = req | buf[:n] if mwHasHeadEnd(req) { return req } } } // mwAccept accepts one connection and completes the upgrade. A negative // return means the listener timed out or failed. func mwAccept(fd int32) (nfd int32) { nfd, _, err := syscall.Accept(fd) if err != nil { return -1 } req := mwReadHead(nfd) head := "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: " | ws.ComputeAccept(mwHeader(req, "Sec-WebSocket-Key: ")) | "\r\n\r\n" syscall.Write(nfd, []byte(head)) return nfd } // mwFrame builds an unmasked server frame. func mwFrame(op byte, payload []byte) (buf []byte) { plen := int32(len(payload)) if plen < 126 { buf = []byte{:2} buf[1] = byte(plen) } else { buf = []byte{:4} buf[1] = 126 buf[2] = byte(plen >> 8) buf[3] = byte(plen) } buf[0] = 0x80 | op return buf | payload } func mwReadN(fd int32, n int32) (out []byte) { out = []byte{:0:n} for len(out) < n { buf := []byte{:n - len(out)} got, err := syscall.Read(fd, buf) if got <= 0 || err != nil { return out } out = out | buf[:got] } return } // mwReadFrame decodes one masked client frame and unmasks it. func mwReadFrame(fd int32) (payload []byte) { hdr := mwReadN(fd, 2) if len(hdr) < 2 { return nil } plen := int32(hdr[1] & 0x7F) if plen == 126 { ext := mwReadN(fd, 2) if len(ext) < 2 { return nil } plen = int32(ext[0])<<8 | int32(ext[1]) } else if plen == 127 { ext := mwReadN(fd, 8) if len(ext) < 8 { return nil } plen = int32(ext[4])<<24 | int32(ext[5])<<16 | int32(ext[6])<<8 | int32(ext[7]) } var mask [4]byte if hdr[1]&0x80 != 0 { m := mwReadN(fd, 4) if len(m) < 4 { return nil } mask = [4]byte{m[0], m[1], m[2], m[3]} } body := mwReadN(fd, plen) for i := int32(0); i < int32(len(body)); i++ { body[i] = body[i] ^ mask[i%4] } return body } // --- frames the client parses --- // mwEventSub builds ["EVENT",{...}], the shape syncOnce's EventSubmission // parser wants. func mwEventSub(kind uint16, pubkeyHex, content string, createdAt int64, tagsJSON string) (raw []byte) { raw = []byte(`["EVENT",{"id":"`) | mxFill(64, '0') raw = raw | `","pubkey":"` | pubkeyHex raw = raw | `","created_at":` | itoa64(createdAt) raw = raw | `,"kind":` | itoa64(int64(kind)) raw = raw | `,"tags":` | tagsJSON raw = raw | `,"content":"` | content | `","sig":""}]` return } // mwEventRes builds ["EVENT","s",{...}], the shape EventResult wants. func mwEventRes(kind uint16, pubkeyHex, content string, createdAt int64, tagsJSON string) (raw []byte) { raw = []byte(`["EVENT","s",{"id":"`) | mxFill(64, '0') raw = raw | `","pubkey":"` | pubkeyHex raw = raw | `","created_at":` | itoa64(createdAt) raw = raw | `,"kind":` | itoa64(int64(kind)) raw = raw | `,"tags":` | tagsJSON raw = raw | `,"content":"` | content | `","sig":""}]` return } func mwEOSE(sub string) (raw []byte) { e := &envelope.EOSE{Subscription: []byte(sub)} return e.Marshal(nil) } func mwOKFrame(id []byte) (raw []byte) { o := &envelope.OK{EventID: id, OK: true, Reason: []byte("ok")} return o.Marshal(nil) } func mwHasSub(hay, needle []byte) (ok bool) { n := int32(len(needle)) if n == 0 || int32(len(hay)) < n { return false } for i := int32(0); i+n <= int32(len(hay)); i++ { if string(hay[i:i+n]) == string(needle) { return true } } return false } // mwParseSub parses an EVENT submission and returns its event. func mwParseSub(raw []byte) (e *envelope.EventSubmission) { _, rem, err := envelope.Identify(raw) if err != nil { return nil } var es envelope.EventSubmission if _, uerr := es.Unmarshal(rem); uerr != nil || es.E == nil { return nil } return &es } // --- scripted peers --- // mwServeSyncRemote answers a sync subscription: one event, EOSE, then close. // When wantA/wantB are non-empty the REQ must contain both before the event is // sent, so the test can assert the filter the client built. func mwServeSyncRemote(fd int32, raw []byte, wantA, wantB []byte) { nfd := mwAccept(fd) if nfd < 0 { syscall.Close(fd) return } req := mwReadFrame(nfd) if len(wantA) > 0 && !mwHasSub(req, wantA) { syscall.Close(nfd) syscall.Close(fd) return } if len(wantB) > 0 && !mwHasSub(req, wantB) { syscall.Close(nfd) syscall.Close(fd) return } syscall.Write(nfd, mwFrame(ws.OpText, raw)) syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("sync"))) syscall.Close(nfd) syscall.Close(fd) } // mwServeSyncLocal consumes the forwarded event and acknowledges it. The OK is // only sent for the expected kind and content, so a client that forwards // nothing (or the wrong thing) fails its drain read. func mwServeSyncLocal(fd int32, wantKind uint16, wantContent []byte) { nfd := mwAccept(fd) if nfd < 0 { syscall.Close(fd) return } raw := mwReadFrame(nfd) es := mwParseSub(raw) if es != nil && es.E.Kind == wantKind && string(es.E.Content) == string(wantContent) { syscall.Write(nfd, mwFrame(ws.OpText, mwOKFrame(es.E.ID))) } syscall.Close(nfd) syscall.Close(fd) } // mwServeOutbox answers the two REQs outboxDiscover sends. func mwServeOutbox(fd int32, k3raw, rlraw []byte) { nfd := mwAccept(fd) if nfd < 0 { syscall.Close(fd) return } mwReadFrame(nfd) syscall.Write(nfd, mwFrame(ws.OpText, k3raw)) syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-k3"))) mwReadFrame(nfd) if rlraw != nil { syscall.Write(nfd, mwFrame(ws.OpText, rlraw)) } syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-rl"))) syscall.Close(nfd) syscall.Close(fd) } // mwServeIndex answers count REQs with one event each and closes. func mwServeIndex(fd int32, raw []byte, count int32) { nfd := mwAccept(fd) if nfd < 0 { syscall.Close(fd) return } for i := int32(0); i < count; i++ { mwReadFrame(nfd) syscall.Write(nfd, mwFrame(ws.OpText, raw)) syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("cr"))) } syscall.Close(nfd) syscall.Close(fd) } // mwServeAckN reads count published events and acknowledges each. func mwServeAckN(fd int32, count int32) { nfd := mwAccept(fd) if nfd < 0 { syscall.Close(fd) return } ids := [][]byte{:count} for i := int32(0); i < count; i++ { raw := mwReadFrame(nfd) es := mwParseSub(raw) if es != nil { ids[i] = es.E.ID } } for i := int32(0); i < count; i++ { syscall.Write(nfd, mwFrame(ws.OpText, mwOKFrame(ids[i]))) } syscall.Close(nfd) syscall.Close(fd) } // mwCrawlLog opens the throwaway log the crawler writes to. func mwCrawlLog(t *testing.T) (out *os.File) { t.Helper() f, err := os.OpenFile("/dev/null", os.O_WRONLY, 0) if err != nil { return os.Stderr } return f } func mwSameStrings(a, b []string) (ok bool) { if len(a) != len(b) { return false } for i := int32(0); i < int32(len(a)); i++ { if a[i] != b[i] { return false } } return true } // --- tests --- func TestSyncOnceForwardsRemoteEvents(t *testing.T) { rfd, rbase := mwListen(t) if rfd == 0 { return } lfd, lbase := mwListen(t) if lfd == 0 { syscall.Close(rfd) return } raw := mwEventSub(1, mxHex64, "mw-sync", 1700000500, "[]") rdone := spawn(mwServeSyncRemote, rfd, raw, []byte(nil), []byte(nil)) ldone := spawn(mwServeSyncLocal, lfd, uint16(1), []byte("mw-sync")) got := syncOnce(rbase, lbase, "", 0) if got != 1700000500 { t.Fatalf("syncOnce = %d, want the forwarded event's timestamp", got) } <-rdone <-ldone } func TestSyncOnceBuildsTheAuthorAndSinceFilter(t *testing.T) { rfd, rbase := mwListen(t) if rfd == 0 { return } lfd, lbase := mwListen(t) if lfd == 0 { syscall.Close(rfd) return } raw := mwEventSub(1, mxHex64, "mw-filter", 1700000600, "[]") // sinceTs is 1700003600, so the filter carries since=1700000000; the // authors list keeps only the 64-hex keys. wantAuthors := []byte(`"authors":["` | mxHex64 | `","` | mxHex64 | `"]`) wantSince := []byte(`"since":1700000000`) rdone := spawn(mwServeSyncRemote, rfd, raw, wantAuthors, wantSince) ldone := spawn(mwServeSyncLocal, lfd, uint16(1), []byte("mw-filter")) got := syncOnce(rbase, lbase, mxHex64|",nothex,"|mxHex64, 1700003600) if got != 1700000600 { t.Fatalf("syncOnce = %d, want 1700000600", got) } <-rdone <-ldone } func TestSyncOnceDialFailures(t *testing.T) { // The remote is refused: the checkpoint is returned unchanged. if got := syncOnce("ws://127.0.0.1:1", "ws://127.0.0.1:1", "", 42); got != 42 { t.Fatalf("remote dial failure = %d, want 42", got) } // The remote answers but the local relay is refused. rfd, rbase := mwListen(t) if rfd == 0 { return } rdone := spawn(mwServeSyncRemote, rfd, mwEventSub(1, mxHex64, "x", 1, "[]"), []byte(nil), []byte(nil)) if got := syncOnce(rbase, "ws://127.0.0.1:1", "", 42); got != 42 { t.Fatalf("local dial failure = %d, want 42", got) } <-rdone } func TestOutboxDiscoverBuildsWriteRelayMap(t *testing.T) { fd, base := mwListen(t) if fd == 0 { return } pkA := mxFill(64, 'a') pkB := mxFill(64, 'b') owner := mxFill(64, 'c') k3 := mwEventRes(3, owner, "", 1700000000, `[["p","`|pkA|`"],["p","`|pkB|`"],["p","short"]]`) rl := mwEventRes(10002, pkA, "", 1700000000, `[["r","wss://a.example"],["r","wss://b.example","write"],["r","wss://c.example","read"],["r","http://bad.example"],["r"]]`) done := spawn(mwServeOutbox, fd, k3, rl) m := outboxDiscover(owner, base) if m == nil { t.Fatal("outboxDiscover returned nil") } if len(m) != 3 { t.Fatalf("write relay buckets = %d, want 3", int32(len(m))) } if !mwSameStrings(m["wss://a.example"], []string{pkA}) { t.Fatal("an unmarked r tag must be a write relay") } if !mwSameStrings(m["wss://b.example"], []string{pkA}) { t.Fatal("a write-marked r tag must be a write relay") } if _, ok := m["wss://c.example"]; ok { t.Fatal("a read-marked r tag must not be a write relay") } if _, ok := m["http://bad.example"]; ok { t.Fatal("a non-ws relay URL must be ignored") } if !mwSameStrings(m["wss://relay.damus.io"], []string{pkB, owner}) { t.Fatal("follows without a kind 10002 (and the owner) must fall back to the default relay") } <-done } func TestOutboxDiscoverDialError(t *testing.T) { if m := outboxDiscover(mxHex64, "ws://127.0.0.1:1"); m != nil { t.Fatal("a refused local relay must return nil") } } func TestCrawlRelayAndPublish(t *testing.T) { out := mwCrawlLog(t) raw := mwEventRes(1, mxHex64, "mw-crawl", 1700000000, "[]") // One relay: crawlRelay returns the event it saw. rfd, rbase := mwListen(t) if rfd == 0 { return } rdone := spawn(mwServeIndex, rfd, raw, int32(1)) events := crawlRelay(rbase, out) if len(events) != 1 { t.Fatalf("crawlRelay returned %d events", int32(len(events))) } if es := mwParseSub(events[0]); es == nil || string(es.E.Content) != "mw-crawl" { t.Fatal("crawlRelay must return the event as an EVENT submission") } <-rdone // crawlPublishBatch forwards them to the local relay and counts the OKs. lfd, lbase := mwListen(t) if lfd == 0 { return } ldone := spawn(mwServeAckN, lfd, int32(1)) if n := crawlPublishBatch(lbase, events, out); n != 1 { t.Fatalf("crawlPublishBatch published %d, want 1", n) } <-ldone // A full pass over one live relay and one dead one. rfd2, rbase2 := mwListen(t) if rfd2 == 0 { return } rdone2 := spawn(mwServeIndex, rfd2, raw, int32(1)) lfd2, lbase2 := mwListen(t) if lfd2 == 0 { return } ldone2 := spawn(mwServeAckN, lfd2, int32(1)) db := newRelayDB() db.add("ws://127.0.0.1:1", 5) db.add(rbase2, 100) if !crawlPass(lbase2, db, out) { t.Fatal("crawlPass must report a completed pass") } <-rdone2 <-ldone2 // No relays at all is a failed pass. if crawlPass(lbase2, newRelayDB(), out) { t.Fatal("an empty relay database must fail the pass") } } // mwServeCloseAfterUpgrade accepts the upgrade and closes, so the client's // first write fails. func mwServeCloseAfterUpgrade(fd int32) { nfd := mwAccept(fd) if nfd >= 0 { syscall.Close(nfd) } syscall.Close(fd) } // mwServeDropAck reads one forwarded event and closes without acknowledging. func mwServeDropAck(fd int32) { nfd := mwAccept(fd) if nfd < 0 { syscall.Close(fd) return } mwReadFrame(nfd) syscall.Close(nfd) syscall.Close(fd) } func TestSyncOnceWriteAndDrainFailures(t *testing.T) { // The remote closes right after the upgrade: the subscribe write fails and // the checkpoint is returned unchanged. rfd, rbase := mwListen(t) if rfd == 0 { return } rdone := spawn(mwServeCloseAfterUpgrade, rfd) if got := syncOnce(rbase, "ws://127.0.0.1:1", "", 7); got != 7 { t.Fatalf("subscribe failure = %d, want 7", got) } <-rdone // The local relay drops the connection instead of acknowledging: the // forwarded event's timestamp is still returned. rfd2, rbase2 := mwListen(t) if rfd2 == 0 { return } lfd, lbase := mwListen(t) if lfd == 0 { syscall.Close(rfd2) return } raw := mwEventSub(1, mxHex64, "mw-drop", 1700000700, "[]") rdone2 := spawn(mwServeSyncRemote, rfd2, raw, []byte(nil), []byte(nil)) ldone := spawn(mwServeDropAck, lfd) // The event was forwarded but never acknowledged, so the checkpoint does // not advance and the next run re-sends it. if got := syncOnce(rbase2, lbase, "", 0); got != 0 { t.Fatalf("drain failure = %d, want the unadvanced checkpoint 0", got) } <-rdone2 <-ldone } func TestOutboxDiscoverWriteError(t *testing.T) { fd, base := mwListen(t) if fd == 0 { return } done := spawn(mwServeCloseAfterUpgrade, fd) if m := outboxDiscover(mxHex64, base); m != nil { t.Fatal("a REQ that cannot be written must return nil") } <-done } func TestOutboxBootstrapRelayLists(t *testing.T) { sfd, sbase := mwListen(t) if sfd == 0 { return } lfd, lbase := mwListen(t) if lfd == 0 { syscall.Close(sfd) return } pkA := mxFill(64, 'a') pkB := mxFill(64, 'b') raw := mwEventRes(10002, pkA, "", 1700000000, `[["r","wss://a.example"]]`) want := []byte(`"authors":["` | pkA | `","` | pkB | `"]`) sdone := spawn(mwServeSeedRelayLists, sfd, raw, want) ldone := spawn(mwServeAckN, lfd, int32(1)) outboxBootstrapRelayLists([]string{pkA, pkB}, lbase, sbase) <-sdone <-ldone // No follows is a no-op; a dead seed returns quietly before dialing local. outboxBootstrapRelayLists([]string{}, lbase, sbase) outboxBootstrapRelayLists([]string{pkA}, lbase, "ws://127.0.0.1:1") // A dead local relay returns quietly too, after the seed has accepted. sfd2, sbase2 := mwListen(t) if sfd2 == 0 { return } sdone2 := spawn(mwServeSeedRelayLists, sfd2, raw, want) outboxBootstrapRelayLists([]string{pkA}, "ws://127.0.0.1:1", sbase2) <-sdone2 } // mwServeSeedRelayLists answers the bootstrap REQ with one kind 10002 event. func mwServeSeedRelayLists(fd int32, raw, wantAuthors []byte) { nfd := mwAccept(fd) if nfd < 0 { syscall.Close(fd) return } req := mwReadFrame(nfd) if !mwHasSub(req, wantAuthors) { syscall.Close(nfd) syscall.Close(fd) return } syscall.Write(nfd, mwFrame(ws.OpText, raw)) syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-rl2"))) syscall.Close(nfd) syscall.Close(fd) } func TestCrawlRelayDialFailureAndShortAck(t *testing.T) { out := mwCrawlLog(t) if got := crawlRelay("ws://127.0.0.1:1", out); got != nil { t.Fatal("a refused relay must return no events") } // The local relay acknowledges one of two events and drops: the count is // what was acknowledged, not what was sent. raw := mwEventRes(1, mxHex64, "mw-two", 1700000000, "[]") lfd, lbase := mwListen(t) if lfd == 0 { return } done := spawn(mwServeAckN, lfd, int32(1)) if n := crawlPublishBatch(lbase, [][]byte{raw, raw}, out); n != 1 { t.Fatalf("crawlPublishBatch = %d, want 1", n) } <-done } // mwServeOutboxSkip answers the first REQ with frames that must be skipped // plus one usable follow, then answers the second REQ with a skipped 10002 and // closes. A usable follow keeps the bootstrap (and its real seed relays) out of // the test. func mwServeOutboxSkip(fd int32, pkA string) { nfd := mwAccept(fd) if nfd < 0 { syscall.Close(fd) return } mwReadFrame(nfd) // A non-EVENT frame, a malformed EVENT, a non-kind-3 event, a kind 3 // without tags, then a usable kind 3; then EOSE. syscall.Write(nfd, mwFrame(ws.OpText, []byte(`["NOTICE","hello"]`))) syscall.Write(nfd, mwFrame(ws.OpText, []byte(`["EVENT","s",{"bad":]`))) syscall.Write(nfd, mwFrame(ws.OpText, mwEventRes(1, mxHex64, "", 1, "[]"))) syscall.Write(nfd, mwFrame(ws.OpText, []byte(`["EVENT","s",{"id":"`|mxFill(64, '0')|`","pubkey":"`|mxHex64|`","created_at":1,"kind":3,"content":"","sig":""}]`))) syscall.Write(nfd, mwFrame(ws.OpText, mwEventRes(3, mxHex64, "", 1, `[["p","`|pkA|`"]]`))) syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-k3"))) mwReadFrame(nfd) // A non-kind-10002 event is skipped; no write relay is discovered, so the // follow falls back to the default relay. syscall.Write(nfd, mwFrame(ws.OpText, mwEventRes(1, pkA, "", 1, `[["r","wss://skip.example"]]`))) syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-rl"))) syscall.Close(nfd) syscall.Close(fd) } func TestOutboxDiscoverSkipsBadEvents(t *testing.T) { fd, base := mwListen(t) if fd == 0 { return } pkA := mxFill(64, 'a') done := spawn(mwServeOutboxSkip, fd, pkA) m := outboxDiscover(mxHex64, base) if m == nil { t.Fatal("a usable follow must produce a write relay map") } if !mwSameStrings(m["wss://relay.damus.io"], []string{pkA, mxHex64}) { t.Fatal("the follow and the owner must fall back to the default relay") } if _, ok := m["wss://skip.example"]; ok { t.Fatal("a relay URL on a non-10002 event must be ignored") } <-done } // mwServeEventThenClose sends one event and then a close frame. func mwServeEventThenClose(fd int32, raw []byte) { nfd := mwAccept(fd) if nfd < 0 { syscall.Close(fd) return } mwReadFrame(nfd) syscall.Write(nfd, mwFrame(ws.OpText, raw)) syscall.Write(nfd, mwFrame(ws.OpClose, []byte{0x03, 0xE8})) syscall.Close(nfd) syscall.Close(fd) } func TestSyncOnceRemoteCloseFrame(t *testing.T) { rfd, rbase := mwListen(t) if rfd == 0 { return } lfd, lbase := mwListen(t) if lfd == 0 { syscall.Close(rfd) return } raw := mwEventSub(1, mxHex64, "mw-close", 1700000800, "[]") rdone := spawn(mwServeEventThenClose, rfd, raw) ldone := spawn(mwServeSyncLocal, lfd, uint16(1), []byte("mw-close")) // A close frame ends the loop like the socket close does. if got := syncOnce(rbase, lbase, "", 0); got != 1700000800 { t.Fatalf("close frame = %d, want 1700000800", got) } <-rdone <-ldone } func TestSyncOnceLocalPublishFailure(t *testing.T) { rfd, rbase := mwListen(t) if rfd == 0 { return } lfd, lbase := mwListen(t) if lfd == 0 { syscall.Close(rfd) return } raw := mwEventSub(1, mxHex64, "mw-localfail", 1700000900, "[]") rdone := spawn(mwServeSyncRemote, rfd, raw, []byte(nil), []byte(nil)) ldone := spawn(mwServeCloseAfterUpgrade, lfd) // The local relay is gone before the forward: nothing is acknowledged, so // the checkpoint stays put. if got := syncOnce(rbase, lbase, "", 0); got != 0 { t.Fatalf("local publish failure = %d, want 0", got) } <-rdone <-ldone }