返回 DeepSeek-Reasonix
broadcaster_test.go
根目录 / internal / serve / broadcaster_test.go
1 package serve
2
3 import (
4 "encoding/json"
5 "strings"
6 "testing"
7
8 "reasonix/internal/event"
9 "reasonix/internal/eventwire"
10 )
11
12 func TestBroadcasterFanOut(t *testing.T) {
13 b := NewBroadcaster()
14 a, ca := b.Subscribe()
15 d, cd := b.Subscribe()
16 defer ca()
17 defer cd()
18
19 if got := b.Subscribers(); got != 2 {
20 t.Fatalf("subscribers = %d, want 2", got)
21 }
22
23 b.Emit(event.Event{Kind: event.Text, Text: "hi"})
24
25 for i, ch := range []<-chan []byte{a, d} {
26 var w eventwire.Event
27 if err := json.Unmarshal(<-ch, &w); err != nil {
28 t.Fatalf("subscriber %d: %v", i, err)
29 }
30 if w.Kind != "text" || w.Text != "hi" {
31 t.Errorf("subscriber %d got %+v", i, w)
32 }
33 }
34 }
35
36 func TestBroadcasterEmitsRetryingJSON(t *testing.T) {
37 b := NewBroadcaster()
38 ch, cancel := b.Subscribe()
39 defer cancel()
40
41 b.Emit(event.Event{Kind: event.Retrying, RetryAttempt: 3, RetryMax: 10})
42
43 s := string(<-ch)
44 for _, want := range []string{`"kind":"retrying"`, `"retryAttempt":3`, `"retryMax":10`} {
45 if !strings.Contains(s, want) {
46 t.Fatalf("retrying broadcast JSON = %s, want it to contain %s", s, want)
47 }
48 }
49 }
50
51 func TestBroadcasterUnsubscribe(t *testing.T) {
52 b := NewBroadcaster()
53 _, cancel := b.Subscribe()
54 if b.Subscribers() != 1 {
55 t.Fatalf("want 1 subscriber")
56 }
57 cancel()
58 if b.Subscribers() != 0 {
59 t.Fatalf("unsubscribe should drop to 0, got %d", b.Subscribers())
60 }
61 // Emitting with no subscribers must not panic.
62 b.Emit(event.Event{Kind: event.TurnDone})
63 }
64
65 func TestBroadcasterDropsSlowSubscriber(t *testing.T) {
66 b := NewBroadcaster()
67 ch, cancel := b.Subscribe()
68 defer cancel()
69 // Overfill far past the 64-slot buffer without reading; Emit must not block.
70 for i := 0; i < 1000; i++ {
71 b.Emit(event.Event{Kind: event.Text, Text: "x"})
72 }
73 if len(ch) == 0 {
74 t.Error("expected some buffered frames")
75 }
76 }
77
77 lines GO