package ws import ( "syscall" "testing" "git.smesh.lol/smesh/pkg/nostr/envelope" "git.smesh.lol/smesh/pkg/nostr/event" "git.smesh.lol/smesh/pkg/nostr/filter" ) // The client state machine: RunReadLoop against a scripted relay, Publish and // its OK, Unsubscribe, CLOSED, and every early return of dispatch and the four // handlers. The existing tests drive Connect/Subscribe and hand-feed one // EOSE frame; these go through the loop itself. func clientEvent(fill byte, content string) (ev *event.E) { ev = event.New() ev.ID = []byte{:32} ev.Pubkey = []byte{:32} ev.Sig = []byte{:64} for i := 0; i < 32; i++ { ev.ID[i] = fill ev.Pubkey[i] = fill + 1 } ev.Kind = 1 ev.CreatedAt = 1700000000 ev.Content = []byte(content) return } // serveAfterReq handshakes, consumes the REQ, writes the script, then consumes // the CLOSE the client sends on Close. func serveAfterReq(fd int32, script []byte) { nfd, _, err := syscall.Accept(fd) if err != nil { syscall.Close(fd) return } req := readUpgradeRequest(nfd) syscall.Write(nfd, []byte(upgradeHead(req))) readClientFrame(nfd) syscall.Write(nfd, script) readClientFrame(nfd) syscall.Close(nfd) syscall.Close(fd) } // serveAckFirstFrame handshakes, consumes the client's first frame (a publish) // and replies with the given bytes. func serveAckFirstFrame(fd int32, reply []byte) { nfd, _, err := syscall.Accept(fd) if err != nil { syscall.Close(fd) return } req := readUpgradeRequest(nfd) syscall.Write(nfd, []byte(upgradeHead(req))) readClientFrame(nfd) syscall.Write(nfd, reply) readClientFrame(nfd) syscall.Close(nfd) syscall.Close(fd) } // serveReqThenClosed handshakes, consumes the REQ and replies with the script. func serveReqThenClosed(fd int32, script []byte) { nfd, _, err := syscall.Accept(fd) if err != nil { syscall.Close(fd) return } req := readUpgradeRequest(nfd) syscall.Write(nfd, []byte(upgradeHead(req))) readClientFrame(nfd) syscall.Write(nfd, script) syscall.Close(nfd) syscall.Close(fd) } func TestClientReadLoopDispatch(t *testing.T) { fd, port, err := listenLoopback() if err != nil { t.Fatalf("listen: %v", err) return } ev := clientEvent(0x11, "loop event") payload := []byte("[\"EVENT\",\"s1\",") | ev.Marshal(nil) | []byte("]") eose := &envelope.EOSE{Subscription: []byte("s1")} script := serverFrame(OpText, payload) script = script | serverFrame(OpText, eose.Marshal(nil)) script = script | rawFrame(OpClose, 2, nil, []byte{0x03, 0xE8}) done := spawn(serveAfterReq, fd, script) c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/") if cerr != nil { t.Fatalf("connect: %v", cerr) return } sub, serr := c.Subscribe(&filter.F{}) if serr != nil { t.Fatalf("subscribe: %v", serr) c.Close() return } c.RunReadLoop() if c.Err != nil { t.Fatalf("a relay close frame must end the loop cleanly: %v", c.Err) } got := <-sub.Events if got == nil || string(got.Content) != "loop event" { t.Fatal("the loop did not deliver the EVENT to the subscription") } if string(got.ID) != string(ev.ID) { t.Fatal("the delivered event is not the one the relay sent") } select { case <-sub.EOSE: default: t.Fatal("the loop did not deliver EOSE") } // A second EOSE for the same subscription is dropped by the eosed guard. c.dispatch(eose.Marshal(nil)) c.Close() <-done } func TestClientPublishAndOK(t *testing.T) { ev := clientEvent(0x22, "publish me") ack := &envelope.OK{EventID: ev.ID, OK: true, Reason: []byte("saved")} reply := serverFrame(OpText, ack.Marshal(nil)) fd, port, err := listenLoopback() if err != nil { t.Fatalf("listen: %v", err) return } done := spawn(serveAckFirstFrame, fd, reply) c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/") if cerr != nil { t.Fatalf("connect: %v", cerr) return } if perr := c.Publish(ev); perr != nil { t.Fatalf("publish: %v", perr) return } op, payload, rerr := c.ws.ReadMessage() if rerr != nil { t.Fatalf("read OK: %v", rerr) return } if op != OpText { t.Fatalf("op = %d", int32(op)) } c.dispatch(payload) ok := <-c.OKs if ok == nil { t.Fatal("no OK was queued") } if !ok.OK { t.Fatal("the OK says the relay rejected the event") } if string(ok.EventID) != string(ev.ID) { t.Fatal("the OK names a different event id") } if string(ok.Reason) != "saved" { t.Fatalf("reason = %s", ok.Reason) } c.Close() <-done } func TestClientUnsubscribe(t *testing.T) { fd, port, err := listenLoopback() if err != nil { t.Fatalf("listen: %v", err) return } done := spawn(serveAfterReq, fd, []byte(nil)) c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/") if cerr != nil { t.Fatalf("connect: %v", cerr) return } sub, serr := c.Subscribe(&filter.F{}) if serr != nil { t.Fatalf("subscribe: %v", serr) c.Close() return } if len(c.subs) != 1 { t.Fatal("the subscription was not registered") } if uerr := c.Unsubscribe(sub); uerr != nil { t.Fatalf("unsubscribe: %v", uerr) } if len(c.subs) != 0 { t.Fatal("Unsubscribe must drop the subscription") } c.Close() <-done } func TestClientClosedRemovesSubscription(t *testing.T) { fd, port, err := listenLoopback() if err != nil { t.Fatalf("listen: %v", err) return } closed := &envelope.Closed{Subscription: []byte("s1"), Reason: []byte("error: too many")} script := serverFrame(OpText, closed.Marshal(nil)) done := spawn(serveReqThenClosed, fd, script) c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/") if cerr != nil { t.Fatalf("connect: %v", cerr) return } sub, serr := c.Subscribe(&filter.F{}) if serr != nil { t.Fatalf("subscribe: %v", serr) c.Close() return } op, payload, rerr := c.ws.ReadMessage() if rerr != nil { t.Fatalf("read CLOSED: %v", rerr) return } if op != OpText { t.Fatalf("op = %d", int32(op)) } c.dispatch(payload) if len(c.subs) != 0 { t.Fatal("CLOSED must remove the subscription") } if _, still := <-sub.Events; still { t.Fatal("CLOSED must close the subscription's event channel") } c.Close() <-done } func TestClientDispatchEarlyReturns(t *testing.T) { // No connection: dispatch only touches the subscription map, the OK queue // and the done channel, so malformed input is safe to drive directly. c := &Client{ subs: map[string]*Sub{}, done: chan struct{}{}, OKs: chan *envelope.OK{4}, } c.subs["s1"] = &Sub{ID: "s1", Events: chan *event.E{4}, EOSE: chan struct{}{1}} bad := [][]byte{ nil, []byte(""), []byte("not json"), []byte("[]"), []byte("[\"EVENT\"]"), []byte("[\"EVENT\",\"s1\"]"), []byte("[\"EVENT\",\"s1\",{\"bad\":]"), []byte("[\"EOSE\"]"), []byte("[\"EOSE\",\"nope\"]"), []byte("[\"OK\"]"), []byte("[\"OK\",\"00\"]"), []byte("[\"CLOSED\"]"), []byte("[\"CLOSED\",\"nope\",\"why\"]"), []byte("[\"NOTICE\",\"hello\"]"), } for _, msg := range bad { c.dispatch(msg) } // A well-formed EVENT for an unknown subscription is parsed and dropped. ev := clientEvent(0x33, "orphan") orphan := []byte("[\"EVENT\",\"nope\",") | ev.Marshal(nil) | []byte("]") c.dispatch(orphan) // An OK with the wrong id length is rejected before it reaches the queue. shortID := &envelope.OK{EventID: []byte("short"), OK: true} c.dispatch(shortID.Marshal(nil)) select { case <-c.OKs: t.Fatal("malformed messages must not queue an OK") default: } // The CLOSED for an unknown subscription leaves the map alone. if len(c.subs) != 1 { t.Fatal("an unknown CLOSED must not touch other subscriptions") } // The known subscription is still usable. good := &envelope.EOSE{Subscription: []byte("s1")} c.dispatch(good.Marshal(nil)) select { case <-c.subs["s1"].EOSE: default: t.Fatal("a well-formed EOSE must still be delivered") } } func TestClientReadLoopReportsReadError(t *testing.T) { fd, port, err := listenLoopback() if err != nil { t.Fatalf("listen: %v", err) return } // The peer vanishes without a close frame: the loop must record the read // error, not treat the drop as a relay shutdown. The server consumes the // REQ before closing, so the Subscribe below is not a broken pipe. done := spawn(serveReqThenClosed, fd, []byte(nil)) c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/") if cerr != nil { t.Fatalf("connect: %v", cerr) return } if _, serr := c.Subscribe(&filter.F{}); serr != nil { t.Fatalf("subscribe: %v", serr) } c.RunReadLoop() if c.Err == nil { t.Fatal("a dropped connection must set Err") } c.Close() <-done } func TestClientConnectRefusesBadURL(t *testing.T) { if _, err := Connect("ftp://127.0.0.1:1/"); err == nil { t.Fatal("Connect must reject a non-ws scheme") } if _, err := Connect("not a url"); err == nil { t.Fatal("Connect must reject an unparseable URL") } } func TestClientWriteErrorsOnADeadSocket(t *testing.T) { fd, port, err := listenLoopback() if err != nil { t.Fatalf("listen: %v", err) return } done := spawn(serveScript, fd, []byte(nil)) c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/") if cerr != nil { t.Fatalf("connect: %v", cerr) return } // Closing the socket directly leaves the subscription map intact, so each // write path reports the failure instead of panicking on a nil map. c.ws.Close() if _, serr := c.Subscribe(&filter.F{}); serr == nil { t.Fatal("Subscribe must report a write on a dead socket") } if len(c.subs) != 0 { t.Fatal("a failed Subscribe must not leave its subscription behind") } sub := &Sub{ID: "gone", Events: chan *event.E{1}, EOSE: chan struct{}{1}} if uerr := c.Unsubscribe(sub); uerr == nil { t.Fatal("Unsubscribe must report a write on a dead socket") } ev := clientEvent(0x44, "dead") if perr := c.Publish(ev); perr == nil { t.Fatal("Publish must report a write on a dead socket") } <-done }