返回 DeepSeek-Reasonix
session_source_compatibility.go
根目录 / desktop / session_source_compatibility.go
1 package main
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/hex"
7 "errors"
8 "fmt"
9 "io"
10 "os"
11 "path/filepath"
12 "slices"
13 "sort"
14 "strings"
15
16 "reasonix/desktop/internal/workspacestate"
17 "reasonix/internal/agent"
18 "reasonix/internal/config"
19 "reasonix/internal/fileutil"
20 "reasonix/internal/session"
21 "reasonix/internal/topicstate"
22 )
23
24 func desktopSourceKey(path, head string) string {
25 // Migration sources include both transcript files and canonical prototype
26 // directories. Keep their persisted key independent from the runtime
27 // session locator, which intentionally accepts transcript paths only.
28 pathKey := agent.CanonicalSessionPath(cleanDesktopPath(path))
29 return agent.SessionSourceKeyFromIdentity(pathKey, head)
30 }
31
32 func (source desktopMigrationSource) mappingKey(path string) string {
33 if source.registeredSourceKey != "" {
34 return source.registeredSourceKey
35 }
36 key := desktopSourceKey(path, source.headID)
37 if source.versionFingerprint != "" {
38 key += ":review:" + source.versionFingerprint
39 }
40 return key
41 }
42
43 // Fingerprints cover source bytes, not the destination's evolving projection.
44 // A continued canonical session must never be replaced by its frozen import.
45 func desktopSourceFingerprint(path string) (string, error) {
46 info, err := os.Lstat(path)
47 if err != nil {
48 return "", err
49 }
50 if info.Mode()&os.ModeSymlink != 0 {
51 return "", errors.New("session source is a symbolic link")
52 }
53 paths := []string{}
54 if info.IsDir() {
55 for _, name := range []string{"manifest.json", "header.json", "events.frames", "events.jsonl"} {
56 candidate := filepath.Join(path, name)
57 if _, err := os.Lstat(candidate); err == nil {
58 paths = append(paths, candidate)
59 } else if !os.IsNotExist(err) {
60 return "", err
61 }
62 }
63 } else {
64 paths = append(paths, path)
65 for _, artifact := range sessionTrashArtifacts(path, filepath.Base(path)) {
66 if artifact.src == path {
67 continue
68 }
69 // DAG selection lives in the event log. Branch .meta also contains
70 // mutable catalog/title projections; hashing it would reject a source
71 // simply because background indexing refreshed its display metadata.
72 if strings.HasSuffix(artifact.name, ".events.jsonl") {
73 if _, err := os.Lstat(artifact.src); err == nil {
74 paths = append(paths, artifact.src)
75 } else if !os.IsNotExist(err) {
76 return "", err
77 }
78 }
79 }
80 }
81 if len(paths) == 0 {
82 return "", errors.New("session source has no durable records")
83 }
84 sort.Strings(paths)
85 h := sha256.New()
86 for _, file := range paths {
87 info, err := os.Lstat(file)
88 if err != nil {
89 return "", err
90 }
91 if !info.Mode().IsRegular() {
92 return "", errors.New("session source contains a non-regular file")
93 }
94 fmt.Fprintf(h, "%s\x00%d\x00", filepath.Base(file), info.Size())
95 f, err := os.Open(file)
96 if err != nil {
97 return "", err
98 }
99 _, copyErr := io.Copy(h, f)
100 closeErr := f.Close()
101 if copyErr != nil {
102 return "", copyErr
103 }
104 if closeErr != nil {
105 return "", closeErr
106 }
107 }
108 return hex.EncodeToString(h.Sum(nil)), nil
109 }
110
111 func (a *App) recordDesktopSource(ctx context.Context, path, format, fingerprint, targetID, workspaceID string) error {
112 presentation := workspacestate.Presentation{SortOrder: -1, TopicID: legacySessionTopicID(path)}
113 if meta, ok, err := agent.LoadBranchMeta(path); err == nil && ok {
114 if meta.TopicID != "" {
115 presentation.TopicID = meta.TopicID
116 }
117 presentation.Title = meta.TopicTitle
118 }
119 presentation = a.historicalTopicPresentation(ctx, workspaceID, presentation)
120 return a.workspaceRegistry().RecordSource(ctx, workspacestate.SourceMapping{
121 SourceKey: desktopSourceKey(path, ""), Path: path, Format: format, Fingerprint: fingerprint,
122 SessionID: targetID, WorkspaceID: workspaceID,
123 }, presentation)
124 }
125
126 func (a *App) commitDesktopImport(ctx context.Context, source desktopMigrationSource, path, format, fingerprint, targetID, workspaceID string) error {
127 current, err := desktopSourceFingerprint(path)
128 if err != nil {
129 return err
130 }
131 if current != fingerprint {
132 return workspacestate.ErrMutationConflict
133 }
134 ref := session.SessionRef{HostID: localDesktopHostID, SessionID: targetID}
135 if err := a.validateDesktopWorkspaceMembership(ctx, workspaceID, ref); err != nil {
136 return err
137 }
138 if _, err := a.desktopSessionService("").Query().Snapshot(ctx, ref); err != nil {
139 return err
140 }
141 key := source.mappingKey(path)
142 mapping := workspacestate.SourceMapping{SourceKey: key, Path: path, HeadID: source.headID, Format: format, Fingerprint: fingerprint, SessionID: targetID, WorkspaceID: workspaceID}
143 mapping.RetainedArtifacts, err = retainedDesktopArtifacts(path)
144 if err != nil {
145 return err
146 }
147 presentation := workspacestate.Presentation{SortOrder: -1}
148 if format == "legacy" {
149 presentation.TopicID = legacySessionTopicID(path)
150 }
151 if meta, ok, err := agent.LoadBranchMeta(path); err == nil && ok {
152 if meta.TopicID != "" {
153 presentation.TopicID = meta.TopicID
154 }
155 presentation.Title = meta.TopicTitle
156 }
157 presentation = a.historicalTopicPresentation(ctx, workspaceID, presentation)
158 opID, err := a.prepareDesktopImport(ctx, source, path, fingerprint, targetID, workspaceID)
159 if err != nil {
160 return err
161 }
162 if err := a.workspaceRegistry().PrepareOperationContent(ctx, opID, []string{targetID}, &mapping, &presentation); err != nil {
163 return err
164 }
165 if source.deferArchive {
166 return nil
167 }
168 return a.workspaceRegistry().CommitOperation(ctx, opID)
169 }
170
171 func (a *App) historicalTopicPresentation(ctx context.Context, workspaceID string, presentation workspacestate.Presentation) workspacestate.Presentation {
172 projects := loadProjectsFile()
173 if workspaceID == workspacestate.GlobalWorkspaceID {
174 return historicalTopicPresentationFrom(projects, workspaceID, "", presentation)
175 }
176 workspaceRoot := ""
177 if state, err := a.workspaceRegistry().Load(ctx); err == nil {
178 workspaceRoot = state.Workspaces[workspaceID].Root
179 }
180 return historicalTopicPresentationFrom(projects, workspaceID, workspaceRoot, presentation)
181 }
182
183 func historicalTopicPresentationFrom(projects desktopProjectFile, workspaceID, workspaceRoot string, presentation workspacestate.Presentation) workspacestate.Presentation {
184 topics, pinned := projects.GlobalTopics, projects.GlobalPinnedTopics
185 if workspaceID != workspacestate.GlobalWorkspaceID {
186 topics, pinned = nil, nil
187 for _, project := range projects.Projects {
188 if (workspaceRoot != "" && sameProjectRoot(project.Root, workspaceRoot)) ||
189 (workspaceRoot == "" && desktopWorkspaceID("project", project.Root) == workspaceID) {
190 topics, pinned = project.Topics, project.PinnedTopics
191 break
192 }
193 }
194 }
195 presentation.Pinned = containsDesktopString(pinned, presentation.TopicID)
196 for rank, topic := range pinnedTopicIDs(topics, pinned) {
197 if topic == presentation.TopicID {
198 presentation.SortOrder = rank
199 break
200 }
201 }
202 return presentation
203 }
204
205 // These references are local recovery evidence, never telemetry. Directories
206 // retain their complete subtree; import success does not authorize cleanup.
207 func retainedDesktopArtifacts(path string) ([]string, error) {
208 info, err := os.Lstat(path)
209 if err != nil {
210 return nil, err
211 }
212 if info.IsDir() {
213 return []string{path}, nil
214 }
215 retained := []string{}
216 for _, artifact := range sessionTrashArtifacts(path, filepath.Base(path)) {
217 if _, err := os.Lstat(artifact.src); err == nil {
218 retained = append(retained, artifact.src)
219 } else if !os.IsNotExist(err) {
220 return nil, err
221 }
222 }
223 return retained, nil
224 }
225
226 // Reserve the destination before publishing content, so restart reconciliation
227 // cannot mistake an interrupted import for an ordinary unregistered session.
228 func (a *App) prepareDesktopImport(ctx context.Context, source desktopMigrationSource, path, fingerprint, targetID, workspaceID string) (string, error) {
229 opID := source.operationID
230 if opID == "" {
231 opID = "import-" + desktopSourceKey(path, source.headID) + "-" + fingerprint
232 if source.deferArchive {
233 opID = "archive-" + opID
234 }
235 state, err := a.workspaceRegistry().Load(ctx)
236 if err != nil {
237 return "", err
238 }
239 if source.versionFingerprint != "" {
240 // Older builds used the ordinary import ID for versioned mappings.
241 // Resume that exact reservation when present, but do not collide
242 // with an ordinary import of the same source fingerprint.
243 previous, exists := state.PendingOperations[opID]
244 if !exists || previous.Mapping == nil || previous.Mapping.SourceKey != source.mappingKey(path) {
245 opID = "review-" + opID
246 }
247 }
248 lifecycle := workspacestate.Active
249 kind := "import"
250 if source.deferArchive {
251 kind, lifecycle = "archive-import", workspacestate.Archived
252 }
253 if previous, ok := state.SessionStates[targetID]; ok {
254 lifecycle = previous.Lifecycle
255 }
256 mapping := &workspacestate.SourceMapping{SourceKey: source.mappingKey(path), Path: path, HeadID: source.headID, Fingerprint: fingerprint, SessionID: targetID, WorkspaceID: workspaceID}
257 if err := a.workspaceRegistry().BeginOperation(ctx, workspacestate.Operation{ID: opID, Kind: kind, WorkspaceID: workspaceID, SessionIDs: []string{targetID}, Mapping: mapping, Lifecycle: lifecycle, ExpectedGeneration: state.Generation}); err != nil {
258 return "", err
259 }
260 }
261 format := "legacy"
262 if info, err := os.Stat(path); err != nil {
263 return "", err
264 } else if info.IsDir() {
265 format = "canonical"
266 }
267 mapping := &workspacestate.SourceMapping{SourceKey: source.mappingKey(path), Path: path, HeadID: source.headID, Format: format, Fingerprint: fingerprint, SessionID: targetID, WorkspaceID: workspaceID}
268 return opID, a.workspaceRegistry().ReserveOperationTargets(ctx, opID, []string{targetID}, mapping)
269 }
270
271 func (a *App) resolveDesktopImportTarget(ctx context.Context, query *session.Query, preferredID, key, mappingKey, contentDigest, path, fingerprint string, heads ...string) (string, bool, error) {
272 headID := ""
273 if len(heads) > 0 {
274 headID = heads[0]
275 }
276 state, err := a.workspaceRegistry().Load(ctx)
277 if err != nil {
278 return "", false, err
279 }
280 for _, op := range state.PendingOperations {
281 // Explicit source versions reserve their own mapping key. Recover the
282 // exact reservation even when the imported manifest has no provenance.
283 if op.Mapping == nil || !slices.Contains(state.SourceKeys(op.Mapping.SourceKey), mappingKey) || op.Mapping.Fingerprint != fingerprint || len(op.SessionIDs) != 1 {
284 continue
285 }
286 id := op.SessionIDs[0]
287 if state.SessionStates[id].Lifecycle == workspacestate.Deleted {
288 // A completed import receipt can outlive purge. It is not a live
289 // reservation and cannot authorize reusing that identity.
290 if op.Phase == "committed" {
291 continue
292 }
293 return "", false, workspacestate.ErrMutationConflict
294 }
295 digest, err := canonicalMigrationDigest(ctx, query, session.SessionRef{HostID: localDesktopHostID, SessionID: id})
296 if errors.Is(err, session.ErrSessionNotFound) {
297 return id, true, nil
298 }
299 if err != nil {
300 return "", false, err
301 }
302 if digest != contentDigest {
303 return "", false, workspacestate.ErrMutationConflict
304 }
305 return id, false, nil
306 }
307 remappedRetired := state.SessionStates[preferredID].Lifecycle == workspacestate.Deleted
308 if remappedRetired {
309 if err := a.proveRetiredImportOrigin(ctx, state, path, headID, preferredID); err != nil {
310 return "", false, err
311 }
312 digest := sha256.Sum256([]byte(mappingKey + "\x00" + preferredID + "\x00" + contentDigest))
313 preferredID = "migr-" + hex.EncodeToString(digest[:12])
314 }
315 id, needsImport, err := resolveMigrationTarget(ctx, query, preferredID, key, contentDigest, path, headID)
316 if err != nil {
317 return "", false, err
318 }
319 if state.SessionStates[id].Lifecycle == workspacestate.Deleted {
320 return "", false, workspacestate.ErrMutationConflict
321 }
322 for _, op := range state.PendingOperations {
323 if (remappedRetired || op.Kind == "archive-import") && op.Phase != "committed" && slices.Contains(op.SessionIDs, id) &&
324 (op.Mapping == nil || op.Mapping.SourceKey != mappingKey || op.Mapping.Fingerprint != fingerprint) {
325 return "", false, workspacestate.ErrMutationConflict
326 }
327 }
328 return id, needsImport, nil
329 }
330
331 func (a *App) legacyCanonicalRef(ctx context.Context, path string) (session.SessionRef, bool, error) {
332 snapshot, err := a.workspaceRegistry().VerifySnapshot(ctx)
333 if err != nil {
334 return session.SessionRef{}, false, err
335 }
336 mapping, adopted, err := snapshot.ResolveSource(desktopSourceKey(path, ""))
337 if err != nil {
338 return session.SessionRef{}, false, err
339 }
340 if !adopted {
341 // DAG migration records each head separately. A path-only legacy tab
342 // still refers to the selected head, not a new import of that path.
343 hasHeads, err := snapshot.HasHeadSource(path)
344 if err != nil {
345 return session.SessionRef{}, false, err
346 }
347 if hasHeads {
348 heads, err := agent.ListSessionHeads(path)
349 if err != nil {
350 return session.SessionRef{}, false, err
351 }
352 for _, head := range heads {
353 if head.Selected && !head.Retired {
354 mapping, adopted, err = snapshot.ResolveSource(desktopSourceKey(path, head.ID))
355 if err != nil {
356 return session.SessionRef{}, false, err
357 }
358 break
359 }
360 }
361 }
362 }
363 if adopted {
364 if snapshot.Session(mapping.SessionID).State.Lifecycle == workspacestate.Deleted {
365 return session.SessionRef{}, true, session.ErrSessionNotFound
366 }
367 // Adoption is durable. Opening the new conversation must not hash or
368 // depend on the retained source, which another CLI may still be using.
369 return session.SessionRef{HostID: localDesktopHostID, SessionID: mapping.SessionID}, true, nil
370 }
371 a.mu.RLock()
372 for _, tab := range a.runtimeTabsLocked() {
373 if tab == nil || tab.Ctrl == nil {
374 continue
375 }
376 if _, runtime, exclusive := exclusiveSessionBinding(tab.Ctrl); exclusive {
377 source := runtime.Session().Manifest().Source
378 if source != nil && sessionRuntimeKey(source.Path) == sessionRuntimeKey(path) {
379 ref := runtime.Ref()
380 a.mu.RUnlock()
381 return ref, true, nil
382 }
383 }
384 }
385 a.mu.RUnlock()
386 return session.SessionRef{}, false, nil
387 }
388
389 // The source manifest remains untouched. Metadata backups are captured before
390 // any registry upgrade and are content-addressed so subsequent starts preserve
391 // every distinct pre-upgrade snapshot.
392 func (a *App) backupDesktopUpgradeMetadata(ctx context.Context) error {
393 return backupDesktopUpgradeMetadataAt(ctx, a.workspaceRegistry().Path())
394 }
395
396 func backupDesktopUpgradeMetadataAt(ctx context.Context, registryPath string) error {
397 dir := filepath.Join(desktopConfigDir(), "desktop", "upgrade-backups")
398 paths := []string{
399 registryPath, filepath.Join(desktopConfigDir(), desktopProjectsFile),
400 filepath.Join(desktopConfigDir(), tabsFileName), desktopMigrationLedgerPath(),
401 }
402 roots := []string{""}
403 for _, project := range loadProjectsFile().Projects {
404 roots = append(roots, project.Root)
405 }
406 for _, root := range roots {
407 for _, path := range legacyTopicPaths(root) {
408 paths = append(paths, path)
409 }
410 databasePath := config.DesktopTopicStatePath(root)
411 if _, err := os.Lstat(databasePath); err == nil {
412 if err := os.MkdirAll(dir, 0700); err != nil {
413 return err
414 }
415 tmp, err := os.MkdirTemp(dir, ".topic-snapshot-")
416 if err != nil {
417 return err
418 }
419 snapshot := filepath.Join(tmp, "snapshot.sqlite")
420 if err := topicstate.BackupExisting(ctx, databasePath, snapshot); err != nil {
421 _ = os.RemoveAll(tmp)
422 return err
423 }
424 body, readErr := os.ReadFile(snapshot)
425 _ = os.RemoveAll(tmp)
426 if readErr != nil {
427 return readErr
428 }
429 digest := sha256.Sum256(body)
430 dest := filepath.Join(dir, "topics-"+hex.EncodeToString(digest[:])+".sqlite")
431 if saved, err := os.ReadFile(dest); err == nil {
432 if sha256.Sum256(saved) != digest {
433 return errors.New("topic backup integrity failed")
434 }
435 } else if !os.IsNotExist(err) {
436 return err
437 } else if err := fileutil.AtomicWriteFileStrict(dest, body, 0600); err != nil {
438 return err
439 }
440 } else if !os.IsNotExist(err) {
441 return err
442 }
443 }
444 for _, path := range paths {
445 if err := ctx.Err(); err != nil {
446 return err
447 }
448 body, err := os.ReadFile(path)
449 if os.IsNotExist(err) {
450 continue
451 }
452 if err != nil {
453 return err
454 }
455 digest := sha256.Sum256(body)
456 dest := filepath.Join(dir, filepath.Base(path)+"-"+hex.EncodeToString(digest[:])+".bak")
457 if saved, err := os.ReadFile(dest); err == nil {
458 if sha256.Sum256(saved) != digest {
459 return errors.New("session upgrade backup integrity failed")
460 }
461 continue
462 } else if !os.IsNotExist(err) {
463 return err
464 }
465 if err := os.MkdirAll(dir, 0700); err != nil {
466 return err
467 }
468 if err := fileutil.AtomicWriteFileStrict(dest, body, 0600); err != nil {
469 return err
470 }
471 }
472 return nil
473 }
474
475 func (a *App) sourceRecovery(ctx context.Context, path, format, reason, scope, root string, heads ...string) error {
476 headID := ""
477 if len(heads) > 0 {
478 headID = heads[0]
479 }
480 key := desktopSourceKey(path, headID)
481 fingerprint, _ := desktopSourceFingerprint(path)
482 return a.workspaceRegistry().RecordRecovery(ctx, workspacestate.RecoveryEntry{
483 ID: desktopRecoveryID(key, fingerprint), SourceKey: key, Path: path, HeadID: headID, Format: format, Reason: reason,
484 Status: "pending", Scope: scope, WorkspaceRoot: root, Fingerprint: fingerprint,
485 })
486 }
487
488 // Preserve alternate DAG heads as independently addressable recovery choices.
489 // Importing one head never selects, retires or rewrites a head in the original.
490 func (a *App) discoverLegacyHeads(ctx context.Context, path, format, scope, root string) error {
491 heads, err := agent.ListSessionHeads(path)
492 if err != nil {
493 return errors.Join(err, a.sourceRecovery(ctx, path, format, "head_scan_failed", scope, root))
494 }
495 var joined error
496 for _, head := range heads {
497 if err := ctx.Err(); err != nil {
498 return err
499 }
500 if head.Selected || head.Retired {
501 continue
502 }
503 joined = errors.Join(joined, a.sourceRecovery(ctx, path, format, "alternate_head", scope, root, head.ID))
504 }
505 return joined
506 }
507
508 func desktopRecoveryID(key, fingerprint string) string {
509 return "legacy-" + key + "-" + fingerprint
510 }
511
512 // Read legacy ledger evidence without altering it or losing unknown fields.
513
513 lines GO