返回 DeepSeek-Reasonix
session_v5_migration_checkpoint.go
根目录 / desktop / session_v5_migration_checkpoint.go
1 package main
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/hex"
7 "encoding/json"
8 "errors"
9 "fmt"
10 "os"
11 "path/filepath"
12 "sort"
13
14 "reasonix/internal/store"
15 )
16
17 // Migration may add members to an existing workspace, but it does not own the
18 // user's title or visibility. Avoid rewriting presentation for each new source.
19 func (a *App) ensureDesktopMigrationWorkspace(ctx context.Context, source desktopMigrationSource) (string, error) {
20 state, err := a.workspaceRegistry().Load(ctx)
21 if err != nil {
22 return "", err
23 }
24 id := desktopWorkspaceOwnerID(state, source.scope, source.workspaceRoot)
25 if _, exists := state.Workspaces[id]; exists {
26 return id, nil
27 }
28 return a.ensureDesktopWorkspace(ctx, source.scope, source.workspaceRoot)
29 }
30
31 // A completed ledger entry records adoption of a source, not equality with
32 // today's target. The target can be continued, archived or deliberately deleted.
33 // None of those actions authorize importing the old conversation again.
34 type desktopMigrationCheckpoint struct {
35 key string
36 files []string
37 revision string
38 record desktopMigrationRecord
39 refresh bool
40 inputDigest string
41 }
42
43 type desktopMigrationReceipt struct {
44 extra map[string]json.RawMessage
45 TargetSessionID string `json:"targetSessionId"`
46 ContentDigest string `json:"contentDigest,omitempty"`
47 SourceRevision string `json:"sourceRevision,omitempty"`
48 }
49
50 func newDesktopMigrationCheckpoint(source desktopMigrationSource, key string, files []string) (desktopMigrationCheckpoint, error) {
51 records := source.records
52 if records == nil {
53 ledger, err := readDesktopMigrationLedger()
54 if err != nil {
55 return desktopMigrationCheckpoint{}, err
56 }
57 records = ledger.Records
58 }
59 revision, err := desktopMigrationSourceRevision(files)
60 if err != nil {
61 return desktopMigrationCheckpoint{}, errors.Join(err, updateDesktopMigrationLedger(key, records[key].TargetSessionID, "failed", "source_stat"))
62 }
63 record := records[key]
64 refresh := record.Status != "completed" && record.PreviousCompletion != nil
65 if refresh {
66 previous := record.PreviousCompletion
67 record.Status, record.TargetSessionID = "completed", previous.TargetSessionID
68 record.ContentDigest, record.SourceRevision = previous.ContentDigest, previous.SourceRevision
69 }
70 checkpoint := desktopMigrationCheckpoint{key: key, files: files, revision: revision, record: record, refresh: refresh}
71 if !checkpoint.unchanged() {
72 checkpoint.inputDigest, err = desktopMigrationInputDigest(files)
73 if err != nil {
74 return desktopMigrationCheckpoint{}, err
75 }
76 }
77 return checkpoint, nil
78 }
79
80 func (c desktopMigrationCheckpoint) completed() bool {
81 return c.record.Status == "completed" && c.record.TargetSessionID != ""
82 }
83
84 func (c desktopMigrationCheckpoint) unchanged() bool {
85 return c.completed() && c.record.SourceRevision == c.revision
86 }
87
88 func (c desktopMigrationCheckpoint) skip() error {
89 if c.refresh {
90 return c.complete(c.record.TargetSessionID, c.record.ContentDigest)
91 }
92 return nil
93 }
94
95 // Old ledgers have no source revision. Compare their recorded source digest
96 // once, without comparing to a target that may already contain newer work.
97 func (c desktopMigrationCheckpoint) matchesCompletedContent(digest string) bool {
98 return c.completed() && digest != "" && c.record.ContentDigest == digest
99 }
100
101 func (c desktopMigrationCheckpoint) complete(targetID, digest string) error {
102 revision, err := c.verifiedRevision()
103 if err != nil {
104 return errors.Join(err, updateDesktopMigrationLedger(c.key, targetID, "failed", "source_changed", digest))
105 }
106 return updateDesktopMigrationLedger(c.key, targetID, "completed", "", digest, revision)
107 }
108
109 func (c desktopMigrationCheckpoint) verify() error {
110 _, err := c.verifiedRevision()
111 return err
112 }
113
114 func (c desktopMigrationCheckpoint) verifiedRevision() (string, error) {
115 revision, err := desktopMigrationSourceRevision(c.files)
116 if err == nil && revision != c.revision && c.inputDigest != "" {
117 // A catalog repair can rewrite identical JSONL bytes and refresh only
118 // derived listing fields. Verify the frozen input, not inode/mtime alone.
119 current, digestErr := desktopMigrationInputDigest(c.files)
120 if digestErr == nil && current == c.inputDigest {
121 c.revision = revision
122 }
123 err = digestErr
124 }
125 if err != nil || revision != c.revision {
126 return "", errors.Join(errors.New("desktop migration source changed during import"), err)
127 }
128 return revision, nil
129 }
130
131 func canonicalMigrationSourceFiles(root, sessionID string) []string {
132 dir := filepath.Join(root, sessionID)
133 // Content blobs are immutable and referenced by the event log. Disposable
134 // indexes, checkpoints and lock files do not change the migration input.
135 return []string{filepath.Join(dir, "manifest.json"), filepath.Join(dir, "events.frames"),
136 filepath.Join(dir, "events.jsonl"), filepath.Join(dir, "header.json")}
137 }
138
139 func desktopCanonicalMigrationKey(root, sessionID string) string {
140 digest := sha256.Sum256([]byte(canonicalRuntimeRoot(root) + "\x00" + sessionID))
141 return hex.EncodeToString(digest[:])
142 }
143
144 func legacyMigrationSourceFiles(path string) []string {
145 // Catalog reads may rebuild disposable indexes while an import is frozen.
146 // Their publication is not a change to the legacy conversation. Keep all
147 // durable sidecars (including ancestry/ownership metadata) in the stamp.
148 files := []string{path}
149 for _, sidecar := range store.SessionSidecarFiles(path) {
150 if sidecar == store.SessionEventIndex(path) || sidecar == store.SessionDisplayIndex(path) || sidecar == store.SessionTranscriptProjection(path) {
151 continue
152 }
153 files = append(files, sidecar)
154 }
155 return files
156 }
157
158 func desktopLegacyMigrationKey(path string) string {
159 digest := sha256.Sum256([]byte(canonicalRuntimeRoot(path)))
160 return hex.EncodeToString(digest[:])
161 }
162
163 // This is a cheap filesystem revision, not a content integrity checksum.
164 // Normal source writes change size or mtime; no history bodies are read here.
165 // Include absent files so creation/removal of a sidecar invalidates the stamp.
166 func desktopMigrationSourceRevision(files []string) (string, error) {
167 paths := append([]string(nil), files...)
168 sort.Strings(paths)
169 hash := sha256.New()
170 for _, path := range paths {
171 info, err := os.Stat(path)
172 if os.IsNotExist(err) {
173 fmt.Fprintf(hash, "%q:missing\n", path)
174 continue
175 }
176 if err != nil {
177 return "", err
178 }
179 fmt.Fprintf(hash, "%q:%d:%d:%d\n", path, info.Mode(), info.Size(), info.ModTime().UnixNano())
180 }
181 return "stat-v1-" + hex.EncodeToString(hash.Sum(nil)), nil
182 }
183
184 func readDesktopMigrationLedger() (desktopMigrationLedger, error) {
185 desktopMigrationMu.Lock()
186 defer desktopMigrationMu.Unlock()
187 ledger, _, err := readDesktopMigrationLedgerFile()
188 return ledger, err
189 }
190
191 // Callers serialize reads and atomic replacement with desktopMigrationMu.
192 func readDesktopMigrationLedgerFile() (desktopMigrationLedger, []byte, error) {
193 ledger := desktopMigrationLedger{Version: 1, Records: map[string]desktopMigrationRecord{}}
194 body, err := os.ReadFile(desktopMigrationLedgerPath())
195 if os.IsNotExist(err) {
196 return ledger, nil, nil
197 }
198 if err != nil {
199 return ledger, nil, err
200 }
201 if err := json.Unmarshal(body, &ledger); err != nil {
202 return ledger, nil, err
203 }
204 if ledger.Version > 1 {
205 return ledger, nil, fmt.Errorf("desktop migration ledger version %d is unsupported", ledger.Version)
206 }
207 if ledger.Records == nil {
208 ledger.Records = map[string]desktopMigrationRecord{}
209 }
210 return ledger, body, nil
211 }
212
213 // Preserve unknown root/record fields when adding an optional revision to a v1
214 // ledger. Older writers can drop the revision; the digest fallback remains safe.
215 func marshalDesktopMigrationRecord(original []byte, ledger desktopMigrationLedger, key string) ([]byte, error) {
216 root := map[string]json.RawMessage{}
217 if len(original) > 0 {
218 if err := json.Unmarshal(original, &root); err != nil {
219 return nil, err
220 }
221 }
222 if root == nil {
223 root = map[string]json.RawMessage{}
224 }
225 records := map[string]map[string]json.RawMessage{}
226 if body := root["records"]; len(body) > 0 {
227 if err := json.Unmarshal(body, &records); err != nil {
228 return nil, err
229 }
230 }
231 if records == nil {
232 records = map[string]map[string]json.RawMessage{}
233 }
234 fields := records[key]
235 if fields == nil {
236 fields = map[string]json.RawMessage{}
237 }
238 for _, name := range []string{"sourceKey", "targetSessionId", "contentDigest", "sourceRevision", "previousCompletion", "status", "errorCode", "attempts", "legacyHeads", "legacySelectedHead", "legacyPrimaryHead", "legacyHeadsRevision", "legacyAdoption", "legacyConversions"} {
239 delete(fields, name)
240 }
241 body, err := json.Marshal(ledger.Records[key])
242 if err != nil {
243 return nil, err
244 }
245 if err := json.Unmarshal(body, &fields); err != nil {
246 return nil, err
247 }
248 records[key] = fields
249 root["records"], err = json.Marshal(records)
250 if err != nil {
251 return nil, err
252 }
253 root["version"], err = json.Marshal(ledger.Version)
254 if err != nil {
255 return nil, err
256 }
257 return json.MarshalIndent(root, "", " ")
258 }
259
259 lines GO