返回 DeepSeek-Reasonix
deferred_rebuild.go
根目录 / desktop / deferred_rebuild.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "maps"
9 "strings"
10 "sync"
11 "time"
12
13 "reasonix/internal/agent"
14 "reasonix/internal/control"
15 "reasonix/internal/secrets"
16 )
17
18 // deferredRebuildRetryInterval is how often the retry loop probes a held
19 // session lease. Package-level so tests can shorten it.
20 var deferredRebuildRetryInterval = 2 * time.Second
21
22 const deferredStartupBuildLabel = "__startup__"
23
24 // deferredRuntimeReloadLabel marks a queued ReloadRuntime in the pending map.
25 // Like the startup label it is not a user setting name; the retry loop routes
26 // it to the boot.Rebuild reload path instead of a settings rebuild.
27 const deferredRuntimeReloadLabel = "__reload__"
28
29 // deferredRebuildState tracks tabs whose settings were saved to disk but whose
30 // runtime could not refresh, plus tabs whose initial startup failed, because
31 // the session lease was held by another Reasonix process. A single background
32 // loop probes the lease and replays the rebuild once the other side releases
33 // it. The loop only runs after enableDeferredRebuildRetry (the startup
34 // hook); tests that never call it get the pending bookkeeping without a
35 // background goroutine.
36 type deferredRebuildState struct {
37 mu sync.Mutex
38 pending map[string]deferredRebuildRequest
39 next uint64
40 enabled bool
41 running bool
42 stopped bool
43 stop chan struct{}
44 }
45
46 type deferredRebuildReason uint8
47
48 const (
49 deferredSettingsReason deferredRebuildReason = iota
50 deferredStartupReason
51 deferredReloadReason
52 )
53
54 type deferredRebuildRequest struct {
55 reason deferredRebuildReason
56 label string
57 target *WorkspaceTab
58 runtimeID string
59 revision string
60 sequence uint64
61 }
62
63 // enableDeferredRebuildRetry arms the retry loop; called from startup.
64 func (a *App) enableDeferredRebuildRetry() {
65 d := &a.deferredRebuild
66 d.mu.Lock()
67 defer d.mu.Unlock()
68 d.enabled = true
69 a.startDeferredRebuildLoopLocked()
70 }
71
72 // startDeferredRebuildLoopLocked starts the loop when it is armed, idle, and
73 // has work. Callers must hold d.mu.
74 func (a *App) startDeferredRebuildLoopLocked() {
75 d := &a.deferredRebuild
76 if !d.enabled || d.running || d.stopped || len(d.pending) == 0 {
77 return
78 }
79 d.running = true
80 if d.stop == nil {
81 d.stop = make(chan struct{})
82 }
83 go a.deferredRebuildLoop(d.stop)
84 }
85
86 // scheduleDeferredRebuild records that tabID needs a runtime refresh for
87 // setting and starts the retry loop if it is not running yet. Repeated calls
88 // for the same tab collapse into one retry carrying the latest label.
89 func (a *App) scheduleDeferredRebuild(tabID, setting string) {
90 tabID = strings.TrimSpace(tabID)
91 if tabID == "" {
92 return
93 }
94 a.mu.RLock()
95 tab := a.tabs[tabID]
96 request := deferredRebuildRequest{label: setting, target: tab}
97 var snapshot modelSettingsSnapshot
98 if tab != nil {
99 request.runtimeID = tab.runtimeID
100 snapshot, _ = tab.Ctrl.(modelSettingsSnapshot)
101 }
102 a.mu.RUnlock()
103 if snapshot != nil {
104 _, request.revision, _ = snapshot.ModelSettingsState()
105 }
106 switch setting {
107 case deferredStartupBuildLabel:
108 request.reason = deferredStartupReason
109 case deferredRuntimeReloadLabel:
110 request.reason = deferredReloadReason
111 }
112 d := &a.deferredRebuild
113 d.mu.Lock()
114 defer d.mu.Unlock()
115 if d.stopped {
116 return
117 }
118 if d.pending == nil {
119 d.pending = map[string]deferredRebuildRequest{}
120 }
121 d.next++
122 request.sequence = d.next
123 d.pending[tabID] = request
124 a.startDeferredRebuildLoopLocked()
125 }
126
127 func (a *App) scheduleDeferredStartupBuild(tabID string) {
128 a.scheduleDeferredRebuild(tabID, deferredStartupBuildLabel)
129 }
130
131 func (a *App) deferredRebuildSequence(tabID string) uint64 {
132 d := &a.deferredRebuild
133 d.mu.Lock()
134 defer d.mu.Unlock()
135 return d.pending[tabID].sequence
136 }
137
138 func (a *App) clearDeferredRebuildVersion(tabID string, sequence uint64) {
139 d := &a.deferredRebuild
140 d.mu.Lock()
141 defer d.mu.Unlock()
142 if sequence != 0 && d.pending[tabID].sequence == sequence {
143 delete(d.pending, tabID)
144 }
145 }
146
147 func (a *App) deferredRebuildPending(tabID string) bool {
148 d := &a.deferredRebuild
149 d.mu.Lock()
150 defer d.mu.Unlock()
151 _, ok := d.pending[tabID]
152 return ok
153 }
154
155 // stopDeferredRebuildRetry permanently stops the retry loop; used on shutdown
156 // and by tests.
157 func (a *App) stopDeferredRebuildRetry() {
158 d := &a.deferredRebuild
159 d.mu.Lock()
160 defer d.mu.Unlock()
161 if d.stopped {
162 return
163 }
164 d.stopped = true
165 if d.stop != nil {
166 close(d.stop)
167 }
168 }
169
170 func (a *App) deferredRebuildLoop(stop <-chan struct{}) {
171 ticker := time.NewTicker(deferredRebuildRetryInterval)
172 defer ticker.Stop()
173 for {
174 select {
175 case <-stop:
176 d := &a.deferredRebuild
177 d.mu.Lock()
178 d.running = false
179 d.mu.Unlock()
180 return
181 case <-ticker.C:
182 }
183 if a.deferredRebuildTickDone() {
184 return
185 }
186 }
187 }
188
189 // deferredRebuildTickDone runs one retry pass and reports true when the loop
190 // should exit because nothing is pending anymore.
191 func (a *App) deferredRebuildTickDone() bool {
192 return a.deferredRebuildTick(true)
193 }
194
195 func (a *App) deferredRebuildTick(markIdle bool) bool {
196 d := &a.deferredRebuild
197 d.mu.Lock()
198 if d.stopped || len(d.pending) == 0 {
199 if markIdle {
200 d.running = false
201 }
202 d.mu.Unlock()
203 return true
204 }
205 pending := make(map[string]deferredRebuildRequest, len(d.pending))
206 maps.Copy(pending, d.pending)
207 d.mu.Unlock()
208
209 for tabID, request := range pending {
210 a.retryDeferredRebuild(tabID, request)
211 }
212 return false
213 }
214
215 func (a *App) kickDeferredRebuildRetry() {
216 if a.ctx == nil {
217 return
218 }
219 a.goSafe("deferredRebuildKick", func() {
220 _ = a.deferredRebuildTick(false)
221 })
222 }
223
224 func (a *App) retryDeferredRebuild(tabID string, request deferredRebuildRequest) {
225 if a.ctx == nil {
226 return
227 }
228 tab := a.tabByID(tabID)
229 a.mu.RLock()
230 valid := tab != nil && tab == request.target && tab.runtimeID == request.runtimeID
231 a.mu.RUnlock()
232 if !valid {
233 // The tab is gone; nothing left to refresh.
234 a.clearDeferredRebuildVersion(tabID, request.sequence)
235 return
236 }
237 if request.reason == deferredStartupReason {
238 a.retryDeferredStartupBuild(tabID, tab, request.sequence)
239 return
240 }
241 if request.reason == deferredReloadReason {
242 a.retryDeferredRuntimeReload(tabID, tab, request.sequence)
243 return
244 }
245 // Hold the rebuild mutex across probe + rebuild: the probe briefly acquires
246 // the session lease, and a concurrent manual rebuild's ensureSessionLease
247 // would see that probe as "held by another runtime" and spuriously defer.
248 a.runtimeRebuildMu.Lock()
249 defer a.runtimeRebuildMu.Unlock()
250 tab.turnStartMu.Lock()
251 defer tab.turnStartMu.Unlock()
252 ctrl := a.controllerForTab(tab)
253 if ctrl == nil {
254 // Mid-(re)build on another path (provider retarget, workspace repair);
255 // racing a second build+swap against it is what this loop must avoid.
256 return
257 }
258 if request.label == "saved model settings" {
259 pending, err := modelSettingsNeedApply(ctrl)
260 if err == nil && !pending {
261 a.clearDeferredRebuildVersion(tabID, request.sequence)
262 return
263 }
264 }
265 if (request.label == "saved model settings" && control.ModelReplacementBlocked(ctrl)) || (request.label != "saved model settings" && controllerHasActiveRuntimeWork(ctrl)) {
266 return
267 }
268 if !a.deferredRebuildLeaseLooksFree(tab) {
269 return
270 }
271 setting := request.label
272 err := a.rebuildSettingTurnLocked(setting, tab, false, setting == "saved model settings")
273 if err == nil {
274 // rebuildSettingLocked already cleared the pending entry for the tab it
275 // refreshed; just announce it.
276 a.noticeForTab(tabID, fmt.Sprintf("%s applied: session refreshed after the lease was released", setting))
277 return
278 }
279 if errors.Is(err, agent.ErrSessionLeaseHeld) {
280 return // grabbed back before we could rebuild; keep waiting
281 }
282 var busy *rebuildBusyError
283 if errors.As(err, &busy) {
284 return // a turn started meanwhile; retry once it finishes
285 }
286 // Anything else will not resolve by waiting; give up loudly instead of
287 // retrying forever.
288 a.clearDeferredRebuildVersion(tabID, request.sequence)
289 if setting == "saved model settings" {
290 a.mu.Lock()
291 if a.ownsRuntimeTabLocked(tab) && tab.Ctrl == ctrl {
292 tab.modelApplication.failure = &modelSettingsApplyFailure{ctrl, request.revision, modelSettingsIssue("apply_failed", err).Message}
293 }
294 a.mu.Unlock()
295 }
296 slog.Warn("desktop: deferred settings rebuild failed", "setting", setting, "tab", tabID, "err", err)
297 a.warnForTab(tabID, fmt.Sprintf("%s was saved but the session could not refresh: %s", setting, err.Error()))
298 }
299
300 // retryDeferredRuntimeReload drives one queued ReloadRuntime pass. The
301 // probing contract mirrors retryDeferredRebuild: wait for the tab to be
302 // active and idle and for its lease to look free, then run the boot.Rebuild
303 // reload; busy/lease answers keep waiting, anything else gives up loudly.
304 func (a *App) retryDeferredRuntimeReload(tabID string, tab *WorkspaceTab, queuedSequence ...uint64) {
305 sequence := a.deferredRebuildSequence(tabID)
306 if len(queuedSequence) > 0 {
307 sequence = queuedSequence[0]
308 }
309 // Hold the rebuild mutex across probe + reload: the probe briefly
310 // acquires the session lease, and a concurrent rebuild's ensure lease
311 // would read that probe as "held by another runtime" and spuriously
312 // defer (same contract as retryDeferredRebuild).
313 a.runtimeRebuildMu.Lock()
314 defer a.runtimeRebuildMu.Unlock()
315 ctrl := a.controllerForTab(tab)
316 if ctrl == nil {
317 // Mid-(re)build on another path; racing a second build+swap against it
318 // is what this loop must avoid.
319 return
320 }
321 if controllerHasActiveRuntimeWork(ctrl) {
322 return
323 }
324 if !a.deferredRebuildLeaseLooksFree(tab) {
325 return
326 }
327 err := a.reloadRuntimeTurnLocked(tab)
328 if err == nil {
329 // rebuildSettingTurnLocked already cleared the pending entry for the
330 // tab it refreshed; just announce it.
331 a.noticeForTab(tabID, "runtime reloaded after the session went idle")
332 return
333 }
334 if errors.Is(err, agent.ErrSessionLeaseHeld) {
335 return // grabbed back before we could reload; keep waiting
336 }
337 var busy *rebuildBusyError
338 if errors.As(err, &busy) {
339 return // a turn started meanwhile; retry once it finishes
340 }
341 // Anything else will not resolve by waiting; give up loudly instead of
342 // retrying forever. The error may come from provider/config plumbing and
343 // carry credential-shaped values (passwords, resolved API keys) — the
344 // tested helper redacts before the text reaches logs or the frontend.
345 a.clearDeferredRebuildVersion(tabID, sequence)
346 failure := deferredReloadFailedText(err)
347 slog.Warn("desktop: "+failure, "tab", tabID)
348 a.warnForTab(tabID, failure)
349 }
350
351 // deferredReloadFailedText is the failure line for an unrecoverable deferred
352 // reload — the single formatter both the log and the tab warning use, so a
353 // credential-shaped error can never reach either sink unredacted.
354 func deferredReloadFailedText(err error) string {
355 return "runtime reload failed: " + secrets.RedactCredentials(err.Error())
356 }
357
358 func (a *App) retryDeferredStartupBuild(tabID string, tab *WorkspaceTab, queuedSequence ...uint64) {
359 sequence := a.deferredRebuildSequence(tabID)
360 if len(queuedSequence) > 0 {
361 sequence = queuedSequence[0]
362 }
363 a.runtimeRebuildMu.Lock()
364 defer a.runtimeRebuildMu.Unlock()
365 if !a.tabHasRetryableStartupLeaseError(tab) {
366 a.clearDeferredRebuildVersion(tabID, sequence)
367 return
368 }
369 a.mu.RLock()
370 path := strings.TrimSpace(tab.SessionPath)
371 a.mu.RUnlock()
372 if path != "" && a.attachExistingSessionRuntime(tab, path, a.ctx) {
373 a.clearDeferredRebuildVersion(tabID, sequence)
374 return
375 }
376 if !a.deferredRebuildLeaseLooksFree(tab) {
377 return
378 }
379 err := a.rebuildStartupTabLocked(tab)
380 if err == nil {
381 a.clearDeferredRebuildVersion(tabID, sequence)
382 return
383 }
384 if errors.Is(err, agent.ErrSessionLeaseHeld) {
385 return
386 }
387 a.clearDeferredRebuildVersion(tabID, sequence)
388 slog.Warn("desktop: deferred session startup failed", "tab", tabID, "err", err)
389 }
390
391 func (a *App) tabHasRetryableStartupLeaseError(tab *WorkspaceTab) bool {
392 if tab == nil {
393 return false
394 }
395 a.mu.RLock()
396 defer a.mu.RUnlock()
397 return a.tabs[tab.ID] == tab && !tab.removed && tab.Ctrl == nil && (tab.StartupErrLeaseHeld || tab.modelApplication.startupRetry)
398 }
399
400 func (a *App) rebuildStartupTabLocked(tab *WorkspaceTab) error {
401 buildCtx, cancel := context.WithCancel(a.bootContext())
402 a.mu.Lock()
403 if tab == nil || a.tabs[tab.ID] != tab || tab.removed {
404 a.mu.Unlock()
405 cancel()
406 return nil
407 }
408 if tab.Ctrl != nil {
409 a.mu.Unlock()
410 cancel()
411 return nil
412 }
413 if !tab.StartupErrLeaseHeld && !tab.modelApplication.startupRetry {
414 a.mu.Unlock()
415 cancel()
416 return nil
417 }
418 tab.buildGeneration++
419 generation := tab.buildGeneration
420 if tab.buildCancel != nil {
421 tab.buildCancel()
422 }
423 tab.buildCancel = cancel
424 tab.Ready = false
425 clearTabStartupError(tab)
426 a.setSessionRuntimePhaseLocked(tab, sessionRuntimeStarting, nil)
427 tab.ActivityStatus = ""
428 if tab.sink == nil {
429 tab.sink = &tabEventSink{tabID: tab.ID, app: a, ctx: a.ctx}
430 }
431 a.saveTabsLocked()
432 a.mu.Unlock()
433
434 a.buildTabControllerWithContext(tab, loadedTabSession{}, buildCtx, generation, cancel)
435
436 a.mu.RLock()
437 stillCurrent := false
438 var ctrl control.SessionAPI
439 startupErr := ""
440 leaseHeld := false
441 if tab != nil {
442 stillCurrent = a.tabs[tab.ID] == tab && !tab.removed
443 ctrl = tab.Ctrl
444 startupErr = tab.StartupErr
445 leaseHeld = tab.StartupErrLeaseHeld
446 }
447 a.mu.RUnlock()
448 if !stillCurrent || ctrl != nil {
449 return nil
450 }
451 if leaseHeld {
452 return agent.ErrSessionLeaseHeld
453 }
454 if strings.TrimSpace(startupErr) != "" {
455 return fmt.Errorf("session startup: %s", startupErr)
456 }
457 return fmt.Errorf("session startup: controller was not built")
458 }
459
460 func (a *App) tryRecoverStartupLeaseHeldTab(tab *WorkspaceTab) bool {
461 if a.ctx == nil || !a.tabHasRetryableStartupLeaseError(tab) {
462 return false
463 }
464 a.runtimeRebuildMu.Lock()
465 defer a.runtimeRebuildMu.Unlock()
466 sequence := a.deferredRebuildSequence(tab.ID)
467 if !a.tabHasRetryableStartupLeaseError(tab) {
468 return a.controllerForTab(tab) != nil
469 }
470 if !a.deferredRebuildLeaseLooksFree(tab) {
471 return false
472 }
473 err := a.rebuildStartupTabLocked(tab)
474 if err == nil {
475 a.clearDeferredRebuildVersion(tab.ID, sequence)
476 return a.controllerForTab(tab) != nil
477 }
478 if errors.Is(err, agent.ErrSessionLeaseHeld) {
479 a.scheduleDeferredStartupBuild(tab.ID)
480 } else {
481 a.clearDeferredRebuildVersion(tab.ID, sequence)
482 }
483 return false
484 }
485
486 // deferredRebuildLeaseLooksFree cheaply probes whether the tab's session lease
487 // could be acquired right now, without touching tab.sessionLease (only the
488 // serialized rebuild paths may mutate that). The probe path can lag the
489 // reconciled path the rebuild will use; a stale answer either re-defers on the
490 // next tick or lets the rebuild fail back into the pending set, so a mismatch
491 // only delays the retry.
492 func (a *App) deferredRebuildLeaseLooksFree(tab *WorkspaceTab) bool {
493 a.mu.RLock()
494 ctrl := tab.Ctrl
495 path := strings.TrimSpace(tab.SessionPath)
496 a.mu.RUnlock()
497 if ctrl != nil {
498 if p := strings.TrimSpace(ctrl.SessionPath()); p != "" {
499 path = p
500 }
501 }
502 if path == "" {
503 return true // nothing to probe; let the rebuild decide
504 }
505 key := sessionRuntimeKey(path)
506 a.mu.RLock()
507 rt := a.runtimeBySessionKey[key]
508 ownedByTab := rt != nil && rt.Owner == tab
509 ownedByOther := rt != nil && rt.Owner != nil && rt.Owner != tab && a.runtimeOwnerLiveLocked(rt)
510 a.mu.RUnlock()
511 if ownedByOther {
512 return false
513 }
514 if ownedByTab && tab.sessionLeaseRuntimeKey() == key {
515 return true
516 }
517 lease, err := agent.TryAcquireSessionLease(key)
518 if err != nil {
519 if sameCurrentProcessLease(err) {
520 // The registry ruled out a live sibling owner above, so this is an
521 // orphaned current-process lease. Let the rebuild helper reclaim it
522 // under runtimeRebuildMu instead of looping against our own marker.
523 return true
524 }
525 return !errors.Is(err, agent.ErrSessionLeaseHeld)
526 }
527 lease.Release()
528 return true
529 }
530
530 lines GO