返回 DeepSeek-Reasonix
buffer.go
根目录 / internal / transcript / buffer.go
1 package transcript
2
3 import (
4 "fmt"
5 "reasonix/internal/event"
6 "reasonix/internal/eventwire"
7 "reasonix/internal/provider"
8 "slices"
9 "strings"
10 )
11
12 // displayTextAccumulator retains provider chunks without repeatedly copying
13 // the complete prefix. A turn only materializes the final string when its
14 // display-only history is persisted; successful executor turns are discarded
15 // without ever joining their chunks.
16 type displayTextAccumulator struct {
17 parts []string
18 size int
19 }
20
21 func (a *displayTextAccumulator) append(text string) {
22 if text == "" {
23 return
24 }
25 a.parts = append(a.parts, text)
26 a.size += len(text)
27 }
28
29 func (a *displayTextAccumulator) replace(text string) {
30 a.parts = nil
31 a.size = 0
32 a.append(text)
33 }
34
35 func (a *displayTextAccumulator) hasNonWhitespace() bool {
36 for _, part := range a.parts {
37 if strings.TrimSpace(part) != "" {
38 return true
39 }
40 }
41 return false
42 }
43
44 func (a *displayTextAccumulator) string() string {
45 switch len(a.parts) {
46 case 0:
47 return ""
48 case 1:
49 return a.parts[0]
50 }
51 var out strings.Builder
52 out.Grow(a.size)
53 for _, part := range a.parts {
54 out.WriteString(part)
55 }
56 return out.String()
57 }
58
59 type bufferedMessage struct {
60 message Message
61 content displayTextAccumulator
62 reasoning displayTextAccumulator
63 }
64
65 func (m *bufferedMessage) materialize() Message {
66 out := m.message
67 if out.Role == "assistant" {
68 out.Content = m.content.string()
69 out.Reasoning = m.reasoning.string()
70 }
71 if len(out.MemoryCitations) > 0 {
72 out.MemoryCitations = append([]provider.MemoryCitation(nil), out.MemoryCitations...)
73 }
74 if len(out.ToolCalls) > 0 {
75 out.ToolCalls = append([]ToolCall(nil), out.ToolCalls...)
76 }
77 return out
78 }
79
80 type Buffer struct {
81 Format Formatter
82 messages []*bufferedMessage
83 byMessageID map[string]*bufferedMessage
84 tools map[string]string
85 completion *eventwire.CompletionSummary
86 userTurns int
87 }
88
89 func (buffer *Buffer) Reset() {
90 buffer.messages = nil
91 buffer.byMessageID = nil
92 buffer.tools = nil
93 buffer.completion = nil
94 buffer.userTurns = 0
95 }
96
97 func (buffer *Buffer) ResultMessages() []Message {
98 var out []Message
99 for _, m := range buffer.messages {
100 if m.message.Code == "turn_result" {
101 out = append(out, m.materialize())
102 }
103 }
104 return out
105 }
106
107 func (buffer *Buffer) Messages() []Message {
108 if len(buffer.messages) == 0 {
109 return nil
110 }
111 out := make([]Message, 0, len(buffer.messages))
112 for _, message := range buffer.messages {
113 out = append(out, message.materialize())
114 }
115 return out
116 }
117
118 func (buffer *Buffer) attachTurnStats(turnID string, usage *TurnUsage, durationMs, completedAt int64, finalMessageID string) {
119 for _, row := range slices.Backward(buffer.messages) {
120 if row.message.TurnID != turnID || row.message.Role != "assistant" {
121 continue
122 }
123 if finalMessageID != "" && row.message.MessageID != finalMessageID {
124 continue
125 }
126 if !row.content.hasNonWhitespace() && !row.reasoning.hasNonWhitespace() {
127 continue
128 }
129 if usage != nil {
130 copy := *usage
131 copy.Routes = append([]string(nil), usage.Routes...)
132 if usage.CacheReadTokens != nil {
133 value := *usage.CacheReadTokens
134 copy.CacheReadTokens = &value
135 }
136 if usage.ReasoningTokens != nil {
137 value := *usage.ReasoningTokens
138 copy.ReasoningTokens = &value
139 }
140 row.message.TurnUsage = &copy
141 }
142 row.message.TurnDurationMs = durationMs
143 if row.message.CreatedAt == 0 {
144 row.message.CreatedAt = completedAt
145 }
146 return
147 }
148 }
149
150 func (buffer *Buffer) Apply(e event.Event) {
151 start := len(buffer.messages)
152 defer buffer.stampAppliedMessages(e, start)
153 switch e.Kind {
154 case event.UserMessage:
155 buffer.applyUserMessage(e)
156 case event.StreamAttempt:
157 buffer.applyStreamAttempt(e)
158 case event.CompletionSummary:
159 buffer.completion = eventwire.ToWire(e).Completion
160 case event.TurnDone:
161 buffer.applyTurnDone(e)
162 case event.Phase:
163 if strings.TrimSpace(e.Text) != "" {
164 buffer.messages = append(buffer.messages, &bufferedMessage{message: Message{Role: "phase", Content: e.Text}})
165 }
166 case event.Reasoning:
167 if e.Text != "" {
168 ensureIdentifiedDisplayAssistant(buffer, e.MessageID, false).reasoning.append(e.Text)
169 }
170 case event.Text:
171 if e.Text != "" {
172 ensureIdentifiedDisplayAssistant(buffer, e.MessageID, false).content.append(e.Text)
173 }
174 case event.Message:
175 buffer.applyAssistantMessage(e)
176 case event.ToolDispatch:
177 recordHistoryToolDispatch(buffer, e)
178 case event.ToolResult:
179 buffer.applyToolResult(e)
180 case event.Notice:
181 buffer.applyNotice(e)
182 }
183 }
184
185 func (buffer *Buffer) stampAppliedMessages(e event.Event, start int) {
186 for i := start; i < len(buffer.messages); i++ {
187 m := &buffer.messages[i].message
188 m.Source, m.TurnID = e.Source, e.TurnID
189 if m.RecordID == "" {
190 switch {
191 case m.Role == "tool" && m.ToolCallID != "":
192 m.RecordID = "tool:" + m.ToolCallID
193 case m.MessageID != "":
194 m.RecordID = "m:" + m.MessageID
195 case e.Sequence > 0:
196 m.RecordID = fmt.Sprintf("e:%s:%d:%d", e.TurnID, e.Sequence, i-start)
197 }
198 }
199 }
200 }
201
202 func (buffer *Buffer) applyUserMessage(e event.Event) {
203 if e.Source != "" && e.Source != event.UsageSourceExecutor {
204 return
205 }
206 if e.MessageID != "" && buffer.byMessageID[e.MessageID] == nil {
207 if buffer.byMessageID == nil {
208 buffer.byMessageID = make(map[string]*bufferedMessage)
209 }
210 m := &bufferedMessage{message: Message{Role: "user", MessageID: e.MessageID, Content: e.Text}}
211 buffer.userTurns++
212 m.message.HistoryTurn = buffer.userTurns
213 buffer.messages = append(buffer.messages, m)
214 buffer.byMessageID[e.MessageID] = m
215 }
216 }
217
218 func (buffer *Buffer) applyStreamAttempt(e event.Event) {
219 if e.StreamAttempt.Action == event.StreamAttemptDiscard && e.MessageID != "" {
220 kept := buffer.messages[:0]
221 for _, message := range buffer.messages {
222 if message.message.MessageID != e.MessageID {
223 kept = append(kept, message)
224 }
225 }
226 clear(buffer.messages[len(kept):])
227 buffer.messages = kept
228 delete(buffer.byMessageID, e.MessageID)
229 }
230 if e.MessageID != "" && e.StreamAttempt.Action == event.StreamAttemptBegin {
231 m := ensureIdentifiedDisplayAssistant(buffer, e.MessageID, false)
232 m.message.AttemptID, m.message.Pending = e.AttemptID, true
233 }
234 if e.MessageID != "" && e.StreamAttempt.Action == event.StreamAttemptCommit {
235 if m := buffer.byMessageID[e.MessageID]; m != nil {
236 m.message.Pending = false
237 }
238 }
239 }
240
241 func (buffer *Buffer) applyTurnDone(e event.Event) {
242 for _, m := range buffer.messages {
243 m.message.Pending = false
244 if m.message.Role == "user" && m.message.TurnID == e.TurnID {
245 m.message.CheckpointTurn = e.CheckpointTurn
246 }
247 for i := range m.message.ToolCalls {
248 m.message.ToolCalls[i].Pending = false
249 }
250 }
251 wire := eventwire.ToWire(e)
252 if wire.Receipt != nil || buffer.completion != nil {
253 buffer.messages = append(buffer.messages, &bufferedMessage{message: Message{
254 Role: "notice", Code: "turn_result", Level: "info", TurnID: e.TurnID,
255 CompletionReceipt: wire.Receipt, CompletionSummary: buffer.completion, CheckpointTurn: e.CheckpointTurn,
256 }})
257 }
258 }
259
260 func (buffer *Buffer) applyAssistantMessage(e event.Event) {
261 if e.Text == "" && e.Reasoning == "" && len(e.MemoryCitations) == 0 {
262 return
263 }
264 hm := ensureIdentifiedDisplayAssistant(buffer, e.MessageID, false)
265 if e.Text != "" {
266 hm.content.replace(e.Text)
267 }
268 if e.Reasoning != "" {
269 hm.reasoning.replace(e.Reasoning)
270 }
271 if len(e.MemoryCitations) > 0 {
272 hm.message.MemoryCitations = append([]provider.MemoryCitation(nil), e.MemoryCitations...)
273 }
274 }
275
276 func (buffer *Buffer) applyToolResult(e event.Event) {
277 callID := strings.TrimSpace(e.Tool.ID)
278 content := firstNonEmpty(e.Tool.Output, e.Tool.Err)
279 display, errPreview := buffer.Format.result(content, e.Tool.Err != "")
280 if callID != "" {
281 updateBufferedToolCallSummary(buffer, callID, content)
282 for _, row := range buffer.messages {
283 for i := range row.message.ToolCalls {
284 if row.message.ToolCalls[i].ID == callID {
285 row.message.ToolCalls[i].Pending = false
286 }
287 }
288 }
289 }
290 toolName := e.Tool.Name
291 if toolName == "" && buffer.tools != nil {
292 toolName = buffer.tools[callID]
293 }
294 result := Message{
295 Role: "tool",
296 ToolCallID: callID,
297 ToolName: toolName,
298 Content: display,
299 ToolResultError: errPreview,
300 PresentedFiles: append([]provider.PresentedFile(nil), e.Tool.PresentedFiles...),
301 }
302 if callID != "" {
303 for _, row := range buffer.messages {
304 if row.message.Role == "tool" && row.message.ToolCallID == callID {
305 result.RecordID = row.message.RecordID
306 result.Source, result.TurnID = row.message.Source, row.message.TurnID
307 row.message = result
308 return
309 }
310 }
311 }
312 buffer.messages = append(buffer.messages, &bufferedMessage{message: result})
313 }
314
315 func (buffer *Buffer) applyNotice(e event.Event) {
316 if strings.TrimSpace(e.Text) == "" {
317 return
318 }
319 level := "info"
320 if e.Level == event.LevelWarn {
321 level = "warn"
322 }
323 buffer.messages = append(buffer.messages, &bufferedMessage{message: Message{
324 Role: "notice",
325 MessageID: e.MessageID,
326 Level: level,
327 Content: e.Text,
328 Detail: e.Detail,
329 Code: e.Code,
330 DecisionReceipt: cloneDecisionReceipt(e.DecisionReceipt),
331 }})
332 }
333
334 func recordHistoryToolDispatch(buffer *Buffer, e event.Event) {
335 if strings.TrimSpace(e.Tool.Name) == "" && e.Tool.ID == "" {
336 return
337 }
338 hm := ensureIdentifiedDisplayAssistant(buffer, e.MessageID, true)
339 resolvedReadOnly := e.Tool.ReadOnly
340 call := ToolCall{
341 Partial: e.Tool.Partial,
342 ArgChars: e.Tool.ArgChars,
343 Pending: true,
344 ParentID: e.Tool.ParentID,
345 StartedAt: e.Tool.StartedAt,
346 ID: e.Tool.ID,
347 Name: e.Tool.Name,
348 Arguments: e.Tool.Args,
349 ResolvedName: e.Tool.ResolvedName,
350 CapabilityID: e.Tool.CapabilityID,
351 ResolvedReadOnly: &resolvedReadOnly,
352 Subject: buffer.Format.subject(e.Tool.Name, e.Tool.Args),
353 Summary: buffer.Format.summary(e.Tool.Name, e.Tool.Args, ""),
354 Diff: e.Tool.Diff,
355 Added: e.Tool.Added,
356 Removed: e.Tool.Removed,
357 }
358 replaced := false
359 if call.ID != "" {
360 for i := range hm.message.ToolCalls {
361 if hm.message.ToolCalls[i].ID == call.ID {
362 if call.Partial {
363 previous := hm.message.ToolCalls[i]
364 // A progress update cannot reopen a committed invocation.
365 if !previous.Partial {
366 return
367 }
368 if call.Name == "" {
369 call.Name = previous.Name
370 }
371 if call.Arguments == "" {
372 call.Arguments = previous.Arguments
373 }
374 }
375 hm.message.ToolCalls[i] = call
376 replaced = true
377 break
378 }
379 }
380 if buffer.tools == nil {
381 buffer.tools = map[string]string{}
382 }
383 buffer.tools[call.ID] = call.Name
384 }
385 if !replaced {
386 hm.message.ToolCalls = append(hm.message.ToolCalls, call)
387 }
388 }
389
390 func ensureIdentifiedDisplayAssistant(buffer *Buffer, messageID string, tool bool) *bufferedMessage {
391 if messageID == "" {
392 if tool {
393 return ensureDisplayAssistantForTool(buffer)
394 }
395 return ensureDisplayAssistant(buffer)
396 }
397 if existing := buffer.byMessageID[messageID]; existing != nil {
398 return existing
399 }
400 if buffer.byMessageID == nil {
401 buffer.byMessageID = make(map[string]*bufferedMessage)
402 }
403 message := &bufferedMessage{message: Message{MessageID: messageID, Role: "assistant"}}
404 buffer.messages = append(buffer.messages, message)
405 buffer.byMessageID[messageID] = message
406 return message
407 }
408 func ensureDisplayAssistant(buffer *Buffer) *bufferedMessage {
409 if n := len(buffer.messages); n > 0 && buffer.messages[n-1].message.Role == "assistant" {
410 return buffer.messages[n-1]
411 }
412 message := &bufferedMessage{message: Message{Role: "assistant"}}
413 buffer.messages = append(buffer.messages, message)
414 return message
415 }
416
417 func ensureDisplayAssistantForTool(buffer *Buffer) *bufferedMessage {
418 if n := len(buffer.messages); n > 0 && buffer.messages[n-1].message.Role == "assistant" && !buffer.messages[n-1].content.hasNonWhitespace() {
419 return buffer.messages[n-1]
420 }
421 message := &bufferedMessage{message: Message{Role: "assistant"}}
422 buffer.messages = append(buffer.messages, message)
423 return message
424 }
425
426 func updateBufferedToolCallSummary(buffer *Buffer, callID, output string) {
427 if callID == "" {
428 return
429 }
430 for _, v := range slices.Backward(buffer.messages) {
431 for j := range v.message.ToolCalls {
432 call := &v.message.ToolCalls[j]
433 if call.ID != callID {
434 continue
435 }
436 if call.Summary == "" {
437 call.Summary = buffer.Format.summary(call.Name, call.Arguments, output)
438 }
439 return
440 }
441 }
442 }
443
444 type Formatter struct {
445 ToolSubject func(name, args string) string
446 ToolSummary func(name, args, output string) string
447 ToolResult func(content string, failed bool) (string, string)
448 }
449
450 func (f Formatter) subject(name, args string) string {
451 if f.ToolSubject != nil {
452 return f.ToolSubject(name, args)
453 }
454 return ""
455 }
456 func (f Formatter) summary(name, args, output string) string {
457 if f.ToolSummary != nil {
458 return f.ToolSummary(name, args, output)
459 }
460 return ""
461 }
462 func (f Formatter) result(content string, failed bool) (string, string) {
463 if f.ToolResult != nil {
464 return f.ToolResult(content, failed)
465 }
466 if failed {
467 return content, content
468 }
469 return content, ""
470 }
471 func (buffer *Buffer) ResetToolsIfEmpty() {
472 if len(buffer.messages) == 0 {
473 buffer.tools = nil
474 }
475 }
476 func firstNonEmpty(values ...string) string {
477 for _, value := range values {
478 if value != "" {
479 return value
480 }
481 }
482 return ""
483 }
484 func cloneDecisionReceipt(in *provider.DecisionReceipt) *provider.DecisionReceipt {
485 if in == nil {
486 return nil
487 }
488 out := *in
489 return &out
490 }
491
491 lines GO