返回 DeepSeek-Reasonix
projection.go
根目录 / internal / session / projection.go
1 package session
2
3 import (
4 "bytes"
5 "encoding/json"
6 "fmt"
7 "io"
8 "maps"
9 "slices"
10 "strings"
11
12 "reasonix/internal/event"
13 "reasonix/internal/permissionpreset"
14 "reasonix/internal/provider"
15 )
16
17 type Projection struct {
18 Submissions SubmissionIndex
19 TranscriptInputs []transcriptInput
20 HiddenTurns map[string]bool
21 RetractedInputs map[string]string
22 // RejectedToolResults retains local execution evidence independently of the
23 // wire projection. Result message IDs scope evidence even when call IDs repeat.
24 RejectedToolResults map[string]rejectedToolResult `json:"rejectedToolResults,omitempty"`
25 CommittedSequence uint64
26 TurnID string
27 TurnStatus event.TurnStatus
28 CurrentTurnStart uint64
29 CurrentTurnStartedAt int64
30 // CurrentTurnMessageID is the stable identity of the newest assistant
31 // message committed inside the open turn. It becomes the turn's final reply
32 // identity when the turn closes.
33 CurrentTurnMessageID string
34 CurrentAttempts map[string]bool
35 CurrentCalls map[string]bool
36 Turns []TurnBoundary
37 Messages []provider.Message
38 // ModelMessages is the exact provider-visible projection. Canonical Messages
39 // remains the complete UI/history transcript; compaction replaces only this
40 // view and never deletes the underlying business history.
41 ModelMessages []provider.Message
42 Todos []event.Todo
43 TodoWritten bool
44 Interactions map[string]string
45 ActiveTools map[string]string
46 StartedTools map[string]bool
47 ActiveSteps map[string]bool
48 Recovery *event.RecoveryStatus
49 PlanState json.RawMessage
50 GoalState json.RawMessage
51 PermissionPreset string
52 // One-based event identity; zero means no explicit preset in the log.
53 PermissionPresetSequence uint64
54 Title string
55 // TitleSequence is the sequence of the latest accepted session/title event.
56 // It is independent from CommittedSequence so ordinary chat appends do not
57 // conflict with a delayed title mutation.
58 TitleSequence uint64
59 ModelRef string
60 ModelIdentity string
61 // StreamCheckpoint is the open turn's newest uncommitted stream output;
62 // it is replaced, never mutated, so shallow projection copies may share it.
63 StreamCheckpoint *provider.Message `json:"streamCheckpoint,omitempty"`
64 }
65
66 type TurnBoundary struct {
67 SamplingCount int `json:"samplingCount,omitempty"`
68 ToolCount int `json:"toolCount,omitempty"`
69 DurationMs int64 `json:"durationMs,omitempty"`
70 TurnID string `json:"turnId"`
71 StartSequence uint64 `json:"startSequence"`
72 EndSequence uint64 `json:"endSequence"`
73 Status event.TurnStatus `json:"status"`
74 // BoundarySequence is the last sequence of the commit that closed this turn.
75 // A cut may only land here: a turn end and the state ending with it can share
76 // one commit, and a cut inside that commit inherits half an operation.
77 BoundarySequence uint64 `json:"boundarySequence"`
78 // Availability is fixed from the complete commit that closed the turn. It
79 // must not be recomputed from the latest projection: a later commit may
80 // resolve authority that the earlier fork prefix would still inherit.
81 Availability ForkAvailability `json:"availability"`
82 // MessageID is the stable transcript identity of the turn's final reply, empty
83 // when the turn committed none. Surfaces match turns to messages through this
84 // identity, never through an array position.
85 MessageID string `json:"messageId,omitempty"`
86 }
87
88 var ProjectionKinds = map[string]bool{
89 "message/complete": true, "message/upsert": true, "message/retract": true, "assistant/attempt": true,
90 "tool/call": true, "tool/start": true, "tool/result": true,
91 "turn/start": true, "turn/end": true, "step/start": true, "step/end": true,
92 "todo/write": true, "interaction/created": true, "interaction/resolved": true,
93 "plan/state": true, "goal/state": true, "session/title": true, "session/config": true, "session/permission-preset": true,
94 "model/context-replace": true, "history/replace": true,
95 "compaction": true, "runtime/recovery": true, "legacy/import": true,
96 "diagnostic": true,
97 }
98
99 var PrototypeProjectionKinds = func() map[string]bool {
100 kinds := make(map[string]bool, len(ProjectionKinds)+1)
101 maps.Copy(kinds, ProjectionKinds)
102 kinds["context/replace"] = true
103 return kinds
104 }()
105
106 func Project(commits []Commit) (Projection, error) {
107 projection := Projection{Todos: []event.Todo{}, Interactions: map[string]string{}, ActiveTools: map[string]string{}}
108 for _, commit := range commits {
109 if err := applyProjectionCommit(&projection, commit); err != nil {
110 return Projection{}, err
111 }
112 }
113 return projection, nil
114 }
115
116 func applyProjectionCommit(projection *Projection, commit Commit) error {
117 initializeProjectionMaps(projection)
118 return applyProjectionEvents(projection, commit)
119 }
120
121 func initializeProjectionMaps(projection *Projection) {
122 if projection.Interactions == nil {
123 projection.Interactions = map[string]string{}
124 }
125 if projection.ActiveTools == nil {
126 projection.ActiveTools = map[string]string{}
127 }
128 if projection.StartedTools == nil {
129 projection.StartedTools = map[string]bool{}
130 }
131 if projection.ActiveSteps == nil {
132 projection.ActiveSteps = map[string]bool{}
133 }
134 }
135
136 func applyProjectionEvents(projection *Projection, commit Commit) error {
137 closedBefore := len(projection.Turns)
138 for _, ev := range commit.Events {
139 projection.CommittedSequence = ev.Sequence
140 supersedeStreamCheckpoint(projection, ev.Kind)
141 var err error
142 switch ev.Kind {
143 case StreamCheckpointKind:
144 projectStreamCheckpoint(projection, ev)
145 case "submission/accepted":
146 err = projectSubmission(projection, commit, ev)
147 case "legacy/import":
148 err = projectLegacyImport(projection, commit, ev)
149 case "message/complete":
150 err = projectMessageComplete(projection, commit, ev)
151 case "message/upsert":
152 err = projectMessageUpsert(projection, commit, ev)
153 case "message/retract":
154 err = projectMessageRetract(projection, commit, ev)
155 case "assistant/attempt":
156 err = projectAssistantAttempt(projection, commit, ev)
157 case "history/replace":
158 err = projectHistoryReplace(projection, commit, ev)
159 case "model/context-replace":
160 err = projectModelContextReplace(projection, commit, ev)
161 case "session/title":
162 err = projectSessionTitle(projection, commit, ev)
163 case "session/config", "session/permission-preset":
164 err = projectSessionSettings(projection, commit, ev)
165 case "compaction":
166 err = projectCompaction(projection, commit, ev)
167 case "turn/start":
168 err = projectTurnStart(projection, commit, ev)
169 case "step/start", "step/end":
170 err = projectStepStart(projection, commit, ev)
171 case "tool/call":
172 err = projectToolCall(projection, commit, ev)
173 case "tool/start":
174 err = projectToolStart(projection, commit, ev)
175 case "tool/result":
176 err = projectToolResult(projection, commit, ev)
177 case "todo/write":
178 err = projectTodoWrite(projection, commit, ev)
179 case "interaction/created":
180 err = projectInteractionCreated(projection, commit, ev)
181 case "interaction/resolved":
182 err = projectInteractionResolved(projection, commit, ev)
183 case "runtime/recovery":
184 err = projectRuntimeRecovery(projection, commit, ev)
185 case "plan/state":
186 err = projectPlanState(projection, commit, ev)
187 case "goal/state":
188 err = projectGoalState(projection, commit, ev)
189 case "diagnostic":
190 err = projectDiagnostic(projection, commit, ev)
191 case "turn/end":
192 err = projectTurnEnd(projection, commit, ev)
193 }
194 if err != nil {
195 return err
196 }
197 if ev.Kind != "message/complete" {
198 applyTranscriptMetadata(projection, commit, ev)
199 }
200 }
201 // turn/end can be followed by more events in the same atomic commit. Only
202 // after the whole commit is projected do we know whether its cut leaves a
203 // turn, interaction, or tool authority open.
204 for index := closedBefore; index < len(projection.Turns); index++ {
205 projection.Turns[index].Availability = forkProjectionAvailability(*projection, projection.Turns[index].BoundarySequence)
206 }
207 return nil
208 }
209
210 func projectLegacyImport(projection *Projection, commit Commit, ev Event) error {
211 var body legacyImportPayload
212 if err := strictPayload(ev.Payload, &body); err != nil || body.Messages == nil {
213 return damagedPayload(ev, err)
214 }
215 messages := firstMessageOccurrences(body.Messages)
216 projection.Messages = append([]provider.Message{}, messages...)
217 projection.ModelMessages = append([]provider.Message{}, provider.ModelMessages(messages)...)
218 projection.GoalState = cloneRaw(body.Goal)
219 projection.ModelRef = strings.TrimSpace(body.ModelRef)
220 projection.ModelIdentity = strings.TrimSpace(body.ModelIdentity)
221 return nil
222 }
223
224 // projectMessageComplete keeps an id's first message. The writer refuses a
225 // repeat; one already on disk is skipped, metadata included, so the log opens.
226 func projectMessageComplete(projection *Projection, commit Commit, ev Event) error {
227 var body struct {
228 Message *provider.Message `json:"message"`
229 }
230 if err := strictPayload(ev.Payload, &body); err != nil || body.Message == nil || body.Message.ID == "" {
231 return damagedPayload(ev, err)
232 }
233 if projectionMessageIndex(projection.Messages, body.Message.ID) >= 0 {
234 return nil
235 }
236 projection.Messages = append(projection.Messages, *body.Message)
237 projection.ModelMessages = append(projection.ModelMessages, provider.ModelMessages([]provider.Message{*body.Message})...)
238 projection.recordTurnReply(*body.Message)
239 applyTranscriptMetadata(projection, commit, ev)
240 return nil
241 }
242
243 // recordTurnReply keeps the open turn's final answer identity. A fork entry
244 // belongs on the turn's answer, so a trailing tool call, a retried attempt, or a
245 // host-generated protocol message must not take the anchor away from the text a
246 // user actually reads.
247 func (projection *Projection) recordTurnReply(message provider.Message) {
248 if projection.TurnID == "" || message.Role != provider.RoleAssistant || message.LocalOnly {
249 return
250 }
251 if strings.TrimSpace(message.RawContent) == "" && strings.TrimSpace(message.Content) == "" {
252 return
253 }
254 projection.CurrentTurnMessageID = message.ID
255 }
256
257 func projectMessageUpsert(projection *Projection, commit Commit, ev Event) error {
258 var body struct {
259 Message *provider.Message `json:"message"`
260 }
261 if err := strictPayload(ev.Payload, &body); err != nil || body.Message == nil || body.Message.ID == "" {
262 return damagedPayload(ev, err)
263 }
264 if !replaceProjectionMessage(projection.Messages, *body.Message) {
265 projection.Messages = append(projection.Messages, *body.Message)
266 }
267 visible := provider.ModelMessages([]provider.Message{*body.Message})
268 modelIndex := projectionMessageIndex(projection.ModelMessages, body.Message.ID)
269 if modelIndex >= 0 {
270 if len(visible) == 0 {
271 projection.ModelMessages = append(projection.ModelMessages[:modelIndex], projection.ModelMessages[modelIndex+1:]...)
272 } else {
273 projection.ModelMessages[modelIndex] = visible[0]
274 }
275 } else if len(visible) > 0 {
276 // Upserts are normally metadata changes to an existing message. The
277 // append case is retained for explicitly-created records.
278 projection.ModelMessages = append(projection.ModelMessages, visible[0])
279 }
280 if projection.CurrentTurnMessageID == body.Message.ID && (body.Message.LocalOnly || body.Message.Role != provider.RoleAssistant) {
281 projection.CurrentTurnMessageID = ""
282 }
283 projection.recordTurnReply(*body.Message)
284 return nil
285 }
286
287 func projectMessageRetract(projection *Projection, commit Commit, ev Event) error {
288 ids, err := retractedMessageIDs(ev, ev.Payload)
289 if err != nil {
290 return err
291 }
292 removed := make(map[string]bool, len(ids))
293 for _, id := range ids {
294 removed[id] = true
295 }
296 filter := func(messages []provider.Message) []provider.Message {
297 return slices.DeleteFunc(messages, func(message provider.Message) bool { return removed[message.ID] })
298 }
299 projection.Messages = filter(projection.Messages)
300 projection.ModelMessages = filter(projection.ModelMessages)
301 if removed[projection.CurrentTurnMessageID] {
302 projection.CurrentTurnMessageID = ""
303 }
304 return nil
305 }
306
307 func projectAssistantAttempt(projection *Projection, commit Commit, ev Event) error {
308 var body struct {
309 ID string `json:"id"`
310 MessageID string `json:"messageId,omitempty"`
311 Action string `json:"action"`
312 Attempt int `json:"attempt,omitempty"`
313 Max int `json:"max,omitempty"`
314 Reason string `json:"reason,omitempty"`
315 }
316 if err := strictPayload(ev.Payload, &body); err != nil || body.ID == "" || (body.Action != "begin" && body.Action != "discard" && body.Action != "commit") {
317 return damagedPayload(ev, err)
318 }
319 if projection.TurnID != "" && body.Action == "begin" {
320 if projection.CurrentAttempts == nil {
321 projection.CurrentAttempts = map[string]bool{}
322 }
323 projection.CurrentAttempts[body.ID] = true
324 }
325 return nil
326 }
327
328 func projectHistoryReplace(projection *Projection, commit Commit, ev Event) error {
329 var body historyReplacePayload
330 if err := strictPayload(ev.Payload, &body); err != nil || body.Messages == nil {
331 return damagedPayload(ev, err)
332 }
333 messages := firstMessageOccurrences(body.Messages)
334 projection.Messages = append([]provider.Message(nil), messages...)
335 projection.ModelMessages = append([]provider.Message(nil), provider.ModelMessages(messages)...)
336 return nil
337 }
338
339 func projectModelContextReplace(projection *Projection, commit Commit, ev Event) error {
340 var body historyReplacePayload
341 if err := strictPayload(ev.Payload, &body); err != nil || body.Messages == nil {
342 return damagedPayload(ev, err)
343 }
344 projection.ModelMessages = append([]provider.Message(nil), body.Messages...)
345 return nil
346 }
347
348 func projectSessionTitle(projection *Projection, commit Commit, ev Event) error {
349 var body struct {
350 Title string `json:"title"`
351 }
352 if err := strictPayload(ev.Payload, &body); err != nil {
353 return damagedPayload(ev, err)
354 }
355 projection.Title = body.Title
356 // Sequence zero is a valid first event, so store the one-based identity;
357 // zero remains the durable "no title event yet" revision.
358 projection.TitleSequence = ev.Sequence + 1
359 return nil
360 }
361
362 func projectSessionConfig(projection *Projection, commit Commit, ev Event) error {
363 var body struct {
364 ModelRef string `json:"modelRef"`
365 ModelIdentity string `json:"modelIdentity,omitempty"`
366 }
367 if err := strictPayload(ev.Payload, &body); err != nil || strings.TrimSpace(body.ModelRef) == "" {
368 return damagedPayload(ev, err)
369 }
370 projection.ModelRef = strings.TrimSpace(body.ModelRef)
371 projection.ModelIdentity = strings.TrimSpace(body.ModelIdentity)
372 return nil
373 }
374
375 func projectSessionSettings(projection *Projection, commit Commit, ev Event) error {
376 if ev.Kind == "session/config" {
377 return projectSessionConfig(projection, commit, ev)
378 }
379 return projectSessionPermissionPreset(projection, commit, ev)
380 }
381
382 func projectSessionPermissionPreset(projection *Projection, _ Commit, ev Event) error {
383 var body struct {
384 Preset string `json:"preset"`
385 }
386 if err := strictPayload(ev.Payload, &body); err != nil {
387 return damagedPayload(ev, err)
388 }
389 if !permissionpreset.Valid(body.Preset) {
390 return damagedPayload(ev, fmt.Errorf("invalid permission preset %q", body.Preset))
391 }
392 projection.PermissionPreset = body.Preset
393 projection.PermissionPresetSequence = ev.Sequence
394 return nil
395 }
396
397 func projectCompaction(projection *Projection, commit Commit, ev Event) error {
398 var body struct {
399 Messages []provider.Message `json:"messages"`
400 Trigger string `json:"trigger,omitempty"`
401 Sources []uint64 `json:"sourceSequences,omitempty"`
402 }
403 if err := strictPayload(ev.Payload, &body); err != nil || body.Messages == nil {
404 return damagedPayload(ev, err)
405 }
406 projection.ModelMessages = append([]provider.Message(nil), body.Messages...)
407 return nil
408 }
409
410 func projectTurnStart(projection *Projection, commit Commit, ev Event) error {
411 projection.TurnID = commit.TurnID
412 projection.TurnStatus = event.TurnInProgress
413 projection.CurrentTurnStart = ev.Sequence
414 projection.CurrentTurnStartedAt = commit.CreatedAt.UnixMilli()
415 projection.CurrentTurnMessageID = ""
416 projection.CurrentAttempts, projection.CurrentCalls = map[string]bool{}, map[string]bool{}
417 projection.Todos, projection.TodoWritten = []event.Todo{}, false
418 projection.Recovery = nil
419 return nil
420 }
421
422 func projectStepStart(projection *Projection, commit Commit, ev Event) error {
423 var body struct {
424 ID string `json:"id"`
425 Status string `json:"status,omitempty"`
426 }
427 if err := strictPayload(ev.Payload, &body); err != nil || body.ID == "" {
428 return damagedPayload(ev, err)
429 }
430 if ev.Kind == "step/start" {
431 projection.ActiveSteps[body.ID] = true
432 } else {
433 delete(projection.ActiveSteps, body.ID)
434 }
435 return nil
436 }
437
438 func projectToolCall(projection *Projection, commit Commit, ev Event) error {
439 var body struct {
440 ID string `json:"id"`
441 Name string `json:"name"`
442 Args string `json:"args,omitempty"`
443 RunState provider.ToolRunState `json:"runState,omitempty"`
444 Diagnostic json.RawMessage `json:"diagnostic,omitempty"`
445 ResolvedName string `json:"resolvedName,omitempty"`
446 CapabilityID string `json:"capabilityId,omitempty"`
447 ReadOnly bool `json:"readOnly,omitempty"`
448 Truncated bool `json:"truncated,omitempty"`
449 DurationMs int64 `json:"durationMs,omitempty"`
450 StartedAt int64 `json:"startedAt,omitempty"`
451 EndedAt int64 `json:"endedAt,omitempty"`
452 Partial bool `json:"partial,omitempty"`
453 ArgChars int `json:"argChars,omitempty"`
454 Refreshed bool `json:"refreshed,omitempty"`
455 ParentID string `json:"parentId,omitempty"`
456 AttemptID string `json:"attemptId,omitempty"`
457 SubagentRef string `json:"subagentRef,omitempty"`
458 SubagentStatus string `json:"subagentStatus,omitempty"`
459 SubagentErrorCode string `json:"subagentErrorCode,omitempty"`
460 SubagentRetryable bool `json:"subagentRetryable,omitempty"`
461 Diff string `json:"diff,omitempty"`
462 Added int `json:"added,omitempty"`
463 Removed int `json:"removed,omitempty"`
464 Profile json.RawMessage `json:"profile,omitempty"`
465 Execution json.RawMessage `json:"execution,omitempty"`
466 PresentedFiles []provider.PresentedFile `json:"presentedFiles,omitempty"`
467 WorkspaceMutation bool `json:"workspaceMutation,omitempty"`
468 WorkspacePaths []string `json:"workspacePaths,omitempty"`
469 WorkspaceAllPaths bool `json:"workspaceAllPaths,omitempty"`
470 }
471 if err := strictPayload(ev.Payload, &body); err != nil || body.ID == "" || body.Name == "" {
472 return damagedPayload(ev, err)
473 }
474 projection.ActiveTools[body.ID] = body.Name
475 if projection.TurnID != "" {
476 if projection.CurrentCalls == nil {
477 projection.CurrentCalls = map[string]bool{}
478 }
479 projection.CurrentCalls[body.ID] = true
480 }
481 return nil
482 }
483
484 func projectToolStart(projection *Projection, commit Commit, ev Event) error {
485 var body struct {
486 ID string `json:"id"`
487 Name string `json:"name"`
488 }
489 if err := strictPayload(ev.Payload, &body); err != nil || body.ID == "" || body.Name == "" {
490 return damagedPayload(ev, err)
491 }
492 projection.ActiveTools[body.ID] = body.Name
493 projection.StartedTools[body.ID] = true
494 return nil
495 }
496
497 func projectToolResult(projection *Projection, commit Commit, ev Event) error {
498 var body struct {
499 ID string `json:"id"`
500 Name string `json:"name"`
501 Args string `json:"args,omitempty"`
502 Error string `json:"error,omitempty"`
503 Output string `json:"output,omitempty"`
504 State string `json:"state,omitempty"`
505 RunState provider.ToolRunState `json:"runState,omitempty"`
506 Diagnostic json.RawMessage `json:"diagnostic,omitempty"`
507 ResolvedName string `json:"resolvedName,omitempty"`
508 CapabilityID string `json:"capabilityId,omitempty"`
509 ReadOnly bool `json:"readOnly,omitempty"`
510 Truncated bool `json:"truncated,omitempty"`
511 DurationMs int64 `json:"durationMs,omitempty"`
512 StartedAt int64 `json:"startedAt,omitempty"`
513 EndedAt int64 `json:"endedAt,omitempty"`
514 Partial bool `json:"partial,omitempty"`
515 ArgChars int `json:"argChars,omitempty"`
516 Refreshed bool `json:"refreshed,omitempty"`
517 ParentID string `json:"parentId,omitempty"`
518 AttemptID string `json:"attemptId,omitempty"`
519 SubagentRef string `json:"subagentRef,omitempty"`
520 SubagentStatus string `json:"subagentStatus,omitempty"`
521 SubagentErrorCode string `json:"subagentErrorCode,omitempty"`
522 SubagentRetryable bool `json:"subagentRetryable,omitempty"`
523 Diff string `json:"diff,omitempty"`
524 Added int `json:"added,omitempty"`
525 Removed int `json:"removed,omitempty"`
526 Profile json.RawMessage `json:"profile,omitempty"`
527 Execution json.RawMessage `json:"execution,omitempty"`
528 PresentedFiles []provider.PresentedFile `json:"presentedFiles,omitempty"`
529 Todos []event.Todo `json:"todos,omitempty"`
530 TodoWritten bool `json:"todoWritten,omitempty"`
531 WorkspaceMutation bool `json:"workspaceMutation,omitempty"`
532 WorkspacePaths []string `json:"workspacePaths,omitempty"`
533 WorkspaceAllPaths bool `json:"workspaceAllPaths,omitempty"`
534 }
535 if err := strictPayload(ev.Payload, &body); err != nil || body.ID == "" || body.Name == "" {
536 return damagedPayload(ev, err)
537 }
538 delete(projection.ActiveTools, body.ID)
539 delete(projection.StartedTools, body.ID)
540 return nil
541 }
542
543 func projectTodoWrite(projection *Projection, commit Commit, ev Event) error {
544 var body struct {
545 Todos []event.Todo `json:"todos"`
546 }
547 if err := strictPayload(ev.Payload, &body); err != nil || validateTodos(body.Todos) != nil {
548 return damagedPayload(ev, err)
549 }
550 projection.Todos, projection.TodoWritten = append([]event.Todo(nil), body.Todos...), true
551 return nil
552 }
553
554 func projectInteractionCreated(projection *Projection, commit Commit, ev Event) error {
555 var body struct {
556 ID string `json:"id"`
557 ToolCallID string `json:"toolCallId,omitempty"`
558 Kind string `json:"kind,omitempty"`
559 State string `json:"state,omitempty"`
560 SessionID string `json:"sessionId,omitempty"`
561 HeadID string `json:"headId,omitempty"`
562 TurnID string `json:"turnId,omitempty"`
563 RuntimeEpoch string `json:"runtimeEpoch,omitempty"`
564 }
565 if err := strictPayload(ev.Payload, &body); err != nil || body.ID == "" || (body.State != "" && body.State != "pending") {
566 return damagedPayload(ev, err)
567 }
568 projection.Interactions[body.ID] = "pending"
569 return nil
570 }
571
572 func projectInteractionResolved(projection *Projection, commit Commit, ev Event) error {
573 var body struct {
574 ID string `json:"id"`
575 State string `json:"state"`
576 }
577 if err := strictPayload(ev.Payload, &body); err != nil || body.ID == "" || !terminalInteractionState(body.State) {
578 return damagedPayload(ev, err)
579 }
580 delete(projection.Interactions, body.ID)
581 return nil
582 }
583
584 func projectRuntimeRecovery(projection *Projection, commit Commit, ev Event) error {
585 var body event.RecoveryStatus
586 if err := strictPayload(ev.Payload, &body); err != nil {
587 return damagedPayload(ev, err)
588 }
589 projection.Recovery = &body
590 return nil
591 }
592
593 func projectPlanState(projection *Projection, commit Commit, ev Event) error {
594 if !validJSONObject(ev.Payload) {
595 return damagedPayload(ev, nil)
596 }
597 projection.PlanState = cloneRaw(ev.Payload)
598 return nil
599 }
600
601 func projectGoalState(projection *Projection, commit Commit, ev Event) error {
602 if !validJSONObject(ev.Payload) {
603 return damagedPayload(ev, nil)
604 }
605 projection.GoalState = cloneRaw(ev.Payload)
606 return nil
607 }
608
609 func projectDiagnostic(projection *Projection, commit Commit, ev Event) error {
610 if len(ev.Payload) > 0 && !json.Valid(ev.Payload) {
611 return damagedPayload(ev, nil)
612 }
613 return nil
614 }
615
616 func projectTurnEnd(projection *Projection, commit Commit, ev Event) error {
617 var body struct {
618 Status event.TurnStatus `json:"status"`
619 }
620 if err := strictPayload(ev.Payload, &body); err != nil || !body.Status.Terminal() {
621 return damagedPayload(ev, err)
622 }
623 if projection.TurnID != "" && projection.CurrentTurnStart != 0 {
624 projection.Turns = append(projection.Turns, TurnBoundary{
625 SamplingCount: len(projection.CurrentAttempts), ToolCount: len(projection.CurrentCalls),
626 DurationMs: max(0, commit.CreatedAt.UnixMilli()-projection.CurrentTurnStartedAt),
627 TurnID: projection.TurnID, StartSequence: projection.CurrentTurnStart,
628 EndSequence: ev.Sequence, Status: body.Status,
629 BoundarySequence: commit.LastSequence(),
630 MessageID: projection.CurrentTurnMessageID,
631 })
632 }
633 projection.TurnID = ""
634 projection.CurrentTurnStart = 0
635 projection.CurrentTurnStartedAt = 0
636 projection.CurrentTurnMessageID = ""
637 projection.TurnStatus = body.Status
638 return nil
639 }
640
641 func replaceProjectionMessage(messages []provider.Message, replacement provider.Message) bool {
642 if index := projectionMessageIndex(messages, replacement.ID); index >= 0 {
643 messages[index] = replacement
644 return true
645 }
646 return false
647 }
648
649 func projectionMessageIndex(messages []provider.Message, id string) int {
650 for i := range messages {
651 if messages[i].ID == id {
652 return i
653 }
654 }
655 return -1
656 }
657
658 func cloneProjection(projection Projection) Projection {
659 projection.TranscriptInputs = append([]transcriptInput(nil), projection.TranscriptInputs...)
660 projection.HiddenTurns = maps.Clone(projection.HiddenTurns)
661 projection.RetractedInputs = maps.Clone(projection.RetractedInputs)
662 projection.RejectedToolResults = maps.Clone(projection.RejectedToolResults)
663 projection.CurrentAttempts = maps.Clone(projection.CurrentAttempts)
664 projection.CurrentCalls = maps.Clone(projection.CurrentCalls)
665 projection.Messages = append([]provider.Message(nil), projection.Messages...)
666 projection.ModelMessages = append([]provider.Message(nil), projection.ModelMessages...)
667 projection.Turns = append([]TurnBoundary(nil), projection.Turns...)
668 projection.Todos = append([]event.Todo(nil), projection.Todos...)
669 interactions := make(map[string]string, len(projection.Interactions))
670 maps.Copy(interactions, projection.Interactions)
671 projection.Interactions = interactions
672 tools := make(map[string]string, len(projection.ActiveTools))
673 maps.Copy(tools, projection.ActiveTools)
674 projection.ActiveTools = tools
675 projection.StartedTools = maps.Clone(projection.StartedTools)
676 projection.ActiveSteps = maps.Clone(projection.ActiveSteps)
677 if projection.Recovery != nil {
678 recovery := *projection.Recovery
679 projection.Recovery = &recovery
680 }
681 projection.PlanState = cloneRaw(projection.PlanState)
682 projection.GoalState = cloneRaw(projection.GoalState)
683 return projection
684 }
685
686 func strictPayload(payload json.RawMessage, target any) error {
687 if len(payload) == 0 {
688 return io.ErrUnexpectedEOF
689 }
690 decoder := json.NewDecoder(bytes.NewReader(payload))
691 decoder.DisallowUnknownFields()
692 if err := decoder.Decode(target); err != nil {
693 return err
694 }
695 var extra any
696 if err := decoder.Decode(&extra); err != io.EOF {
697 if err == nil {
698 return fmt.Errorf("multiple JSON values")
699 }
700 return err
701 }
702 return nil
703 }
704
705 func validateTodos(todos []event.Todo) error {
706 if todos == nil {
707 return fmt.Errorf("todos must be an array")
708 }
709 seen := make(map[string]bool, len(todos))
710 for i, todo := range todos {
711 content := strings.TrimSpace(todo.Content)
712 if content == "" || content != todo.Content || seen[content] {
713 return fmt.Errorf("todos[%d].content is invalid", i)
714 }
715 seen[content] = true
716 switch todo.Status {
717 case "pending", "in_progress", "completed":
718 default:
719 return fmt.Errorf("todos[%d].status is invalid", i)
720 }
721 }
722 return nil
723 }
724
725 func terminalInteractionState(state string) bool {
726 switch state {
727 case "answered", "rejected", "cancelled", "unavailable":
728 return true
729 default:
730 return false
731 }
732 }
733
734 func validJSONObject(raw json.RawMessage) bool {
735 var object map[string]json.RawMessage
736 return len(raw) > 0 && json.Unmarshal(raw, &object) == nil && object != nil
737 }
738
739 func damagedPayload(event Event, cause error) error {
740 if cause != nil {
741 return fmt.Errorf("%w: invalid %s payload at %d: %w", ErrDamagedStore, event.Kind, event.Sequence, cause)
742 }
743 return fmt.Errorf("%w: invalid %s payload at %d", ErrDamagedStore, event.Kind, event.Sequence)
744 }
745
746 func cloneRaw(raw json.RawMessage) json.RawMessage { return append(json.RawMessage(nil), raw...) }
747
747 lines GO