返回 DeepSeek-Reasonix
session_v5_migration.go
根目录 / desktop / session_v5_migration.go
1 package main
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/hex"
7 "encoding/json"
8 "errors"
9 "log/slog"
10 "os"
11 "path/filepath"
12 "strings"
13 "sync"
14
15 "reasonix/desktop/internal/workspacestate"
16 "reasonix/internal/agent"
17 "reasonix/internal/config"
18 "reasonix/internal/fileutil"
19 "reasonix/internal/session"
20 "reasonix/internal/store"
21 )
22
23 type desktopMigrationRecord struct {
24 SourceKey string `json:"sourceKey"`
25 TargetSessionID string `json:"targetSessionId"`
26 ContentDigest string `json:"contentDigest,omitempty"`
27 SourceRevision string `json:"sourceRevision,omitempty"`
28 Status string `json:"status"`
29 ErrorCode string `json:"errorCode,omitempty"`
30 Attempts int `json:"attempts"`
31 PreviousCompletion *desktopMigrationReceipt `json:"previousCompletion,omitempty"`
32 LegacyHeads []string `json:"legacyHeads,omitempty"`
33 LegacySelectedHead string `json:"legacySelectedHead,omitempty"`
34 LegacyPrimaryHead string `json:"legacyPrimaryHead,omitempty"`
35 LegacyHeadsRevision string `json:"legacyHeadsRevision,omitempty"`
36 LegacyAdoption *desktopMigrationReceipt `json:"legacyAdoption,omitempty"`
37 LegacyConversions []desktopMigrationConversion `json:"legacyConversions,omitempty"`
38 }
39
40 type desktopMigrationLedger struct {
41 Version int `json:"version"`
42 Records map[string]desktopMigrationRecord `json:"records"`
43 }
44
45 var desktopMigrationMu sync.Mutex
46
47 func desktopMigrationLedgerPath() string {
48 return filepath.Join(desktopConfigDir(), "desktop", "session-migration-v5.json")
49 }
50
51 func (a *App) startDesktopSessionMigration(ctx context.Context) {
52 if a == nil {
53 return
54 }
55 // Capture reservations before startup admits renderer requests. Background
56 // replay must not abort a new create whose body has not been published yet.
57 startupState, err := a.workspaceRegistry().Load(ctx)
58 if err != nil {
59 a.desktopMigrationFailed.Store(true)
60 slogWarnDesktopMigration(err)
61 close(a.desktopMigrationDone)
62 return
63 }
64 c := &a.historicalImports
65 c.mu.Lock()
66 c.initialize(ctx)
67 if c.stopped || a.shuttingDown.Load() {
68 c.mu.Unlock()
69 close(a.desktopMigrationDone)
70 return
71 }
72 ctx = c.ctx
73 c.catalogEnabled = true
74 c.discoveryPending = true
75 c.workers.Add(1)
76 c.mu.Unlock()
77 go func() {
78 defer c.workers.Done()
79 admitted := false
80 defer func() {
81 if !admitted {
82 close(a.desktopMigrationDone)
83 }
84 }()
85 // Reservation recovery is independent of historical discovery. The
86 // catalog starts metadata discovery after the restored shell exists.
87 if err := a.recoverDesktopPendingCreateSnapshot(ctx, startupState.PendingCreates); err != nil {
88 a.desktopMigrationFailed.Store(true)
89 slogWarnDesktopMigration(err)
90 }
91 // Historical content waits for an explicit request. Keep prepared
92 // reservations intact for the on-demand importer.
93 if err := a.recoverDesktopOperations(ctx, false); err != nil {
94 slogWarnDesktopMigration(err)
95 }
96 // Saved-tab reconciliation waits for reservation recovery. It must be
97 // released before waiting for the shell that reconciliation will create.
98 close(a.desktopMigrationDone)
99 admitted = true
100 select {
101 case <-a.tabsRestoredSignal():
102 case <-ctx.Done():
103 return
104 }
105 c.mu.Lock()
106 c.discoveryPending = false
107 c.mu.Unlock()
108 a.requestHistoricalCatalog()
109 // The watcher owns initial discovery. Recovery changes the registry;
110 // invalidating every legacy root here queues a duplicate full scan.
111 a.emitProjectTreeMetadataChanged()
112 }()
113 }
114
115 // recoverDesktopPendingCreates completes the registry half of a create that
116 // reached durable session publication before the process stopped. A missing
117 // target is safe to forget: no canonical content exists for the pending ID and
118 // the UI can retry creation without inventing a replacement identity.
119 func (a *App) recoverDesktopPendingCreates(ctx context.Context) error {
120 state, err := a.workspaceRegistry().Load(ctx)
121 if err != nil {
122 return err
123 }
124 return a.recoverDesktopPendingCreateSnapshot(ctx, state.PendingCreates)
125 }
126
127 func (a *App) recoverDesktopPendingCreateSnapshot(ctx context.Context, pendingCreates map[string]workspacestate.PendingCreate) error {
128 service := a.desktopSessionService("")
129 var joined error
130 for sessionID, pending := range pendingCreates {
131 ref := session.SessionRef{HostID: localDesktopHostID, SessionID: sessionID}
132 if _, err := service.Query().Snapshot(ctx, ref); err == nil {
133 var attachErr error
134 if strings.HasPrefix(pending.OperationID, "rotate-") {
135 attachErr = a.workspaceRegistry().CommitRotation(ctx, pending.OperationID, pending.WorkspaceID, sessionID, "", pending.ArchiveSource)
136 } else {
137 attachErr = a.workspaceRegistry().AttachSession(ctx, pending.OperationID, pending.WorkspaceID, sessionID, "")
138 }
139 if attachErr != nil {
140 joined = errors.Join(joined, attachErr)
141 } else {
142 a.desktopSessions.pendingCreateRecovered.Add(1)
143 }
144 } else if errors.Is(err, session.ErrSessionNotFound) {
145 if abortErr := a.workspaceRegistry().AbortCreate(ctx, sessionID); abortErr != nil {
146 joined = errors.Join(joined, abortErr)
147 }
148 } else {
149 joined = errors.Join(joined, err)
150 }
151 }
152 return joined
153 }
154
155 func slogWarnDesktopMigration(err error) {
156 // Keep migration logs content- and path-free. Detailed per-source state is
157 // available through the local ledger and UI health row.
158 if err != nil {
159 slog.Warn("desktop session migration incomplete")
160 }
161 }
162
163 type desktopMigrationSource struct {
164 operationID string
165 registeredSourceKey string
166 headID string
167 deferArchive bool
168 versionFingerprint string
169 root string
170 scope string
171 workspaceRoot string
172 exact map[string]bool
173 records map[string]desktopMigrationRecord
174 pairedRoot string
175 pairedIDs map[string]bool
176 legacyAdoption *desktopMigrationReceipt
177 pairedAdoption *desktopMigrationReceipt
178 conversions map[string][]desktopMigrationConversion
179 headConversions []desktopMigrationConversion
180 handledStores map[string]bool
181 conversionAdoptions []*desktopMigrationReceipt
182 }
183
184 func (a *App) migrateDesktopSessionsV5(ctx context.Context) error {
185 replayErr := a.recoverDesktopSessionOperations(ctx)
186 sources, legacySources := a.desktopHistoricalRoots()
187 joined := replayErr
188 conversions, handledStores, conversionErr := discoverDesktopMigrationConversions(ctx, sources)
189 joined = errors.Join(joined, conversionErr)
190 for _, source := range legacySources {
191 source.conversions, source.handledStores = conversions, handledStores
192 joined = errors.Join(joined, markDesktopMigrationPairedStores(source, sources))
193 if err := a.migrateLegacyDirectory(ctx, source); err != nil {
194 joined = errors.Join(joined, err)
195 }
196 }
197 joined = errors.Join(joined, a.migrateStoredConversionLineages(ctx, conversions, sources, handledStores))
198 for _, source := range sources {
199 source.handledStores = handledStores
200 if err := a.migrateCanonicalStore(ctx, *source); err != nil {
201 joined = errors.Join(joined, err)
202 }
203 }
204 return errors.Join(joined, a.discoverHistoricalTrash(ctx), a.reconcileUnregisteredSessions(ctx))
205 }
206
207 // Enumerate trusted storage roots without reading or converting transcripts.
208 func (a *App) desktopHistoricalRoots() (map[string]*desktopMigrationSource, map[string]desktopMigrationSource) {
209 tabs := loadTabsFile()
210 projects := loadProjectsFile()
211 sources := map[string]*desktopMigrationSource{}
212 add := func(scope, workspaceRoot, root string) *desktopMigrationSource {
213 root = filepath.Clean(strings.TrimSpace(root))
214 if root == "." || root == "" || sameDesktopPath(root, a.desktopSessions.root) {
215 return nil
216 }
217 key := canonicalRuntimeRoot(root)
218 if current := sources[key]; current != nil {
219 return current
220 }
221 source := &desktopMigrationSource{root: root, scope: scope, workspaceRoot: workspaceRoot, exact: map[string]bool{}}
222 sources[key] = source
223 return source
224 }
225 addStores := func(scope, workspaceRoot, root string) {
226 if strings.TrimSpace(root) == "" {
227 return
228 }
229 for _, candidate := range desktopLegacyStoreRoots(root) {
230 add(scope, workspaceRoot, candidate)
231 }
232 }
233 addStores("global", "", config.SessionStoreDir())
234 addStores("global", "", config.ProjectSessionStoreDir(globalWorkspaceRoot()))
235 for _, project := range projects.Projects {
236 addStores("project", project.Root, config.ProjectSessionStoreDir(project.Root))
237 }
238 for _, tab := range tabs.Tabs {
239 if strings.TrimSpace(tab.SessionID) == "" {
240 continue
241 }
242 root := config.ProjectSessionStoreDir(globalWorkspaceRoot())
243 if tab.Scope == "project" {
244 root = config.ProjectSessionStoreDir(tab.WorkspaceRoot)
245 }
246 if source := add(tab.Scope, tab.WorkspaceRoot, root); source != nil {
247 source.exact[tab.SessionID] = true
248 }
249 addStores(tab.Scope, tab.WorkspaceRoot, root)
250 }
251 legacySources := map[string]desktopMigrationSource{}
252 addLegacy := func(scope, workspaceRoot, dir string) {
253 dir = filepath.Clean(strings.TrimSpace(dir))
254 if dir == "." || dir == "" {
255 return
256 }
257 key := canonicalRuntimeRoot(dir)
258 if _, ok := legacySources[key]; !ok {
259 legacySources[key] = desktopMigrationSource{root: dir, scope: scope, workspaceRoot: workspaceRoot, exact: map[string]bool{}, pairedRoot: filepath.Join(filepath.Dir(dir), "sessions-v4")}
260 }
261 }
262 addLegacy("global", "", config.SessionDir())
263 addLegacy("global", "", desktopSessionDir(globalWorkspaceRoot()))
264 for _, project := range projects.Projects {
265 addLegacy("project", project.Root, desktopSessionDir(project.Root))
266 }
267 for _, tab := range tabs.Tabs {
268 path := filepath.Clean(strings.TrimSpace(tab.SessionPath))
269 if path == "." || path == "" {
270 continue
271 }
272 dir := filepath.Dir(path)
273 key := canonicalRuntimeRoot(dir)
274 source, ok := legacySources[key]
275 if !ok {
276 scope, workspaceRoot := desktopTabLegacyScope(tab, path, dir)
277 source = desktopMigrationSource{root: dir, scope: scope, workspaceRoot: workspaceRoot, exact: map[string]bool{}, pairedRoot: filepath.Join(filepath.Dir(dir), "sessions-v4")}
278 }
279 source.exact[path] = true
280 legacySources[key] = source
281 }
282 // Include retired roots discovered only through saved legacy tabs before
283 // indexing conversion provenance.
284 for _, source := range legacySources {
285 addStores(source.scope, source.workspaceRoot, source.pairedRoot)
286 }
287 return sources, legacySources
288 }
289
290 // A directory no registered scope owns belongs to the project whose session
291 // directory it is. The session's own recorded root proves that; the saved tab's
292 // scope does not, since the tab may have been rebound after the file was written.
293 func desktopTabLegacyScope(tab desktopTabEntry, path, dir string) (scope, workspaceRoot string) {
294 meta, exists, err := agent.LoadBranchMetaBounded(context.Background(), path)
295 if err == nil && exists && strings.TrimSpace(meta.WorkspaceRoot) != "" && sameDesktopPath(desktopSessionDir(meta.WorkspaceRoot), dir) {
296 return "project", meta.WorkspaceRoot
297 }
298 return tab.Scope, tab.WorkspaceRoot
299 }
300
301 // A paired checkpoint and event store are one migration decision. Never
302 // publish the sidecar independently after a conflict or source failure.
303 func markDesktopMigrationPairedStores(source desktopMigrationSource, sources map[string]*desktopMigrationSource) error {
304 entries, err := os.ReadDir(source.root)
305 if os.IsNotExist(err) {
306 return nil
307 }
308 if err != nil {
309 return err
310 }
311 var joined error
312 for _, entry := range entries {
313 if entry.IsDir() || !store.IsSessionTranscriptName(entry.Name()) || strings.HasPrefix(entry.Name(), ".") {
314 continue
315 }
316 path := filepath.Join(source.root, entry.Name())
317 pairedRoot, err := desktopLegacyPairedRoot(path, source.pairedRoot)
318 if err != nil {
319 joined = errors.Join(joined, err)
320 continue
321 }
322 if paired := sources[canonicalRuntimeRoot(pairedRoot)]; paired != nil {
323 if paired.pairedIDs == nil {
324 paired.pairedIDs = map[string]bool{}
325 }
326 paired.pairedIDs[agent.BranchID(path)] = true
327 }
328 }
329 return joined
330 }
331
332 func (a *App) migrateCanonicalStore(ctx context.Context, source desktopMigrationSource) (retErr error) {
333 if _, err := os.Stat(source.root); os.IsNotExist(err) {
334 return nil
335 } else if err != nil {
336 return err
337 }
338 persistence := session.NewFilesystemPersistence(source.root)
339 old, err := session.NewService("migration-source", persistence)
340 if err != nil {
341 return err
342 }
343 defer func() { retErr = errors.Join(retErr, old.Shutdown(context.Background())) }()
344 ledger, err := readDesktopMigrationLedger()
345 if err != nil {
346 return err
347 }
348 source.records = ledger.Records
349 var joined error
350 var cursor string
351 for {
352 // A disposable catalog cache cannot decide whether durable history
353 // exists. Enumerate every identity, including entries with failed or
354 // missing metadata, without scheduling writes to the source cache.
355 page, err := persistence.List(ctx, cursor, 100)
356 if err != nil {
357 return errors.Join(joined, err)
358 }
359 for _, info := range page.Sessions {
360 if ctx.Err() != nil {
361 return errors.Join(joined, ctx.Err())
362 }
363 if source.pairedIDs[info.SessionID] || source.handledStores[canonicalRuntimeRoot(filepath.Join(source.root, info.SessionID))] {
364 continue
365 }
366 if info.Codec == session.PrototypeCodec || info.Codec == session.LegacyLinearCodec || info.Codec == session.FinalV31Codec {
367 joined = errors.Join(joined, a.migratePreviewSession(ctx, source, info.SessionID))
368 continue
369 }
370 if err := a.migrateCanonicalSession(ctx, old, source, "", info.SessionID); err != nil {
371 joined = errors.Join(joined, err)
372 }
373 }
374 if page.NextCursor == "" {
375 return joined
376 }
377 if page.NextCursor <= cursor {
378 return errors.Join(joined, errors.New("desktop migration source cursor did not advance"))
379 }
380 cursor = page.NextCursor
381 }
382 }
383
384 func (a *App) migrateCanonicalSession(ctx context.Context, old *session.Service, source desktopMigrationSource, workspaceID, sessionID string) error {
385 preview, err := isDesktopStoredPreview(filepath.Join(source.root, sessionID))
386 if err != nil {
387 return errors.Join(err, updateDesktopMigrationLedger(desktopCanonicalMigrationKey(source.root, sessionID), sessionID, "failed", "source_read"))
388 }
389 if preview {
390 return a.migratePreviewSession(ctx, source, sessionID)
391 }
392 key := desktopCanonicalMigrationKey(source.root, sessionID)
393 if source.versionFingerprint != "" {
394 key += ":review:" + source.versionFingerprint
395 }
396 checkpoint, err := newDesktopMigrationCheckpoint(source, key, canonicalMigrationSourceFiles(source.root, sessionID))
397 if err != nil {
398 return err
399 }
400 if handled, err := a.checkAdoptedMigrationSource(ctx, source, checkpoint); handled || err != nil {
401 return err
402 }
403 if checkpoint.unchanged() {
404 return a.completeRegisteredMigration(ctx, source, checkpoint, checkpoint.record.TargetSessionID, checkpoint.record.ContentDigest)
405 }
406 oldRef := session.SessionRef{HostID: "migration-source", SessionID: sessionID}
407 release, err := acquireHistoricalCanonicalRead(filepath.Join(source.root, sessionID))
408 if err != nil {
409 return err
410 }
411 defer release()
412 contentDigest, err := canonicalMigrationDigest(ctx, old.Query(), oldRef)
413 if err != nil {
414 return errors.Join(err, updateDesktopMigrationLedger(key, sessionID, "failed", "source_read"))
415 }
416 if handled, err := a.quarantineChangedMigration(ctx, source, checkpoint, contentDigest); handled || err != nil {
417 return err
418 }
419 if checkpoint.matchesCompletedContent(contentDigest) {
420 return a.completeRegisteredMigration(ctx, source, checkpoint, checkpoint.record.TargetSessionID, contentDigest)
421 }
422 if handled, err := a.reconcileCanonicalConversion(ctx, source, checkpoint, contentDigest); handled || err != nil {
423 return err
424 }
425 if workspaceID == "" {
426 workspaceID, err = a.ensureDesktopMigrationWorkspace(ctx, source)
427 if err != nil {
428 return err
429 }
430 }
431 target := a.desktopSessionService("")
432 targetID, needsImport, err := a.resolveRegisteredMigrationTarget(ctx, source, checkpoint, sessionID, contentDigest)
433 if err != nil {
434 _ = updateDesktopMigrationLedger(key, sessionID, "failed", "target_conflict", contentDigest)
435 return err
436 }
437 if err := updateDesktopMigrationLedger(key, targetID, "pending", "", contentDigest); err != nil {
438 return err
439 }
440 if err := a.prepareRegisteredMigration(ctx, source, checkpoint, targetID, workspaceID); err != nil {
441 return err
442 }
443 if needsImport {
444 tmp, err := os.MkdirTemp("", "reasonix-session-v5-export-")
445 if err != nil {
446 return err
447 }
448 bundle := filepath.Join(tmp, "bundle")
449 defer os.RemoveAll(tmp)
450 if err := old.TryExportCold(ctx, oldRef, bundle); err != nil {
451 _ = updateDesktopMigrationLedger(key, targetID, "failed", "export", contentDigest)
452 return err
453 }
454 if _, err := target.ImportWithHeader(ctx, bundle, session.CreateOptions{
455 SessionID: targetID, CWD: desktopWorkspaceRoot(source.scope, source.workspaceRoot), Origin: session.SessionOriginCanonicalImport,
456 }); err != nil {
457 _ = updateDesktopMigrationLedger(key, targetID, "failed", "import", contentDigest)
458 return err
459 }
460 }
461
462 return a.completeRegisteredMigration(ctx, source, checkpoint, targetID, contentDigest)
463 }
464
465 func (a *App) migrateLegacyDirectory(ctx context.Context, source desktopMigrationSource) error {
466 entries, err := os.ReadDir(source.root)
467 if os.IsNotExist(err) {
468 return nil
469 }
470 if err != nil {
471 return err
472 }
473 ledger, err := readDesktopMigrationLedger()
474 if err != nil {
475 return err
476 }
477 source.records = ledger.Records
478 var joined error
479 for _, entry := range entries {
480 if ctx.Err() != nil {
481 return ctx.Err()
482 }
483 if entry.IsDir() || !store.IsSessionTranscriptName(entry.Name()) || strings.HasPrefix(entry.Name(), ".") {
484 continue
485 }
486 path := filepath.Join(source.root, entry.Name())
487 if desktopMigrationAutomaticRecovery(path) {
488 continue
489 }
490 perFile := source
491 perFile.pairedRoot, err = desktopLegacyPairedRoot(path, source.pairedRoot)
492 if err != nil {
493 joined = errors.Join(joined, err, updateDesktopMigrationLedger(desktopLegacyMigrationKey(path), "", "failed", "source_stat"))
494 continue
495 }
496 if err := a.migrateLegacyHeads(ctx, path, perFile); err != nil {
497 joined = errors.Join(joined, err)
498 }
499 }
500 return joined
501 }
502
503 func (a *App) migrateLegacySession(ctx context.Context, path string, source desktopMigrationSource, workspaceID string) (retErr error) {
504 key := desktopLegacyMigrationKey(path)
505 if source.headID != "" {
506 key = desktopLegacyHeadKey(path, source.headID)
507 }
508 return a.migrateLegacyHead(ctx, path, source, workspaceID, source.headID, key, true)
509 }
510
511 func (a *App) migrateLegacyHead(ctx context.Context, path string, source desktopMigrationSource, workspaceID, headID, key string, paired bool) (retErr error) {
512 source.headID = headID
513 if source.operationID == "" {
514 meta, exists, err := agent.LoadBranchMeta(path)
515 if err != nil {
516 return errors.Join(err, a.sourceRecovery(ctx, path, "legacy", "metadata_unreadable", source.scope, source.workspaceRoot, headID))
517 }
518 if exists && meta.WorkspaceRoot != "" && !sameDesktopPath(meta.WorkspaceRoot, desktopWorkspaceRoot(source.scope, source.workspaceRoot)) {
519 return errors.Join(errSessionWorkspaceConflict, a.sourceRecovery(ctx, path, "legacy", "workspace_conflict", source.scope, source.workspaceRoot, headID))
520 }
521 }
522 if source.versionFingerprint != "" {
523 key += ":review:" + source.versionFingerprint
524 }
525 checkpoint, err := newDesktopMigrationCheckpoint(source, key, desktopLegacyMigrationFiles(path, source))
526 if err != nil {
527 return err
528 }
529 if handled, err := a.checkAdoptedMigrationSource(ctx, source, checkpoint); handled || err != nil {
530 return err
531 }
532 if checkpoint.unchanged() {
533 return a.completeRegisteredMigration(ctx, source, checkpoint, checkpoint.record.TargetSessionID, checkpoint.record.ContentDigest)
534 }
535 if len(source.headConversions) > 0 {
536 if err := a.migrateConversionLineage(ctx, path, headID, source, &checkpoint, workspaceID); err != nil {
537 return errors.Join(err, updateDesktopMigrationLedger(key, checkpoint.record.TargetSessionID, "failed", "conversion_import"))
538 }
539 return nil
540 }
541 if paired && source.pairedRoot != "" {
542 // Previous v5 builds may already have adopted the paired canonical
543 // store before discovering its metadata-less legacy checkpoint.
544 pairedKey := desktopCanonicalMigrationKey(source.pairedRoot, agent.BranchID(path))
545 record := source.records[pairedKey]
546 if record.Status == "completed" {
547 source.pairedAdoption = &desktopMigrationReceipt{TargetSessionID: record.TargetSessionID, ContentDigest: record.ContentDigest}
548 } else {
549 source.pairedAdoption = record.PreviousCompletion
550 }
551 }
552 stageRoot, err := os.MkdirTemp("", "reasonix-legacy-import-")
553 if err != nil {
554 return err
555 }
556 defer os.RemoveAll(stageRoot)
557 stage, err := session.NewService("migration-stage", session.NewFilesystemPersistence(filepath.Join(stageRoot, "sessions-v4")))
558 if err != nil {
559 return err
560 }
561 defer func() { retErr = errors.Join(retErr, stage.Shutdown(context.Background())) }()
562 var runtime *session.Runtime
563 if paired && source.pairedRoot != "" {
564 runtime, _, err = stage.ContinueImportedFrom(ctx, path, source.pairedRoot, headID)
565 } else {
566 runtime, _, err = stage.ContinueImported(ctx, path, headID)
567 }
568 if err != nil {
569 _ = updateDesktopMigrationLedger(key, "", "failed", "legacy_import")
570 var diagnostic *session.TranscriptInitializationError
571 if errors.As(err, &diagnostic) {
572 // The source key matches the local migration ledger. Never log the
573 // source path or the unrestricted error string from imported data.
574 slog.Warn("desktop session migration transcript initialization failed", "source_key", key,
575 "stage", "legacy_import", "diagnostic", diagnostic)
576 queueTranscriptInitializationFailure(diagnostic)
577 }
578 return err
579 }
580 if err := stage.Close(ctx, runtime.Ref()); err != nil {
581 return err
582 }
583 return a.publishStagedMigration(ctx, source, checkpoint, stage, runtime.Ref(), workspaceID, session.SessionOriginLegacyImport)
584 }
585
586 func (a *App) publishStagedMigration(ctx context.Context, source desktopMigrationSource, checkpoint desktopMigrationCheckpoint, stage *session.Service, ref session.SessionRef, workspaceID string, origin session.SessionOrigin) error {
587 key := checkpoint.key
588 target := a.desktopSessionService("")
589 contentDigest, err := canonicalMigrationDigest(ctx, stage.Query(), ref)
590 if err != nil {
591 return err
592 }
593 if handled, err := a.quarantineChangedMigration(ctx, source, checkpoint, contentDigest); handled || err != nil {
594 return err
595 }
596 if checkpoint.matchesCompletedContent(contentDigest) {
597 return a.completeRegisteredMigration(ctx, source, checkpoint, checkpoint.record.TargetSessionID, contentDigest)
598 }
599 // An older migrator adopted only the selected DAG head under the path key.
600 // Recognize that receipt when adding explicit head identities, even if the
601 // selected head has since changed and the v5 target has been continued.
602 if source.legacyAdoption != nil && source.legacyAdoption.ContentDigest == contentDigest {
603 return a.completeRegisteredMigration(ctx, source, checkpoint, source.legacyAdoption.TargetSessionID, contentDigest)
604 }
605 if source.pairedAdoption != nil && source.pairedAdoption.ContentDigest == contentDigest {
606 return a.completeRegisteredMigration(ctx, source, checkpoint, source.pairedAdoption.TargetSessionID, contentDigest)
607 }
608 for _, receipt := range source.conversionAdoptions {
609 if receipt.ContentDigest == contentDigest {
610 return a.completeRegisteredMigration(ctx, source, checkpoint, receipt.TargetSessionID, contentDigest)
611 }
612 }
613 if workspaceID == "" {
614 workspaceID, err = a.ensureDesktopMigrationWorkspace(ctx, source)
615 if err != nil {
616 return err
617 }
618 }
619 targetID, needsImport, err := a.resolveRegisteredMigrationTarget(ctx, source, checkpoint, ref.SessionID, contentDigest)
620 if err != nil {
621 _ = updateDesktopMigrationLedger(key, ref.SessionID, "failed", "target_conflict", contentDigest)
622 return err
623 }
624 if err := updateDesktopMigrationLedger(key, targetID, "pending", "", contentDigest); err != nil {
625 return err
626 }
627 if err := a.prepareRegisteredMigration(ctx, source, checkpoint, targetID, workspaceID); err != nil {
628 return err
629 }
630 if needsImport {
631 tmp, err := os.MkdirTemp("", "reasonix-v5-bundle-")
632 if err != nil {
633 return err
634 }
635 defer os.RemoveAll(tmp)
636 bundle := filepath.Join(tmp, "bundle")
637 if err := stage.Export(ctx, ref, bundle); err != nil {
638 return errors.Join(err, updateDesktopMigrationLedger(key, targetID, "failed", "export", contentDigest))
639 }
640 if _, err := target.ImportWithHeader(ctx, bundle, session.CreateOptions{
641 SessionID: targetID, CWD: desktopWorkspaceRoot(source.scope, source.workspaceRoot), Origin: origin,
642 }); err != nil {
643 _ = updateDesktopMigrationLedger(key, targetID, "failed", "import", contentDigest)
644 return err
645 }
646 }
647
648 return a.completeRegisteredMigration(ctx, source, checkpoint, targetID, contentDigest)
649 }
650
651 func canonicalMigrationDigest(ctx context.Context, query *session.Query, ref session.SessionRef) (string, error) {
652 messages, err := query.History(ctx, ref)
653 if err != nil {
654 return "", err
655 }
656 return agent.ContentDigestForMessages(messages)
657 }
658
659 func resolveMigrationTarget(ctx context.Context, query *session.Query, preferredID, sourceKey, contentDigest string, sourcePaths ...string) (string, bool, error) {
660 check := func(sessionID string) (bool, error) {
661 digest, err := canonicalMigrationDigest(ctx, query, session.SessionRef{HostID: localDesktopHostID, SessionID: sessionID})
662 if errors.Is(err, session.ErrSessionNotFound) {
663 return false, nil
664 }
665 if err != nil {
666 return false, err
667 }
668 if digest != contentDigest {
669 return false, nil
670 }
671 if len(sourcePaths) > 0 {
672 info, err := query.Stat(ctx, session.SessionRef{HostID: localDesktopHostID, SessionID: sessionID})
673 if err != nil {
674 return false, err
675 }
676 body, err := os.ReadFile(filepath.Join(info.Path, "manifest.json"))
677 if err != nil {
678 return false, err
679 }
680 var manifest session.Manifest
681 if err := json.Unmarshal(body, &manifest); err != nil {
682 return false, err
683 }
684 if manifest.Source == nil || !sameDesktopPath(manifest.Source.Path, sourcePaths[0]) {
685 return false, nil
686 }
687 if len(sourcePaths) > 1 && sourcePaths[1] != "" && manifest.Source.LegacyHeadID != sourcePaths[1] {
688 return false, nil
689 }
690 }
691 return true, nil
692 }
693 if identical, err := check(preferredID); err != nil {
694 return "", false, err
695 } else if identical {
696 return preferredID, false, nil
697 } else if _, err := query.Snapshot(ctx, session.SessionRef{HostID: localDesktopHostID, SessionID: preferredID}); errors.Is(err, session.ErrSessionNotFound) {
698 return preferredID, true, nil
699 } else if err != nil {
700 return "", false, err
701 }
702 digest := sha256.Sum256([]byte(sourceKey + "\x00" + contentDigest))
703 conflictID := "migr-" + hex.EncodeToString(digest[:12])
704 if identical, err := check(conflictID); err != nil {
705 return "", false, err
706 } else if identical {
707 return conflictID, false, nil
708 } else if _, err := query.Snapshot(ctx, session.SessionRef{HostID: localDesktopHostID, SessionID: conflictID}); errors.Is(err, session.ErrSessionNotFound) {
709 return conflictID, true, nil
710 } else if err != nil {
711 return "", false, err
712 }
713 return "", false, errors.New("migration target identity collision")
714 }
715
716 // sourceState optionally supplies the content digest followed by a file revision.
717 func updateDesktopMigrationLedger(sourceKey, targetID, status, errorCode string, sourceState ...string) error {
718 desktopMigrationMu.Lock()
719 defer desktopMigrationMu.Unlock()
720 path := desktopMigrationLedgerPath()
721 release, lockErr := lockDesktopMigrationLedger()
722 if lockErr != nil {
723 return lockErr
724 }
725 defer release()
726 ledger, original, err := readDesktopMigrationLedgerFile()
727 if err != nil {
728 return err
729 }
730 record := ledger.Records[sourceKey]
731 // A failed scan or interrupted update must not erase proof that an older
732 // source revision was already adopted (and may have been continued).
733 if record.Status == "completed" && status != "completed" {
734 record.PreviousCompletion = &desktopMigrationReceipt{
735 TargetSessionID: record.TargetSessionID, ContentDigest: record.ContentDigest, SourceRevision: record.SourceRevision,
736 }
737 }
738 record.SourceKey, record.TargetSessionID, record.Status, record.ErrorCode = sourceKey, targetID, status, errorCode
739 if len(sourceState) > 0 {
740 record.ContentDigest = sourceState[0]
741 }
742 record.SourceRevision = ""
743 if status == "completed" && len(sourceState) > 1 {
744 record.SourceRevision = sourceState[1]
745 }
746 if status == "completed" {
747 record.PreviousCompletion = nil
748 }
749 if status == "pending" {
750 record.Attempts++
751 }
752 ledger.Records[sourceKey] = record
753 body, err := marshalDesktopMigrationRecord(original, ledger, sourceKey)
754 if err != nil {
755 return err
756 }
757 if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil {
758 return err
759 }
760 return fileutil.AtomicWriteFileStrict(path, append(body, '\n'), 0o600)
761 }
762
762 lines GO