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