返回 DeepSeek-Reasonix
historical_reconciliation.go
根目录 / desktop / historical_reconciliation.go
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
255 lines GO