main_session_test.mx raw
1 // The relay command's networked helpers, driven against loopback WebSocket
2 // servers instead of real relays: syncOnce (two peers: a remote feed and the
3 // local relay it forwards into), outboxDiscover (the two REQs that build the
4 // write-relay map), and the crawler (crawlRelay, crawlPublishBatch,
5 // crawlPass). The remaining main.mx blocks are the subcommands that loop
6 // forever or spawn the binary, which stay integration-only.
7 //
8 // The server side speaks just enough of the protocol: the 101 upgrade with a
9 // correct Sec-WebSocket-Accept, unmasked text frames, and masked client frames
10 // read back. It validates what the client sends and withholds the reply when
11 // the request is wrong, so a passing assertion means the client really sent it.
12 package main
13
14 import (
15 "os"
16 "syscall"
17 "testing"
18
19 "git.smesh.lol/nostr/pkg/envelope"
20 "git.smesh.lol/nostr/pkg/ws"
21 )
22
23 // --- loopback WebSocket server ---
24
25 func mwBind() (fd int32, port int32, err error) {
26 fd, err = syscall.Socket(syscall.AF_INET, syscall.SOCK_STREAM, 0)
27 if err != nil {
28 return 0, 0, err
29 }
30 syscall.SetsockoptInt(fd, syscall.SOL_SOCKET, syscall.SO_REUSEADDR, 1)
31 tv := syscall.Timeval{Sec: 5}
32 syscall.SetsockoptTimeval(fd, syscall.SOL_SOCKET, syscall.SO_RCVTIMEO, &tv)
33 sa := &syscall.SockaddrInet4{Port: 0, Addr: [4]byte{127, 0, 0, 1}}
34 if err = syscall.Bind(fd, sa); err != nil {
35 syscall.Close(fd)
36 return 0, 0, err
37 }
38 if err = syscall.Listen(fd, 8); err != nil {
39 syscall.Close(fd)
40 return 0, 0, err
41 }
42 got, gerr := syscall.Getsockname(fd)
43 if gerr != nil {
44 syscall.Close(fd)
45 return 0, 0, gerr
46 }
47 sa4, ok := got.(*syscall.SockaddrInet4)
48 if !ok {
49 syscall.Close(fd)
50 return 0, 0, syscall.EINVAL
51 }
52 return fd, sa4.Port, nil
53 }
54
55 func mwPortStr(port int32) (s string) {
56 if port == 0 {
57 return "0"
58 }
59 buf := []byte{:6}
60 n := int32(6)
61 for port > 0 {
62 n--
63 buf[n] = byte('0' + port%10)
64 port = port / 10
65 }
66 return string(buf[n:])
67 }
68
69 func mwListen(t *testing.T) (fd int32, base string) {
70 t.Helper()
71 fd, port, err := mwBind()
72 if err != nil {
73 t.Fatalf("listen: %s", err.Error())
74 return 0, ""
75 }
76 return fd, "ws://127.0.0.1:" | mwPortStr(port)
77 }
78
79 // mwHeader returns the value of an HTTP header, tested case-insensitively.
80 func mwHeader(req []byte, name string) (v string) {
81 n := int32(len(name))
82 for i := int32(0); i+n <= int32(len(req)); i++ {
83 if string(req[i:i+n]) == name {
84 j := i + n
85 for j < int32(len(req)) && req[j] != '\r' && req[j] != '\n' {
86 j++
87 }
88 return string(req[i+n : j])
89 }
90 }
91 return ""
92 }
93
94 func mwHasHeadEnd(b []byte) (ok bool) {
95 for i := int32(0); i+3 < int32(len(b)); i++ {
96 if b[i] == '\r' && b[i+1] == '\n' && b[i+2] == '\r' && b[i+3] == '\n' {
97 return true
98 }
99 }
100 return false
101 }
102
103 func mwReadHead(fd int32) (req []byte) {
104 buf := []byte{:2048}
105 for {
106 n, err := syscall.Read(fd, buf)
107 if n <= 0 || err != nil {
108 return req
109 }
110 req = req | buf[:n]
111 if mwHasHeadEnd(req) {
112 return req
113 }
114 }
115 }
116
117 // mwAccept accepts one connection and completes the upgrade. A negative
118 // return means the listener timed out or failed.
119 func mwAccept(fd int32) (nfd int32) {
120 nfd, _, err := syscall.Accept(fd)
121 if err != nil {
122 return -1
123 }
124 req := mwReadHead(nfd)
125 head := "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: " |
126 ws.ComputeAccept(mwHeader(req, "Sec-WebSocket-Key: ")) | "\r\n\r\n"
127 syscall.Write(nfd, []byte(head))
128 return nfd
129 }
130
131 // mwFrame builds an unmasked server frame.
132 func mwFrame(op byte, payload []byte) (buf []byte) {
133 plen := int32(len(payload))
134 if plen < 126 {
135 buf = []byte{:2}
136 buf[1] = byte(plen)
137 } else {
138 buf = []byte{:4}
139 buf[1] = 126
140 buf[2] = byte(plen >> 8)
141 buf[3] = byte(plen)
142 }
143 buf[0] = 0x80 | op
144 return buf | payload
145 }
146
147 func mwReadN(fd int32, n int32) (out []byte) {
148 out = []byte{:0:n}
149 for len(out) < n {
150 buf := []byte{:n - len(out)}
151 got, err := syscall.Read(fd, buf)
152 if got <= 0 || err != nil {
153 return out
154 }
155 out = out | buf[:got]
156 }
157 return
158 }
159
160 // mwReadFrame decodes one masked client frame and unmasks it.
161 func mwReadFrame(fd int32) (payload []byte) {
162 hdr := mwReadN(fd, 2)
163 if len(hdr) < 2 {
164 return nil
165 }
166 plen := int32(hdr[1] & 0x7F)
167 if plen == 126 {
168 ext := mwReadN(fd, 2)
169 if len(ext) < 2 {
170 return nil
171 }
172 plen = int32(ext[0])<<8 | int32(ext[1])
173 } else if plen == 127 {
174 ext := mwReadN(fd, 8)
175 if len(ext) < 8 {
176 return nil
177 }
178 plen = int32(ext[4])<<24 | int32(ext[5])<<16 | int32(ext[6])<<8 | int32(ext[7])
179 }
180 var mask [4]byte
181 if hdr[1]&0x80 != 0 {
182 m := mwReadN(fd, 4)
183 if len(m) < 4 {
184 return nil
185 }
186 mask = [4]byte{m[0], m[1], m[2], m[3]}
187 }
188 body := mwReadN(fd, plen)
189 for i := int32(0); i < int32(len(body)); i++ {
190 body[i] = body[i] ^ mask[i%4]
191 }
192 return body
193 }
194
195 // --- frames the client parses ---
196
197 // mwEventSub builds ["EVENT",{...}], the shape syncOnce's EventSubmission
198 // parser wants.
199 func mwEventSub(kind uint16, pubkeyHex, content string, createdAt int64, tagsJSON string) (raw []byte) {
200 raw = []byte(`["EVENT",{"id":"`) | mxFill(64, '0')
201 raw = raw | `","pubkey":"` | pubkeyHex
202 raw = raw | `","created_at":` | itoa64(createdAt)
203 raw = raw | `,"kind":` | itoa64(int64(kind))
204 raw = raw | `,"tags":` | tagsJSON
205 raw = raw | `,"content":"` | content | `","sig":""}]`
206 return
207 }
208
209 // mwEventRes builds ["EVENT","s",{...}], the shape EventResult wants.
210 func mwEventRes(kind uint16, pubkeyHex, content string, createdAt int64, tagsJSON string) (raw []byte) {
211 raw = []byte(`["EVENT","s",{"id":"`) | mxFill(64, '0')
212 raw = raw | `","pubkey":"` | pubkeyHex
213 raw = raw | `","created_at":` | itoa64(createdAt)
214 raw = raw | `,"kind":` | itoa64(int64(kind))
215 raw = raw | `,"tags":` | tagsJSON
216 raw = raw | `,"content":"` | content | `","sig":""}]`
217 return
218 }
219
220 func mwEOSE(sub string) (raw []byte) {
221 e := &envelope.EOSE{Subscription: []byte(sub)}
222 return e.Marshal(nil)
223 }
224
225 func mwOKFrame(id []byte) (raw []byte) {
226 o := &envelope.OK{EventID: id, OK: true, Reason: []byte("ok")}
227 return o.Marshal(nil)
228 }
229
230 func mwHasSub(hay, needle []byte) (ok bool) {
231 n := int32(len(needle))
232 if n == 0 || int32(len(hay)) < n {
233 return false
234 }
235 for i := int32(0); i+n <= int32(len(hay)); i++ {
236 if string(hay[i:i+n]) == string(needle) {
237 return true
238 }
239 }
240 return false
241 }
242
243 // mwParseSub parses an EVENT submission and returns its event.
244 func mwParseSub(raw []byte) (e *envelope.EventSubmission) {
245 _, rem, err := envelope.Identify(raw)
246 if err != nil {
247 return nil
248 }
249 var es envelope.EventSubmission
250 if _, uerr := es.Unmarshal(rem); uerr != nil || es.E == nil {
251 return nil
252 }
253 return &es
254 }
255
256 // --- scripted peers ---
257
258 // mwServeSyncRemote answers a sync subscription: one event, EOSE, then close.
259 // When wantA/wantB are non-empty the REQ must contain both before the event is
260 // sent, so the test can assert the filter the client built.
261 func mwServeSyncRemote(fd int32, raw []byte, wantA, wantB []byte) {
262 nfd := mwAccept(fd)
263 if nfd < 0 {
264 syscall.Close(fd)
265 return
266 }
267 req := mwReadFrame(nfd)
268 if len(wantA) > 0 && !mwHasSub(req, wantA) {
269 syscall.Close(nfd)
270 syscall.Close(fd)
271 return
272 }
273 if len(wantB) > 0 && !mwHasSub(req, wantB) {
274 syscall.Close(nfd)
275 syscall.Close(fd)
276 return
277 }
278 syscall.Write(nfd, mwFrame(ws.OpText, raw))
279 syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("sync")))
280 syscall.Close(nfd)
281 syscall.Close(fd)
282 }
283
284 // mwServeSyncLocal consumes the forwarded event and acknowledges it. The OK is
285 // only sent for the expected kind and content, so a client that forwards
286 // nothing (or the wrong thing) fails its drain read.
287 func mwServeSyncLocal(fd int32, wantKind uint16, wantContent []byte) {
288 nfd := mwAccept(fd)
289 if nfd < 0 {
290 syscall.Close(fd)
291 return
292 }
293 raw := mwReadFrame(nfd)
294 es := mwParseSub(raw)
295 if es != nil && es.E.Kind == wantKind && string(es.E.Content) == string(wantContent) {
296 syscall.Write(nfd, mwFrame(ws.OpText, mwOKFrame(es.E.ID)))
297 }
298 syscall.Close(nfd)
299 syscall.Close(fd)
300 }
301
302 // mwServeOutbox answers the two REQs outboxDiscover sends.
303 func mwServeOutbox(fd int32, k3raw, rlraw []byte) {
304 nfd := mwAccept(fd)
305 if nfd < 0 {
306 syscall.Close(fd)
307 return
308 }
309 mwReadFrame(nfd)
310 syscall.Write(nfd, mwFrame(ws.OpText, k3raw))
311 syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-k3")))
312 mwReadFrame(nfd)
313 if rlraw != nil {
314 syscall.Write(nfd, mwFrame(ws.OpText, rlraw))
315 }
316 syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-rl")))
317 syscall.Close(nfd)
318 syscall.Close(fd)
319 }
320
321 // mwServeIndex answers count REQs with one event each and closes.
322 func mwServeIndex(fd int32, raw []byte, count int32) {
323 nfd := mwAccept(fd)
324 if nfd < 0 {
325 syscall.Close(fd)
326 return
327 }
328 for i := int32(0); i < count; i++ {
329 mwReadFrame(nfd)
330 syscall.Write(nfd, mwFrame(ws.OpText, raw))
331 syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("cr")))
332 }
333 syscall.Close(nfd)
334 syscall.Close(fd)
335 }
336
337 // mwServeAckN reads count published events and acknowledges each.
338 func mwServeAckN(fd int32, count int32) {
339 nfd := mwAccept(fd)
340 if nfd < 0 {
341 syscall.Close(fd)
342 return
343 }
344 ids := [][]byte{:count}
345 for i := int32(0); i < count; i++ {
346 raw := mwReadFrame(nfd)
347 es := mwParseSub(raw)
348 if es != nil {
349 ids[i] = es.E.ID
350 }
351 }
352 for i := int32(0); i < count; i++ {
353 syscall.Write(nfd, mwFrame(ws.OpText, mwOKFrame(ids[i])))
354 }
355 syscall.Close(nfd)
356 syscall.Close(fd)
357 }
358
359 // mwCrawlLog opens the throwaway log the crawler writes to.
360 func mwCrawlLog(t *testing.T) (out *os.File) {
361 t.Helper()
362 f, err := os.OpenFile("/dev/null", os.O_WRONLY, 0)
363 if err != nil {
364 return os.Stderr
365 }
366 return f
367 }
368
369 func mwSameStrings(a, b []string) (ok bool) {
370 if len(a) != len(b) {
371 return false
372 }
373 for i := int32(0); i < int32(len(a)); i++ {
374 if a[i] != b[i] {
375 return false
376 }
377 }
378 return true
379 }
380
381 // --- tests ---
382
383 func TestSyncOnceForwardsRemoteEvents(t *testing.T) {
384 rfd, rbase := mwListen(t)
385 if rfd == 0 {
386 return
387 }
388 lfd, lbase := mwListen(t)
389 if lfd == 0 {
390 syscall.Close(rfd)
391 return
392 }
393 raw := mwEventSub(1, mxHex64, "mw-sync", 1700000500, "[]")
394 rdone := spawn(mwServeSyncRemote, rfd, raw, []byte(nil), []byte(nil))
395 ldone := spawn(mwServeSyncLocal, lfd, uint16(1), []byte("mw-sync"))
396
397 got := syncOnce(rbase, lbase, "", 0)
398 if got != 1700000500 {
399 t.Fatalf("syncOnce = %d, want the forwarded event's timestamp", got)
400 }
401 <-rdone
402 <-ldone
403 }
404
405 func TestSyncOnceBuildsTheAuthorAndSinceFilter(t *testing.T) {
406 rfd, rbase := mwListen(t)
407 if rfd == 0 {
408 return
409 }
410 lfd, lbase := mwListen(t)
411 if lfd == 0 {
412 syscall.Close(rfd)
413 return
414 }
415 raw := mwEventSub(1, mxHex64, "mw-filter", 1700000600, "[]")
416 // sinceTs is 1700003600, so the filter carries since=1700000000; the
417 // authors list keeps only the 64-hex keys.
418 wantAuthors := []byte(`"authors":["` | mxHex64 | `","` | mxHex64 | `"]`)
419 wantSince := []byte(`"since":1700000000`)
420 rdone := spawn(mwServeSyncRemote, rfd, raw, wantAuthors, wantSince)
421 ldone := spawn(mwServeSyncLocal, lfd, uint16(1), []byte("mw-filter"))
422
423 got := syncOnce(rbase, lbase, mxHex64|",nothex,"|mxHex64, 1700003600)
424 if got != 1700000600 {
425 t.Fatalf("syncOnce = %d, want 1700000600", got)
426 }
427 <-rdone
428 <-ldone
429 }
430
431 func TestSyncOnceDialFailures(t *testing.T) {
432 // The remote is refused: the checkpoint is returned unchanged.
433 if got := syncOnce("ws://127.0.0.1:1", "ws://127.0.0.1:1", "", 42); got != 42 {
434 t.Fatalf("remote dial failure = %d, want 42", got)
435 }
436 // The remote answers but the local relay is refused.
437 rfd, rbase := mwListen(t)
438 if rfd == 0 {
439 return
440 }
441 rdone := spawn(mwServeSyncRemote, rfd, mwEventSub(1, mxHex64, "x", 1, "[]"), []byte(nil), []byte(nil))
442 if got := syncOnce(rbase, "ws://127.0.0.1:1", "", 42); got != 42 {
443 t.Fatalf("local dial failure = %d, want 42", got)
444 }
445 <-rdone
446 }
447
448 func TestOutboxDiscoverBuildsWriteRelayMap(t *testing.T) {
449 fd, base := mwListen(t)
450 if fd == 0 {
451 return
452 }
453 pkA := mxFill(64, 'a')
454 pkB := mxFill(64, 'b')
455 owner := mxFill(64, 'c')
456 k3 := mwEventRes(3, owner, "", 1700000000,
457 `[["p","`|pkA|`"],["p","`|pkB|`"],["p","short"]]`)
458 rl := mwEventRes(10002, pkA, "", 1700000000,
459 `[["r","wss://a.example"],["r","wss://b.example","write"],["r","wss://c.example","read"],["r","http://bad.example"],["r"]]`)
460 done := spawn(mwServeOutbox, fd, k3, rl)
461
462 m := outboxDiscover(owner, base)
463 if m == nil {
464 t.Fatal("outboxDiscover returned nil")
465 }
466 if len(m) != 3 {
467 t.Fatalf("write relay buckets = %d, want 3", int32(len(m)))
468 }
469 if !mwSameStrings(m["wss://a.example"], []string{pkA}) {
470 t.Fatal("an unmarked r tag must be a write relay")
471 }
472 if !mwSameStrings(m["wss://b.example"], []string{pkA}) {
473 t.Fatal("a write-marked r tag must be a write relay")
474 }
475 if _, ok := m["wss://c.example"]; ok {
476 t.Fatal("a read-marked r tag must not be a write relay")
477 }
478 if _, ok := m["http://bad.example"]; ok {
479 t.Fatal("a non-ws relay URL must be ignored")
480 }
481 if !mwSameStrings(m["wss://relay.damus.io"], []string{pkB, owner}) {
482 t.Fatal("follows without a kind 10002 (and the owner) must fall back to the default relay")
483 }
484 <-done
485 }
486
487 func TestOutboxDiscoverDialError(t *testing.T) {
488 if m := outboxDiscover(mxHex64, "ws://127.0.0.1:1"); m != nil {
489 t.Fatal("a refused local relay must return nil")
490 }
491 }
492
493 func TestCrawlRelayAndPublish(t *testing.T) {
494 out := mwCrawlLog(t)
495 raw := mwEventRes(1, mxHex64, "mw-crawl", 1700000000, "[]")
496
497 // One relay: crawlRelay returns the event it saw.
498 rfd, rbase := mwListen(t)
499 if rfd == 0 {
500 return
501 }
502 rdone := spawn(mwServeIndex, rfd, raw, int32(1))
503 events := crawlRelay(rbase, out)
504 if len(events) != 1 {
505 t.Fatalf("crawlRelay returned %d events", int32(len(events)))
506 }
507 if es := mwParseSub(events[0]); es == nil || string(es.E.Content) != "mw-crawl" {
508 t.Fatal("crawlRelay must return the event as an EVENT submission")
509 }
510 <-rdone
511
512 // crawlPublishBatch forwards them to the local relay and counts the OKs.
513 lfd, lbase := mwListen(t)
514 if lfd == 0 {
515 return
516 }
517 ldone := spawn(mwServeAckN, lfd, int32(1))
518 if n := crawlPublishBatch(lbase, events, out); n != 1 {
519 t.Fatalf("crawlPublishBatch published %d, want 1", n)
520 }
521 <-ldone
522
523 // A full pass over one live relay and one dead one.
524 rfd2, rbase2 := mwListen(t)
525 if rfd2 == 0 {
526 return
527 }
528 rdone2 := spawn(mwServeIndex, rfd2, raw, int32(1))
529 lfd2, lbase2 := mwListen(t)
530 if lfd2 == 0 {
531 return
532 }
533 ldone2 := spawn(mwServeAckN, lfd2, int32(1))
534 db := newRelayDB()
535 db.add("ws://127.0.0.1:1", 5)
536 db.add(rbase2, 100)
537 if !crawlPass(lbase2, db, out) {
538 t.Fatal("crawlPass must report a completed pass")
539 }
540 <-rdone2
541 <-ldone2
542
543 // No relays at all is a failed pass.
544 if crawlPass(lbase2, newRelayDB(), out) {
545 t.Fatal("an empty relay database must fail the pass")
546 }
547 }
548
549 // mwServeCloseAfterUpgrade accepts the upgrade and closes, so the client's
550 // first write fails.
551 func mwServeCloseAfterUpgrade(fd int32) {
552 nfd := mwAccept(fd)
553 if nfd >= 0 {
554 syscall.Close(nfd)
555 }
556 syscall.Close(fd)
557 }
558
559 // mwServeDropAck reads one forwarded event and closes without acknowledging.
560 func mwServeDropAck(fd int32) {
561 nfd := mwAccept(fd)
562 if nfd < 0 {
563 syscall.Close(fd)
564 return
565 }
566 mwReadFrame(nfd)
567 syscall.Close(nfd)
568 syscall.Close(fd)
569 }
570
571 func TestSyncOnceWriteAndDrainFailures(t *testing.T) {
572 // The remote closes right after the upgrade: the subscribe write fails and
573 // the checkpoint is returned unchanged.
574 rfd, rbase := mwListen(t)
575 if rfd == 0 {
576 return
577 }
578 rdone := spawn(mwServeCloseAfterUpgrade, rfd)
579 if got := syncOnce(rbase, "ws://127.0.0.1:1", "", 7); got != 7 {
580 t.Fatalf("subscribe failure = %d, want 7", got)
581 }
582 <-rdone
583
584 // The local relay drops the connection instead of acknowledging: the
585 // forwarded event's timestamp is still returned.
586 rfd2, rbase2 := mwListen(t)
587 if rfd2 == 0 {
588 return
589 }
590 lfd, lbase := mwListen(t)
591 if lfd == 0 {
592 syscall.Close(rfd2)
593 return
594 }
595 raw := mwEventSub(1, mxHex64, "mw-drop", 1700000700, "[]")
596 rdone2 := spawn(mwServeSyncRemote, rfd2, raw, []byte(nil), []byte(nil))
597 ldone := spawn(mwServeDropAck, lfd)
598 // The event was forwarded but never acknowledged, so the checkpoint does
599 // not advance and the next run re-sends it.
600 if got := syncOnce(rbase2, lbase, "", 0); got != 0 {
601 t.Fatalf("drain failure = %d, want the unadvanced checkpoint 0", got)
602 }
603 <-rdone2
604 <-ldone
605 }
606
607 func TestOutboxDiscoverWriteError(t *testing.T) {
608 fd, base := mwListen(t)
609 if fd == 0 {
610 return
611 }
612 done := spawn(mwServeCloseAfterUpgrade, fd)
613 if m := outboxDiscover(mxHex64, base); m != nil {
614 t.Fatal("a REQ that cannot be written must return nil")
615 }
616 <-done
617 }
618
619 func TestOutboxBootstrapRelayLists(t *testing.T) {
620 sfd, sbase := mwListen(t)
621 if sfd == 0 {
622 return
623 }
624 lfd, lbase := mwListen(t)
625 if lfd == 0 {
626 syscall.Close(sfd)
627 return
628 }
629 pkA := mxFill(64, 'a')
630 pkB := mxFill(64, 'b')
631 raw := mwEventRes(10002, pkA, "", 1700000000, `[["r","wss://a.example"]]`)
632 want := []byte(`"authors":["` | pkA | `","` | pkB | `"]`)
633 sdone := spawn(mwServeSeedRelayLists, sfd, raw, want)
634 ldone := spawn(mwServeAckN, lfd, int32(1))
635
636 outboxBootstrapRelayLists([]string{pkA, pkB}, lbase, sbase)
637 <-sdone
638 <-ldone
639
640 // No follows is a no-op; a dead seed returns quietly before dialing local.
641 outboxBootstrapRelayLists([]string{}, lbase, sbase)
642 outboxBootstrapRelayLists([]string{pkA}, lbase, "ws://127.0.0.1:1")
643
644 // A dead local relay returns quietly too, after the seed has accepted.
645 sfd2, sbase2 := mwListen(t)
646 if sfd2 == 0 {
647 return
648 }
649 sdone2 := spawn(mwServeSeedRelayLists, sfd2, raw, want)
650 outboxBootstrapRelayLists([]string{pkA}, "ws://127.0.0.1:1", sbase2)
651 <-sdone2
652 }
653
654 // mwServeSeedRelayLists answers the bootstrap REQ with one kind 10002 event.
655 func mwServeSeedRelayLists(fd int32, raw, wantAuthors []byte) {
656 nfd := mwAccept(fd)
657 if nfd < 0 {
658 syscall.Close(fd)
659 return
660 }
661 req := mwReadFrame(nfd)
662 if !mwHasSub(req, wantAuthors) {
663 syscall.Close(nfd)
664 syscall.Close(fd)
665 return
666 }
667 syscall.Write(nfd, mwFrame(ws.OpText, raw))
668 syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-rl2")))
669 syscall.Close(nfd)
670 syscall.Close(fd)
671 }
672
673 func TestCrawlRelayDialFailureAndShortAck(t *testing.T) {
674 out := mwCrawlLog(t)
675 if got := crawlRelay("ws://127.0.0.1:1", out); got != nil {
676 t.Fatal("a refused relay must return no events")
677 }
678 // The local relay acknowledges one of two events and drops: the count is
679 // what was acknowledged, not what was sent.
680 raw := mwEventRes(1, mxHex64, "mw-two", 1700000000, "[]")
681 lfd, lbase := mwListen(t)
682 if lfd == 0 {
683 return
684 }
685 done := spawn(mwServeAckN, lfd, int32(1))
686 if n := crawlPublishBatch(lbase, [][]byte{raw, raw}, out); n != 1 {
687 t.Fatalf("crawlPublishBatch = %d, want 1", n)
688 }
689 <-done
690 }
691
692 // mwServeOutboxSkip answers the first REQ with frames that must be skipped
693 // plus one usable follow, then answers the second REQ with a skipped 10002 and
694 // closes. A usable follow keeps the bootstrap (and its real seed relays) out of
695 // the test.
696 func mwServeOutboxSkip(fd int32, pkA string) {
697 nfd := mwAccept(fd)
698 if nfd < 0 {
699 syscall.Close(fd)
700 return
701 }
702 mwReadFrame(nfd)
703 // A non-EVENT frame, a malformed EVENT, a non-kind-3 event, a kind 3
704 // without tags, then a usable kind 3; then EOSE.
705 syscall.Write(nfd, mwFrame(ws.OpText, []byte(`["NOTICE","hello"]`)))
706 syscall.Write(nfd, mwFrame(ws.OpText, []byte(`["EVENT","s",{"bad":]`)))
707 syscall.Write(nfd, mwFrame(ws.OpText, mwEventRes(1, mxHex64, "", 1, "[]")))
708 syscall.Write(nfd, mwFrame(ws.OpText, []byte(`["EVENT","s",{"id":"`|mxFill(64, '0')|`","pubkey":"`|mxHex64|`","created_at":1,"kind":3,"content":"","sig":""}]`)))
709 syscall.Write(nfd, mwFrame(ws.OpText, mwEventRes(3, mxHex64, "", 1, `[["p","`|pkA|`"]]`)))
710 syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-k3")))
711 mwReadFrame(nfd)
712 // A non-kind-10002 event is skipped; no write relay is discovered, so the
713 // follow falls back to the default relay.
714 syscall.Write(nfd, mwFrame(ws.OpText, mwEventRes(1, pkA, "", 1, `[["r","wss://skip.example"]]`)))
715 syscall.Write(nfd, mwFrame(ws.OpText, mwEOSE("ob-rl")))
716 syscall.Close(nfd)
717 syscall.Close(fd)
718 }
719
720 func TestOutboxDiscoverSkipsBadEvents(t *testing.T) {
721 fd, base := mwListen(t)
722 if fd == 0 {
723 return
724 }
725 pkA := mxFill(64, 'a')
726 done := spawn(mwServeOutboxSkip, fd, pkA)
727 m := outboxDiscover(mxHex64, base)
728 if m == nil {
729 t.Fatal("a usable follow must produce a write relay map")
730 }
731 if !mwSameStrings(m["wss://relay.damus.io"], []string{pkA, mxHex64}) {
732 t.Fatal("the follow and the owner must fall back to the default relay")
733 }
734 if _, ok := m["wss://skip.example"]; ok {
735 t.Fatal("a relay URL on a non-10002 event must be ignored")
736 }
737 <-done
738 }
739
740 // mwServeEventThenClose sends one event and then a close frame.
741 func mwServeEventThenClose(fd int32, raw []byte) {
742 nfd := mwAccept(fd)
743 if nfd < 0 {
744 syscall.Close(fd)
745 return
746 }
747 mwReadFrame(nfd)
748 syscall.Write(nfd, mwFrame(ws.OpText, raw))
749 syscall.Write(nfd, mwFrame(ws.OpClose, []byte{0x03, 0xE8}))
750 syscall.Close(nfd)
751 syscall.Close(fd)
752 }
753
754 func TestSyncOnceRemoteCloseFrame(t *testing.T) {
755 rfd, rbase := mwListen(t)
756 if rfd == 0 {
757 return
758 }
759 lfd, lbase := mwListen(t)
760 if lfd == 0 {
761 syscall.Close(rfd)
762 return
763 }
764 raw := mwEventSub(1, mxHex64, "mw-close", 1700000800, "[]")
765 rdone := spawn(mwServeEventThenClose, rfd, raw)
766 ldone := spawn(mwServeSyncLocal, lfd, uint16(1), []byte("mw-close"))
767 // A close frame ends the loop like the socket close does.
768 if got := syncOnce(rbase, lbase, "", 0); got != 1700000800 {
769 t.Fatalf("close frame = %d, want 1700000800", got)
770 }
771 <-rdone
772 <-ldone
773 }
774
775 func TestSyncOnceLocalPublishFailure(t *testing.T) {
776 rfd, rbase := mwListen(t)
777 if rfd == 0 {
778 return
779 }
780 lfd, lbase := mwListen(t)
781 if lfd == 0 {
782 syscall.Close(rfd)
783 return
784 }
785 raw := mwEventSub(1, mxHex64, "mw-localfail", 1700000900, "[]")
786 rdone := spawn(mwServeSyncRemote, rfd, raw, []byte(nil), []byte(nil))
787 ldone := spawn(mwServeCloseAfterUpgrade, lfd)
788 // The local relay is gone before the forward: nothing is acknowledged, so
789 // the checkpoint stays put.
790 if got := syncOnce(rbase, lbase, "", 0); got != 0 {
791 t.Fatalf("local publish failure = %d, want 0", got)
792 }
793 <-rdone
794 <-ldone
795 }
796