返回 DeepSeek-Reasonix
turn_admission.go
根目录 / desktop / turn_admission.go
1 package main
2
3 import (
4 "fmt"
5
6 "reasonix/internal/control"
7 )
8
9 type imageCapabilitySnapshot interface{ ImageCapabilityChanged() bool }
10
11 // Use the existing build/swap/lease boundary before accepting a new turn.
12 // A failed rebuild leaves the previous snapshot visible and rejects this turn.
13 func (a *App) refreshTabImageCapability(tab *WorkspaceTab) error {
14 a.runtimeRebuildMu.Lock()
15 defer a.runtimeRebuildMu.Unlock()
16 tab.turnStartMu.Lock()
17 defer tab.turnStartMu.Unlock()
18 current, ok := a.controllerForTab(tab).(imageCapabilitySnapshot)
19 if !ok || !current.ImageCapabilityChanged() {
20 return nil
21 }
22 if err := a.rebuildSettingTurnLocked("image input", tab, false, false); err != nil {
23 return fmt.Errorf("refresh image input configuration: %w", err)
24 }
25 return nil
26 }
27
28 // tabTurnAdmission owns both locks acquired while a foreground turn starts.
29 type tabTurnAdmission struct {
30 app *App
31 tab *WorkspaceTab
32 released bool
33 }
34
35 type turnFinishingWaiter interface {
36 TurnFinishingDone() (<-chan struct{}, bool)
37 }
38
39 func (admission *tabTurnAdmission) finish(ctrl control.SessionAPI) bool {
40 if admission == nil || admission.released {
41 return false
42 }
43 admission.released = true
44 tab := admission.tab
45 if tab != nil {
46 // Defers preserve lock release if RuntimeStatus panics.
47 defer admission.app.runtimeAdmissionMu.RUnlock()
48 defer tab.turnStartMu.Unlock()
49 }
50 started := ctrl != nil && ctrl.RuntimeStatus().Running
51 if !started && tab != nil && tab.sink != nil {
52 tab.sink.cancelTurnStart()
53 }
54 return started
55 }
56
57 func (admission *tabTurnAdmission) abort() {
58 admission.finish(nil)
59 }
60
61 // beginTabTurn reserves one tab until its TurnDone fan-out completes.
62 func (a *App) beginTabTurn(tabID string, reclaim bool, submissionID ...string) (*tabTurnAdmission, control.SessionAPI, error) {
63 return a.beginRuntimeTurnChecked(tabID, reclaim, false, identifiedSubmissionCheck(submissionID), submissionID...)
64 }
65
66 func (a *App) beginRuntimeTurn(tabID string, reclaim, detached bool, submissionID ...string) (*tabTurnAdmission, control.SessionAPI, error) {
67 return a.beginRuntimeTurnChecked(tabID, reclaim, detached, nil, submissionID...)
68 }
69
70 func checkRuntimeAdmission(ctrl control.SessionAPI, check func(control.SessionAPI) error) error {
71 if check == nil {
72 return nil
73 }
74 return check(ctrl)
75 }
76
77 func controllerAuthenticationError(ctrl control.SessionAPI) error {
78 authentication, ok := ctrl.(interface {
79 AuthenticationState() control.AuthenticationState
80 })
81 if !ok {
82 return nil
83 }
84 state := authentication.AuthenticationState()
85 if state.Ready() {
86 return nil
87 }
88 return &control.AuthenticationError{State: state}
89 }
90
91 // check observes the selected controller under the same admission locks as the
92 // running check and submit. A refusal must carry this owner's identity with it.
93 func (a *App) beginRuntimeTurnChecked(tabID string, reclaim, detached bool, check func(control.SessionAPI) error, submissionID ...string) (*tabTurnAdmission, control.SessionAPI, error) {
94 return a.beginRuntimeTurnWithModelChoice(tabID, reclaim, detached, check, nil, submissionID...)
95 }
96
97 func (a *App) beginRuntimeTurnWithModelChoice(tabID string, reclaim, detached bool, check func(control.SessionAPI) error, choice *control.ModelApplicationChoice, submissionID ...string) (*tabTurnAdmission, control.SessionAPI, error) {
98 for {
99 tab, ctrl := a.tabAndCtrlByID(tabID)
100 if detached {
101 a.mu.RLock()
102 tab = a.tabByEventSinkIDLocked(tabID)
103 ctrl = nil
104 if tab != nil {
105 ctrl = tab.Ctrl
106 }
107 a.mu.RUnlock()
108 }
109 if a.tabIsReadOnly(tab) {
110 return nil, nil, readOnlyChannelErr()
111 }
112 if err := a.workspaceRuntimeAdmissionErr(tab, ctrl); err != nil {
113 return nil, nil, err
114 }
115 // Slow workspace repair stays outside the runtime admission barrier.
116 if err := a.ensureTabControllerWorkspace(tab); err != nil {
117 return nil, nil, err
118 }
119
120 a.runtimeAdmissionMu.RLock()
121 abort := func() {
122 tab.turnStartMu.Unlock()
123 a.runtimeAdmissionMu.RUnlock()
124 }
125 tab.turnStartMu.Lock()
126 if err := a.validateDraftAdmission(tab, firstSubmissionID(submissionID)); err != nil {
127 abort()
128 return nil, nil, err
129 }
130 if a.tabIsReadOnly(tab) {
131 abort()
132 return nil, nil, readOnlyChannelErr()
133 }
134 if reclaim && a.botBridge != nil {
135 a.botBridge.reclaimFromDesktop(tab.ID)
136 }
137 ctrl = a.controllerForTab(tab)
138 if err := a.workspaceRuntimeAdmissionErr(tab, ctrl); err != nil {
139 abort()
140 return nil, nil, err
141 }
142 ctrl = a.controllerForTab(tab)
143 if err := a.workspaceRuntimeAdmissionErr(tab, ctrl); err != nil {
144 abort()
145 return nil, nil, err
146 }
147 running := ctrl.RuntimeStatus().Running
148 if err := checkRuntimeAdmission(ctrl, check); err != nil {
149 abort()
150 return nil, nil, err
151 }
152 if running {
153 if waiter, ok := ctrl.(turnFinishingWaiter); ok {
154 if done, finishing := waiter.TurnFinishingDone(); finishing {
155 // Re-resolve after waiting so close/switch cannot misroute retry.
156 abort()
157 <-done
158 continue
159 }
160 // Fan-out can end between RuntimeStatus and TurnFinishingDone.
161 // Re-check before reporting busy so that completed boundary retries
162 // instead of preserving the original false rejection window.
163 if !ctrl.RuntimeStatus().Running {
164 abort()
165 continue
166 }
167 }
168 abort()
169 return nil, nil, control.ErrTurnRunning
170 }
171 usingApplied, choiceErr := validateModelApplicationChoice(ctrl, choice)
172 if choiceErr != nil {
173 abort()
174 return nil, nil, choiceErr
175 }
176 if retry, err := a.applyTurnModelSettings(tab, ctrl, usingApplied, abort); retry || err != nil {
177 if err != nil {
178 return nil, nil, err
179 }
180 continue
181 }
182 if snapshot, ok := ctrl.(imageCapabilitySnapshot); a.ctx != nil && !usingApplied && ok && snapshot.ImageCapabilityChanged() {
183 a.mu.RLock()
184 draftPending := tab.PendingCreateOperationID != ""
185 a.mu.RUnlock()
186 abort()
187 if draftPending {
188 return nil, nil, fmt.Errorf("image configuration changed before draft admission")
189 }
190 if err := a.refreshTabImageCapability(tab); err != nil {
191 return nil, nil, err
192 }
193 continue
194 }
195 if err := controllerAuthenticationError(ctrl); err != nil {
196 abort()
197 return nil, nil, err
198 }
199 if tab.sink != nil && !tab.sink.tryBeginTurn(submissionID...) {
200 abort()
201 return nil, nil, control.ErrTurnRunning
202 }
203 if metadata, ok := ctrl.(interface {
204 SetTurnEventRoutingMetadata(runtimeEpoch, submissionID string)
205 }); ok {
206 epoch := ""
207 if tab.sink != nil {
208 epoch = tab.sink.runtimeEpochSnapshot()
209 }
210 metadata.SetTurnEventRoutingMetadata(epoch, firstSubmissionID(submissionID))
211 }
212 return &tabTurnAdmission{app: a, tab: tab}, ctrl, nil
213 }
214 }
215
216 // Durable follow-ups are new runs even when a controller dispatches them on
217 // its own after TurnDone. Resolve the runtime owner again after detach/attach.
218 func (a *App) beforeInboxDispatch(ctrl *control.Controller) (func(), error) {
219 a.mu.RLock()
220 var owner *WorkspaceTab
221 for _, tab := range a.runtimeTabsLocked() {
222 if tab.Ctrl == ctrl {
223 owner = tab
224 break
225 }
226 }
227 a.mu.RUnlock()
228 if owner == nil {
229 return nil, control.ErrInboxRuntimeUnpublished
230 }
231 admission, current, err := a.beginRuntimeTurn(owner.ID, false, true)
232 if err != nil {
233 return nil, err
234 }
235 if current != ctrl {
236 admission.abort()
237 if replacement, ok := current.(*control.Controller); ok {
238 go replacement.NotifyInboxRuntimeReady()
239 }
240 return nil, control.ErrInboxRuntimeUnpublished
241 }
242 return func() { admission.finish(ctrl) }, nil
243 }
244
244 lines GO