返回 DeepSeek-Reasonix
coalesce.go
根目录 / internal / event / coalesce.go
1 package event
2
3 import (
4 "reflect"
5 "strings"
6 "sync"
7 "time"
8
9 "reasonix/internal/evidence"
10 "reasonix/internal/nilutil"
11 )
12
13 // coalesceMaxBytes bounds a merged delta so one event never carries an
14 // unbounded payload across a frontend bridge.
15 const coalesceMaxBytes = 16 << 10
16
17 // DefaultStreamDeltaWindow caps how often coalesced streaming deltas cross
18 // into a frontend: at most one merged event per window under load — about one
19 // per display frame — so a fast provider (hundreds of chunks/sec) cannot
20 // flood a webview bridge, an SSE stream, or a terminal redraw loop.
21 const DefaultStreamDeltaWindow = 16 * time.Millisecond
22
23 // Coalesce wraps inner so bursts of consecutive streaming deltas — Text or
24 // Reasoning events carrying nothing but a Text payload, or ToolProgress events
25 // carrying nothing but one tool's Output — merge into one event.
26 // The first delta of a burst forwards immediately (time-to-first-token is
27 // unchanged); later deltas buffer at most window, flushing earlier on any
28 // other event (total order preserved), a kind switch, or coalesceMaxBytes.
29 func Coalesce(inner Sink, window time.Duration) Sink {
30 if nilutil.IsNil(inner) {
31 return Discard
32 }
33 if window <= 0 {
34 return inner
35 }
36 return &coalescer{inner: inner, window: window}
37 }
38
39 type coalescer struct {
40 inner Sink
41 window time.Duration
42
43 // mu guards buffering state and the outbound queue; inner.Emit is never
44 // called under mu. A single drainer forwards FIFO, so a sink that
45 // synchronously re-enters Emit enqueues and returns instead of deadlocking.
46 mu sync.Mutex
47 key deltaKey
48 buf strings.Builder
49 pending bool
50 timer *time.Timer
51 lastForward time.Time
52 queue []coalescedEvent
53 draining bool
54 }
55
56 type coalescedEvent struct {
57 event Event
58 done chan error
59 forward func()
60 }
61
62 var _ OptionalSinkCapabilities = (*coalescer)(nil)
63 var _ CheckedSink = (*coalescer)(nil)
64
65 // deltaKey is the identity a buffered burst merges under; any change flushes.
66 type deltaKey struct {
67 kind Kind
68 source, messageID, attemptID string
69 toolID string
70 }
71
72 // streamDelta reports whether e is a pure streaming delta, with its merge key
73 // and payload: merging is only safe when no other field carries meaning. The
74 // zero-probe comparison keeps this true by construction as Event grows fields.
75 func streamDelta(e Event) (deltaKey, string, bool) {
76 key := deltaKey{kind: e.Kind, source: e.Source, messageID: e.MessageID, attemptID: e.AttemptID}
77 probe := e
78 probe.Source, probe.MessageID, probe.AttemptID = "", "", ""
79 payload := e.Text
80 switch e.Kind {
81 case Text, Reasoning:
82 probe.Text = ""
83 case ToolProgress:
84 key.toolID, payload = e.Tool.ID, e.Tool.Output
85 probe.Tool.ID, probe.Tool.Output = "", ""
86 default:
87 return deltaKey{}, "", false
88 }
89 if payload == "" || (e.Kind == ToolProgress && key.toolID == "") || !reflect.DeepEqual(probe, Event{Kind: e.Kind}) {
90 return deltaKey{}, "", false
91 }
92 return key, payload, true
93 }
94
95 func (k deltaKey) event(payload string) Event {
96 e := Event{Kind: k.kind, Source: k.source, MessageID: k.messageID, AttemptID: k.attemptID}
97 if k.kind == ToolProgress {
98 e.Tool = Tool{ID: k.toolID, Output: payload}
99 } else {
100 e.Text = payload
101 }
102 return e
103 }
104
105 func (c *coalescer) Emit(e Event) {
106 _ = c.enqueue(e, false)
107 }
108
109 // EmitChecked is a synchronous ordering barrier. Buffered deltas are written
110 // before e, and it returns only after the durable inner sink has acknowledged
111 // e. Regular streaming Emit calls remain non-blocking while a drainer is
112 // active; they learn asynchronous failures through the lifecycle sink's
113 // poisoned-ledger state.
114 func (c *coalescer) EmitChecked(e Event) error {
115 return c.enqueue(e, true)
116 }
117
118 func (c *coalescer) enqueue(e Event, checked bool) error {
119 var done chan error
120 if checked {
121 done = make(chan error, 1)
122 }
123 key, payload, delta := streamDelta(e)
124 c.mu.Lock()
125 if checked && delta {
126 c.enqueueFlushLocked()
127 c.queue = append(c.queue, coalescedEvent{event: e, done: done})
128 c.drainAndUnlock()
129 return <-done
130 }
131 if !delta {
132 c.enqueueFlushLocked()
133 c.queue = append(c.queue, coalescedEvent{event: e, done: done})
134 c.drainAndUnlock()
135 if done != nil {
136 return <-done
137 }
138 return nil
139 }
140 if c.pending && c.key != key {
141 c.enqueueFlushLocked()
142 }
143 if !c.pending && time.Since(c.lastForward) >= c.window {
144 c.lastForward = time.Now()
145 c.queue = append(c.queue, coalescedEvent{event: e, done: done})
146 c.drainAndUnlock()
147 if done != nil {
148 return <-done
149 }
150 return nil
151 }
152 if !c.pending {
153 c.pending = true
154 c.key = key
155 if c.timer == nil {
156 c.timer = time.AfterFunc(c.window, c.flush)
157 } else {
158 c.timer.Reset(c.window)
159 }
160 }
161 c.buf.WriteString(payload)
162 if c.buf.Len() >= coalesceMaxBytes {
163 c.enqueueFlushLocked()
164 }
165 c.drainAndUnlock()
166 if done != nil {
167 return <-done
168 }
169 return nil
170 }
171
172 func (c *coalescer) flush() {
173 c.mu.Lock()
174 c.enqueueFlushLocked()
175 c.drainAndUnlock()
176 }
177
178 // enqueueFlushLocked moves the buffered delta (if any) onto the outbound queue.
179 func (c *coalescer) enqueueFlushLocked() {
180 if !c.pending {
181 return
182 }
183 c.timer.Stop()
184 c.queue = append(c.queue, coalescedEvent{event: c.key.event(c.buf.String())})
185 c.buf.Reset()
186 c.pending = false
187 c.key = deltaKey{}
188 c.lastForward = time.Now()
189 }
190
191 // drainAndUnlock forwards queued events in FIFO order and releases mu. Exactly
192 // one goroutine drains at a time; others enqueue and return.
193 func (c *coalescer) drainAndUnlock() {
194 if c.draining || len(c.queue) == 0 {
195 c.mu.Unlock()
196 return
197 }
198 c.draining = true
199 for len(c.queue) > 0 {
200 batch := c.queue
201 c.queue = nil
202 c.mu.Unlock()
203 for _, item := range batch {
204 if item.forward != nil {
205 item.forward()
206 continue
207 }
208 err := EmitChecked(c.inner, item.event)
209 if item.done != nil {
210 item.done <- err
211 close(item.done)
212 }
213 }
214 c.mu.Lock()
215 }
216 c.draining = false
217 c.mu.Unlock()
218 }
219
220 // Optional sink capabilities flush first so audits never overtake a buffered
221 // delta, then forward to inner sinks that opt in.
222
223 func (c *coalescer) enqueueCapability(forward func()) {
224 c.mu.Lock()
225 c.enqueueFlushLocked()
226 c.queue = append(c.queue, coalescedEvent{forward: forward})
227 c.drainAndUnlock()
228 }
229
230 func (c *coalescer) RecordDelegationAudit(a evidence.DelegationAudit) {
231 c.enqueueCapability(func() { RecordDelegationAudit(c.inner, a) })
232 }
233
234 func (c *coalescer) RecordReadinessAudit(a evidence.ReadinessAudit) {
235 c.enqueueCapability(func() { RecordReadinessAudit(c.inner, a) })
236 }
237
238 func (c *coalescer) RecordAnchorSafetyAudit(a AnchorSafetyAudit) {
239 c.enqueueCapability(func() { RecordAnchorSafetyAudit(c.inner, a) })
240 }
241
242 func (c *coalescer) RecordTurnCompletion() {
243 c.enqueueCapability(func() { RecordTurnCompletion(c.inner) })
244 }
245
246 func (c *coalescer) RecordProtocolRecovery(a ProtocolRecoveryAudit) {
247 c.enqueueCapability(func() { RecordProtocolRecovery(c.inner, a) })
248 }
249
250 func (c *coalescer) RecordContractShadow(a ContractShadowAudit) {
251 c.enqueueCapability(func() { RecordContractShadow(c.inner, a) })
252 }
253
254 func (c *coalescer) RecordCompletionReport(a CompletionReportAudit) {
255 c.enqueueCapability(func() { RecordCompletionReport(c.inner, a) })
256 }
257
258 func (c *coalescer) RecordOutcomeProgress(sample evidence.OutcomeSample) {
259 c.enqueueCapability(func() { RecordOutcomeProgress(c.inner, sample) })
260 }
261
262 func (c *coalescer) RecordMemoryRecall(a MemoryRecallAudit) {
263 c.enqueueCapability(func() { RecordMemoryRecall(c.inner, a) })
264 }
265
266 func (c *coalescer) RecordDelegationAdmission(a DelegationAdmissionAudit) {
267 c.enqueueCapability(func() { RecordDelegationAdmission(c.inner, a) })
268 }
269
270 func (c *coalescer) RecordWorkspaceMutation(m WorkspaceMutation) {
271 c.enqueueCapability(func() { RecordWorkspaceMutation(c.inner, m) })
272 }
273
274 func (c *coalescer) RecordRunBudget(sample RunBudgetSample) {
275 c.enqueueCapability(func() { RecordRunBudget(c.inner, sample) })
276 }
277
278 func (c *coalescer) RecordSubagentLifecycle(info SubagentLifecycleInfo) {
279 c.enqueueCapability(func() { RecordSubagentLifecycle(c.inner, info) })
280 }
281
281 lines GO