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/smesh/pkg/nostr/envelope"
  20  	"git.smesh.lol/smesh/pkg/nostr/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