| 1 | package control |
| 2 | |
| 3 | import ( |
| 4 | "reasonix/internal/agent" |
| 5 | "reasonix/internal/provider" |
| 6 | "strings" |
| 7 | "time" |
| 8 | ) |
| 9 | |
| 10 | // planCancelledMessages preserves the existing provider replay policy without |
| 11 | // mutating the executor, storage, or UI. Its result belongs to the caller. |
| 12 | func planCancelledMessages(msgs []provider.Message, idx int, fallback provider.Message, startedAt time.Time, canReplay func(provider.Message) bool, evidence *interruptedTailEvidence) []provider.Message { |
| 13 | if start, ok := resolveInterruptedTurnStart(msgs, idx, true, startedAt, fallback); ok { |
| 14 | idx = start |
| 15 | } |
| 16 | if idx < 0 { |
| 17 | idx = 0 |
| 18 | } |
| 19 | if idx > len(msgs) { |
| 20 | idx = len(msgs) |
| 21 | } |
| 22 | next := append([]provider.Message{}, msgs[:idx]...) |
| 23 | keptUser := false |
| 24 | userEnd := idx |
| 25 | for i, m := range msgs[idx:] { |
| 26 | if !agent.IsUserAuthoredTurnMessage(m) { |
| 27 | continue |
| 28 | } |
| 29 | m.Content = StripComposePrefixes(m.Content) |
| 30 | next = append(next, m) |
| 31 | keptUser = true |
| 32 | userEnd = idx + i + 1 |
| 33 | break |
| 34 | } |
| 35 | if !keptUser && agent.IsUserAuthoredTurnMessage(fallback) { |
| 36 | fallback.Content = StripComposePrefixes(fallback.Content) |
| 37 | if strings.TrimSpace(fallback.Content) != "" { |
| 38 | fallback.Images = append([]string(nil), fallback.Images...) |
| 39 | next = append(next, fallback) |
| 40 | keptUser = true |
| 41 | userEnd = idx |
| 42 | } |
| 43 | } |
| 44 | if !keptUser && len(msgs) <= idx { |
| 45 | return nil |
| 46 | } |
| 47 | recovery := &provider.InterruptedTurnRecovery{Pending: true} |
| 48 | localIndexes := make([]int, 0, 1) |
| 49 | for i := userEnd; i < len(msgs); { |
| 50 | m := msgs[i] |
| 51 | if m.LocalOnly { |
| 52 | m.Role = provider.RoleTool |
| 53 | m.ToolCallID = provider.LocalOnlyToolID |
| 54 | m.Name = provider.LocalOnlyToolName |
| 55 | previousRecovery := m.InterruptedTurn |
| 56 | m.InterruptedTurn = nil |
| 57 | next = append(next, m) |
| 58 | localIndexes = append(localIndexes, len(next)-1) |
| 59 | recovery.DroppedPartialText = recovery.DroppedPartialText || strings.TrimSpace(m.Content) != "" |
| 60 | recovery.DroppedPartialReasoning = recovery.DroppedPartialReasoning || strings.TrimSpace(m.ReasoningContent) != "" |
| 61 | if previousRecovery != nil { |
| 62 | mergeInterruptedRecovery(recovery, previousRecovery) |
| 63 | } else { |
| 64 | for _, call := range m.ToolCalls { |
| 65 | provider.RecordToolRecovery(recovery, interruptedToolSummary(call), provider.ToolRunUnknown) |
| 66 | } |
| 67 | } |
| 68 | i++ |
| 69 | continue |
| 70 | } |
| 71 | // Keep compaction digests between the pinned input and recent tool tail; |
| 72 | // their summarized work is no longer available verbatim. |
| 73 | if agent.IsCompactionSummary(m) { |
| 74 | next = append(next, m) |
| 75 | i++ |
| 76 | continue |
| 77 | } |
| 78 | if m.Role == provider.RoleAssistant { |
| 79 | recordInterruptedAssistantRecovery(recovery, msgs, i, evidence) |
| 80 | } |
| 81 | if end, ok := completeToolTurnEnd(msgs, i); ok && canReplay(m) { |
| 82 | next = append(next, msgs[i:end]...) |
| 83 | i = end |
| 84 | continue |
| 85 | } |
| 86 | switch m.Role { |
| 87 | case provider.RoleAssistant: |
| 88 | local := m |
| 89 | local.Role = provider.RoleTool |
| 90 | local.LocalOnly = true |
| 91 | local.ToolCallID = provider.LocalOnlyToolID |
| 92 | local.Name = provider.LocalOnlyToolName |
| 93 | local.InterruptedTurn = nil |
| 94 | next = append(next, local) |
| 95 | localIndexes = append(localIndexes, len(next)-1) |
| 96 | recovery.DroppedPartialText = recovery.DroppedPartialText || strings.TrimSpace(local.Content) != "" |
| 97 | recovery.DroppedPartialReasoning = recovery.DroppedPartialReasoning || strings.TrimSpace(local.ReasoningContent) != "" |
| 98 | case provider.RoleTool: |
| 99 | local := m |
| 100 | local.LocalOnly = true |
| 101 | local.ToolCalls = []provider.ToolCall{{ID: m.ToolCallID, Name: m.Name}} |
| 102 | local.ToolCallID = provider.LocalOnlyToolID |
| 103 | local.Name = provider.LocalOnlyToolName |
| 104 | next = append(next, local) |
| 105 | localIndexes = append(localIndexes, len(next)-1) |
| 106 | } |
| 107 | i++ |
| 108 | } |
| 109 | if len(localIndexes) == 0 { |
| 110 | next = append(next, provider.Message{ |
| 111 | Role: provider.RoleTool, ToolCallID: provider.LocalOnlyToolID, |
| 112 | Name: provider.LocalOnlyToolName, LocalOnly: true, |
| 113 | }) |
| 114 | localIndexes = append(localIndexes, len(next)-1) |
| 115 | } |
| 116 | if evidence != nil { |
| 117 | recovery.Cause = "runtime_restart" |
| 118 | recovery.TurnID = evidence.turnID |
| 119 | recovery.SilentInterruption = len(recovery.ToolCalls) == 0 && len(recovery.CompletedTools) == 0 && !recovery.DroppedPartialText && !recovery.DroppedPartialReasoning |
| 120 | } |
| 121 | next[localIndexes[len(localIndexes)-1]].InterruptedTurn = recovery |
| 122 | return next |
| 123 | } |
| 124 | |
| 125 | func mergeInterruptedRecovery(dst, src *provider.InterruptedTurnRecovery) { |
| 126 | if src.TerminalStatus != "" { |
| 127 | dst.TerminalStatus = src.TerminalStatus |
| 128 | dst.FailureDiagnostic = nil |
| 129 | } |
| 130 | if src.FailureDiagnostic != nil { |
| 131 | diagnostic := *src.FailureDiagnostic |
| 132 | dst.FailureDiagnostic = &diagnostic |
| 133 | } |
| 134 | dst.CompletedTools = append(dst.CompletedTools, src.CompletedTools...) |
| 135 | dst.InterruptedTools = append(dst.InterruptedTools, src.InterruptedTools...) |
| 136 | dst.NotStartedTools = append(dst.NotStartedTools, src.NotStartedTools...) |
| 137 | dst.UnknownTools = append(dst.UnknownTools, src.UnknownTools...) |
| 138 | } |
| 139 |