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