返回 DeepSeek-Reasonix
session_migration_registry.go
根目录 / desktop / session_migration_registry.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "log/slog"
7 "os"
8 "path/filepath"
9
10 "reasonix/desktop/internal/legacycleanup"
11 "reasonix/desktop/internal/workspacestate"
12 filelock "reasonix/internal/identitylock"
13 "reasonix/internal/session"
14 )
15
16 func migrationCheckpointPath(cp desktopMigrationCheckpoint) (string, string) {
17 path := cp.files[0]
18 if filepath.Base(path) == "manifest.json" {
19 return filepath.Dir(path), "canonical"
20 }
21 return path, "legacy"
22 }
23
24 // Reconcile adoption independently of content comparison: the destination may
25 // have been continued, archived or purged since this receipt was written.
26 func (a *App) completeRegisteredMigration(ctx context.Context, source desktopMigrationSource, cp desktopMigrationCheckpoint, id, digest string) error {
27 path, format := migrationCheckpointPath(cp)
28 state, err := a.workspaceRegistry().Load(ctx)
29 if err != nil {
30 return err
31 }
32 if lifecycle := state.SessionStates[id].Lifecycle; lifecycle == workspacestate.Deleted || lifecycle == workspacestate.Archived {
33 if lifecycle == workspacestate.Archived {
34 if _, err := a.desktopSessionService("").Query().Snapshot(ctx, session.SessionRef{HostID: localDesktopHostID, SessionID: id}); err != nil {
35 return err
36 }
37 }
38 fingerprint, err := desktopSourceFingerprint(path)
39 if err != nil {
40 return err
41 }
42 // Verify the frozen input before publishing a receipt. The final ledger
43 // update remains independently retryable after the registry commit.
44 if err := cp.verify(); err != nil {
45 return err
46 }
47 mapping := workspacestate.SourceMapping{SourceKey: source.mappingKey(path), Path: path, HeadID: source.headID,
48 Format: format, Fingerprint: fingerprint, SessionID: id, WorkspaceID: desktopWorkspaceOwnerID(state, source.scope, source.workspaceRoot)}
49 if old, ok, err := state.ResolveSource(mapping.SourceKey); err != nil {
50 return err
51 } else if ok {
52 if old.SessionID != id {
53 return workspacestate.ErrMutationConflict
54 }
55 return cp.complete(id, digest)
56 }
57 mapping.RetainedArtifacts, err = retainedDesktopArtifacts(path)
58 if err != nil {
59 return err
60 }
61 if err := a.workspaceRegistry().RecordRetiredSource(ctx, mapping, state.Generation); err != nil {
62 return err
63 }
64 return cp.complete(id, digest)
65 }
66 ref := session.SessionRef{HostID: localDesktopHostID, SessionID: id}
67 if _, err := a.desktopSessionService("").Query().Stat(ctx, ref); err != nil {
68 if errors.Is(err, session.ErrSessionNotFound) {
69 return a.sourceRecovery(ctx, path, format, "adopted_target_missing", source.scope, source.workspaceRoot, source.headID)
70 }
71 return err
72 }
73 workspace, err := a.ensureDesktopMigrationWorkspace(ctx, source)
74 if err != nil {
75 return err
76 }
77 fingerprint, err := desktopSourceFingerprint(path)
78 if err != nil {
79 return err
80 }
81 key := source.mappingKey(path)
82 mapping, exists, err := state.ResolveSource(key)
83 if err != nil {
84 return err
85 }
86 if exists && source.operationID == "" {
87 if mapping.SessionID != id {
88 return workspacestate.ErrMutationConflict
89 }
90 if _, err := a.canonicalSessionWorkspace(ctx, ref); err == nil {
91 return cp.complete(id, digest)
92 }
93 }
94 if hook := a.desktopSessions.beforeMigrationRegistryCommit; hook != nil {
95 if err := hook(); err != nil {
96 return errors.Join(err, updateDesktopMigrationLedger(cp.key, id, "failed", "registry", digest))
97 }
98 }
99 if err := a.commitDesktopImport(ctx, source, path, format, fingerprint, id, workspace); err != nil {
100 return err
101 }
102 if format == "legacy" {
103 if err := a.bindLegacyCleanupMigration(ctx, path, source.headID, id, workspace); err != nil &&
104 !errors.Is(err, legacycleanup.ErrNotInitialized) && !errors.Is(err, errLegacyCleanupStateChanged) {
105 slog.Warn("desktop: legacy cleanup migration binding unavailable")
106 }
107 }
108 return cp.complete(id, digest)
109 }
110
111 func (a *App) prepareRegisteredMigration(ctx context.Context, source desktopMigrationSource, cp desktopMigrationCheckpoint, id, workspace string) error {
112 state, err := a.workspaceRegistry().Load(ctx)
113 if err != nil {
114 return err
115 }
116 if _, exists := state.Workspaces[workspace]; !exists {
117 return errors.Join(workspacestate.ErrWorkspaceNotFound, updateDesktopMigrationLedger(cp.key, id, "failed", "registry"))
118 }
119 path, _ := migrationCheckpointPath(cp)
120 fingerprint, err := desktopSourceFingerprint(path)
121 if err != nil {
122 return err
123 }
124 _, err = a.prepareDesktopImport(ctx, source, path, fingerprint, id, workspace)
125 return err
126 }
127
128 func (a *App) resolveRegisteredMigrationTarget(ctx context.Context, source desktopMigrationSource, cp desktopMigrationCheckpoint, preferredID, digest string) (string, bool, error) {
129 path, _ := migrationCheckpointPath(cp)
130 fingerprint, err := desktopSourceFingerprint(path)
131 if err != nil {
132 return "", false, err
133 }
134 return a.resolveDesktopImportTarget(ctx, a.desktopSessionService("").Query(), preferredID, cp.key, source.mappingKey(path), digest, path, fingerprint, source.headID)
135 }
136
137 func (a *App) quarantineChangedMigration(ctx context.Context, source desktopMigrationSource, cp desktopMigrationCheckpoint, digest string) (bool, error) {
138 if source.operationID != "" || !cp.completed() || cp.matchesCompletedContent(digest) {
139 return false, nil
140 }
141 path, format := migrationCheckpointPath(cp)
142 // The old path receipt adopted the selected head, not every head. Adding
143 // an independently identified head is discovery, not a changed adoption.
144 if source.headID != "" && source.legacyAdoption != nil && cp.record.TargetSessionID == source.legacyAdoption.TargetSessionID {
145 state, err := a.workspaceRegistry().Load(ctx)
146 if err != nil {
147 return true, err
148 }
149 if _, exists, err := state.ResolveSource(desktopSourceKey(path, source.headID)); err != nil {
150 return true, err
151 } else if !exists {
152 return false, nil
153 }
154 }
155 return true, a.sourceRecovery(ctx, path, format, "source_changed_after_adoption", source.scope, source.workspaceRoot, source.headID)
156 }
157
158 func lockDesktopMigrationLedger() (func(), error) {
159 path := desktopMigrationLedgerPath()
160 if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
161 return nil, err
162 }
163 return filelock.TryAcquire(path + ".lock")
164 }
165
166 // Registry fingerprints distinguish a metadata-only stat change from a new
167 // historical version without converting or publishing the source again.
168 func (a *App) checkAdoptedMigrationSource(ctx context.Context, source desktopMigrationSource, cp desktopMigrationCheckpoint) (bool, error) {
169 if !cp.completed() || source.operationID != "" {
170 return false, nil
171 }
172 path, format := migrationCheckpointPath(cp)
173 state, err := a.workspaceRegistry().Load(ctx)
174 if err != nil {
175 return true, err
176 }
177 mapping, exists, err := state.ResolveSource(source.mappingKey(path))
178 if err != nil {
179 return true, err
180 }
181 if !exists {
182 return false, nil
183 }
184 fingerprint, err := desktopSourceFingerprint(path)
185 if err != nil {
186 return true, err
187 }
188 if mapping.Fingerprint == fingerprint {
189 return true, a.completeRegisteredMigration(ctx, source, cp, mapping.SessionID, cp.record.ContentDigest)
190 }
191 return true, a.sourceRecovery(ctx, path, format, "source_changed_after_adoption", source.scope, source.workspaceRoot, source.headID)
192 }
193
193 lines GO