| 1 | package agent |
| 2 | |
| 3 | import ( |
| 4 | "strings" |
| 5 | |
| 6 | "reasonix/internal/event" |
| 7 | ) |
| 8 | |
| 9 | // deferredStreamSink keeps selected stream events local until the caller |
| 10 | // chooses which provider response to adopt. On an ordinary healthy DeepSeek |
| 11 | // turn, reasoning arrives before tool calls and unlocks live tool-card events. |
| 12 | // On the rare malformed turn with no reasoning, only the speculative partial |
| 13 | // tool cards remain buffered, so retrying does not flash duplicate cards in the |
| 14 | // UI. A recovery attempt buffers everything because it may be discarded. |
| 15 | type deferredStreamSink struct { |
| 16 | inner event.Sink |
| 17 | deferAll bool |
| 18 | waitingForReasoning bool |
| 19 | sawReasoning bool |
| 20 | events []event.Event |
| 21 | } |
| 22 | |
| 23 | func newReasoningAwareStreamSink(inner event.Sink) *deferredStreamSink { |
| 24 | return &deferredStreamSink{inner: inner, waitingForReasoning: true} |
| 25 | } |
| 26 | |
| 27 | func newDeferredStreamSink(inner event.Sink) *deferredStreamSink { |
| 28 | return &deferredStreamSink{inner: inner, deferAll: true} |
| 29 | } |
| 30 | |
| 31 | func (s *deferredStreamSink) Emit(e event.Event) { |
| 32 | if s == nil { |
| 33 | return |
| 34 | } |
| 35 | if s.deferAll { |
| 36 | s.events = append(s.events, e) |
| 37 | return |
| 38 | } |
| 39 | if s.waitingForReasoning && e.Kind == event.Reasoning && strings.TrimSpace(e.Text) != "" { |
| 40 | s.sawReasoning = true |
| 41 | s.inner.Emit(e) |
| 42 | s.flushBuffered() |
| 43 | return |
| 44 | } |
| 45 | if s.waitingForReasoning && !s.sawReasoning { |
| 46 | switch e.Kind { |
| 47 | case event.ToolDispatch, event.ToolResult, event.Text, event.Message: |
| 48 | // Keep every user-visible speculative event private until reasoning |
| 49 | // proves the turn replayable. Healthy DeepSeek responses emit |
| 50 | // reasoning first, so their live-streaming fast path is unchanged. |
| 51 | s.events = append(s.events, e) |
| 52 | return |
| 53 | } |
| 54 | } |
| 55 | s.inner.Emit(e) |
| 56 | } |
| 57 | |
| 58 | func (s *deferredStreamSink) flushBuffered() { |
| 59 | if s == nil { |
| 60 | return |
| 61 | } |
| 62 | for _, e := range s.events { |
| 63 | s.inner.Emit(e) |
| 64 | } |
| 65 | s.events = nil |
| 66 | } |
| 67 | |
| 68 | func (s *deferredStreamSink) Flush() { |
| 69 | if s == nil { |
| 70 | return |
| 71 | } |
| 72 | s.flushBuffered() |
| 73 | } |
| 74 | |
| 75 | func (s *deferredStreamSink) Discard() { |
| 76 | if s != nil { |
| 77 | s.events = nil |
| 78 | } |
| 79 | } |
| 80 |