返回 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 controllerHasActiveRuntimeWork(ctrl) {
259 return
260 }
261 if !a.deferredRebuildLeaseLooksFree(tab) {
262 return
263 }
264 setting := request.label
265 err := a.rebuildSettingTurnLocked(setting, tab, false, false)
266 if err == nil {
267 // rebuildSettingLocked already cleared the pending entry for the tab it
268 // refreshed; just announce it.
269 a.noticeForTab(tabID, fmt.Sprintf("%s applied: session refreshed after the lease was released", setting))
270 return
271 }
272 if errors.Is(err, agent.ErrSessionLeaseHeld) {
273 return // grabbed back before we could rebuild; keep waiting
274 }
275 var busy *rebuildBusyError
276 if errors.As(err, &busy) {
277 return // a turn started meanwhile; retry once it finishes
278 }
279 // Anything else will not resolve by waiting; give up loudly instead of
280 // retrying forever.
281 a.clearDeferredRebuildVersion(tabID, request.sequence)
282 slog.Warn("desktop: deferred settings rebuild failed", "setting", setting, "tab", tabID, "err", err)
283 a.warnForTab(tabID, fmt.Sprintf("%s was saved but the session could not refresh: %s", setting, err.Error()))
284 }
285
286 // retryDeferredRuntimeReload drives one queued ReloadRuntime pass. The
287 // probing contract mirrors retryDeferredRebuild: wait for the tab to be
288 // active and idle and for its lease to look free, then run the boot.Rebuild
289 // reload; busy/lease answers keep waiting, anything else gives up loudly.
290 func (a *App) retryDeferredRuntimeReload(tabID string, tab *WorkspaceTab, queuedSequence ...uint64) {
291 sequence := a.deferredRebuildSequence(tabID)
292 if len(queuedSequence) > 0 {
293 sequence = queuedSequence[0]
294 }
295 // Hold the rebuild mutex across probe + reload: the probe briefly
296 // acquires the session lease, and a concurrent rebuild's ensure lease
297 // would read that probe as "held by another runtime" and spuriously
298 // defer (same contract as retryDeferredRebuild).
299 a.runtimeRebuildMu.Lock()
300 defer a.runtimeRebuildMu.Unlock()
301 ctrl := a.controllerForTab(tab)
302 if ctrl == nil {
303 // Mid-(re)build on another path; racing a second build+swap against it
304 // is what this loop must avoid.
305 return
306 }
307 if controllerHasActiveRuntimeWork(ctrl) {
308 return
309 }
310 if !a.deferredRebuildLeaseLooksFree(tab) {
311 return
312 }
313 err := a.reloadRuntimeTurnLocked(tab)
314 if err == nil {
315 // rebuildSettingTurnLocked already cleared the pending entry for the
316 // tab it refreshed; just announce it.
317 a.noticeForTab(tabID, "runtime reloaded after the session went idle")
318 return
319 }
320 if errors.Is(err, agent.ErrSessionLeaseHeld) {
321 return // grabbed back before we could reload; keep waiting
322 }
323 var busy *rebuildBusyError
324 if errors.As(err, &busy) {
325 return // a turn started meanwhile; retry once it finishes
326 }
327 // Anything else will not resolve by waiting; give up loudly instead of
328 // retrying forever. The error may come from provider/config plumbing and
329 // carry credential-shaped values (passwords, resolved API keys) — the
330 // tested helper redacts before the text reaches logs or the frontend.
331 a.clearDeferredRebuildVersion(tabID, sequence)
332 failure := deferredReloadFailedText(err)
333 slog.Warn("desktop: "+failure, "tab", tabID)
334 a.warnForTab(tabID, failure)
335 }
336
337 // deferredReloadFailedText is the failure line for an unrecoverable deferred
338 // reload — the single formatter both the log and the tab warning use, so a
339 // credential-shaped error can never reach either sink unredacted.
340 func deferredReloadFailedText(err error) string {
341 return "runtime reload failed: " + secrets.RedactCredentials(err.Error())
342 }
343
344 func (a *App) retryDeferredStartupBuild(tabID string, tab *WorkspaceTab, queuedSequence ...uint64) {
345 sequence := a.deferredRebuildSequence(tabID)
346 if len(queuedSequence) > 0 {
347 sequence = queuedSequence[0]
348 }
349 a.runtimeRebuildMu.Lock()
350 defer a.runtimeRebuildMu.Unlock()
351 if !a.tabHasRetryableStartupLeaseError(tab) {
352 a.clearDeferredRebuildVersion(tabID, sequence)
353 return
354 }
355 a.mu.RLock()
356 path := strings.TrimSpace(tab.SessionPath)
357 a.mu.RUnlock()
358 if path != "" && a.attachExistingSessionRuntime(tab, path, a.ctx) {
359 a.clearDeferredRebuildVersion(tabID, sequence)
360 return
361 }
362 if !a.deferredRebuildLeaseLooksFree(tab) {
363 return
364 }
365 err := a.rebuildStartupTabLocked(tab)
366 if err == nil {
367 a.clearDeferredRebuildVersion(tabID, sequence)
368 return
369 }
370 if errors.Is(err, agent.ErrSessionLeaseHeld) {
371 return
372 }
373 a.clearDeferredRebuildVersion(tabID, sequence)
374 slog.Warn("desktop: deferred session startup failed", "tab", tabID, "err", err)
375 }
376
377 func (a *App) tabHasRetryableStartupLeaseError(tab *WorkspaceTab) bool {
378 if tab == nil {
379 return false
380 }
381 a.mu.RLock()
382 defer a.mu.RUnlock()
383 return a.tabs[tab.ID] == tab && !tab.removed && tab.Ctrl == nil && (tab.StartupErrLeaseHeld || tab.modelApplication.startupRetry)
384 }
385
386 func (a *App) rebuildStartupTabLocked(tab *WorkspaceTab) error {
387 buildCtx, cancel := context.WithCancel(a.bootContext())
388 a.mu.Lock()
389 if tab == nil || a.tabs[tab.ID] != tab || tab.removed {
390 a.mu.Unlock()
391 cancel()
392 return nil
393 }
394 if tab.Ctrl != nil {
395 a.mu.Unlock()
396 cancel()
397 return nil
398 }
399 if !tab.StartupErrLeaseHeld && !tab.modelApplication.startupRetry {
400 a.mu.Unlock()
401 cancel()
402 return nil
403 }
404 tab.buildGeneration++
405 generation := tab.buildGeneration
406 if tab.buildCancel != nil {
407 tab.buildCancel()
408 }
409 tab.buildCancel = cancel
410 tab.Ready = false
411 clearTabStartupError(tab)
412 a.setSessionRuntimePhaseLocked(tab, sessionRuntimeStarting, nil)
413 tab.ActivityStatus = ""
414 if tab.sink == nil {
415 tab.sink = &tabEventSink{tabID: tab.ID, app: a, ctx: a.ctx}
416 }
417 a.saveTabsLocked()
418 a.mu.Unlock()
419
420 a.buildTabControllerWithContext(tab, loadedTabSession{}, buildCtx, generation, cancel)
421
422 a.mu.RLock()
423 stillCurrent := false
424 var ctrl control.SessionAPI
425 startupErr := ""
426 leaseHeld := false
427 if tab != nil {
428 stillCurrent = a.tabs[tab.ID] == tab && !tab.removed
429 ctrl = tab.Ctrl
430 startupErr = tab.StartupErr
431 leaseHeld = tab.StartupErrLeaseHeld
432 }
433 a.mu.RUnlock()
434 if !stillCurrent || ctrl != nil {
435 return nil
436 }
437 if leaseHeld {
438 return agent.ErrSessionLeaseHeld
439 }
440 if strings.TrimSpace(startupErr) != "" {
441 return fmt.Errorf("session startup: %s", startupErr)
442 }
443 return fmt.Errorf("session startup: controller was not built")
444 }
445
446 func (a *App) tryRecoverStartupLeaseHeldTab(tab *WorkspaceTab) bool {
447 if a.ctx == nil || !a.tabHasRetryableStartupLeaseError(tab) {
448 return false
449 }
450 a.runtimeRebuildMu.Lock()
451 defer a.runtimeRebuildMu.Unlock()
452 sequence := a.deferredRebuildSequence(tab.ID)
453 if !a.tabHasRetryableStartupLeaseError(tab) {
454 return a.controllerForTab(tab) != nil
455 }
456 if !a.deferredRebuildLeaseLooksFree(tab) {
457 return false
458 }
459 err := a.rebuildStartupTabLocked(tab)
460 if err == nil {
461 a.clearDeferredRebuildVersion(tab.ID, sequence)
462 return a.controllerForTab(tab) != nil
463 }
464 if errors.Is(err, agent.ErrSessionLeaseHeld) {
465 a.scheduleDeferredStartupBuild(tab.ID)
466 } else {
467 a.clearDeferredRebuildVersion(tab.ID, sequence)
468 }
469 return false
470 }
471
472 // deferredRebuildLeaseLooksFree cheaply probes whether the tab's session lease
473 // could be acquired right now, without touching tab.sessionLease (only the
474 // serialized rebuild paths may mutate that). The probe path can lag the
475 // reconciled path the rebuild will use; a stale answer either re-defers on the
476 // next tick or lets the rebuild fail back into the pending set, so a mismatch
477 // only delays the retry.
478 func (a *App) deferredRebuildLeaseLooksFree(tab *WorkspaceTab) bool {
479 a.mu.RLock()
480 ctrl := tab.Ctrl
481 path := strings.TrimSpace(tab.SessionPath)
482 a.mu.RUnlock()
483 if ctrl != nil {
484 if p := strings.TrimSpace(ctrl.SessionPath()); p != "" {
485 path = p
486 }
487 }
488 if path == "" {
489 return true // nothing to probe; let the rebuild decide
490 }
491 key := sessionRuntimeKey(path)
492 a.mu.RLock()
493 rt := a.runtimeBySessionKey[key]
494 ownedByTab := rt != nil && rt.Owner == tab
495 ownedByOther := rt != nil && rt.Owner != nil && rt.Owner != tab && a.runtimeOwnerLiveLocked(rt)
496 a.mu.RUnlock()
497 if ownedByOther {
498 return false
499 }
500 if ownedByTab && tab.sessionLeaseRuntimeKey() == key {
501 return true
502 }
503 lease, err := agent.TryAcquireSessionLease(key)
504 if err != nil {
505 if sameCurrentProcessLease(err) {
506 // The registry ruled out a live sibling owner above, so this is an
507 // orphaned current-process lease. Let the rebuild helper reclaim it
508 // under runtimeRebuildMu instead of looping against our own marker.
509 return true
510 }
511 return !errors.Is(err, agent.ErrSessionLeaseHeld)
512 }
513 lease.Release()
514 return true
515 }
516
516 lines GO