relay_test.go raw

   1  package relay
   2  
   3  import (
   4  	"context"
   5  	"encoding/json"
   6  	"net/http"
   7  	"runtime"
   8  	"strings"
   9  	"testing"
  10  	"time"
  11  
  12  	"git.mleku.dev/mleku/dendrite/pkg/nostr"
  13  )
  14  
  15  // TestPublishAndSubscribe verifies the basic relay contract:
  16  // publish events, subscribe, receive them back.
  17  func TestPublishAndSubscribe(t *testing.T) {
  18  	r := New("test-relay")
  19  	addr, err := r.Listen("127.0.0.1:0")
  20  	if err != nil {
  21  		t.Fatal(err)
  22  	}
  23  	defer r.Shutdown(context.Background())
  24  
  25  	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
  26  	defer cancel()
  27  
  28  	// Create identity and compose events.
  29  	id, err := nostr.NewIdentity()
  30  	if err != nil {
  31  		t.Fatal(err)
  32  	}
  33  
  34  	meta, err := id.ComposeMetadata(0, 1, "deadbeef")
  35  	if err != nil {
  36  		t.Fatal(err)
  37  	}
  38  	status, err := id.ComposeStatus(0, 1, 2922, 980, 0.335, 0.284, 3)
  39  	if err != nil {
  40  		t.Fatal(err)
  41  	}
  42  
  43  	// Connect and publish.
  44  	client, err := nostr.Connect(ctx, "ws://"+addr)
  45  	if err != nil {
  46  		t.Fatal(err)
  47  	}
  48  	defer client.Disconnect()
  49  	go client.Listen(ctx)
  50  
  51  	for _, ev := range []*nostr.Event{meta, status} {
  52  		if err := client.Publish(ctx, ev); err != nil {
  53  			t.Fatalf("publish: %v", err)
  54  		}
  55  	}
  56  
  57  	// Wait for OK responses.
  58  	for i := 0; i < 2; i++ {
  59  		select {
  60  		case ok := <-client.OKs:
  61  			if !ok.Accepted {
  62  				t.Fatalf("relay rejected event %s: %s", ok.EventID[:16], ok.Message)
  63  			}
  64  		case <-time.After(2 * time.Second):
  65  			t.Fatal("timeout waiting for OK")
  66  		}
  67  	}
  68  
  69  	if r.EventCount() != 2 {
  70  		t.Fatalf("expected 2 stored events, got %d", r.EventCount())
  71  	}
  72  
  73  	// Subscribe from a second client.
  74  	sub, err := nostr.Connect(ctx, "ws://"+addr)
  75  	if err != nil {
  76  		t.Fatal(err)
  77  	}
  78  	defer sub.Disconnect()
  79  	go sub.Listen(ctx)
  80  
  81  	limit := 10
  82  	if err := sub.Subscribe(ctx, "test-sub", nostr.Filter{
  83  		Authors: []string{id.PubKey},
  84  		Limit:   &limit,
  85  	}); err != nil {
  86  		t.Fatal(err)
  87  	}
  88  
  89  	// Collect replayed events.
  90  	received := make(map[string]bool)
  91  	deadline := time.After(2 * time.Second)
  92  	for len(received) < 2 {
  93  		select {
  94  		case ev := <-sub.Events:
  95  			received[ev.ID] = true
  96  		case <-deadline:
  97  			t.Fatalf("timeout: received %d/2 events", len(received))
  98  		}
  99  	}
 100  
 101  	if !received[meta.ID] {
 102  		t.Fatal("missing metadata event")
 103  	}
 104  	if !received[status.ID] {
 105  		t.Fatal("missing status event")
 106  	}
 107  }
 108  
 109  // TestBroadcast verifies that new events are pushed to active subscribers.
 110  func TestBroadcast(t *testing.T) {
 111  	r := New("test-relay")
 112  	addr, err := r.Listen("127.0.0.1:0")
 113  	if err != nil {
 114  		t.Fatal(err)
 115  	}
 116  	defer r.Shutdown(context.Background())
 117  
 118  	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
 119  	defer cancel()
 120  
 121  	id, _ := nostr.NewIdentity()
 122  
 123  	// Subscriber connects first and sets up subscription.
 124  	sub, err := nostr.Connect(ctx, "ws://"+addr)
 125  	if err != nil {
 126  		t.Fatal(err)
 127  	}
 128  	defer sub.Disconnect()
 129  	go sub.Listen(ctx)
 130  
 131  	if err := sub.Subscribe(ctx, "live", nostr.Filter{
 132  		Kinds: []int{1},
 133  	}); err != nil {
 134  		t.Fatal(err)
 135  	}
 136  
 137  	// Wait for EOSE (subscription registered).
 138  	time.Sleep(100 * time.Millisecond)
 139  
 140  	// Publisher connects and publishes.
 141  	pub, err := nostr.Connect(ctx, "ws://"+addr)
 142  	if err != nil {
 143  		t.Fatal(err)
 144  	}
 145  	defer pub.Disconnect()
 146  	go pub.Listen(ctx)
 147  
 148  	status, _ := id.ComposeStatus(0, 1, 2922, 980, 0.335, 0.284, 3)
 149  	if err := pub.Publish(ctx, status); err != nil {
 150  		t.Fatal(err)
 151  	}
 152  
 153  	// Subscriber should receive the event via broadcast.
 154  	select {
 155  	case ev := <-sub.Events:
 156  		if ev.ID != status.ID {
 157  			t.Fatalf("wrong event: got %s, want %s", ev.ID[:16], status.ID[:16])
 158  		}
 159  	case <-time.After(2 * time.Second):
 160  		t.Fatal("timeout waiting for broadcast event")
 161  	}
 162  }
 163  
 164  // TestInvalidEventRejected verifies that events with bad signatures are rejected.
 165  func TestInvalidEventRejected(t *testing.T) {
 166  	r := New("test-relay")
 167  	addr, err := r.Listen("127.0.0.1:0")
 168  	if err != nil {
 169  		t.Fatal(err)
 170  	}
 171  	defer r.Shutdown(context.Background())
 172  
 173  	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
 174  	defer cancel()
 175  
 176  	id, _ := nostr.NewIdentity()
 177  	ev, _ := id.ComposeStatus(0, 1, 100, 50, 0.5, 0.1, 1)
 178  
 179  	// Tamper with the content (invalidates signature).
 180  	ev.Content = "tampered content"
 181  
 182  	client, err := nostr.Connect(ctx, "ws://"+addr)
 183  	if err != nil {
 184  		t.Fatal(err)
 185  	}
 186  	defer client.Disconnect()
 187  	go client.Listen(ctx)
 188  
 189  	if err := client.Publish(ctx, ev); err != nil {
 190  		t.Fatal(err)
 191  	}
 192  
 193  	select {
 194  	case ok := <-client.OKs:
 195  		if ok.Accepted {
 196  			t.Fatal("relay accepted invalid event")
 197  		}
 198  	case <-time.After(2 * time.Second):
 199  		t.Fatal("timeout waiting for OK")
 200  	}
 201  
 202  	if r.EventCount() != 0 {
 203  		t.Fatal("invalid event was stored")
 204  	}
 205  }
 206  
 207  // TestFilterMatching verifies kind and author filter logic.
 208  func TestFilterMatching(t *testing.T) {
 209  	r := New("test-relay")
 210  	addr, err := r.Listen("127.0.0.1:0")
 211  	if err != nil {
 212  		t.Fatal(err)
 213  	}
 214  	defer r.Shutdown(context.Background())
 215  
 216  	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
 217  	defer cancel()
 218  
 219  	id1, _ := nostr.NewIdentity()
 220  	id2, _ := nostr.NewIdentity()
 221  
 222  	// Publish events from two different authors.
 223  	pub, err := nostr.Connect(ctx, "ws://"+addr)
 224  	if err != nil {
 225  		t.Fatal(err)
 226  	}
 227  	defer pub.Disconnect()
 228  	go pub.Listen(ctx)
 229  
 230  	ev1, _ := id1.ComposeStatus(0, 1, 100, 50, 0.5, 0.1, 1)
 231  	ev2, _ := id2.ComposeStatus(1, 1, 200, 100, 0.5, 0.2, 2)
 232  	meta1, _ := id1.ComposeMetadata(0, 1, "aabbccdd")
 233  
 234  	for _, ev := range []*nostr.Event{ev1, ev2, meta1} {
 235  		pub.Publish(ctx, ev)
 236  	}
 237  
 238  	// Wait for all OKs.
 239  	for i := 0; i < 3; i++ {
 240  		select {
 241  		case <-pub.OKs:
 242  		case <-time.After(2 * time.Second):
 243  			t.Fatal("timeout waiting for OK")
 244  		}
 245  	}
 246  
 247  	// Subscribe for only id1's kind 1 events.
 248  	sub, err := nostr.Connect(ctx, "ws://"+addr)
 249  	if err != nil {
 250  		t.Fatal(err)
 251  	}
 252  	defer sub.Disconnect()
 253  	go sub.Listen(ctx)
 254  
 255  	limit := 10
 256  	if err := sub.Subscribe(ctx, "filter-test", nostr.Filter{
 257  		Authors: []string{id1.PubKey},
 258  		Kinds:   []int{1},
 259  		Limit:   &limit,
 260  	}); err != nil {
 261  		t.Fatal(err)
 262  	}
 263  
 264  	// Should receive only ev1 (id1's kind 1).
 265  	var received []*nostr.Event
 266  	deadline := time.After(2 * time.Second)
 267  loop:
 268  	for {
 269  		select {
 270  		case ev := <-sub.Events:
 271  			received = append(received, ev)
 272  		case <-deadline:
 273  			break loop
 274  		case <-time.After(500 * time.Millisecond):
 275  			break loop
 276  		}
 277  	}
 278  
 279  	if len(received) != 1 {
 280  		t.Fatalf("expected 1 event, got %d", len(received))
 281  	}
 282  	if received[0].ID != ev1.ID {
 283  		t.Fatalf("wrong event: got %s, want %s", received[0].ID[:16], ev1.ID[:16])
 284  	}
 285  }
 286  
 287  // TestRoundTrip verifies byte-level event fidelity through the relay.
 288  func TestRoundTrip(t *testing.T) {
 289  	r := New("test-relay")
 290  	addr, err := r.Listen("127.0.0.1:0")
 291  	if err != nil {
 292  		t.Fatal(err)
 293  	}
 294  	defer r.Shutdown(context.Background())
 295  
 296  	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
 297  	defer cancel()
 298  
 299  	id, _ := nostr.NewIdentity()
 300  	original, _ := id.ComposeStatus(0, 1, 2922, 980, 0.335, 0.284, 3)
 301  	origJSON, _ := json.Marshal(original)
 302  
 303  	// Publish.
 304  	pub, err := nostr.Connect(ctx, "ws://"+addr)
 305  	if err != nil {
 306  		t.Fatal(err)
 307  	}
 308  	defer pub.Disconnect()
 309  	go pub.Listen(ctx)
 310  
 311  	pub.Publish(ctx, original)
 312  	select {
 313  	case ok := <-pub.OKs:
 314  		if !ok.Accepted {
 315  			t.Fatalf("rejected: %s", ok.Message)
 316  		}
 317  	case <-time.After(2 * time.Second):
 318  		t.Fatal("timeout")
 319  	}
 320  
 321  	// Subscribe and receive.
 322  	sub, err := nostr.Connect(ctx, "ws://"+addr)
 323  	if err != nil {
 324  		t.Fatal(err)
 325  	}
 326  	defer sub.Disconnect()
 327  	go sub.Listen(ctx)
 328  
 329  	limit := 1
 330  	sub.Subscribe(ctx, "rt", nostr.Filter{
 331  		IDs:   []string{original.ID},
 332  		Limit: &limit,
 333  	})
 334  
 335  	select {
 336  	case ev := <-sub.Events:
 337  		recvJSON, _ := json.Marshal(ev)
 338  		if string(origJSON) != string(recvJSON) {
 339  			t.Fatalf("DIVERGENCE:\norig: %s\nrecv: %s", origJSON, recvJSON)
 340  		}
 341  		t.Log("round-trip: ZERO DIVERGENCE")
 342  	case <-time.After(2 * time.Second):
 343  		t.Fatal("timeout waiting for event")
 344  	}
 345  }
 346  
 347  // TestNoGoroutineLeak verifies that the relay doesn't leak goroutines
 348  // when clients connect and disconnect.
 349  func TestNoGoroutineLeak(t *testing.T) {
 350  	r := New("test-relay")
 351  	addr, err := r.Listen("127.0.0.1:0")
 352  	if err != nil {
 353  		t.Fatal(err)
 354  	}
 355  
 356  	ctx := context.Background()
 357  
 358  	// Baseline goroutine count.
 359  	runtime.GC()
 360  	time.Sleep(50 * time.Millisecond)
 361  	baseline := runtime.NumGoroutine()
 362  
 363  	// Connect and disconnect 10 clients.
 364  	for i := 0; i < 10; i++ {
 365  		cctx, cancel := context.WithTimeout(ctx, time.Second)
 366  		c, err := nostr.Connect(cctx, "ws://"+addr)
 367  		if err != nil {
 368  			cancel()
 369  			t.Fatal(err)
 370  		}
 371  		go c.Listen(cctx)
 372  		time.Sleep(10 * time.Millisecond)
 373  		c.Disconnect()
 374  		cancel()
 375  	}
 376  
 377  	r.Shutdown(ctx)
 378  
 379  	// Wait for goroutines to settle.
 380  	time.Sleep(200 * time.Millisecond)
 381  	runtime.GC()
 382  	time.Sleep(50 * time.Millisecond)
 383  
 384  	after := runtime.NumGoroutine()
 385  	leaked := after - baseline
 386  	if leaked > 5 {
 387  		t.Fatalf("goroutine leak: %d goroutines before, %d after (%d leaked)", baseline, after, leaked)
 388  	}
 389  }
 390  
 391  // TestNIP11 verifies the relay information document.
 392  func TestNIP11(t *testing.T) {
 393  	r := New("dendrite-test")
 394  	addr, err := r.Listen("127.0.0.1:0")
 395  	if err != nil {
 396  		t.Fatal(err)
 397  	}
 398  	defer r.Shutdown(context.Background())
 399  
 400  	// Fetch NIP-11 document via HTTP.
 401  	req, _ := http.NewRequest("GET", "http://"+addr+"/", nil)
 402  	req.Header.Set("Accept", "application/nostr+json")
 403  
 404  	resp, err := http.DefaultClient.Do(req)
 405  	if err != nil {
 406  		t.Fatal(err)
 407  	}
 408  	defer resp.Body.Close()
 409  
 410  	var info map[string]any
 411  	if err := json.NewDecoder(resp.Body).Decode(&info); err != nil {
 412  		t.Fatal(err)
 413  	}
 414  
 415  	if info["name"] != "dendrite-test" {
 416  		t.Fatalf("wrong name: %v", info["name"])
 417  	}
 418  	if info["software"] != "dendrite" {
 419  		t.Fatalf("wrong software: %v", info["software"])
 420  	}
 421  }
 422  
 423  // TestNIP09Deletion verifies that kind 5 events delete referenced events
 424  // by the same author, and that cross-author deletion is rejected.
 425  func TestNIP09Deletion(t *testing.T) {
 426  	r := New("test-relay")
 427  	addr, err := r.Listen("127.0.0.1:0")
 428  	if err != nil {
 429  		t.Fatal(err)
 430  	}
 431  	defer r.Shutdown(context.Background())
 432  
 433  	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
 434  	defer cancel()
 435  
 436  	author, _ := nostr.NewIdentity()
 437  	other, _ := nostr.NewIdentity()
 438  
 439  	// Publish an event from author.
 440  	client, err := nostr.Connect(ctx, "ws://"+addr)
 441  	if err != nil {
 442  		t.Fatal(err)
 443  	}
 444  	defer client.Disconnect()
 445  	go client.Listen(ctx)
 446  
 447  	ev, _ := author.ComposeStatus(0, 1, 100, 50, 0.5, 0.1, 1)
 448  	if err := client.Publish(ctx, ev); err != nil {
 449  		t.Fatal(err)
 450  	}
 451  	<-client.OKs
 452  
 453  	if r.EventCount() != 1 {
 454  		t.Fatalf("expected 1 event, got %d", r.EventCount())
 455  	}
 456  
 457  	// Other author tries to delete it — should have no effect.
 458  	delByOther := &nostr.Event{
 459  		CreatedAt: time.Now().Unix(),
 460  		Kind:      5,
 461  		Tags:      [][]string{{"e", ev.ID}},
 462  		Content:   "",
 463  	}
 464  	delByOther.Sign(other.PrivKeyHex())
 465  	client.Publish(ctx, delByOther)
 466  	<-client.OKs
 467  
 468  	// Event count should be 2 (original + deletion event), original not removed.
 469  	if r.EventCount() != 2 {
 470  		t.Fatalf("after cross-author deletion attempt: expected 2 events, got %d", r.EventCount())
 471  	}
 472  
 473  	// Author deletes their own event.
 474  	delBySelf := &nostr.Event{
 475  		CreatedAt: time.Now().Unix(),
 476  		Kind:      5,
 477  		Tags:      [][]string{{"e", ev.ID}},
 478  		Content:   "",
 479  	}
 480  	delBySelf.Sign(author.PrivKeyHex())
 481  	client.Publish(ctx, delBySelf)
 482  	<-client.OKs
 483  
 484  	// Original event should be gone: 2 deletion events + 1 other's event attempt = 3 minus 1 deleted = 2.
 485  	// Actually: original deleted, other's deletion event stored, self's deletion event stored = 2.
 486  	if r.EventCount() != 2 {
 487  		t.Fatalf("after self-deletion: expected 2 events (deletion events only), got %d", r.EventCount())
 488  	}
 489  }
 490  
 491  // TestNIP42Auth verifies that an authenticated relay rejects events
 492  // from unauthenticated clients and accepts after AUTH.
 493  func TestNIP42Auth(t *testing.T) {
 494  	r := New("auth-relay")
 495  	r.RequireAuth = true
 496  	addr, err := r.Listen("127.0.0.1:0")
 497  	if err != nil {
 498  		t.Fatal(err)
 499  	}
 500  	defer r.Shutdown(context.Background())
 501  
 502  	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
 503  	defer cancel()
 504  
 505  	id, _ := nostr.NewIdentity()
 506  
 507  	client, err := nostr.Connect(ctx, "ws://"+addr)
 508  	if err != nil {
 509  		t.Fatal(err)
 510  	}
 511  	defer client.Disconnect()
 512  	go client.Listen(ctx)
 513  
 514  	// The relay sends an AUTH challenge on connect. Read it.
 515  	// Give the relay a moment to send the challenge.
 516  	time.Sleep(100 * time.Millisecond)
 517  
 518  	// Try to publish without authenticating — should be rejected.
 519  	ev, _ := id.ComposeStatus(0, 1, 100, 50, 0.5, 0.1, 1)
 520  	client.Publish(ctx, ev)
 521  
 522  	select {
 523  	case ok := <-client.OKs:
 524  		if ok.Accepted {
 525  			t.Fatal("relay accepted event from unauthenticated client")
 526  		}
 527  		if !strings.Contains(ok.Message, "auth-required") {
 528  			t.Fatalf("expected auth-required message, got: %s", ok.Message)
 529  		}
 530  	case <-time.After(2 * time.Second):
 531  		t.Fatal("timeout")
 532  	}
 533  
 534  	if r.EventCount() != 0 {
 535  		t.Fatal("unauthenticated event was stored")
 536  	}
 537  }
 538  
 539  // TestCoherenceFilter verifies that the relay rejects events whose
 540  // content doesn't bond to any known lattice structure.
 541  func TestCoherenceFilter(t *testing.T) {
 542  	r := New("filter-relay")
 543  	// Only admit events that decompose into "hashtag" elements.
 544  	r.ContentFilter = &CoherenceFilter{
 545  		KnownTypes: map[string]bool{"hashtag": true},
 546  	}
 547  
 548  	addr, err := r.Listen("127.0.0.1:0")
 549  	if err != nil {
 550  		t.Fatal(err)
 551  	}
 552  	defer r.Shutdown(context.Background())
 553  
 554  	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
 555  	defer cancel()
 556  
 557  	id, _ := nostr.NewIdentity()
 558  
 559  	client, err := nostr.Connect(ctx, "ws://"+addr)
 560  	if err != nil {
 561  		t.Fatal(err)
 562  	}
 563  	defer client.Disconnect()
 564  	go client.Listen(ctx)
 565  
 566  	// Event WITHOUT hashtag — should be rejected.
 567  	evNoHash := &nostr.Event{
 568  		CreatedAt: 1700000000,
 569  		Kind:      1,
 570  		Content:   "just plain text without any tags",
 571  	}
 572  	evNoHash.Sign(id.PrivKeyHex())
 573  	client.Publish(ctx, evNoHash)
 574  
 575  	select {
 576  	case ok := <-client.OKs:
 577  		if ok.Accepted {
 578  			t.Fatal("relay accepted incoherent event")
 579  		}
 580  	case <-time.After(2 * time.Second):
 581  		t.Fatal("timeout")
 582  	}
 583  
 584  	// Event WITH hashtag — should be accepted.
 585  	evHash := &nostr.Event{
 586  		CreatedAt: 1700000000,
 587  		Kind:      1,
 588  		Content:   "talking about nostr",
 589  		Tags:      [][]string{{"t", "nostr"}},
 590  	}
 591  	evHash.Sign(id.PrivKeyHex())
 592  	client.Publish(ctx, evHash)
 593  
 594  	select {
 595  	case ok := <-client.OKs:
 596  		if !ok.Accepted {
 597  			t.Fatalf("relay rejected coherent event: %s", ok.Message)
 598  		}
 599  	case <-time.After(2 * time.Second):
 600  		t.Fatal("timeout")
 601  	}
 602  }
 603