package app import ( "testing" "time" "git.smesh.lol/orly/pkg/nostr/encoders/event" ) // TestPubSubEnqueueNeverBlocks is a regression test for the subscription // delivery stall: a subscriber that is not draining its channel used to block // the publisher actor for the full write timeout per delivered event, which // stalled delivery for every other subscriber on the relay. // // The receiver here is deliberately unbuffered with nothing reading it, so the // pump can push at most one event and then blocks. enqueue must still return // immediately for every call. func TestPubSubEnqueueNeverBlocks(t *testing.T) { receiver := make(event.C) // unbuffered, never read q := newPubSub(receiver) defer q.close() // Keep offering events until the queue bound rejects one. The exact count is // not asserted: the pump drains concurrently, which is the point. deadline := time.Now().Add(10 * time.Second) dropped := false for i := 0; i < q.maxQueue*4; i++ { if time.Now().After(deadline) { t.Fatal("enqueue stopped making progress: it appears to block") } if !q.enqueue(&event.E{Kind: 1}) { dropped = true break } } if !dropped { t.Fatalf("queue bound of %d was never enforced", q.maxQueue) } if q.dropped.Load() == 0 { t.Fatal("expected overflow to be counted in dropped") } // Enqueue must stay non-blocking even once full. start := time.Now() for i := 0; i < 1000; i++ { q.enqueue(&event.E{Kind: 1}) } if elapsed := time.Since(start); elapsed > time.Second { t.Fatalf("enqueue blocked once full: 1000 calls took %s", elapsed) } } // TestPubSubPumpDeliversThenClosesReceiver verifies the normal delivery path and // that stopping a pump closes the receiver channel, which is what lets the // per-subscription consumer goroutine exit instead of leaking. func TestPubSubPumpDeliversThenClosesReceiver(t *testing.T) { receiver := make(event.C, 8) q := newPubSub(receiver) ev := &event.E{Kind: 1} if !q.enqueue(ev) { t.Fatal("enqueue rejected an event on an empty queue") } select { case got := <-receiver: if got != ev { t.Fatalf("pump delivered the wrong event: %p != %p", got, ev) } case <-time.After(2 * time.Second): t.Fatal("pump did not deliver the queued event") } // close() must release the pump and close the receiver so the consumer // terminates, even though nothing more is read. q.close() select { case _, ok := <-receiver: if ok { // An event queued before the close may still be delivered; drain // until the channel reports closed. select { case _, ok = <-receiver: if ok { t.Fatal("receiver still open after close()") } case <-time.After(2 * time.Second): t.Fatal("pump did not close the receiver channel after close()") } } case <-time.After(2 * time.Second): t.Fatal("pump did not close the receiver channel after close()") } // enqueue after close must be rejected, not panic. if q.enqueue(&event.E{Kind: 1}) { t.Fatal("enqueue accepted an event after close") } }