返回 DeepSeek-Reasonix
session_lifecycle.go
根目录 / desktop / session_lifecycle.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "slices"
8 "sort"
9 "strings"
10
11 "reasonix/desktop/internal/workspacestate"
12 "reasonix/internal/agent"
13 "reasonix/internal/control"
14 "reasonix/internal/session"
15 )
16
17 func (a *App) archiveSessionRefsWithOperation(refs []session.SessionRef, operationID string, dependencies ...string) error {
18 return a.archiveSessionRefsWithOperationConditional(refs, operationID, nil, dependencies...)
19 }
20
21 // archiveSessionRefsWithOperationConditional runs verify after runtime and
22 // filesystem maintenance ownership has been acquired, while session removal is
23 // still serialized. It is used by maintenance jobs whose read decision must be
24 // fenced from a concurrent title/content/runtime mutation.
25 //
26 // Archiving the last visible session leaves the surface empty on purpose: the
27 // frontend lands on the workspace draft instead of a replacement blank session.
28 func (a *App) archiveSessionRefsWithOperationConditional(refs []session.SessionRef, operationID string, verify func(context.Context, workspacestate.State) error, dependencies ...string) error {
29 a.sessionRemovalMu.Lock()
30 defer a.sessionRemovalMu.Unlock()
31 cleanup := archivedRuntimeCleanup{app: a}
32 defer cleanup.finish()
33 return a.archiveSessionRefsRemovalHeld(refs, operationID, verify, &cleanup, dependencies...)
34 }
35
36 // archiveSessionRefsRemovalHeld commits and detaches only. The caller holds
37 // sessionRemovalMu and must finish cleanup after releasing title/index locks.
38 func (a *App) archiveSessionRefsRemovalHeld(refs []session.SessionRef, operationID string, verify func(context.Context, workspacestate.State) error, cleanup *archivedRuntimeCleanup, dependencies ...string) error {
39 ctx := a.bootContext()
40 service := a.desktopSessionService("")
41 unique := map[string]session.SessionRef{}
42 for _, ref := range refs {
43 if err := validateLocalSessionRef(ref); err != nil {
44 return err
45 }
46 a.cancelAISessionTitle((SessionTarget{SessionRef: ref}).key())
47 unique[ref.SessionID] = ref
48 }
49 if len(unique) == 0 {
50 return errors.New("no sessions to archive")
51 }
52 ids := make([]string, 0, len(unique))
53 for id := range unique {
54 ids = append(ids, id)
55 }
56 sort.Strings(ids)
57 state, err := a.workspaceRegistry().Load(ctx)
58 if err != nil {
59 return err
60 }
61 legacyTargets := map[string]bool{}
62 for _, dependency := range dependencies {
63 if mapping := state.PendingOperations[dependency].Mapping; mapping != nil {
64 legacyTargets[sessionRuntimeKey(mapping.Path)] = true
65 }
66 }
67 removed, err := a.idleArchiveRuntimes(ids, legacyTargets)
68 if err != nil {
69 return err
70 }
71 guards := []func(){}
72 staged := map[string]string{}
73 for _, dependency := range dependencies {
74 op, ok := state.PendingOperations[dependency]
75 if !ok || op.Kind != "archive-import" || op.Phase != "content_ready" {
76 return workspacestate.ErrMutationConflict
77 }
78 for _, id := range op.SessionIDs {
79 staged[id] = op.WorkspaceID
80 }
81 }
82 defer func() {
83 for _, release := range slices.Backward(guards) {
84 release()
85 }
86 }()
87 for _, id := range ids {
88 ref := unique[id]
89 if err := a.validateConditionalArchiveWorkspace(ctx, state, ref, staged[id], verify != nil); err != nil {
90 return err
91 }
92 if runtime, live := service.Runtime(ref); live {
93 phase := runtime.StateSnapshot().Phase
94 if phase != session.RuntimeIdle && phase != session.RuntimeRecoveryRequired {
95 return errTopicHasActiveWork
96 }
97 } else {
98 guard, err := session.NewFilesystemPersistence(a.desktopSessions.root).AcquireMaintenance(id)
99 if err != nil {
100 return userFacingSessionLeaseError("", err)
101 }
102 guards = append(guards, guard)
103 }
104 if err := requireArchivableSession(ctx, service, ref); err != nil {
105 return err
106 }
107 }
108 state, err = a.workspaceRegistry().Load(ctx)
109 if err != nil {
110 return err
111 }
112 op := workspacestate.Operation{ID: operationID, Kind: "archive", Lifecycle: workspacestate.Archived, SessionIDs: ids, ExpectedGeneration: state.Generation, Dependencies: dependencies}
113 if err := a.beginConditionalArchiveOperation(ctx, state, op, verify); err != nil {
114 return err
115 }
116 if err := a.workspaceRegistry().PrepareOperationContent(ctx, op.ID, ids, nil, nil); err != nil {
117 return err
118 }
119 a.lifecycleCheckpoint("before-archive-commit")
120 if err := a.workspaceRegistry().CommitOperation(ctx, op.ID); err != nil {
121 return err
122 }
123 a.detachArchivedRuntimeBindings(removed)
124 cleanup.removed = append(cleanup.removed, removed...)
125 for _, ref := range unique {
126 cleanup.refs = append(cleanup.refs, ref)
127 }
128 return nil
129 }
130
131 func (a *App) validateConditionalArchiveWorkspace(ctx context.Context, state workspacestate.State, ref session.SessionRef, stagedWorkspaceID string, maintenance bool) error {
132 if stagedWorkspaceID != "" {
133 return a.validateDesktopWorkspaceMembership(ctx, stagedWorkspaceID, ref)
134 }
135 if !maintenance {
136 _, err := a.canonicalSessionWorkspace(ctx, ref)
137 return err
138 }
139 // A final content/lifecycle CAS lets maintenance archive proven-empty
140 // sessions even when immutable and registry workspace ownership disagree.
141 if !conditionalArchiveRegistryHasSingleActiveOwner(state, ref.SessionID) {
142 return workspacestate.ErrMutationConflict
143 }
144 return nil
145 }
146
147 func conditionalArchiveRegistryHasSingleActiveOwner(state workspacestate.State, sessionID string) bool {
148 if state.SessionStates[sessionID].Lifecycle != workspacestate.Active {
149 return false
150 }
151 owners := 0
152 for _, workspace := range state.Workspaces {
153 if slices.Contains(workspace.SessionIDs, sessionID) {
154 owners++
155 }
156 }
157 return owners == 1
158 }
159
160 func (a *App) beginConditionalArchiveOperation(ctx context.Context, state workspacestate.State, op workspacestate.Operation, verify func(context.Context, workspacestate.State) error) error {
161 if verify != nil {
162 if err := verify(ctx, state); err != nil {
163 return err
164 }
165 }
166 return a.workspaceRegistry().BeginOperation(ctx, op)
167 }
168
169 // Called only after durable commit, with runtime mutation admission held.
170 func (a *App) finishArchivedRuntimeBindings(removed []removedSessionRuntime) {
171 a.detachArchivedRuntimeBindings(removed)
172 a.finalizeRemovedTopicRuntimes(removed)
173 a.closeRemainingRemovedSessionRuntimesAdmissionHeld(removed, map[control.SessionAPI]bool{})
174 }
175
176 func (a *App) detachArchivedRuntimeBindings(removed []removedSessionRuntime) {
177 a.mu.Lock()
178 for _, item := range removed {
179 tab := item.tab
180 if tab.Ctrl != item.ctrl {
181 continue
182 }
183 a.markTabRemovedLocked(tab)
184 stopTabAutosave(tab)
185 a.releaseSessionRuntimeLocked(tab)
186 a.unregisterDetachedRuntimeLocked(tab)
187 delete(a.tabs, tab.ID)
188 a.removeTabOrderLocked(tab.ID)
189 if a.activeTabID == tab.ID {
190 a.activeTabID = ""
191 }
192 }
193 if a.activeTabID == "" && len(a.tabOrder) > 0 {
194 a.activeTabID = a.tabOrder[0]
195 }
196 var dir, activeID string
197 var entries []desktopTabEntry
198 var version uint64
199 if len(removed) > 0 {
200 dir, entries, activeID, version = a.saveTabsCollectLocked()
201 }
202 a.mu.Unlock()
203 if len(removed) > 0 {
204 a.saveTabsWrite(dir, entries, activeID, version)
205 }
206 }
207
208 func (a *App) archiveCompatibleTopic(topicID string) error {
209 release, ok := a.tryLockRuntimeMutationBounded("archive topic")
210 if !ok {
211 return errTopicArchiveBusy
212 }
213 defer release()
214 return a.archiveCompatibleTopicAdmissionHeld(topicID, "")
215 }
216
217 func (a *App) archiveCompatibleTopicAdmissionHeld(topicID, operationID string) error {
218 return a.archiveCompatibleTopicWithCleanupAdmissionHeld(topicID, operationID, nil)
219 }
220
221 // A non-nil cleanup means the caller already holds sessionRemovalMu, followed
222 // by the title/index locks, for a confirmed session-backed topic removal.
223 func (a *App) archiveCompatibleTopicWithCleanupAdmissionHeld(topicID, operationID string, cleanup *archivedRuntimeCleanup) error {
224 topicID = strings.TrimSpace(topicID)
225 if topicID == "" {
226 return fmt.Errorf("topicID is required")
227 }
228 if a.topicHasActiveRuntimeWork(topicID) {
229 return errTopicHasActiveWork
230 }
231 refs := map[string]session.SessionRef{}
232 state, err := a.workspaceRegistry().Load(a.bootContext())
233 if err != nil {
234 return err
235 }
236 for id, presentation := range state.Presentation {
237 if presentation.TopicID == topicID {
238 refs[id] = session.SessionRef{HostID: localDesktopHostID, SessionID: id}
239 }
240 }
241 if id, ok := strings.CutPrefix(topicID, "canonical-"); ok {
242 if _, registered := state.SessionStates[id]; registered {
243 refs[id] = session.SessionRef{HostID: localDesktopHostID, SessionID: id}
244 }
245 }
246 a.mu.RLock()
247 for _, tab := range a.runtimeTabsLocked() {
248 if tab != nil && tab.TopicID == topicID && tab.SessionID != "" {
249 refs[tab.SessionID] = session.SessionRef{HostID: localDesktopHostID, SessionID: tab.SessionID}
250 }
251 }
252 a.mu.RUnlock()
253 owners := a.captureTopicRuntimeBindings(topicID)
254 if err := a.snapshotTopicRuntimeBindings(owners); err != nil {
255 return err
256 }
257 // Snapshot first: a legacy runtime may publish its first durable file here.
258 // Only a topic with neither canonical identities nor legacy content may use
259 // the metadata-only removal path.
260 targets, err := a.topicTrashTargets(topicID)
261 if err != nil {
262 return err
263 }
264 if err := a.validateCompatibleTopicOwner(state, topicID, len(targets) > 0); err != nil {
265 return err
266 }
267 if len(refs) == 0 && len(targets) == 0 {
268 return a.removeEmptyCompatibleTopicAdmissionHeld(topicID, cleanup)
269 }
270 // Originals stay in place, so retain existing leases and acquire only cold
271 // sources. The importer recognizes these same-process owners when freezing.
272 localOwners := topicArchiveLeaseOwners(owners)
273 leases := []*agent.SessionLease{}
274 defer func() {
275 for _, lease := range leases {
276 lease.Release()
277 }
278 }()
279 for _, target := range targets {
280 if localOwners[sessionRuntimeKey(target.sessionPath)] != nil {
281 continue
282 }
283 lease, err := agent.TryAcquireSessionLease(target.sessionPath)
284 if err != nil {
285 if errors.Is(err, agent.ErrSessionLeaseHeld) {
286 return errSessionBusyElsewhere
287 }
288 return err
289 }
290 leases = append(leases, lease)
291 }
292 dependencies := []string{}
293 for _, target := range targets {
294 ref, dependency, err := a.stageArchiveSource(a.bootContext(), target.sessionPath)
295 if err != nil {
296 return err
297 }
298 if dependency != "" {
299 dependencies = append(dependencies, dependency)
300 }
301 refs[ref.SessionID] = ref
302 }
303 list := make([]session.SessionRef, 0, len(refs))
304 for _, ref := range refs {
305 list = append(list, ref)
306 }
307 if operationID == "" {
308 operationID = "archive-" + newTabID()
309 }
310 if cleanup != nil {
311 err = a.archiveSessionRefsRemovalHeld(list, operationID, nil, cleanup, dependencies...)
312 } else {
313 err = a.archiveSessionRefsWithOperation(list, operationID, dependencies...)
314 }
315 if err != nil {
316 return err
317 }
318 for _, lease := range leases {
319 lease.Release()
320 }
321 leases = nil
322 a.emitProjectTreeChanged()
323 return nil
324 }
325
326 func (a *App) removeEmptyCompatibleTopicAdmissionHeld(topicID string, cleanup *archivedRuntimeCleanup) error {
327 // A confirmed session removal must not silently turn into placeholder
328 // removal (which also acquires the removal/title locks).
329 if cleanup != nil {
330 return workspacestate.ErrMutationConflict
331 }
332 return a.removeCompatiblePlaceholderAdmissionHeld(topicID)
333 }
334
335 func (a *App) stageArchiveSource(ctx context.Context, path string) (session.SessionRef, string, error) {
336 if ref, found, err := a.legacyCanonicalRef(ctx, path); found || err != nil {
337 return ref, "", err
338 }
339 meta, _, err := agent.LoadBranchMeta(path)
340 if err != nil {
341 return session.SessionRef{}, "", err
342 }
343 scope, root := "global", ""
344 if meta.WorkspaceRoot != "" && !a.isGlobalWorkspacePath(ctx, meta.WorkspaceRoot) {
345 scope, root = "project", meta.WorkspaceRoot
346 }
347 workspaceID, err := a.ensureDesktopWorkspace(ctx, scope, root)
348 if err != nil {
349 return session.SessionRef{}, "", err
350 }
351 fingerprint, err := desktopSourceFingerprint(path)
352 if err != nil {
353 return session.SessionRef{}, "", err
354 }
355 if err := a.migrateLegacySession(ctx, path, desktopMigrationSource{scope: scope, workspaceRoot: root, deferArchive: true}, workspaceID); err != nil {
356 return session.SessionRef{}, "", err
357 }
358 opID := "archive-import-" + desktopSourceKey(path, "") + "-" + fingerprint
359 state, err := a.workspaceRegistry().Load(ctx)
360 if err != nil {
361 return session.SessionRef{}, "", err
362 }
363 op, exists := state.PendingOperations[opID]
364 if !exists || op.Phase != "content_ready" || len(op.SessionIDs) != 1 {
365 return session.SessionRef{}, "", workspacestate.ErrMutationConflict
366 }
367 return session.SessionRef{HostID: localDesktopHostID, SessionID: op.SessionIDs[0]}, opID, nil
368 }
369
370 func (a *App) restoreCanonicalSession(ctx context.Context, ref session.SessionRef, operationID string, recoveryIDs ...string) (SessionRestoreResult, error) {
371 if err := validateLocalSessionRef(ref); err != nil {
372 return SessionRestoreResult{}, err
373 }
374 workspace, err := a.canonicalSessionWorkspace(ctx, ref)
375 if errors.Is(err, errSessionWorkspaceConflict) && len(recoveryIDs) == 1 {
376 info, readErr := a.desktopSessionService("").Query().Stat(ctx, ref)
377 if readErr != nil {
378 return SessionRestoreResult{}, readErr
379 }
380 if info.CWD == "" || info.Origin == "" {
381 return SessionRestoreResult{}, errSessionWorkspaceConflict
382 }
383 scope, root := "project", info.CWD
384 if a.isGlobalWorkspacePath(ctx, root) {
385 scope, root = "global", ""
386 }
387 workspaceID, ensureErr := a.ensureDesktopWorkspace(ctx, scope, root)
388 if ensureErr != nil {
389 return SessionRestoreResult{}, ensureErr
390 }
391 state, loadErr := a.workspaceRegistry().Load(ctx)
392 if loadErr != nil {
393 return SessionRestoreResult{}, loadErr
394 }
395 workspace, err = state.Workspaces[workspaceID], nil
396 }
397 if err != nil {
398 return SessionRestoreResult{}, err
399 }
400 if err := requireArchivableSession(ctx, a.desktopSessionService(""), ref); err != nil {
401 return SessionRestoreResult{}, err
402 }
403 if operationID == "" {
404 operationID = "restore-" + newTabID()
405 }
406 state, err := a.workspaceRegistry().Load(ctx)
407 if err != nil {
408 return SessionRestoreResult{}, err
409 }
410 op := workspacestate.Operation{ID: operationID, Kind: "restore", WorkspaceID: workspace.ID, Lifecycle: workspacestate.Active, SessionIDs: []string{ref.SessionID}, ExpectedGeneration: state.Generation}
411 if len(recoveryIDs) == 1 {
412 op.RecoveryEntryID = recoveryIDs[0]
413 }
414 if err := a.workspaceRegistry().BeginOperation(ctx, op); err != nil {
415 return SessionRestoreResult{}, err
416 }
417 if err := a.workspaceRegistry().PrepareOperationContent(ctx, op.ID, op.SessionIDs, nil, nil); err != nil {
418 return SessionRestoreResult{}, err
419 }
420 if err := a.workspaceRegistry().CommitOperation(ctx, op.ID); err != nil {
421 return SessionRestoreResult{}, err
422 }
423 a.markLegacyCleanupSessionRestored(ref.SessionID)
424 state, err = a.workspaceRegistry().Load(ctx)
425 if err != nil {
426 return SessionRestoreResult{}, err
427 }
428 a.emitProjectTreeChanged()
429 return SessionRestoreResult{Session: ref, WorkspaceID: workspace.ID, Generation: state.PendingOperations[op.ID].ResultGeneration}, nil
430 }
431
432 type SessionRestoreResult struct {
433 Session session.SessionRef `json:"session"`
434 WorkspaceID string `json:"workspaceId"`
435 Generation uint64 `json:"generation"`
436 }
437
438 func (a *App) idleArchiveRuntimes(ids []string, legacyTargets map[string]bool) ([]removedSessionRuntime, error) {
439 a.mu.RLock()
440 removed := []removedSessionRuntime{}
441 for _, tab := range a.runtimeTabsLocked() {
442 if tab == nil || (!containsDesktopString(ids, tab.SessionID) && !legacyTargets[sessionRuntimeKey(tab.currentSessionPath())]) {
443 continue
444 }
445 item := removedRuntimeFromTab(tab, tabRuntimeSessionDir(tab), tab.currentSessionPath())
446 item.failedStartup = a.suppressTabStartupRestoreLocked(tab)
447 removed = append(removed, item)
448 }
449 a.mu.RUnlock()
450 for _, item := range removed {
451 if item.ctrl != nil && controllerHasActiveRuntimeWork(item.ctrl) {
452 return nil, errTopicHasActiveWork
453 }
454 }
455 if err := a.snapshotTopicRuntimeBindings(removed); err != nil {
456 return nil, err
457 }
458 return removed, nil
459 }
460
461 // requireArchivableSession proves the session exists. Moving a session in or
462 // out of the archive never reads its content, so a damaged store stays movable:
463 // archiving is how a user sets aside a session that no longer opens.
464 func requireArchivableSession(ctx context.Context, service *session.Service, ref session.SessionRef) error {
465 _, err := service.Query().Snapshot(ctx, ref)
466 if errors.Is(err, session.ErrDamagedStore) {
467 return nil
468 }
469 return err
470 }
471
471 lines GO