| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "errors" |
| 5 | "fmt" |
| 6 | "log/slog" |
| 7 | "strings" |
| 8 | |
| 9 | "reasonix/internal/agent" |
| 10 | "reasonix/internal/control" |
| 11 | "reasonix/internal/event" |
| 12 | "reasonix/internal/turnevent" |
| 13 | ) |
| 14 | |
| 15 | // TurnStartView is the synchronous admission receipt for the new Wails turn |
| 16 | // API. Events remain the streaming authority after admission. |
| 17 | type TurnStartView struct { |
| 18 | TurnID string `json:"turnId"` |
| 19 | Status event.TurnStatus `json:"status"` |
| 20 | Disposition control.SubmitDisposition `json:"disposition"` |
| 21 | OperationID string `json:"operationId,omitempty"` |
| 22 | // A management refusal is correlated with its existing operation, not a failed chat turn. |
| 23 | ManagementErrorCode string `json:"managementErrorCode,omitempty"` |
| 24 | RuntimeEpoch string `json:"runtimeEpoch,omitempty"` |
| 25 | SubmissionID string `json:"submissionId,omitempty"` |
| 26 | } |
| 27 | |
| 28 | // validatePromptIdentity fences a decision to the runtime and turn that |
| 29 | // rendered its card. It is intentionally shared by every decision surface; |
| 30 | // callers must still resolve the prompt on the same controller instance. |
| 31 | func (a *App) validatePromptIdentity(tabID, turnID, runtimeEpoch string) (control.SessionAPI, error) { |
| 32 | tab, ctrl := a.tabAndCtrlByID(tabID) |
| 33 | if ctrl == nil { |
| 34 | return nil, a.workspaceNotReadyErr(tab) |
| 35 | } |
| 36 | status := ctrl.RuntimeStatus() |
| 37 | if strings.TrimSpace(turnID) == "" || status.TurnID != strings.TrimSpace(turnID) { |
| 38 | return nil, fmt.Errorf("turn %q is not the active turn for tab %q", turnID, tabID) |
| 39 | } |
| 40 | if epoch := strings.TrimSpace(runtimeEpoch); epoch != "" && tab != nil && tab.sink != nil && tab.sink.runtimeEpochSnapshot() != epoch { |
| 41 | return nil, fmt.Errorf("runtime changed while resolving prompt for tab %q", tabID) |
| 42 | } |
| 43 | return ctrl, nil |
| 44 | } |
| 45 | |
| 46 | // stoppableCtrl resolves and captures the controller a Stop request targets. |
| 47 | // The request is session-scoped and idle cancellation is idempotent. |
| 48 | func (a *App) stoppableCtrl(tabID, turnID string) (control.SessionAPI, error) { |
| 49 | tab, ctrl := a.tabAndCtrlByID(tabID) |
| 50 | if ctrl == nil { |
| 51 | return nil, a.workspaceNotReadyErr(tab) |
| 52 | } |
| 53 | // Logging must never refresh execution state before delivering Stop. An |
| 54 | // empty ID is the ordinary session-scoped command, not a stale turn. |
| 55 | if reader, ok := ctrl.(control.PublishedRuntimeStateReader); ok && strings.TrimSpace(turnID) != "" { |
| 56 | state := reader.PublishedRuntimeStateSnapshot() |
| 57 | if turnID = strings.TrimSpace(turnID); turnID != state.TurnID { |
| 58 | slog.Info("desktop: stop targeted a stale turn id; interrupting the active turn", "tab", tabID, "requested", turnID, "active", state.TurnID) |
| 59 | } |
| 60 | } |
| 61 | return ctrl, nil |
| 62 | } |
| 63 | |
| 64 | // CancelSessionForTab is the protocol-v2 Stop operation. It captures the tab's |
| 65 | // current controller exactly once, so later tab switches cannot retarget it. |
| 66 | func (a *App) CancelSessionForTab(tabID string) (control.CancelReceipt, error) { |
| 67 | ctrl, err := a.stoppableCtrl(tabID, "") |
| 68 | if err != nil { |
| 69 | return control.CancelReceipt{}, err |
| 70 | } |
| 71 | if concrete, ok := ctrl.(*control.Controller); ok { |
| 72 | return concrete.CancelSessionFrom("user_stop"), nil |
| 73 | } |
| 74 | if session, ok := ctrl.(interface{ CancelSession() control.CancelReceipt }); ok { |
| 75 | return session.CancelSession(), nil |
| 76 | } |
| 77 | ctrl.Cancel() |
| 78 | if reader, ok := ctrl.(control.PublishedRuntimeStateReader); ok { |
| 79 | state := reader.PublishedRuntimeStateSnapshot() |
| 80 | return control.CancelReceipt{ |
| 81 | SessionRef: ctrl.SessionPath(), HeadID: agent.BranchID(ctrl.SessionPath()), |
| 82 | RuntimeEpoch: state.RuntimeEpoch, Accepted: true, |
| 83 | AlreadyIdle: !state.Running && !state.PendingPrompt && state.BackgroundJobs == 0, |
| 84 | RecoveryRequired: state.Phase == "recovery_required", |
| 85 | }, nil |
| 86 | } |
| 87 | status := ctrl.RuntimeStatus() |
| 88 | return control.CancelReceipt{SessionRef: ctrl.SessionPath(), HeadID: agent.BranchID(ctrl.SessionPath()), Accepted: true, AlreadyIdle: !status.Running && !status.PendingPrompt}, nil |
| 89 | } |
| 90 | |
| 91 | // StartTurnForTab is the turn-id-aware replacement for SubmitToTab. Existing |
| 92 | // Submit entry points remain compatibility wrappers during the protocol cutover. |
| 93 | func (a *App) StartTurnForTab(tabID, input, submissionID string) (TurnStartView, error) { |
| 94 | if strings.TrimSpace(submissionID) == "" { |
| 95 | return TurnStartView{}, fmt.Errorf("submissionId is required") |
| 96 | } |
| 97 | if _, ctrl := a.tabAndCtrlByID(tabID); ctrl != nil { |
| 98 | if identified, ok := ctrl.(*control.Controller); ok { |
| 99 | receipt, found, err := identified.LookupSubmission(control.SubmissionRequest{ID: submissionID, Input: input, Display: input}) |
| 100 | if err != nil { |
| 101 | return TurnStartView{}, err |
| 102 | } |
| 103 | if found { |
| 104 | return TurnStartView{TurnID: receipt.TurnID, Status: event.TurnQueued, Disposition: control.SubmitTurnStarted, SubmissionID: submissionID}, nil |
| 105 | } |
| 106 | } |
| 107 | } |
| 108 | result, err := a.submitToTabResult(tabID, input, false, true, submissionID) |
| 109 | if err != nil { |
| 110 | if result.Disposition == control.SubmitManagementHandled && result.OperationID != "" { |
| 111 | code := "" |
| 112 | switch { |
| 113 | case errors.Is(err, control.ErrMaintenanceBusy): |
| 114 | code = "maintenance_busy" |
| 115 | case errors.Is(err, control.ErrMaintenanceRecovery): |
| 116 | code = "maintenance_recovery_required" |
| 117 | } |
| 118 | if code != "" { |
| 119 | return TurnStartView{Disposition: result.Disposition, OperationID: result.OperationID, |
| 120 | SubmissionID: submissionID, ManagementErrorCode: code}, nil |
| 121 | } |
| 122 | } |
| 123 | return TurnStartView{}, err |
| 124 | } |
| 125 | if result.Disposition == control.SubmitManagementHandled { |
| 126 | return TurnStartView{Disposition: result.Disposition, OperationID: result.OperationID, SubmissionID: submissionID}, nil |
| 127 | } |
| 128 | tab, ctrl := a.tabAndCtrlByID(tabID) |
| 129 | if ctrl == nil { |
| 130 | return TurnStartView{}, a.workspaceNotReadyErr(tab) |
| 131 | } |
| 132 | // Admission is released by now and the tab may hold a rebuilt controller, |
| 133 | // so the admitting controller's receipt is the authority when it has one. |
| 134 | turnID := result.TurnID |
| 135 | if admitted, ok := ctrl.(interface{ TurnIDForSubmission(string) string }); ok && turnID == "" { |
| 136 | turnID = admitted.TurnIDForSubmission(submissionID) |
| 137 | } |
| 138 | if strings.TrimSpace(turnID) == "" { |
| 139 | return TurnStartView{}, fmt.Errorf("turn admission did not produce a durable turn id") |
| 140 | } |
| 141 | epoch := "" |
| 142 | if tab != nil && tab.sink != nil { |
| 143 | epoch = tab.sink.runtimeEpochSnapshot() |
| 144 | } |
| 145 | // This is an admission receipt, not a potentially raced runtime snapshot. |
| 146 | // Ordered events carry every later transition, including a provider that |
| 147 | // completed before the Wails Promise was delivered. |
| 148 | return TurnStartView{TurnID: turnID, Status: event.TurnQueued, Disposition: control.SubmitTurnStarted, RuntimeEpoch: epoch, SubmissionID: submissionID}, nil |
| 149 | } |
| 150 | |
| 151 | func (a *App) StartTurnForTabWithDrafts(tabID, input, submissionID string, draftIDs []string) (TurnStartView, error) { |
| 152 | if strings.TrimSpace(submissionID) == "" { |
| 153 | return TurnStartView{}, fmt.Errorf("submissionId is required") |
| 154 | } |
| 155 | req := control.SubmissionRequest{ID: submissionID, Input: input, Display: input, DraftIDs: append([]string(nil), draftIDs...)} |
| 156 | if _, ctrl := a.tabAndCtrlByID(tabID); ctrl != nil { |
| 157 | if identified, ok := ctrl.(*control.Controller); ok { |
| 158 | receipt, found, err := identified.LookupSubmission(req) |
| 159 | if err != nil { |
| 160 | return TurnStartView{}, err |
| 161 | } |
| 162 | if found { |
| 163 | return TurnStartView{TurnID: receipt.TurnID, Status: event.TurnQueued, Disposition: control.SubmitTurnStarted, SubmissionID: submissionID}, nil |
| 164 | } |
| 165 | } |
| 166 | } |
| 167 | admission, ctrl, err := a.beginTabTurn(tabID, true, submissionID) |
| 168 | if err != nil { |
| 169 | return TurnStartView{}, a.submissionAdmissionError(tabID, req, err) |
| 170 | } |
| 171 | defer admission.abort() |
| 172 | if err := a.ensureTabTopicIndexedForUserTurn(admission.tab); err != nil { |
| 173 | return TurnStartView{}, err |
| 174 | } |
| 175 | identified, ok := ctrl.(*control.Controller) |
| 176 | if !ok { |
| 177 | return TurnStartView{}, fmt.Errorf("unsupported: attachments-v1") |
| 178 | } |
| 179 | receipt, err := identified.SubmitIdentified(req) |
| 180 | if err != nil { |
| 181 | return TurnStartView{}, inboxBridgeError(err) |
| 182 | } |
| 183 | admission.finish(ctrl) |
| 184 | turnID := receipt.TurnID |
| 185 | if strings.TrimSpace(turnID) == "" { |
| 186 | return TurnStartView{}, fmt.Errorf("turn admission did not produce a durable turn id") |
| 187 | } |
| 188 | epoch := "" |
| 189 | if admission.tab != nil && admission.tab.sink != nil { |
| 190 | epoch = admission.tab.sink.runtimeEpochSnapshot() |
| 191 | } |
| 192 | return TurnStartView{TurnID: turnID, Status: event.TurnQueued, Disposition: control.SubmitTurnStarted, RuntimeEpoch: epoch, SubmissionID: submissionID}, nil |
| 193 | } |
| 194 | |
| 195 | // InterruptTurnForTab stops the tab's active work. Stop is a session-level |
| 196 | // request: a turn id from a stale button still interrupts whatever is running |
| 197 | // now, because an unstoppable turn is worse than stopping its replacement. |
| 198 | func (a *App) InterruptTurnForTab(tabID, turnID string) error { |
| 199 | ctrl, err := a.stoppableCtrl(tabID, turnID) |
| 200 | if err != nil { |
| 201 | return err |
| 202 | } |
| 203 | if session, ok := ctrl.(interface{ CancelSession() control.CancelReceipt }); ok { |
| 204 | session.CancelSession() |
| 205 | } else { |
| 206 | ctrl.Cancel() |
| 207 | } |
| 208 | return nil |
| 209 | } |
| 210 | |
| 211 | // InterruptTurnWithInboxItemsForTab is the receipt-capable Stop used by the |
| 212 | // Composer when it also discards queued follow-ups. |
| 213 | func (a *App) InterruptTurnWithInboxItemsForTab(tabID, turnID string, itemIDs []string) (InboxCancelResultView, error) { |
| 214 | view := InboxCancelResultView{DiscardedItemIDs: []string{}} |
| 215 | ctrl, err := a.stoppableCtrl(tabID, turnID) |
| 216 | if err != nil { |
| 217 | return view, err |
| 218 | } |
| 219 | result, err := ctrl.CancelWithInboxItemsResult(itemIDs, "desktop") |
| 220 | if err != nil { |
| 221 | return view, inboxBridgeError(err) |
| 222 | } |
| 223 | view.DiscardedItemIDs = append(view.DiscardedItemIDs, result.DiscardedItemIDs...) |
| 224 | view.Warning = result.Warning |
| 225 | a.emitInboxChanged(tabID) |
| 226 | return view, nil |
| 227 | } |
| 228 | |
| 229 | // AnswerPromptForTab resolves an Ask only when it belongs to the exact active |
| 230 | // turn. Controller-side prompt ids remain independently idempotent. |
| 231 | func (a *App) AnswerPromptForTab(tabID, turnID, promptID string, answers []QuestionAnswer) error { |
| 232 | tab, ctrl := a.tabAndCtrlByID(tabID) |
| 233 | if ctrl == nil { |
| 234 | return a.workspaceNotReadyErr(tab) |
| 235 | } |
| 236 | status := ctrl.RuntimeStatus() |
| 237 | if strings.TrimSpace(turnID) == "" || status.TurnID != strings.TrimSpace(turnID) { |
| 238 | return fmt.Errorf("turn %q is not the active turn for tab %q", turnID, tabID) |
| 239 | } |
| 240 | // Resolve on the controller instance that passed the turn-id fence. Calling |
| 241 | // the legacy app wrapper here would re-resolve the tab and could deliver a |
| 242 | // late answer to a replacement controller after a runtime rebuild. |
| 243 | out := make([]event.AskAnswer, len(answers)) |
| 244 | for i, answer := range answers { |
| 245 | out[i] = event.AskAnswer{QuestionID: answer.QuestionID, Selected: answer.Selected} |
| 246 | } |
| 247 | if checked, ok := ctrl.(interface { |
| 248 | AnswerQuestionChecked(string, []event.AskAnswer) error |
| 249 | }); ok { |
| 250 | return checked.AnswerQuestionChecked(promptID, out) |
| 251 | } |
| 252 | ctrl.AnswerQuestion(promptID, out) |
| 253 | return nil |
| 254 | } |
| 255 | |
| 256 | type turnEventReader interface { |
| 257 | TurnEventReplay(after uint64) (turnevent.ReplayView, error) |
| 258 | } |
| 259 | |
| 260 | type TurnEventReplayView struct { |
| 261 | Events []turnevent.Envelope `json:"events"` |
| 262 | FloorSequence uint64 `json:"floorSeq"` |
| 263 | LatestSequence uint64 `json:"latestSeq"` |
| 264 | NextAfterSequence uint64 `json:"nextAfterSeq"` |
| 265 | HasMore bool `json:"hasMore"` |
| 266 | ResetRequired bool `json:"resetRequired"` |
| 267 | TranscriptRevision int64 `json:"transcriptRevision,omitempty"` |
| 268 | TranscriptDigest string `json:"transcriptDigest,omitempty"` |
| 269 | HeadID string `json:"headId,omitempty"` |
| 270 | LeafMessageID string `json:"leafMessageId,omitempty"` |
| 271 | RuntimeEpoch string `json:"runtimeEpoch,omitempty"` |
| 272 | } |
| 273 | |
| 274 | // TurnEventsForTab supplies the durable suffix used to repair sequence gaps or |
| 275 | // rebuild after a runtime epoch change. |
| 276 | func (a *App) TurnEventsForTab(tabID string, afterSeq uint64) (TurnEventReplayView, error) { |
| 277 | empty := TurnEventReplayView{Events: []turnevent.Envelope{}} |
| 278 | tab, ctrl := a.tabAndCtrlByID(tabID) |
| 279 | if ctrl == nil { |
| 280 | return empty, a.workspaceNotReadyErr(tab) |
| 281 | } |
| 282 | reader, ok := ctrl.(turnEventReader) |
| 283 | if !ok { |
| 284 | return empty, fmt.Errorf("turn event replay is unavailable") |
| 285 | } |
| 286 | // Re-check the controller under the app lock before sampling the epoch. |
| 287 | // This prevents pairing an old controller with a replacement runtime after |
| 288 | // a session rebind races tabAndCtrlByID. |
| 289 | epoch := "" |
| 290 | a.mu.RLock() |
| 291 | bound := tab != nil && a.tabs[tabID] == tab && tab.Ctrl == ctrl |
| 292 | if bound && tab.sink != nil { |
| 293 | epoch = tab.sink.runtimeEpochSnapshot() |
| 294 | } |
| 295 | a.mu.RUnlock() |
| 296 | if !bound { |
| 297 | return empty, fmt.Errorf("runtime changed while binding turn event replay") |
| 298 | } |
| 299 | replay, err := reader.TurnEventReplay(afterSeq) |
| 300 | if replay.Events == nil { |
| 301 | replay.Events = []turnevent.Envelope{} |
| 302 | } |
| 303 | return TurnEventReplayView{ |
| 304 | Events: replay.Events, FloorSequence: replay.FloorSequence, |
| 305 | LatestSequence: replay.LatestSequence, NextAfterSequence: replay.NextAfterSequence, |
| 306 | HasMore: replay.HasMore, ResetRequired: replay.ResetRequired, |
| 307 | TranscriptRevision: replay.TranscriptRevision, TranscriptDigest: replay.TranscriptDigest, |
| 308 | HeadID: replay.HeadID, LeafMessageID: replay.LeafMessageID, |
| 309 | RuntimeEpoch: epoch, |
| 310 | }, err |
| 311 | } |
| 312 | |
| 313 | var _ control.SessionAPI = (*control.Controller)(nil) |
| 314 |