返回 DeepSeek-Reasonix
turn_runtime_api.go
根目录 / desktop / turn_runtime_api.go
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
314 lines GO