publisher_pubsub_test.go raw

   1  package app
   2  
   3  import (
   4  	"testing"
   5  	"time"
   6  
   7  	"git.smesh.lol/orly/pkg/nostr/encoders/event"
   8  )
   9  
  10  // TestPubSubEnqueueNeverBlocks is a regression test for the subscription
  11  // delivery stall: a subscriber that is not draining its channel used to block
  12  // the publisher actor for the full write timeout per delivered event, which
  13  // stalled delivery for every other subscriber on the relay.
  14  //
  15  // The receiver here is deliberately unbuffered with nothing reading it, so the
  16  // pump can push at most one event and then blocks. enqueue must still return
  17  // immediately for every call.
  18  func TestPubSubEnqueueNeverBlocks(t *testing.T) {
  19  	receiver := make(event.C) // unbuffered, never read
  20  	q := newPubSub(receiver)
  21  	defer q.close()
  22  
  23  	// Keep offering events until the queue bound rejects one. The exact count is
  24  	// not asserted: the pump drains concurrently, which is the point.
  25  	deadline := time.Now().Add(10 * time.Second)
  26  	dropped := false
  27  	for i := 0; i < q.maxQueue*4; i++ {
  28  		if time.Now().After(deadline) {
  29  			t.Fatal("enqueue stopped making progress: it appears to block")
  30  		}
  31  		if !q.enqueue(&event.E{Kind: 1}) {
  32  			dropped = true
  33  			break
  34  		}
  35  	}
  36  
  37  	if !dropped {
  38  		t.Fatalf("queue bound of %d was never enforced", q.maxQueue)
  39  	}
  40  	if q.dropped.Load() == 0 {
  41  		t.Fatal("expected overflow to be counted in dropped")
  42  	}
  43  	// Enqueue must stay non-blocking even once full.
  44  	start := time.Now()
  45  	for i := 0; i < 1000; i++ {
  46  		q.enqueue(&event.E{Kind: 1})
  47  	}
  48  	if elapsed := time.Since(start); elapsed > time.Second {
  49  		t.Fatalf("enqueue blocked once full: 1000 calls took %s", elapsed)
  50  	}
  51  }
  52  
  53  // TestPubSubPumpDeliversThenClosesReceiver verifies the normal delivery path and
  54  // that stopping a pump closes the receiver channel, which is what lets the
  55  // per-subscription consumer goroutine exit instead of leaking.
  56  func TestPubSubPumpDeliversThenClosesReceiver(t *testing.T) {
  57  	receiver := make(event.C, 8)
  58  	q := newPubSub(receiver)
  59  
  60  	ev := &event.E{Kind: 1}
  61  	if !q.enqueue(ev) {
  62  		t.Fatal("enqueue rejected an event on an empty queue")
  63  	}
  64  
  65  	select {
  66  	case got := <-receiver:
  67  		if got != ev {
  68  			t.Fatalf("pump delivered the wrong event: %p != %p", got, ev)
  69  		}
  70  	case <-time.After(2 * time.Second):
  71  		t.Fatal("pump did not deliver the queued event")
  72  	}
  73  
  74  	// close() must release the pump and close the receiver so the consumer
  75  	// terminates, even though nothing more is read.
  76  	q.close()
  77  
  78  	select {
  79  	case _, ok := <-receiver:
  80  		if ok {
  81  			// An event queued before the close may still be delivered; drain
  82  			// until the channel reports closed.
  83  			select {
  84  			case _, ok = <-receiver:
  85  				if ok {
  86  					t.Fatal("receiver still open after close()")
  87  				}
  88  			case <-time.After(2 * time.Second):
  89  				t.Fatal("pump did not close the receiver channel after close()")
  90  			}
  91  		}
  92  	case <-time.After(2 * time.Second):
  93  		t.Fatal("pump did not close the receiver channel after close()")
  94  	}
  95  
  96  	// enqueue after close must be rejected, not panic.
  97  	if q.enqueue(&event.E{Kind: 1}) {
  98  		t.Fatal("enqueue accepted an event after close")
  99  	}
 100  }
 101