返回 DeepSeek-Reasonix
session_listing_repair.go
根目录 / internal / agent / session_listing_repair.go
1 package agent
2
3 import (
4 "context"
5 "crypto/sha256"
6 "errors"
7 "fmt"
8 "os"
9 "strings"
10
11 "reasonix/internal/store"
12 )
13
14 type SessionListingRepairStatus string
15
16 const (
17 SessionListingRepairApplied SessionListingRepairStatus = "applied"
18 SessionListingRepairAlreadyCurrent SessionListingRepairStatus = "already_current"
19 SessionListingRepairSourceChanged SessionListingRepairStatus = "source_changed"
20 SessionListingRepairDamaged SessionListingRepairStatus = "damaged"
21 SessionListingRepairUnsupported SessionListingRepairStatus = "unsupported"
22 )
23
24 var ErrSessionListingRepairBusy = errors.New("session listing repair source is busy")
25
26 // SessionListingRepairResult describes one convergent repair attempt. Preview
27 // and Turns are safe to publish only for applied/already_current results.
28 type SessionListingRepairResult struct {
29 Status SessionListingRepairStatus
30 Preview string
31 Turns int
32 LedgerRepaired bool
33 ContentFingerprint string
34 MetaFingerprint string
35 }
36
37 // SessionListingGeneration identifies the transcript/event-log generation and
38 // its metadata sidecar while all writer locks for the session remain held.
39 type SessionListingGeneration struct {
40 ContentFingerprint string
41 MetaFingerprint string
42 }
43
44 // TryLockSessionListingGeneration fences catalog publication against a
45 // foreground writer. The caller must invoke the returned unlock function.
46 func TryLockSessionListingGeneration(path string) (SessionListingGeneration, func(), error) {
47 path = strings.TrimSpace(path)
48 if path == "" {
49 return SessionListingGeneration{}, nil, fmt.Errorf("empty session path")
50 }
51 unlock, err := lockSessionListingRepair(path)
52 if err != nil {
53 return SessionListingGeneration{}, nil, err
54 }
55 return SessionListingGeneration{
56 ContentFingerprint: sessionListingCatalogContentFingerprint(path),
57 MetaFingerprint: sessionListingCatalogFileFingerprint(BranchMetaPath(path)),
58 }, unlock, nil
59 }
60
61 // RepairSessionListingProjection repairs one session generation while holding
62 // only that session's save/file/meta locks. Foreground saves win immediately;
63 // callers persist a retry instead of waiting behind active work.
64 func RepairSessionListingProjection(ctx context.Context, path string) (result SessionListingRepairResult, resultErr error) {
65 path = strings.TrimSpace(path)
66 if path == "" {
67 return SessionListingRepairResult{}, fmt.Errorf("empty session path")
68 }
69 if err := ctx.Err(); err != nil {
70 return SessionListingRepairResult{}, err
71 }
72 unlock, err := lockSessionListingRepair(path)
73 if err != nil {
74 return SessionListingRepairResult{}, err
75 }
76 defer unlock()
77 defer func() {
78 result.ContentFingerprint = sessionListingCatalogContentFingerprint(path)
79 result.MetaFingerprint = sessionListingCatalogFileFingerprint(BranchMetaPath(path))
80 }()
81 meta, metaOK, err := LoadBranchMeta(path)
82 if err != nil {
83 if isDamagedSessionRepairError(err) {
84 return SessionListingRepairResult{Status: SessionListingRepairDamaged}, nil
85 }
86 return SessionListingRepairResult{}, err
87 }
88 if !metaOK {
89 meta = BranchMeta{ID: BranchID(path)}
90 }
91 if !sessionHeadIndexStale(path) {
92 if result, handled, err := repairSessionListingFromIndex(path, meta); handled || err != nil {
93 return result, err
94 }
95 }
96 return repairSessionListingFromReplay(ctx, path, meta)
97 }
98
99 func lockSessionListingRepair(path string) (func(), error) {
100 unlockSave, ok := tryLockSessionSavePath(path)
101 if !ok {
102 return nil, ErrSessionListingRepairBusy
103 }
104 unlockFile, err := tryLockSessionFile(path)
105 if err != nil {
106 unlockSave()
107 if errors.Is(err, ErrSessionFileLockHeld) {
108 return nil, ErrSessionListingRepairBusy
109 }
110 return nil, err
111 }
112 unlockMeta, ok, err := tryLockSessionMetaPath(path)
113 if err != nil {
114 unlockFile()
115 unlockSave()
116 return nil, err
117 }
118 if !ok {
119 unlockFile()
120 unlockSave()
121 return nil, ErrSessionListingRepairBusy
122 }
123 return func() { unlockMeta(); unlockFile(); unlockSave() }, nil
124 }
125
126 func repairSessionListingFromIndex(path string, meta BranchMeta) (SessionListingRepairResult, bool, error) {
127 preview, turns, ok, unsupported, err := indexedSessionListing(path, meta)
128 if err != nil {
129 return SessionListingRepairResult{}, true, err
130 }
131 if unsupported {
132 return SessionListingRepairResult{Status: SessionListingRepairUnsupported}, true, nil
133 }
134 if !ok {
135 return SessionListingRepairResult{}, false, nil
136 }
137 if sessionListingProjectionFresh(meta.SchemaVersion, meta.Turns, meta.Revision,
138 meta.ListingRevision, meta.ContentDigest, meta.ListingContentDigest) &&
139 meta.Preview == preview && meta.Turns == turns {
140 return SessionListingRepairResult{Status: SessionListingRepairAlreadyCurrent, Preview: preview, Turns: turns}, true, nil
141 }
142 meta.Preview, meta.Turns, meta.SchemaVersion = preview, turns, BranchMetaCountsVersion
143 stampSessionListingProjection(&meta)
144 if err := saveBranchMeta(path, meta, false); err != nil {
145 return SessionListingRepairResult{}, true, err
146 }
147 return SessionListingRepairResult{Status: SessionListingRepairApplied, Preview: preview, Turns: turns}, true, nil
148 }
149
150 func repairSessionListingFromReplay(ctx context.Context, path string, meta BranchMeta) (SessionListingRepairResult, error) {
151 before, err := sessionRepairContentFingerprint(path)
152 if err != nil {
153 return SessionListingRepairResult{}, err
154 }
155 msgs, state, repairable, err := loadSessionDisplayMessagesContextUnlocked(ctx, path)
156 if err != nil {
157 if ctxErr := ctx.Err(); ctxErr != nil {
158 return SessionListingRepairResult{}, ctxErr
159 }
160 if isUnsupportedSessionRepairError(err) {
161 return SessionListingRepairResult{Status: SessionListingRepairUnsupported}, nil
162 }
163 if isDamagedSessionRepairError(err) {
164 return SessionListingRepairResult{Status: SessionListingRepairDamaged}, nil
165 }
166 return SessionListingRepairResult{}, err
167 }
168 if !repairable {
169 return SessionListingRepairResult{Status: SessionListingRepairDamaged}, nil
170 }
171 if err := ctx.Err(); err != nil {
172 return SessionListingRepairResult{}, err
173 }
174 after, err := sessionRepairContentFingerprint(path)
175 if err != nil {
176 return SessionListingRepairResult{}, err
177 }
178 if before != after {
179 return SessionListingRepairResult{Status: SessionListingRepairSourceChanged}, nil
180 }
181
182 preview, turns := SessionPreviewFromMessages(msgs)
183 ledgerCurrent := meta.Revision > 0 && strings.TrimSpace(meta.ContentDigest) == state.DigestHex
184 if meta.Recovered && strings.TrimSpace(meta.RecoveryDigest) != state.DigestHex {
185 ledgerCurrent = false
186 }
187 ledgerRepaired := !ledgerCurrent
188 if ledgerRepaired {
189 meta.Revision = max(int64(1), meta.Revision+1)
190 meta.ContentDigest = state.DigestHex
191 if meta.Recovered || strings.TrimSpace(meta.RecoveryDigest) != "" {
192 meta.RecoveryDigest = state.DigestHex
193 }
194 meta.WriterID = SessionWriterID()
195 }
196 meta.Preview = preview
197 meta.Turns = turns
198 meta.SchemaVersion = BranchMetaCountsVersion
199 stampSessionListingProjection(&meta)
200
201 // The display index ranges describe the compatibility JSONL exactly, so
202 // refresh it from this single authoritative replay before publishing meta.
203 if err := writeSessionMessagesContext(ctx, path, msgs); err != nil {
204 return SessionListingRepairResult{}, fmt.Errorf("write session display read model: %w", err)
205 }
206 if err := refreshSessionEventIndexContext(ctx, path, msgs, state.Digest, meta.Revision); err != nil {
207 return SessionListingRepairResult{}, err
208 }
209 idx, err := BuildSessionDisplayIndexContext(ctx, msgs, meta.Revision, true, state.Digest)
210 if err != nil {
211 return SessionListingRepairResult{}, fmt.Errorf("encode session display index: %w", err)
212 }
213 if err := WriteSessionDisplayIndexContext(ctx, store.SessionDisplayIndex(path), idx); err != nil {
214 return SessionListingRepairResult{}, err
215 }
216 if err := saveBranchMetaContext(ctx, path, meta, false); err != nil {
217 return SessionListingRepairResult{}, err
218 }
219 return SessionListingRepairResult{
220 Status: SessionListingRepairApplied, Preview: preview, Turns: turns, LedgerRepaired: ledgerRepaired,
221 }, nil
222 }
223
224 func indexedSessionListing(path string, meta BranchMeta) (preview string, turns int, ok, unsupported bool, err error) {
225 digest := strings.TrimSpace(meta.ContentDigest)
226 if meta.Revision <= 0 || digest == "" || meta.Recovered && strings.TrimSpace(meta.RecoveryDigest) != digest {
227 return "", 0, false, false, nil
228 }
229 idx, err := LoadSessionDisplayIndex(store.SessionDisplayIndex(path))
230 if err != nil || idx == nil || !idx.ListingPreviewKnown || !idx.RevisionKnown ||
231 idx.Revision != meta.Revision || idx.ContentDigest != digest {
232 return "", 0, false, false, nil
233 }
234 info, err := os.Stat(path)
235 if err != nil {
236 return "", 0, false, false, err
237 }
238 if info.IsDir() || info.Size() != idx.TranscriptSize {
239 return "", 0, false, false, nil
240 }
241 indexInfo, err := os.Stat(store.SessionDisplayIndex(path))
242 if err != nil || !indexInfo.ModTime().After(SessionContentModTime(path)) {
243 return "", 0, false, false, nil
244 }
245 probe, err := probeSessionEventLog(path)
246 if err != nil {
247 return "", 0, false, false, err
248 }
249 if probe.futureSchema {
250 return "", 0, false, true, nil
251 }
252 if probe.native && probe.size > 0 {
253 eventIdx, err := readSessionEventIndex(path)
254 if err != nil || eventIdx == nil || eventIdx.LogSize != probe.size ||
255 eventIdx.Revision != meta.Revision || eventIdx.ContentDigest != digest ||
256 eventIdx.MessageCount != idx.MessageCount {
257 return "", 0, false, false, nil
258 }
259 eventInfo, eventErr := os.Stat(store.SessionEventIndex(path))
260 logInfo, logErr := os.Stat(store.SessionEventLog(path))
261 if eventErr != nil || logErr != nil || !eventInfo.ModTime().After(logInfo.ModTime()) {
262 return "", 0, false, false, nil
263 }
264 }
265 return idx.ListingPreview, idx.AuthoredTurns, true, false, nil
266 }
267
268 func sessionRepairContentFingerprint(path string) (string, error) {
269 h := sha256.New()
270 for _, artifact := range []string{path, store.SessionEventLog(path)} {
271 info, err := os.Stat(artifact)
272 if err != nil {
273 if os.IsNotExist(err) {
274 _, _ = fmt.Fprintf(h, "missing\x00")
275 continue
276 }
277 return "", err
278 }
279 _, _ = fmt.Fprintf(h, "%d\x00%d\x00%d\x00", info.Size(), info.ModTime().UnixNano(), info.Mode())
280 }
281 return fmt.Sprintf("%x", h.Sum(nil)), nil
282 }
283
284 func sessionListingCatalogFileFingerprint(path string) string {
285 info, err := os.Stat(path)
286 if err != nil {
287 return ""
288 }
289 return fmt.Sprintf("%d:%d", info.Size(), info.ModTime().UnixNano())
290 }
291
292 func sessionListingCatalogContentFingerprint(path string) string {
293 return sessionListingCatalogFileFingerprint(path) + "|" + sessionListingCatalogFileFingerprint(store.SessionEventLog(path))
294 }
295
296 func isUnsupportedSessionRepairError(err error) bool {
297 return errors.Is(err, ErrSessionReplayLimitExceeded) ||
298 strings.Contains(err.Error(), "unsupported schema") ||
299 strings.Contains(err.Error(), "supports up to") ||
300 strings.Contains(err.Error(), "unsupported event type")
301 }
302
303 func isDamagedSessionRepairError(err error) bool {
304 text := err.Error()
305 return strings.Contains(text, "decode session transcript") ||
306 strings.Contains(text, "decode session event log") ||
307 strings.Contains(text, "decode ")
308 }
309
309 lines GO