| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "fmt" |
| 7 | "log/slog" |
| 8 | "os" |
| 9 | "path/filepath" |
| 10 | "strings" |
| 11 | |
| 12 | "reasonix/desktop/internal/workspacestate" |
| 13 | "reasonix/internal/historywork" |
| 14 | "reasonix/internal/session" |
| 15 | ) |
| 16 | |
| 17 | // Repair adoption before publishing ordinary rows. Discovery owns this work; |
| 18 | // sidebar pagination only reads its result and never starts a second migration. |
| 19 | func (a *App) reconcileHistoricalCatalog(ctx context.Context, rows []historicalCatalogEntry, sources map[string]historicalSource) error { |
| 20 | for i := range rows { |
| 21 | if err := ctx.Err(); err != nil { |
| 22 | return err |
| 23 | } |
| 24 | entry := &rows[i] |
| 25 | source := sources[entry.node.Source.SourceKey] |
| 26 | release, err := a.historyMaintenance.BackgroundSlice(ctx, false) |
| 27 | if err != nil { |
| 28 | return err |
| 29 | } |
| 30 | err = a.reconcileHistoricalCatalogSource(ctx, entry.node.Source.SourceKey, source, entry.node.Health) |
| 31 | release(historywork.ReadChunk) |
| 32 | if err != nil && !historicalSourceBusyError(err) && !errors.Is(err, os.ErrPermission) && !errors.Is(err, context.Canceled) { |
| 33 | // Unavailable rows remain discoverable for a later retry, but are not |
| 34 | // advertised as usable conversations. No source bytes are removed. |
| 35 | entry.node.Health = "unavailable" |
| 36 | slog.Warn("desktop: historical source recovery deferred", "source_key", entry.node.Source.SourceKey) |
| 37 | } |
| 38 | } |
| 39 | return ctx.Err() |
| 40 | } |
| 41 | |
| 42 | func (a *App) reconcileHistoricalCatalogSource(ctx context.Context, key string, source historicalSource, health string) error { |
| 43 | state, err := a.workspaceRegistry().Load(ctx) |
| 44 | if err != nil { |
| 45 | return err |
| 46 | } |
| 47 | mapping, mapped, resolveErr := state.ResolveSource(key) |
| 48 | if mapped && a.historicalAdoptionReady(ctx, state, mapping, key) { |
| 49 | return nil |
| 50 | } |
| 51 | if resolveErr != nil && !errors.Is(resolveErr, workspacestate.ErrMutationConflict) { |
| 52 | return resolveErr |
| 53 | } |
| 54 | receipts, err := historicalReconciliationReceipts(ctx, source) |
| 55 | if err != nil { |
| 56 | return err |
| 57 | } |
| 58 | if !mapped && resolveErr == nil && len(receipts) == 0 && (health == "ok" || health == "metadata_pending") { |
| 59 | return nil |
| 60 | } |
| 61 | releaseRuntime, ok := a.tryLockRuntimeMutation("reconcile historical source") |
| 62 | if !ok { |
| 63 | return errHistoricalSourceBusy |
| 64 | } |
| 65 | defer releaseRuntime() |
| 66 | release, err := acquireHistoricalSource(ctx, key, source) |
| 67 | if err != nil { |
| 68 | return err |
| 69 | } |
| 70 | defer release() |
| 71 | migration := desktopMigrationSource{scope: source.scope, workspaceRoot: source.root, headID: source.head, registeredSourceKey: key} |
| 72 | cp, err := historicalReconciliationCheckpoint(source, migration) |
| 73 | if err != nil { |
| 74 | return err |
| 75 | } |
| 76 | fingerprint, err := desktopSourceFingerprint(source.path) |
| 77 | if err != nil { |
| 78 | return err |
| 79 | } |
| 80 | old, oldRef, finish, err := openHistoricalReconciliationSource(ctx, source) |
| 81 | if err != nil { |
| 82 | return err |
| 83 | } |
| 84 | defer finish() |
| 85 | digest, err := canonicalMigrationDigest(ctx, old.Query(), oldRef) |
| 86 | if err != nil { |
| 87 | return err |
| 88 | } |
| 89 | if !mapped && resolveErr == nil && len(receipts) == 0 { |
| 90 | return nil |
| 91 | } |
| 92 | target := mapping.SessionID |
| 93 | if !mapped { |
| 94 | target, err = a.selectCanonicalConversionTarget(ctx, migration, source.path, digest, receipts) |
| 95 | if err != nil { |
| 96 | return err |
| 97 | } |
| 98 | if target == "" { |
| 99 | return workspacestate.ErrMutationConflict |
| 100 | } |
| 101 | } |
| 102 | if mapped && fingerprint != mapping.Fingerprint { |
| 103 | return fmt.Errorf("retained source changed: %w", workspacestate.ErrMutationConflict) |
| 104 | } |
| 105 | if err := cp.verify(); err != nil { |
| 106 | return err |
| 107 | } |
| 108 | if err := a.repairMissingHistoricalTarget(ctx, old, oldRef, migration, source.path, target, fingerprint); err != nil { |
| 109 | return fmt.Errorf("repair target: %w", err) |
| 110 | } |
| 111 | return a.recordHistoricalReconciliation(ctx, source, migration, cp, target, digest, fingerprint) |
| 112 | } |
| 113 | |
| 114 | func (a *App) historicalAdoptionReady(ctx context.Context, state workspacestate.State, mapping workspacestate.SourceMapping, key string) bool { |
| 115 | if state.SessionStates[mapping.SessionID].Lifecycle == workspacestate.Deleted { |
| 116 | return true |
| 117 | } |
| 118 | if !containsDesktopString(state.Workspaces[mapping.WorkspaceID].SessionIDs, mapping.SessionID) || pendingHistoricalOperation(state, key) != nil { |
| 119 | return false |
| 120 | } |
| 121 | ref := session.SessionRef{HostID: localDesktopHostID, SessionID: mapping.SessionID} |
| 122 | info, err := a.desktopSessionService("").Query().Stat(ctx, ref) |
| 123 | return err == nil && info.Error == "" && info.MetadataStatus != session.MetadataFailed |
| 124 | } |
| 125 | |
| 126 | func (a *App) recordHistoricalReconciliation(ctx context.Context, source historicalSource, migration desktopMigrationSource, cp desktopMigrationCheckpoint, target, digest, fingerprint string) error { |
| 127 | state, err := a.workspaceRegistry().Load(ctx) |
| 128 | if err != nil { |
| 129 | return err |
| 130 | } |
| 131 | workspace := desktopWorkspaceOwnerID(state, source.scope, source.root) |
| 132 | if state.SessionStates[target].Lifecycle == workspacestate.Deleted || containsDesktopString(state.Workspaces[workspace].SessionIDs, target) { |
| 133 | if actual, err := desktopSourceFingerprint(source.path); err != nil || actual != fingerprint { |
| 134 | return errors.Join(err, workspacestate.ErrMutationConflict) |
| 135 | } |
| 136 | artifacts, err := retainedDesktopArtifacts(source.path) |
| 137 | if err != nil { |
| 138 | return err |
| 139 | } |
| 140 | if err := cp.verify(); err != nil { |
| 141 | return err |
| 142 | } |
| 143 | if err := a.workspaceRegistry().RecordRecoveredSource(ctx, workspacestate.SourceMapping{SourceKey: migration.mappingKey(source.path), Path: source.path, Format: source.format, |
| 144 | HeadID: source.head, SessionID: target, WorkspaceID: workspace, Fingerprint: fingerprint, RetainedArtifacts: artifacts}, state.Generation); err != nil { |
| 145 | return err |
| 146 | } |
| 147 | return cp.complete(target, digest) |
| 148 | } |
| 149 | return a.completeRegisteredMigration(ctx, migration, cp, target, digest) |
| 150 | } |
| 151 | |
| 152 | // Recover only an absent destination, using its original identity. Existing |
| 153 | // content (including a damaged or continued target) is never overwritten. |
| 154 | // An ordinary import journal makes publication restartable and lifecycle-fenced. |
| 155 | func (a *App) repairMissingHistoricalTarget(ctx context.Context, old *session.Service, oldRef session.SessionRef, source desktopMigrationSource, path, id, fingerprint string) error { |
| 156 | state, err := a.workspaceRegistry().Load(ctx) |
| 157 | if err != nil { |
| 158 | return err |
| 159 | } |
| 160 | if state.SessionStates[id].Lifecycle == workspacestate.Deleted { |
| 161 | return nil |
| 162 | } |
| 163 | ref := session.SessionRef{HostID: localDesktopHostID, SessionID: id} |
| 164 | query := a.desktopSessionService("").Query() |
| 165 | if _, err := query.Snapshot(ctx, ref); err == nil { |
| 166 | if op := pendingHistoricalOperation(state, source.mappingKey(path)); op != nil && strings.HasPrefix(op.ID, "repair-") && op.Mapping.SessionID == id { |
| 167 | return a.finishHistoricalTargetRepair(ctx, *op) |
| 168 | } |
| 169 | return nil |
| 170 | } else if !errors.Is(err, session.ErrSessionNotFound) { |
| 171 | return err |
| 172 | } |
| 173 | if _, err := os.Lstat(filepath.Join(a.desktopSessions.root, id)); !os.IsNotExist(err) { |
| 174 | return errors.Join(workspacestate.ErrMutationConflict, err) |
| 175 | } |
| 176 | return a.restoreHistoricalTargetBundle(ctx, old, oldRef, source, path, id, fingerprint, state) |
| 177 | } |
| 178 | |
| 179 | func (a *App) restoreHistoricalTargetBundle(ctx context.Context, old *session.Service, ref session.SessionRef, source desktopMigrationSource, path, id, fingerprint string, state workspacestate.State) error { |
| 180 | workspace, err := a.ensureDesktopMigrationWorkspace(ctx, source) |
| 181 | if err != nil { |
| 182 | return err |
| 183 | } |
| 184 | // This temporary export also validates every referenced blob before any |
| 185 | // registry reservation or native target is published. |
| 186 | tmp, err := os.MkdirTemp("", "reasonix-history-repair-") |
| 187 | if err != nil { |
| 188 | return err |
| 189 | } |
| 190 | defer os.RemoveAll(tmp) |
| 191 | bundle := filepath.Join(tmp, "bundle") |
| 192 | if err := old.TryExportCold(ctx, ref, bundle); err != nil { |
| 193 | return err |
| 194 | } |
| 195 | a.lifecycleCheckpoint("historical-repair-source-exported") |
| 196 | if actual, err := desktopSourceFingerprint(path); err != nil || actual != fingerprint { |
| 197 | return errors.Join(err, workspacestate.ErrMutationConflict) |
| 198 | } |
| 199 | source.operationID = fmt.Sprintf("repair-%s-%d", source.mappingKey(path), state.Generation) |
| 200 | if op := pendingHistoricalOperation(state, source.mappingKey(path)); op != nil && strings.HasPrefix(op.ID, "repair-") && op.Mapping.SessionID == id { |
| 201 | source.operationID = op.ID |
| 202 | } |
| 203 | lifecycle := state.SessionStates[id].Lifecycle |
| 204 | if lifecycle == "" { |
| 205 | lifecycle = workspacestate.Active |
| 206 | } |
| 207 | presentation := state.Presentation[id] |
| 208 | if err := a.workspaceRegistry().BeginOperation(ctx, workspacestate.Operation{ID: source.operationID, Presentation: &presentation, |
| 209 | Kind: "import", WorkspaceID: workspace, SessionIDs: []string{id}, Lifecycle: lifecycle, |
| 210 | ExpectedGeneration: state.Generation}); err != nil { |
| 211 | return err |
| 212 | } |
| 213 | if _, err := a.prepareDesktopImport(ctx, source, path, fingerprint, id, workspace); err != nil { |
| 214 | return fmt.Errorf("reserve repair: %w", err) |
| 215 | } |
| 216 | origin := session.SessionOriginCanonicalImport |
| 217 | if info, err := os.Stat(path); err != nil { |
| 218 | return err |
| 219 | } else if !info.IsDir() { |
| 220 | origin = session.SessionOriginLegacyImport |
| 221 | } |
| 222 | if _, err := a.desktopSessionService("").ImportWithHeader(ctx, bundle, session.CreateOptions{SessionID: id, CWD: desktopWorkspaceRoot(source.scope, source.workspaceRoot), Origin: origin}); err != nil { |
| 223 | return err |
| 224 | } |
| 225 | a.lifecycleCheckpoint("historical-repair-content-published") |
| 226 | current, err := a.workspaceRegistry().Load(ctx) |
| 227 | if err != nil { |
| 228 | return err |
| 229 | } |
| 230 | return a.finishHistoricalTargetRepair(ctx, current.PendingOperations[source.operationID]) |
| 231 | } |
| 232 | |
| 233 | func (a *App) finishHistoricalTargetRepair(ctx context.Context, op workspacestate.Operation) error { |
| 234 | if op.Mapping == nil || len(op.SessionIDs) != 1 || op.Mapping.SessionID != op.SessionIDs[0] { |
| 235 | return workspacestate.ErrMutationConflict |
| 236 | } |
| 237 | path, fingerprint := op.Mapping.Path, op.Mapping.Fingerprint |
| 238 | if actual, err := desktopSourceFingerprint(path); err != nil || actual != fingerprint { |
| 239 | return errors.Join(err, workspacestate.ErrMutationConflict) |
| 240 | } |
| 241 | artifacts, err := retainedDesktopArtifacts(path) |
| 242 | if err != nil { |
| 243 | return err |
| 244 | } |
| 245 | mapping := *op.Mapping |
| 246 | mapping.RetainedArtifacts = artifacts |
| 247 | if err := a.workspaceRegistry().PrepareOperationContent(ctx, op.ID, op.SessionIDs, &mapping, op.Presentation); err != nil { |
| 248 | return err |
| 249 | } |
| 250 | if err := a.workspaceRegistry().CommitOperation(ctx, op.ID); err != nil { |
| 251 | return fmt.Errorf("commit repair: %w", err) |
| 252 | } |
| 253 | return nil |
| 254 | } |
| 255 |