| 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 |