| 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 |