client_session_test.mx raw
1 package ws
2
3 import (
4 "syscall"
5 "testing"
6
7 "git.smesh.lol/nostr/pkg/envelope"
8 "git.smesh.lol/nostr/pkg/event"
9 "git.smesh.lol/nostr/pkg/filter"
10 )
11
12 // The client state machine: RunReadLoop against a scripted relay, Publish and
13 // its OK, Unsubscribe, CLOSED, and every early return of dispatch and the four
14 // handlers. The existing tests drive Connect/Subscribe and hand-feed one
15 // EOSE frame; these go through the loop itself.
16
17 func clientEvent(fill byte, content string) (ev *event.E) {
18 ev = event.New()
19 ev.ID = []byte{:32}
20 ev.Pubkey = []byte{:32}
21 ev.Sig = []byte{:64}
22 for i := 0; i < 32; i++ {
23 ev.ID[i] = fill
24 ev.Pubkey[i] = fill + 1
25 }
26 ev.Kind = 1
27 ev.CreatedAt = 1700000000
28 ev.Content = []byte(content)
29 return
30 }
31
32 // serveAfterReq handshakes, consumes the REQ, writes the script, then consumes
33 // the CLOSE the client sends on Close.
34 func serveAfterReq(fd int32, script []byte) {
35 nfd, _, err := syscall.Accept(fd)
36 if err != nil {
37 syscall.Close(fd)
38 return
39 }
40 req := readUpgradeRequest(nfd)
41 syscall.Write(nfd, []byte(upgradeHead(req)))
42 readClientFrame(nfd)
43 syscall.Write(nfd, script)
44 readClientFrame(nfd)
45 syscall.Close(nfd)
46 syscall.Close(fd)
47 }
48
49 // serveAckFirstFrame handshakes, consumes the client's first frame (a publish)
50 // and replies with the given bytes.
51 func serveAckFirstFrame(fd int32, reply []byte) {
52 nfd, _, err := syscall.Accept(fd)
53 if err != nil {
54 syscall.Close(fd)
55 return
56 }
57 req := readUpgradeRequest(nfd)
58 syscall.Write(nfd, []byte(upgradeHead(req)))
59 readClientFrame(nfd)
60 syscall.Write(nfd, reply)
61 readClientFrame(nfd)
62 syscall.Close(nfd)
63 syscall.Close(fd)
64 }
65
66 // serveReqThenClosed handshakes, consumes the REQ and replies with the script.
67 func serveReqThenClosed(fd int32, script []byte) {
68 nfd, _, err := syscall.Accept(fd)
69 if err != nil {
70 syscall.Close(fd)
71 return
72 }
73 req := readUpgradeRequest(nfd)
74 syscall.Write(nfd, []byte(upgradeHead(req)))
75 readClientFrame(nfd)
76 syscall.Write(nfd, script)
77 syscall.Close(nfd)
78 syscall.Close(fd)
79 }
80
81 func TestClientReadLoopDispatch(t *testing.T) {
82 fd, port, err := listenLoopback()
83 if err != nil {
84 t.Fatalf("listen: %v", err)
85 return
86 }
87 ev := clientEvent(0x11, "loop event")
88 payload := []byte("[\"EVENT\",\"s1\",") | ev.Marshal(nil) | []byte("]")
89 eose := &envelope.EOSE{Subscription: []byte("s1")}
90 script := serverFrame(OpText, payload)
91 script = script | serverFrame(OpText, eose.Marshal(nil))
92 script = script | rawFrame(OpClose, 2, nil, []byte{0x03, 0xE8})
93 done := spawn(serveAfterReq, fd, script)
94
95 c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/")
96 if cerr != nil {
97 t.Fatalf("connect: %v", cerr)
98 return
99 }
100 sub, serr := c.Subscribe(&filter.F{})
101 if serr != nil {
102 t.Fatalf("subscribe: %v", serr)
103 c.Close()
104 return
105 }
106 c.RunReadLoop()
107 if c.Err != nil {
108 t.Fatalf("a relay close frame must end the loop cleanly: %v", c.Err)
109 }
110
111 got := <-sub.Events
112 if got == nil || string(got.Content) != "loop event" {
113 t.Fatal("the loop did not deliver the EVENT to the subscription")
114 }
115 if string(got.ID) != string(ev.ID) {
116 t.Fatal("the delivered event is not the one the relay sent")
117 }
118 select {
119 case <-sub.EOSE:
120 default:
121 t.Fatal("the loop did not deliver EOSE")
122 }
123 // A second EOSE for the same subscription is dropped by the eosed guard.
124 c.dispatch(eose.Marshal(nil))
125 c.Close()
126 <-done
127 }
128
129 func TestClientPublishAndOK(t *testing.T) {
130 ev := clientEvent(0x22, "publish me")
131 ack := &envelope.OK{EventID: ev.ID, OK: true, Reason: []byte("saved")}
132 reply := serverFrame(OpText, ack.Marshal(nil))
133
134 fd, port, err := listenLoopback()
135 if err != nil {
136 t.Fatalf("listen: %v", err)
137 return
138 }
139 done := spawn(serveAckFirstFrame, fd, reply)
140 c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/")
141 if cerr != nil {
142 t.Fatalf("connect: %v", cerr)
143 return
144 }
145 if perr := c.Publish(ev); perr != nil {
146 t.Fatalf("publish: %v", perr)
147 return
148 }
149 op, payload, rerr := c.ws.ReadMessage()
150 if rerr != nil {
151 t.Fatalf("read OK: %v", rerr)
152 return
153 }
154 if op != OpText {
155 t.Fatalf("op = %d", int32(op))
156 }
157 c.dispatch(payload)
158 ok := <-c.OKs
159 if ok == nil {
160 t.Fatal("no OK was queued")
161 }
162 if !ok.OK {
163 t.Fatal("the OK says the relay rejected the event")
164 }
165 if string(ok.EventID) != string(ev.ID) {
166 t.Fatal("the OK names a different event id")
167 }
168 if string(ok.Reason) != "saved" {
169 t.Fatalf("reason = %s", ok.Reason)
170 }
171 c.Close()
172 <-done
173 }
174
175 func TestClientUnsubscribe(t *testing.T) {
176 fd, port, err := listenLoopback()
177 if err != nil {
178 t.Fatalf("listen: %v", err)
179 return
180 }
181 done := spawn(serveAfterReq, fd, []byte(nil))
182 c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/")
183 if cerr != nil {
184 t.Fatalf("connect: %v", cerr)
185 return
186 }
187 sub, serr := c.Subscribe(&filter.F{})
188 if serr != nil {
189 t.Fatalf("subscribe: %v", serr)
190 c.Close()
191 return
192 }
193 if len(c.subs) != 1 {
194 t.Fatal("the subscription was not registered")
195 }
196 if uerr := c.Unsubscribe(sub); uerr != nil {
197 t.Fatalf("unsubscribe: %v", uerr)
198 }
199 if len(c.subs) != 0 {
200 t.Fatal("Unsubscribe must drop the subscription")
201 }
202 c.Close()
203 <-done
204 }
205
206 func TestClientClosedRemovesSubscription(t *testing.T) {
207 fd, port, err := listenLoopback()
208 if err != nil {
209 t.Fatalf("listen: %v", err)
210 return
211 }
212 closed := &envelope.Closed{Subscription: []byte("s1"), Reason: []byte("error: too many")}
213 script := serverFrame(OpText, closed.Marshal(nil))
214 done := spawn(serveReqThenClosed, fd, script)
215 c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/")
216 if cerr != nil {
217 t.Fatalf("connect: %v", cerr)
218 return
219 }
220 sub, serr := c.Subscribe(&filter.F{})
221 if serr != nil {
222 t.Fatalf("subscribe: %v", serr)
223 c.Close()
224 return
225 }
226 op, payload, rerr := c.ws.ReadMessage()
227 if rerr != nil {
228 t.Fatalf("read CLOSED: %v", rerr)
229 return
230 }
231 if op != OpText {
232 t.Fatalf("op = %d", int32(op))
233 }
234 c.dispatch(payload)
235 if len(c.subs) != 0 {
236 t.Fatal("CLOSED must remove the subscription")
237 }
238 if _, still := <-sub.Events; still {
239 t.Fatal("CLOSED must close the subscription's event channel")
240 }
241 c.Close()
242 <-done
243 }
244
245 func TestClientDispatchEarlyReturns(t *testing.T) {
246 // No connection: dispatch only touches the subscription map, the OK queue
247 // and the done channel, so malformed input is safe to drive directly.
248 c := &Client{
249 subs: map[string]*Sub{},
250 done: chan struct{}{},
251 OKs: chan *envelope.OK{4},
252 }
253 c.subs["s1"] = &Sub{ID: "s1", Events: chan *event.E{4}, EOSE: chan struct{}{1}}
254
255 bad := [][]byte{
256 nil,
257 []byte(""),
258 []byte("not json"),
259 []byte("[]"),
260 []byte("[\"EVENT\"]"),
261 []byte("[\"EVENT\",\"s1\"]"),
262 []byte("[\"EVENT\",\"s1\",{\"bad\":]"),
263 []byte("[\"EOSE\"]"),
264 []byte("[\"EOSE\",\"nope\"]"),
265 []byte("[\"OK\"]"),
266 []byte("[\"OK\",\"00\"]"),
267 []byte("[\"CLOSED\"]"),
268 []byte("[\"CLOSED\",\"nope\",\"why\"]"),
269 []byte("[\"NOTICE\",\"hello\"]"),
270 }
271 for _, msg := range bad {
272 c.dispatch(msg)
273 }
274 // A well-formed EVENT for an unknown subscription is parsed and dropped.
275 ev := clientEvent(0x33, "orphan")
276 orphan := []byte("[\"EVENT\",\"nope\",") | ev.Marshal(nil) | []byte("]")
277 c.dispatch(orphan)
278 // An OK with the wrong id length is rejected before it reaches the queue.
279 shortID := &envelope.OK{EventID: []byte("short"), OK: true}
280 c.dispatch(shortID.Marshal(nil))
281 select {
282 case <-c.OKs:
283 t.Fatal("malformed messages must not queue an OK")
284 default:
285 }
286 // The CLOSED for an unknown subscription leaves the map alone.
287 if len(c.subs) != 1 {
288 t.Fatal("an unknown CLOSED must not touch other subscriptions")
289 }
290 // The known subscription is still usable.
291 good := &envelope.EOSE{Subscription: []byte("s1")}
292 c.dispatch(good.Marshal(nil))
293 select {
294 case <-c.subs["s1"].EOSE:
295 default:
296 t.Fatal("a well-formed EOSE must still be delivered")
297 }
298 }
299
300 func TestClientReadLoopReportsReadError(t *testing.T) {
301 fd, port, err := listenLoopback()
302 if err != nil {
303 t.Fatalf("listen: %v", err)
304 return
305 }
306 // The peer vanishes without a close frame: the loop must record the read
307 // error, not treat the drop as a relay shutdown. The server consumes the
308 // REQ before closing, so the Subscribe below is not a broken pipe.
309 done := spawn(serveReqThenClosed, fd, []byte(nil))
310 c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/")
311 if cerr != nil {
312 t.Fatalf("connect: %v", cerr)
313 return
314 }
315 if _, serr := c.Subscribe(&filter.F{}); serr != nil {
316 t.Fatalf("subscribe: %v", serr)
317 }
318 c.RunReadLoop()
319 if c.Err == nil {
320 t.Fatal("a dropped connection must set Err")
321 }
322 c.Close()
323 <-done
324 }
325
326 func TestClientConnectRefusesBadURL(t *testing.T) {
327 if _, err := Connect("ftp://127.0.0.1:1/"); err == nil {
328 t.Fatal("Connect must reject a non-ws scheme")
329 }
330 if _, err := Connect("not a url"); err == nil {
331 t.Fatal("Connect must reject an unparseable URL")
332 }
333 }
334
335 func TestClientWriteErrorsOnADeadSocket(t *testing.T) {
336 fd, port, err := listenLoopback()
337 if err != nil {
338 t.Fatalf("listen: %v", err)
339 return
340 }
341 done := spawn(serveScript, fd, []byte(nil))
342 c, cerr := Connect("ws://127.0.0.1:" | portStr(port) | "/")
343 if cerr != nil {
344 t.Fatalf("connect: %v", cerr)
345 return
346 }
347 // Closing the socket directly leaves the subscription map intact, so each
348 // write path reports the failure instead of panicking on a nil map.
349 c.ws.Close()
350 if _, serr := c.Subscribe(&filter.F{}); serr == nil {
351 t.Fatal("Subscribe must report a write on a dead socket")
352 }
353 if len(c.subs) != 0 {
354 t.Fatal("a failed Subscribe must not leave its subscription behind")
355 }
356 sub := &Sub{ID: "gone", Events: chan *event.E{1}, EOSE: chan struct{}{1}}
357 if uerr := c.Unsubscribe(sub); uerr == nil {
358 t.Fatal("Unsubscribe must report a write on a dead socket")
359 }
360 ev := clientEvent(0x44, "dead")
361 if perr := c.Publish(ev); perr == nil {
362 t.Fatal("Publish must report a write on a dead socket")
363 }
364 <-done
365 }
366