返回 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 // ToolResult is a completion delta, not a replacement canonical
306 // message. Preserve formal identity and update only event-owned
307 // fields; e.MessageID may identify the owning assistant message.
308 row.message.Content = display
309 row.message.ToolResultError = errPreview
310 row.message.Pending = false
311 if toolName != "" {
312 row.message.ToolName = toolName
313 }
314 if len(e.Tool.PresentedFiles) > 0 {
315 row.message.PresentedFiles = append([]provider.PresentedFile(nil), e.Tool.PresentedFiles...)
316 }
317 return
318 }
319 }
320 }
321 buffer.messages = append(buffer.messages, &bufferedMessage{message: result})
322 }
323
324 func (buffer *Buffer) applyNotice(e event.Event) {
325 if strings.TrimSpace(e.Text) == "" {
326 return
327 }
328 recordID := ""
329 if e.Code == event.NoticeCodeUnappliedSteer && e.MessageID != "" {
330 recordID = unappliedSteerRecordID(e.MessageID)
331 for _, row := range buffer.messages {
332 if row.message.RecordID == recordID && row.message.Role == "notice" {
333 return
334 }
335 }
336 }
337 level := "info"
338 if e.Level == event.LevelWarn {
339 level = "warn"
340 }
341 buffer.messages = append(buffer.messages, &bufferedMessage{message: Message{
342 Role: "notice",
343 RecordID: recordID,
344 MessageID: e.MessageID,
345 Level: level,
346 Content: e.Text,
347 Detail: e.Detail,
348 Code: e.Code,
349 DecisionReceipt: cloneDecisionReceipt(e.DecisionReceipt),
350 }})
351 }
352
353 func recordHistoryToolDispatch(buffer *Buffer, e event.Event) {
354 if strings.TrimSpace(e.Tool.Name) == "" && e.Tool.ID == "" {
355 return
356 }
357 hm := ensureIdentifiedDisplayAssistant(buffer, e.MessageID, true)
358 resolvedReadOnly := e.Tool.ReadOnly
359 call := ToolCall{
360 Partial: e.Tool.Partial,
361 ArgChars: e.Tool.ArgChars,
362 Pending: true,
363 ParentID: e.Tool.ParentID,
364 StartedAt: e.Tool.StartedAt,
365 ID: e.Tool.ID,
366 Name: e.Tool.Name,
367 Arguments: e.Tool.Args,
368 ResolvedName: e.Tool.ResolvedName,
369 CapabilityID: e.Tool.CapabilityID,
370 ResolvedReadOnly: &resolvedReadOnly,
371 Subject: buffer.Format.subject(e.Tool.Name, e.Tool.Args),
372 Summary: buffer.Format.summary(e.Tool.Name, e.Tool.Args, ""),
373 Diff: e.Tool.Diff,
374 Added: e.Tool.Added,
375 Removed: e.Tool.Removed,
376 }
377 replaced := false
378 if call.ID != "" {
379 for i := range hm.message.ToolCalls {
380 if hm.message.ToolCalls[i].ID == call.ID {
381 if call.Partial {
382 previous := hm.message.ToolCalls[i]
383 // A progress update cannot reopen a committed invocation.
384 if !previous.Partial {
385 return
386 }
387 if call.Name == "" {
388 call.Name = previous.Name
389 }
390 if call.Arguments == "" {
391 call.Arguments = previous.Arguments
392 }
393 }
394 hm.message.ToolCalls[i] = call
395 replaced = true
396 break
397 }
398 }
399 if buffer.tools == nil {
400 buffer.tools = map[string]string{}
401 }
402 buffer.tools[call.ID] = call.Name
403 }
404 if !replaced {
405 hm.message.ToolCalls = append(hm.message.ToolCalls, call)
406 }
407 }
408
409 func ensureIdentifiedDisplayAssistant(buffer *Buffer, messageID string, tool bool) *bufferedMessage {
410 if messageID == "" {
411 if tool {
412 return ensureDisplayAssistantForTool(buffer)
413 }
414 return ensureDisplayAssistant(buffer)
415 }
416 if existing := buffer.byMessageID[messageID]; existing != nil {
417 return existing
418 }
419 if buffer.byMessageID == nil {
420 buffer.byMessageID = make(map[string]*bufferedMessage)
421 }
422 message := &bufferedMessage{message: Message{MessageID: messageID, Role: "assistant"}}
423 buffer.messages = append(buffer.messages, message)
424 buffer.byMessageID[messageID] = message
425 return message
426 }
427 func ensureDisplayAssistant(buffer *Buffer) *bufferedMessage {
428 if n := len(buffer.messages); n > 0 && buffer.messages[n-1].message.Role == "assistant" {
429 return buffer.messages[n-1]
430 }
431 message := &bufferedMessage{message: Message{Role: "assistant"}}
432 buffer.messages = append(buffer.messages, message)
433 return message
434 }
435
436 func ensureDisplayAssistantForTool(buffer *Buffer) *bufferedMessage {
437 if n := len(buffer.messages); n > 0 && buffer.messages[n-1].message.Role == "assistant" && !buffer.messages[n-1].content.hasNonWhitespace() {
438 return buffer.messages[n-1]
439 }
440 message := &bufferedMessage{message: Message{Role: "assistant"}}
441 buffer.messages = append(buffer.messages, message)
442 return message
443 }
444
445 func updateBufferedToolCallSummary(buffer *Buffer, callID, output string) {
446 if callID == "" {
447 return
448 }
449 for _, v := range slices.Backward(buffer.messages) {
450 for j := range v.message.ToolCalls {
451 call := &v.message.ToolCalls[j]
452 if call.ID != callID {
453 continue
454 }
455 if call.Summary == "" {
456 call.Summary = buffer.Format.summary(call.Name, call.Arguments, output)
457 }
458 return
459 }
460 }
461 }
462
463 type Formatter struct {
464 ToolSubject func(name, args string) string
465 ToolSummary func(name, args, output string) string
466 ToolResult func(content string, failed bool) (string, string)
467 }
468
469 func (f Formatter) subject(name, args string) string {
470 if f.ToolSubject != nil {
471 return f.ToolSubject(name, args)
472 }
473 return ""
474 }
475 func (f Formatter) summary(name, args, output string) string {
476 if f.ToolSummary != nil {
477 return f.ToolSummary(name, args, output)
478 }
479 return ""
480 }
481 func (f Formatter) result(content string, failed bool) (string, string) {
482 if f.ToolResult != nil {
483 return f.ToolResult(content, failed)
484 }
485 if failed {
486 return content, content
487 }
488 return content, ""
489 }
490 func (buffer *Buffer) ResetToolsIfEmpty() {
491 if len(buffer.messages) == 0 {
492 buffer.tools = nil
493 }
494 }
495 func firstNonEmpty(values ...string) string {
496 for _, value := range values {
497 if value != "" {
498 return value
499 }
500 }
501 return ""
502 }
503 func cloneDecisionReceipt(in *provider.DecisionReceipt) *provider.DecisionReceipt {
504 if in == nil {
505 return nil
506 }
507 out := *in
508 return &out
509 }
510
510 lines GO