返回 DeepSeek-Reasonix
stream_sink_test.go
根目录 / internal / agent / stream_sink_test.go
1 package agent
2
3 import (
4 "reflect"
5 "testing"
6
7 "reasonix/internal/event"
8 )
9
10 func sinkKinds(evs []event.Event) []event.Kind {
11 kinds := make([]event.Kind, 0, len(evs))
12 for _, e := range evs {
13 kinds = append(kinds, e.Kind)
14 }
15 return kinds
16 }
17
18 func TestDeferredStreamSinkBuffersSpeculativeEventsUntilReasoning(t *testing.T) {
19 inner := &recordSink{}
20 s := newReasoningAwareStreamSink(inner)
21
22 s.Emit(event.Event{Kind: event.Text, Text: "a"})
23 s.Emit(event.Event{Kind: event.ToolDispatch, Tool: event.Tool{ID: "c1"}})
24 s.Emit(event.Event{Kind: event.ToolResult, Tool: event.Tool{ID: "c1"}})
25 s.Emit(event.Event{Kind: event.Message, Text: "m"})
26 // Usage/turn bookkeeping is never speculative and passes through live, as
27 // does an empty reasoning metadata event (it must not unlock the buffer).
28 s.Emit(event.Event{Kind: event.Usage})
29 s.Emit(event.Event{Kind: event.Reasoning, Text: " "})
30 s.Emit(event.Event{Kind: event.Text, Text: "b"})
31
32 if got, want := sinkKinds(inner.evs), []event.Kind{event.Usage, event.Reasoning}; !reflect.DeepEqual(got, want) {
33 t.Fatalf("pre-reasoning events = %v, want only %v", got, want)
34 }
35
36 s.Emit(event.Event{Kind: event.Reasoning, Text: "thinking"})
37 want := []event.Kind{
38 event.Usage, event.Reasoning, event.Reasoning,
39 event.Text, event.ToolDispatch, event.ToolResult, event.Message, event.Text,
40 }
41 if got := sinkKinds(inner.evs); !reflect.DeepEqual(got, want) {
42 t.Fatalf("unlock flush = %v, want reasoning then buffered order %v", got, want)
43 }
44
45 s.Emit(event.Event{Kind: event.Text, Text: "c"})
46 if got := sinkKinds(inner.evs); !reflect.DeepEqual(got, append(want, event.Text)) {
47 t.Fatalf("post-unlock events = %v, want live passthrough", got)
48 }
49 }
50
51 func TestDeferredStreamSinkDiscardDropsBufferedEvents(t *testing.T) {
52 inner := &recordSink{}
53 s := newReasoningAwareStreamSink(inner)
54 s.Emit(event.Event{Kind: event.Text, Text: "speculative"})
55 s.Emit(event.Event{Kind: event.ToolDispatch, Tool: event.Tool{ID: "c1"}})
56 s.Discard()
57 s.Flush()
58
59 if got := len(inner.evs); got != 0 {
60 t.Fatalf("discarded sink emitted %d buffered events, want 0", got)
61 }
62 s.Emit(event.Event{Kind: event.Reasoning, Text: "thinking"})
63 s.Emit(event.Event{Kind: event.Text, Text: "adopted"})
64 if got, want := sinkKinds(inner.evs), []event.Kind{event.Reasoning, event.Text}; !reflect.DeepEqual(got, want) {
65 t.Fatalf("post-discard events = %v, want %v", got, want)
66 }
67 }
68
69 func TestDeferredStreamSinkDefersEverythingUntilFlush(t *testing.T) {
70 inner := &recordSink{}
71 s := newDeferredStreamSink(inner)
72 s.Emit(event.Event{Kind: event.Reasoning, Text: "thinking"})
73 s.Emit(event.Event{Kind: event.Usage})
74 s.Emit(event.Event{Kind: event.Text, Text: "a"})
75 if got := len(inner.evs); got != 0 {
76 t.Fatalf("defer-all sink leaked %d events before Flush", got)
77 }
78 s.Flush()
79 want := []event.Kind{event.Reasoning, event.Usage, event.Text}
80 if got := sinkKinds(inner.evs); !reflect.DeepEqual(got, want) {
81 t.Fatalf("flushed events = %v, want arrival order %v", got, want)
82 }
83 s.Flush()
84 if got := len(inner.evs); got != len(want) {
85 t.Fatalf("second Flush re-emitted events: %d total", got)
86 }
87 }
88
89 func TestDeferredStreamSinkNilIsInert(t *testing.T) {
90 var s *deferredStreamSink
91 s.Emit(event.Event{Kind: event.Text, Text: "x"})
92 s.Flush()
93 s.Discard()
94 }
95
95 lines GO