返回 DeepSeek-Reasonix
session_v5_migration_lineage.go
根目录 / desktop / session_v5_migration_lineage.go
1 package main
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "os"
8 "path/filepath"
9 "sort"
10 "strings"
11
12 "reasonix/internal/agent"
13 "reasonix/internal/provider"
14 "reasonix/internal/session"
15 "reasonix/internal/store"
16 )
17
18 // A conversion is linked by persisted provenance, never by title or a coincident
19 // ID. LegacyDir identifies the immutable original snapshot of an unnamed head.
20 type desktopMigrationConversion struct {
21 extra map[string]json.RawMessage
22 Root string `json:"root"`
23 SessionID string `json:"sessionId"`
24 HeadID string `json:"headId,omitempty"`
25 LegacyDir string `json:"legacyDir,omitempty"`
26 Codec string `json:"codec"`
27 Depth int `json:"depth"`
28 }
29
30 type desktopMigrationStoredSource struct {
31 dir string
32 manifest session.Manifest
33 }
34
35 func discoverDesktopMigrationConversions(ctx context.Context, sources map[string]*desktopMigrationSource) (map[string][]desktopMigrationConversion, map[string]bool, error) {
36 nodes := map[string]desktopMigrationStoredSource{}
37 var joined error
38 for _, source := range sources {
39 entries, err := os.ReadDir(source.root)
40 if os.IsNotExist(err) {
41 continue
42 }
43 if err != nil {
44 joined = errors.Join(joined, err)
45 continue
46 }
47 for _, entry := range entries {
48 if err := ctx.Err(); err != nil {
49 return nil, nil, err
50 }
51 if !entry.IsDir() || strings.HasPrefix(entry.Name(), ".") {
52 continue
53 }
54 dir := filepath.Join(source.root, entry.Name())
55 manifest, err := readDesktopMigrationManifest(dir)
56 // Normal source migration reports damaged/unsupported stores. A bad
57 // manifest cannot authorize suppressing another source.
58 if err != nil || manifest.SessionID != entry.Name() {
59 continue
60 }
61 nodes[canonicalRuntimeRoot(dir)] = desktopMigrationStoredSource{dir: dir, manifest: manifest}
62 }
63 }
64 var origin func(desktopMigrationStoredSource, map[string]bool) (string, string, string, int, bool)
65 origin = func(node desktopMigrationStoredSource, seen map[string]bool) (string, string, string, int, bool) {
66 key := canonicalRuntimeRoot(node.dir)
67 if seen[key] || len(seen) >= 32 {
68 return "", "", "", 0, false
69 }
70 if node.manifest.Source == nil {
71 return key, "", "", 0, true
72 }
73 seen[key] = true
74 provenance := node.manifest.Source
75 if store.IsSessionTranscriptName(filepath.Base(provenance.Path)) && (provenance.Version == "" || provenance.Version == "legacy") {
76 return canonicalRuntimeRoot(provenance.Path), provenance.LegacyHeadID, filepath.Join(node.dir, "legacy"), 1, true
77 }
78 parent, found := nodes[canonicalRuntimeRoot(provenance.Path)]
79 if !found {
80 // Preview conversion preserves the old manifest inside its archive,
81 // so ancestry remains available after removal of an intermediate store.
82 dir := filepath.Join(node.dir, "legacy", "prototype")
83 manifest, err := readDesktopMigrationManifest(dir)
84 if err != nil {
85 return "", "", "", 0, false
86 }
87 parent = desktopMigrationStoredSource{dir: dir, manifest: manifest}
88 }
89 path, head, legacyDir, depth, ok := origin(parent, seen)
90 if ok && parent.manifest.Source == nil {
91 // Separate conversions archive the same native source in different
92 // directories. Keep its persisted identity after the original is
93 // removed; the archive location is not a new conversation origin.
94 path = canonicalRuntimeRoot(provenance.Path)
95 }
96 return path, head, legacyDir, depth + 1, ok
97 }
98 hidden := desktopMigrationRecoveryStores(nodes)
99 result := map[string][]desktopMigrationConversion{}
100 for key, node := range nodes {
101 if hidden[key] {
102 continue
103 }
104 path, head, legacyDir, depth, ok := origin(node, map[string]bool{})
105 if !ok {
106 continue
107 }
108 result[path] = append(result[path], desktopMigrationConversion{Root: filepath.Dir(node.dir), SessionID: filepath.Base(node.dir), HeadID: head, LegacyDir: legacyDir, Codec: node.manifest.Codec, Depth: depth})
109 }
110 for path := range result {
111 sort.Slice(result[path], func(i, j int) bool {
112 a, b := result[path][i], result[path][j]
113 if a.Depth != b.Depth {
114 return a.Depth > b.Depth
115 }
116 return filepath.Join(a.Root, a.SessionID) < filepath.Join(b.Root, b.SessionID)
117 })
118 }
119 return result, hidden, joined
120 }
121
122 func readDesktopMigrationManifest(dir string) (session.Manifest, error) {
123 body, err := os.ReadFile(filepath.Join(dir, "manifest.json"))
124 if err != nil {
125 return session.Manifest{}, err
126 }
127 var manifest session.Manifest
128 if err := json.Unmarshal(body, &manifest); err != nil {
129 return manifest, err
130 }
131 if manifest.Codec != session.Codec && manifest.Codec != session.PrototypeCodec && manifest.Codec != session.LegacyLinearCodec && manifest.Codec != session.FinalV31Codec {
132 return manifest, session.ErrUnsupportedVersion
133 }
134 return manifest, nil
135 }
136
137 func desktopLegacyMigrationFiles(path string, source desktopMigrationSource) []string {
138 files := desktopLegacySourceFiles(path, source.pairedRoot)
139 for _, converted := range source.conversions[canonicalRuntimeRoot(path)] {
140 files = append(files, canonicalMigrationSourceFiles(converted.Root, converted.SessionID)...)
141 files = append(files, legacyMigrationSourceFiles(filepath.Join(converted.LegacyDir, filepath.Base(path)))...)
142 }
143 return files
144 }
145
146 func resolveDesktopConversionHeads(ctx context.Context, path string, candidates []desktopMigrationConversion) ([]desktopMigrationConversion, error) {
147 resolved := append([]desktopMigrationConversion(nil), candidates...)
148 for index := range resolved {
149 converted := &resolved[index]
150 if converted.HeadID != "" || converted.LegacyDir == "" {
151 continue
152 }
153 // An omitted head means the selected head at conversion time, not the
154 // source's current selection. Read the preserved snapshot to recover it.
155 frozenPath := filepath.Join(converted.LegacyDir, filepath.Base(path))
156 heads, err := session.LegacyMigrationHeads(ctx, frozenPath)
157 if err != nil {
158 return nil, err
159 }
160 for _, head := range heads {
161 if head.Selected && !head.Retired {
162 converted.HeadID = head.ID
163 break
164 }
165 }
166 }
167 return resolved, nil
168 }
169
170 type desktopMigrationStagedNode struct {
171 service *session.Service
172 checkpoint desktopMigrationCheckpoint
173 ref session.SessionRef
174 messages []provider.Message
175 digest string
176 targetID string
177 coveredBy int
178 origin session.SessionOrigin
179 }
180
181 // All generations of one legacy head are staged before publication. Maximal
182 // histories survive; equal/prefix ancestors get receipts for the same target.
183 // Incomparable continuations remain separate conversations.
184 func (a *App) migrateConversionLineage(ctx context.Context, path, headID string, source desktopMigrationSource, cp *desktopMigrationCheckpoint, workspaceID string) (retErr error) {
185 tmp, err := os.MkdirTemp("", "reasonix-lineage-")
186 if err != nil {
187 return err
188 }
189 defer os.RemoveAll(tmp)
190 stages := []*session.Service{}
191 defer func() {
192 for _, stage := range stages {
193 retErr = errors.Join(retErr, stage.Shutdown(context.Background()))
194 }
195 }()
196 newStage := func() (*session.Service, error) {
197 // A continued descendant can have the same ID/source stamp as the
198 // import of its ancestor. Isolate inputs so reuse cannot substitute a
199 // descendant's newer history while computing the ancestor's digest.
200 dir, err := os.MkdirTemp(tmp, "source-")
201 if err != nil {
202 return nil, err
203 }
204 stage, err := session.NewService("migration-stage", session.NewFilesystemPersistence(filepath.Join(dir, "sessions-v4")))
205 if err == nil {
206 stages = append(stages, stage)
207 }
208 return stage, err
209 }
210 nodes := []desktopMigrationStagedNode{}
211 add := func(stage *session.Service, ref session.SessionRef, checkpoint desktopMigrationCheckpoint, origin session.SessionOrigin) error {
212 if err := stage.Close(ctx, ref); err != nil {
213 return err
214 }
215 messages, err := stage.Query().History(ctx, ref)
216 if err != nil {
217 return err
218 }
219 digest, err := agent.ContentDigestForMessages(messages)
220 if err != nil {
221 return err
222 }
223 nodes = append(nodes, desktopMigrationStagedNode{service: stage, checkpoint: checkpoint, ref: ref, messages: messages, digest: digest, coveredBy: -1, origin: origin})
224 return nil
225 }
226 for _, converted := range source.headConversions {
227 key := desktopCanonicalMigrationKey(converted.Root, converted.SessionID)
228 checkpoint, err := newDesktopMigrationCheckpoint(source, key, canonicalMigrationSourceFiles(converted.Root, converted.SessionID))
229 if err != nil {
230 return err
231 }
232 stage, err := newStage()
233 if err != nil {
234 return err
235 }
236 var runtime *session.Runtime
237 if converted.Codec == session.Codec {
238 runtime, _, err = stage.ContinueImportedSource(ctx, path, filepath.Join(converted.Root, converted.SessionID), headID)
239 } else {
240 runtime, _, err = stage.ContinuePrototype(ctx, filepath.Join(converted.Root, converted.SessionID))
241 }
242 if err != nil {
243 return errors.Join(err, updateDesktopMigrationLedger(key, "", "failed", "conversion_read"))
244 }
245 if err := add(stage, runtime.Ref(), checkpoint, session.SessionOriginCanonicalImport); err != nil {
246 return err
247 }
248 }
249 if cp != nil {
250 stage, err := newStage()
251 if err != nil {
252 return err
253 }
254 runtime, _, err := stage.ContinueImported(ctx, path, headID)
255 if err != nil {
256 return err
257 }
258 if err := add(stage, runtime.Ref(), *cp, session.SessionOriginLegacyImport); err != nil {
259 return err
260 }
261 }
262 for _, node := range nodes {
263 if node.checkpoint.completed() {
264 source.conversionAdoptions = append(source.conversionAdoptions, &desktopMigrationReceipt{TargetSessionID: node.checkpoint.record.TargetSessionID, ContentDigest: node.checkpoint.record.ContentDigest})
265 }
266 }
267 reduceDesktopMigrationLineage(nodes)
268 for i := range nodes {
269 if nodes[i].coveredBy >= 0 {
270 continue
271 }
272 if err := a.publishStagedMigration(ctx, source, nodes[i].checkpoint, nodes[i].service, nodes[i].ref, workspaceID, nodes[i].origin); err != nil {
273 return err
274 }
275 ledger, err := readDesktopMigrationLedger()
276 if err != nil {
277 return err
278 }
279 nodes[i].targetID = ledger.Records[nodes[i].checkpoint.key].TargetSessionID
280 }
281 for i := range nodes {
282 if nodes[i].coveredBy < 0 {
283 continue
284 }
285 winner := nodes[i].coveredBy
286 for nodes[winner].coveredBy >= 0 {
287 winner = nodes[winner].coveredBy
288 }
289 // Preserve an older receipt for a target that may already have been
290 // continued/deleted. Do not detach or resurrect that user's target.
291 targetID := nodes[winner].targetID
292 if nodes[i].checkpoint.matchesCompletedContent(nodes[i].digest) {
293 targetID = nodes[i].checkpoint.record.TargetSessionID
294 }
295 if err := a.completeRegisteredMigration(ctx, source, nodes[i].checkpoint, targetID, nodes[i].digest); err != nil {
296 return err
297 }
298 }
299 return nil
300 }
301
302 // Select maximal histories within an already-proven lineage. Equal histories
303 // prefer the earlier (more recently converted) node; prefix edges cannot cycle.
304 func reduceDesktopMigrationLineage(nodes []desktopMigrationStagedNode) {
305 for i := range nodes {
306 for j := range nodes {
307 if i == j || !session.MigrationHistoryContains(nodes[j].messages, nodes[i].messages) {
308 continue
309 }
310 if session.MigrationHistoryContains(nodes[i].messages, nodes[j].messages) && j > i {
311 continue
312 }
313 nodes[i].coveredBy = j
314 break
315 }
316 }
317 }
318
319 // Stored-only chains (for example native v3 -> v4 with both directories left
320 // behind) use the same reduction even when no JSONL source remains on disk.
321 func (a *App) migrateStoredConversionLineages(ctx context.Context, conversions map[string][]desktopMigrationConversion, sources map[string]*desktopMigrationSource, handled map[string]bool) error {
322 var joined error
323 for path, candidates := range conversions {
324 remaining := []desktopMigrationConversion{}
325 for _, candidate := range candidates {
326 if handled[canonicalRuntimeRoot(filepath.Join(candidate.Root, candidate.SessionID))] {
327 continue
328 }
329 owner := sources[canonicalRuntimeRoot(candidate.Root)]
330 if owner != nil && owner.pairedIDs[candidate.SessionID] {
331 continue
332 }
333 remaining = append(remaining, candidate)
334 }
335 if len(remaining) < 2 {
336 continue
337 }
338 // Head recovery can require replaying an archived DAG. Completed
339 // stored-only chains must bypass that work just like live legacy files.
340 ledger, err := readDesktopMigrationLedger()
341 if err != nil {
342 return errors.Join(joined, err)
343 }
344 allUnchanged := true
345 for _, candidate := range remaining {
346 cp, err := newDesktopMigrationCheckpoint(desktopMigrationSource{records: ledger.Records}, desktopCanonicalMigrationKey(candidate.Root, candidate.SessionID), canonicalMigrationSourceFiles(candidate.Root, candidate.SessionID))
347 if err != nil || !cp.unchanged() {
348 allUnchanged = false
349 break
350 }
351 if err := cp.skip(); err != nil {
352 joined = errors.Join(joined, err)
353 allUnchanged = false
354 break
355 }
356 }
357 if allUnchanged {
358 for _, candidate := range remaining {
359 handled[canonicalRuntimeRoot(filepath.Join(candidate.Root, candidate.SessionID))] = true
360 }
361 continue
362 }
363 resolved, err := resolveDesktopConversionHeads(ctx, path, remaining)
364 if err != nil {
365 joined = errors.Join(joined, err)
366 for _, candidate := range remaining {
367 handled[canonicalRuntimeRoot(filepath.Join(candidate.Root, candidate.SessionID))] = true
368 joined = errors.Join(joined, updateDesktopMigrationLedger(desktopCanonicalMigrationKey(candidate.Root, candidate.SessionID), "", "failed", "conversion_heads"))
369 }
370 continue
371 }
372 groups := map[string][]desktopMigrationConversion{}
373 for _, candidate := range resolved {
374 groups[candidate.HeadID] = append(groups[candidate.HeadID], candidate)
375 }
376 for _, group := range groups {
377 if len(group) < 2 {
378 continue
379 }
380 owner := sources[canonicalRuntimeRoot(group[0].Root)]
381 if owner == nil {
382 continue
383 }
384 source := *owner
385 source.headConversions = group
386 ledger, err := readDesktopMigrationLedger()
387 if err != nil {
388 return errors.Join(joined, err)
389 }
390 source.records = ledger.Records
391 unchanged := true
392 for _, candidate := range group {
393 handled[canonicalRuntimeRoot(filepath.Join(candidate.Root, candidate.SessionID))] = true
394 cp, err := newDesktopMigrationCheckpoint(source, desktopCanonicalMigrationKey(candidate.Root, candidate.SessionID), canonicalMigrationSourceFiles(candidate.Root, candidate.SessionID))
395 if err != nil {
396 joined = errors.Join(joined, err)
397 unchanged = false
398 continue
399 }
400 if !cp.unchanged() {
401 unchanged = false
402 } else if err := cp.skip(); err != nil {
403 joined = errors.Join(joined, err)
404 unchanged = false
405 }
406 }
407 if unchanged {
408 continue
409 }
410 if err := a.migrateConversionLineage(ctx, "", "", source, nil, ""); err != nil {
411 joined = errors.Join(joined, err)
412 }
413 }
414 }
415 return joined
416 }
417
417 lines GO