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