返回 DeepSeek-Reasonix
transaction.go
根目录 / internal / checkpoint / transaction.go
1 package checkpoint
2
3 import (
4 "encoding/json"
5 "errors"
6 "fmt"
7 "log/slog"
8 "os"
9 "path/filepath"
10 "slices"
11 "sort"
12 "strings"
13 "time"
14 )
15
16 // InjectFail is a test seam. When set, CommitRewind fails at the named phase
17 // after optionally publishing the first N files. Empty = disabled.
18 //
19 // Known phases: "publish_file", "delete_file", "conversation", "truncate",
20 // "after_conversation_before_finalize", "finalize".
21 type InjectFail struct {
22 Phase string
23 AfterFiles int // fail after successfully handling this many file targets
24 }
25
26 // ConversationApplier applies conversation truncation during commit and restores
27 // forward conversation on compensate. Implemented by the control layer.
28 type ConversationApplier interface {
29 // ApplyConversationTruncate replaces the live message log with msgs[:boundary].
30 // forward is the full pre-truncate snapshot for later restore.
31 ApplyConversationTruncate(boundary int, forward []byte) error
32 // RestoreConversation reinstalls the forward snapshot.
33 RestoreConversation(forward []byte) error
34 // TruncateCheckpoints drops checkpoints at or after turn.
35 TruncateCheckpoints(fromTurn int) error
36 // RestoreCheckpoints reinstalls backed-up future checkpoints.
37 RestoreCheckpoints(backup []byte) error
38 }
39
40 // PrepareRewind builds a plan and optionally a prepared transaction without
41 // mutating workspace or conversation. Conflict detection uses last-owned after
42 // fingerprints when available.
43 func (s *Store) PrepareRewind(turn int, scope RewindScope, sessionRev int64, boundary int, hasBound bool) (RewindPlan, error) {
44 if s == nil {
45 return RewindPlan{}, fmt.Errorf("checkpoints unavailable")
46 }
47 plan := RewindPlan{
48 PlanID: newID("plan"),
49 Turn: turn,
50 Scope: scope,
51 SessionRevision: sessionRev,
52 BoundaryIndex: boundary,
53 HasBoundary: hasBound,
54 CreatedAt: time.Now(),
55 WorkspaceToken: fmt.Sprintf("%d", s.barrier.Generation()),
56 }
57
58 s.mu.Lock()
59 writers := append([]ActiveWriter(nil), s.activeWriters...)
60 plan.ActiveWriters = writers
61 cov, gaps, legacy, expired := s.coverageFromTurnLocked(turn)
62 plan.Coverage = cov
63 plan.CoverageGaps = gaps
64 plan.Legacy = legacy
65 plan.ExpiredFilePayload = expired
66 files := s.filesFromTurnLocked(turn)
67 plan.Files = files
68 plan.FileCount = len(files)
69 s.mu.Unlock()
70
71 wantFiles := scope == RewindCode || scope == RewindBoth
72 wantConv := scope == RewindConversation || scope == RewindBoth
73
74 if len(writers) > 0 {
75 plan.CanFiles = false
76 plan.CanConversation = false
77 plan.DisabledReason = "active background writer"
78 for _, w := range writers {
79 plan.Conflicts = append(plan.Conflicts, RewindConflict{
80 Path: "",
81 Reason: ConflictBusyWriter,
82 })
83 _ = w
84 }
85 return plan, nil
86 }
87
88 if wantConv {
89 if !hasBound {
90 plan.CanConversation = false
91 if scope == RewindConversation || scope == RewindBoth {
92 plan.DisabledReason = "conversation boundary unavailable"
93 }
94 } else {
95 plan.CanConversation = true
96 }
97 }
98
99 if wantFiles {
100 if len(files) == 0 && scope == RewindBoth {
101 // The file half of a combined rewind is an atomic no-op when this
102 // conversation range never touched a tracked file.
103 plan.CanFiles = true
104 } else if cov == CoverageNone {
105 plan.CanFiles = false
106 plan.DisabledReason = "no file captures"
107 } else if expired {
108 plan.CanFiles = false
109 plan.DisabledReason = "file recovery payload expired"
110 plan.Conflicts = append(plan.Conflicts, RewindConflict{Reason: ConflictExpired})
111 } else if legacy {
112 // Legacy: files can be restored only with explicit warning; batch
113 // overwrite without prompt is forbidden. Prepare still reports files
114 // but CanFiles stays false for the unprompted path.
115 plan.CanFiles = false
116 plan.DisabledReason = "legacy checkpoint cannot verify later manual edits"
117 plan.Conflicts = append(plan.Conflicts, RewindConflict{Reason: ConflictCoverageLegacy})
118 } else {
119 conflicts := s.precheckFiles(turn)
120 plan.Conflicts = append(plan.Conflicts, conflicts...)
121 plan.CanFiles = len(conflicts) == 0 && len(files) > 0
122 if len(conflicts) > 0 {
123 plan.DisabledReason = "file conflicts detected"
124 }
125 }
126 }
127
128 // both requires both sides to pass precheck.
129 if scope == RewindBoth {
130 if !plan.CanFiles || !plan.CanConversation {
131 if plan.DisabledReason == "" {
132 plan.DisabledReason = "both scope requires file and conversation precheck"
133 }
134 }
135 }
136
137 // Persist plan token so Commit can verify freshness.
138 s.mu.Lock()
139 if s.plans == nil {
140 s.plans = map[string]preparedPlan{}
141 }
142 s.plans[plan.PlanID] = preparedPlan{plan: plan, created: time.Now()}
143 // Drop stale plans older than 10 minutes.
144 for id, p := range s.plans {
145 if time.Since(p.created) > 10*time.Minute {
146 delete(s.plans, id)
147 }
148 }
149 s.mu.Unlock()
150 return plan, nil
151 }
152
153 type preparedPlan struct {
154 plan RewindPlan
155 created time.Time
156 previewFingerprint *Fingerprint
157 }
158
159 // ValidatePlanSessionRevision binds a preview to the controller's exact
160 // conversation revision. The controller holds its rotation gate while calling
161 // this and committing, so no turn can slip between validation and mutation.
162 func (s *Store) ValidatePlanSessionRevision(planID string, current int64) error {
163 if s == nil {
164 return fmt.Errorf("checkpoints unavailable")
165 }
166 s.mu.Lock()
167 defer s.mu.Unlock()
168 prepared, ok := s.plans[planID]
169 if !ok {
170 return fmt.Errorf("unknown or expired plan %q", planID)
171 }
172 if prepared.plan.SessionRevision != current {
173 return fmt.Errorf("conversation changed since preview")
174 }
175 return nil
176 }
177
178 // CommitRewind executes a previously prepared plan under exclusive barrier.
179 // conversation/checkpoints are applied via applier when non-nil.
180 func (s *Store) CommitRewind(planID string, applier ConversationApplier, inject *InjectFail) (RewindResult, error) {
181 if s == nil {
182 return RewindResult{}, fmt.Errorf("checkpoints unavailable")
183 }
184 s.mu.Lock()
185 pp, ok := s.plans[planID]
186 if ok {
187 delete(s.plans, planID)
188 }
189 s.mu.Unlock()
190 if !ok {
191 return RewindResult{OK: false, Error: "unknown or expired plan"}, fmt.Errorf("unknown or expired plan %q", planID)
192 }
193 plan := pp.plan
194
195 // Re-validate gate conditions before any mutation.
196 if plan.Scope == RewindBoth && (!plan.CanFiles || !plan.CanConversation) {
197 return RewindResult{OK: false, Error: plan.DisabledReason, Conflicts: plan.Conflicts, Coverage: plan.Coverage}, fmt.Errorf("%s", plan.DisabledReason)
198 }
199 if (plan.Scope == RewindCode || plan.Scope == RewindBoth) && !plan.CanFiles && plan.Scope != RewindConversation {
200 if plan.Scope == RewindCode || plan.Scope == RewindBoth {
201 return RewindResult{OK: false, Error: plan.DisabledReason, Conflicts: plan.Conflicts, Coverage: plan.Coverage}, fmt.Errorf("%s", plan.DisabledReason)
202 }
203 }
204 if (plan.Scope == RewindConversation || plan.Scope == RewindBoth) && !plan.CanConversation {
205 return RewindResult{OK: false, Error: plan.DisabledReason}, fmt.Errorf("%s", plan.DisabledReason)
206 }
207
208 // Workspace exclusive barrier.
209 if !s.barrier.TryEnterExclusive() {
210 err := fmt.Errorf("workspace mutation in progress")
211 return RewindResult{OK: false, Error: err.Error(), Conflicts: []RewindConflict{{Reason: ConflictBusyWriter}}}, err
212 }
213 defer s.barrier.ExitExclusive()
214 if conflicts := s.activeWriterConflicts(); len(conflicts) > 0 {
215 err := fmt.Errorf("active background writer")
216 return RewindResult{OK: false, Error: err.Error(), Conflicts: conflicts, Coverage: plan.Coverage}, err
217 }
218
219 if plan.Scope == RewindCode || plan.Scope == RewindBoth {
220 if plan.WorkspaceToken != fmt.Sprintf("%d", s.barrier.Generation()) {
221 conflict := RewindConflict{Reason: ConflictStalePlan}
222 return RewindResult{OK: false, Error: "workspace changed since preview", Conflicts: []RewindConflict{conflict}, Coverage: plan.Coverage}, fmt.Errorf("workspace changed since preview")
223 }
224 conflicts := s.precheckFiles(plan.Turn)
225 if len(conflicts) > 0 {
226 return RewindResult{OK: false, Error: "file conflicts detected", Conflicts: conflicts, Coverage: plan.Coverage}, fmt.Errorf("file conflicts detected")
227 }
228 }
229
230 tx, err := s.prepareTransaction(plan, applier)
231 if err != nil {
232 return RewindResult{OK: false, Error: err.Error()}, err
233 }
234
235 result, err := s.commitTransaction(tx, applier, inject)
236 return result, err
237 }
238
239 // UndoRewind reverses a committed transaction when still available.
240 func (s *Store) UndoRewind(transactionID string, applier ConversationApplier) (RewindResult, error) {
241 if s == nil {
242 return RewindResult{}, fmt.Errorf("checkpoints unavailable")
243 }
244 s.mu.Lock()
245 last := s.lastUndo
246 s.mu.Unlock()
247 if last == nil || last.ID != transactionID || last.State != TxCommitted {
248 return RewindResult{OK: false, Error: "undo not available"}, fmt.Errorf("undo not available for %q", transactionID)
249 }
250
251 if !s.barrier.TryEnterExclusive() {
252 err := fmt.Errorf("workspace mutation in progress")
253 return RewindResult{OK: false, Error: err.Error(), Conflicts: []RewindConflict{{Reason: ConflictBusyWriter}}}, err
254 }
255 defer s.barrier.ExitExclusive()
256 if conflicts := s.activeWriterConflicts(); len(conflicts) > 0 {
257 err := fmt.Errorf("active background writer")
258 return RewindResult{OK: false, Error: err.Error(), Conflicts: conflicts}, err
259 }
260
261 // Precheck that current disk still matches what we published (targets' restore state).
262 for _, t := range last.Targets {
263 fp, err := FingerprintPath(s.root, t.AbsPath)
264 if err != nil && !os.IsNotExist(err) {
265 return RewindResult{OK: false, Error: err.Error()}, fmt.Errorf("fingerprint %s before undo: %w", t.Path, err)
266 }
267 // After commit, disk should match restore image. If it doesn't, refuse.
268 if t.Action == "delete" {
269 if fp.Existed {
270 return RewindResult{OK: false, Error: "file changed since rewind", Conflicts: []RewindConflict{{
271 Path: t.Path, Reason: ConflictManualEdit, CurrentSHA: fp.SHA256,
272 }}}, fmt.Errorf("file changed since rewind: %s", t.Path)
273 }
274 } else {
275 if !fingerprintMatches(fp, t.RestoreExisted, t.RestoreSHA, t.RestoreMode) {
276 restoreExisted := t.RestoreExisted
277 return RewindResult{OK: false, Error: "file changed since rewind", Conflicts: []RewindConflict{{
278 Path: t.Path, Reason: CompareIdentity(fp, t.RestoreSHA, &restoreExisted, t.RestoreMode),
279 CurrentSHA: fp.SHA256, LastOwnedSHA: t.RestoreSHA, CurrentMode: fp.Mode, CheckpointMode: t.RestoreMode,
280 }}}, fmt.Errorf("file changed since rewind: %s", t.Path)
281 }
282 }
283 }
284
285 // Build inverse transaction: restore forward images.
286 undo := &TransactionManifest{
287 SchemaVersion: SchemaV2,
288 ID: newID("tx"),
289 SessionID: last.SessionID,
290 WorkspaceRoot: last.WorkspaceRoot,
291 State: TxPrepared,
292 Kind: "undo",
293 Turn: last.Turn,
294 Scope: last.Scope,
295 CreatedAt: time.Now(),
296 UpdatedAt: time.Now(),
297 SessionRevision: last.SessionRevision,
298 ParentTransaction: last.ID,
299 HasBoundary: last.HasBoundary,
300 BoundaryIndex: last.BoundaryIndex,
301 ConversationForward: last.ConversationForward,
302 CheckpointBackup: last.CheckpointBackup,
303 TruncateFrom: last.TruncateFrom,
304 }
305 for _, t := range last.Targets {
306 inv := TransactionTarget{
307 Path: t.Path,
308 AbsPath: t.AbsPath,
309 RestoreExisted: t.ForwardExisted,
310 RestoreMode: t.ForwardMode,
311 RestoreSHA: t.ForwardSHA,
312 RestoreBlob: t.ForwardBlob,
313 RestoreInline: clonePayload(t.ForwardInline),
314 ForwardExisted: t.RestoreExisted,
315 ForwardMode: t.RestoreMode,
316 ForwardSHA: t.RestoreSHA,
317 ForwardBlob: t.RestoreBlob,
318 ForwardInline: clonePayload(t.RestoreInline),
319 }
320 if t.ForwardExisted {
321 inv.Action = "write"
322 } else {
323 inv.Action = "delete"
324 }
325 inv.PublishTmp, inv.BackupPath = transactionSiblingPaths(inv.AbsPath, undo.ID, len(undo.Targets))
326 if inv.Action != "write" {
327 inv.PublishTmp = ""
328 }
329 undo.Targets = append(undo.Targets, inv)
330 }
331 if err := s.persistTransaction(undo); err != nil {
332 return RewindResult{OK: false, Error: err.Error()}, err
333 }
334
335 // Stage publish temps for write targets.
336 for i := range undo.Targets {
337 t := &undo.Targets[i]
338 if t.Action != "write" {
339 continue
340 }
341 data, err := s.loadBlobOrInline(t.RestoreBlob, t.RestoreInline)
342 if err != nil {
343 s.cleanupPublishTemps(undo.Targets)
344 _ = s.abortTransaction(undo, err)
345 return RewindResult{OK: false, Error: err.Error()}, err
346 }
347 mode := os.FileMode(0o644)
348 if t.RestoreMode != 0 {
349 mode = os.FileMode(t.RestoreMode)
350 }
351 if err := s.writePublishTemp(t.PublishTmp, data, mode); err != nil {
352 s.cleanupPublishTemps(undo.Targets)
353 _ = s.abortTransaction(undo, err)
354 return RewindResult{OK: false, Error: err.Error()}, err
355 }
356 }
357 undo.State = TxPrepared
358 if err := s.persistTransaction(undo); err != nil {
359 err = s.failTransaction(undo, undo.Targets, nil, err)
360 return RewindResult{OK: false, Error: err.Error()}, err
361 }
362
363 // For undo of conversation: restore forward conversation and checkpoints.
364 // Commit path for undo: publish files, then restore conversation/checkpoints.
365 result, err := s.commitUndoTransaction(undo, last, applier)
366 return result, err
367 }
368
369 func (s *Store) commitUndoTransaction(undo, original *TransactionManifest, applier ConversationApplier) (RewindResult, error) {
370 undo.State = TxCommitting
371 undo.UpdatedAt = time.Now()
372 if err := s.persistTransaction(undo); err != nil {
373 err = s.failTransaction(undo, undo.Targets, nil, err)
374 return RewindResult{OK: false, Error: err.Error()}, err
375 }
376
377 result := RewindResult{TransactionID: undo.ID, Coverage: CoverageComplete}
378 var stages []FileStage
379
380 // Publish files (inverse).
381 for i := range undo.Targets {
382 t := &undo.Targets[i]
383 st := FileStage{Path: t.Path, Phase: "commit", Action: t.Action}
384 t.Published = true
385 undo.UpdatedAt = time.Now()
386 if err := s.persistTransaction(undo); err != nil {
387 t.Published = false
388 st.Error = err.Error()
389 stages = append(stages, st)
390 err = s.failTransaction(undo, undo.Targets, stages, err)
391 result.Error = err.Error()
392 result.Files = stages
393 return result, err
394 }
395 if err := s.publishTarget(t); err != nil {
396 st.Error = err.Error()
397 stages = append(stages, st)
398 err = s.failTransaction(undo, undo.Targets[:i+1], stages, err)
399 result.OK = false
400 result.Error = err.Error()
401 result.Files = stages
402 return result, err
403 }
404 undo.UpdatedAt = time.Now()
405 if err := s.persistTransaction(undo); err != nil {
406 st.Error = err.Error()
407 stages = append(stages, st)
408 err = s.failTransaction(undo, undo.Targets[:i+1], stages, err)
409 result.Error = err.Error()
410 result.Files = stages
411 return result, err
412 }
413 st.Phase = "done"
414 stages = append(stages, st)
415 if t.Action == "write" {
416 result.Written = append(result.Written, t.Path)
417 } else {
418 result.Deleted = append(result.Deleted, t.Path)
419 }
420 }
421
422 // Restore conversation and checkpoints to pre-rewind state.
423 if applier != nil && len(original.ConversationForward) > 0 {
424 if err := applier.RestoreConversation(original.ConversationForward); err != nil {
425 restoreErr := s.restoreOriginalRewind(original, applier)
426 err = s.failTransactionAfterStateCompensation(undo, undo.Targets, stages, err, restoreErr)
427 result.OK = false
428 result.Error = err.Error()
429 result.Files = stages
430 return result, err
431 }
432 result.ConversationOK = true
433 }
434 if applier != nil && len(original.CheckpointBackup) > 0 {
435 if err := applier.RestoreCheckpoints(original.CheckpointBackup); err != nil {
436 // Return every side to the original rewind state before compensating
437 // the inverse file publish. Re-restoring the forward conversation here
438 // would leave conversation and files at opposite endpoints.
439 restoreErr := s.restoreOriginalRewind(original, applier)
440 err = s.failTransactionAfterStateCompensation(undo, undo.Targets, stages, err, restoreErr)
441 result.OK = false
442 result.Error = err.Error()
443 result.Files = stages
444 return result, err
445 }
446 }
447
448 // Controller appliers restore this same store and then rebuild their boundary
449 // index. Only the store-only path needs a direct restore here.
450 if applier == nil && len(original.CheckpointBackup) > 0 {
451 if err := s.restoreCheckpointBackup(original.CheckpointBackup); err != nil {
452 err = s.failTransactionAfterStateCompensation(undo, undo.Targets, stages, err, nil)
453 result.OK = false
454 result.Error = err.Error()
455 result.Files = stages
456 return result, err
457 }
458 }
459
460 undo.State = TxCommitted
461 undo.UpdatedAt = time.Now()
462 if err := s.persistTransaction(undo); err != nil {
463 restoreErr := s.restoreOriginalRewind(original, applier)
464 err = s.failTransactionAfterStateCompensation(undo, undo.Targets, stages, err, restoreErr)
465 result.OK = false
466 result.Error = err.Error()
467 result.Files = stages
468 return result, err
469 }
470
471 // Mark original as undone; clear lastUndo.
472 original.State = TxUndone
473 original.UpdatedAt = time.Now()
474 if err := s.persistTransaction(original); err != nil {
475 // The committed undo manifest durably names its parent, so startup will
476 // suppress the stale parent even if this secondary write failed.
477 slog.Warn("checkpoint: persist original transaction as undone", "err", err)
478 }
479 s.mu.Lock()
480 s.lastUndo = nil
481 s.mu.Unlock()
482
483 result.OK = true
484 result.UndoAvailable = false
485 result.Files = stages
486 return result, nil
487 }
488
489 func (s *Store) restoreOriginalRewind(original *TransactionManifest, applier ConversationApplier) error {
490 if original == nil || applier == nil {
491 return nil
492 }
493 if original.Scope != RewindConversation && original.Scope != RewindBoth {
494 return nil
495 }
496 var err error
497 if original.HasBoundary {
498 err = errors.Join(err, applier.ApplyConversationTruncate(original.BoundaryIndex, original.ConversationForward))
499 }
500 err = errors.Join(err, applier.TruncateCheckpoints(original.TruncateFrom))
501 return err
502 }
503
504 func (s *Store) prepareTransaction(plan RewindPlan, applier ConversationApplier) (*TransactionManifest, error) {
505 tx := newRewindTransaction(s.root, plan)
506 prepared := false
507 defer func() {
508 if !prepared {
509 s.cleanupPublishTemps(tx.Targets)
510 }
511 }()
512
513 if plan.Scope == RewindCode || plan.Scope == RewindBoth {
514 earliest := s.earliestRevisions(plan.Turn)
515 // Stable order for deterministic inject tests.
516 paths := make([]string, 0, len(earliest))
517 for p := range earliest {
518 paths = append(paths, p)
519 }
520 sort.Strings(paths)
521 for targetIndex, p := range paths {
522 rev := earliest[p]
523 abs, err := safePath(s.root, p)
524 if err != nil {
525 return nil, err
526 }
527 // Capture forward image.
528 fwd, gap, err := CapturePath(abs, CaptureOptions{WorkspaceRoot: s.root, ReadContent: true})
529 if err != nil && gap != nil {
530 return nil, fmt.Errorf("capture forward %s: %w", p, err)
531 }
532 t := TransactionTarget{
533 Path: p,
534 AbsPath: abs,
535 RestoreExisted: rev.Existed,
536 RestoreMode: rev.Mode,
537 RestoreSHA: rev.SHA256,
538 RestoreBlob: rev.BlobRef,
539 RestoreEncoding: rev.Encoding,
540 ForwardExisted: fwd.Existed,
541 ForwardMode: fwd.Mode,
542 ForwardSHA: fwd.SHA256,
543 }
544 if rev.Existed {
545 t.Action = "write"
546 if t.RestoreBlob == "" && rev.Content == nil {
547 return nil, fmt.Errorf("missing restore payload for %s", p)
548 }
549 // Stage publish temp. Blobs hold raw on-disk bytes; inline
550 // Content is decoded text and must be re-encoded. Legacy v1
551 // snapshots often omit Encoding — fall back to the current
552 // file's encoding (same as the pre-v2 RestoreCode path).
553 var data []byte
554 if rev.BlobRef != "" {
555 var lerr error
556 data, lerr = s.loadRevisionBytes(rev)
557 if lerr != nil {
558 return nil, lerr
559 }
560 } else if rev.Content != nil {
561 var eerr error
562 if data, eerr = s.encodeRevisionContent(rev, abs); eerr != nil {
563 return nil, eerr
564 }
565 } else {
566 return nil, fmt.Errorf("missing restore payload for %s", p)
567 }
568 mode := os.FileMode(0o644)
569 if rev.Mode != 0 {
570 mode = os.FileMode(rev.Mode)
571 }
572 if t.RestoreBlob == "" && s.blobs != nil {
573 ref, err := s.blobs.Put(data)
574 if err != nil {
575 return nil, err
576 }
577 t.RestoreBlob = ref
578 } else if t.RestoreBlob == "" {
579 t.RestoreInline = clonePayload(data)
580 }
581 t.PublishTmp, t.BackupPath = transactionSiblingPaths(abs, tx.ID, targetIndex)
582 if err := s.writePublishTemp(t.PublishTmp, data, mode); err != nil {
583 return nil, err
584 }
585 } else {
586 t.Action = "delete"
587 _, t.BackupPath = transactionSiblingPaths(abs, tx.ID, targetIndex)
588 }
589 if fwd.Existed && s.blobs != nil {
590 ref, err := s.blobs.Put(fwd.Content)
591 if err != nil {
592 if t.PublishTmp != "" {
593 _ = secureRemove(s.root, t.PublishTmp)
594 }
595 return nil, err
596 }
597 t.ForwardBlob = ref
598 } else if fwd.Existed {
599 t.ForwardInline = clonePayload(fwd.Content)
600 }
601 // Backup existing file for delete path (move later at commit).
602 tx.Targets = append(tx.Targets, t)
603 }
604 }
605
606 if shouldTruncateConversation(tx) && applier != nil {
607 // Backup future checkpoints for undo.
608 backup, err := s.backupCheckpointsFrom(plan.Turn)
609 if err != nil {
610 return nil, err
611 }
612 tx.CheckpointBackup = backup
613 }
614
615 if err := s.persistTransaction(tx); err != nil {
616 return nil, err
617 }
618 prepared = true
619 return tx, nil
620 }
621
622 func transactionSiblingPaths(absPath, transactionID string, index int) (publish, backup string) {
623 dir := filepath.Dir(absPath)
624 base := filepath.Base(absPath)
625 prefix := fmt.Sprintf(".%s.reasonix-%s-%d", base, transactionID, index)
626 return filepath.Join(dir, prefix+".tmp"), filepath.Join(dir, prefix+".bak")
627 }
628
629 func (s *Store) writePublishTemp(path string, data []byte, mode os.FileMode) error {
630 if err := secureWriteNew(s.root, path, data, mode); err != nil {
631 return fmt.Errorf("create publish temp: %w", err)
632 }
633 return nil
634 }
635
636 func (s *Store) cleanupPublishTemps(targets []TransactionTarget) {
637 for _, target := range targets {
638 if target.PublishTmp != "" {
639 _ = secureRemove(s.root, target.PublishTmp)
640 }
641 }
642 }
643
644 func (s *Store) commitTransaction(tx *TransactionManifest, applier ConversationApplier, inject *InjectFail) (RewindResult, error) {
645 tx.State = TxCommitting
646 tx.UpdatedAt = time.Now()
647 if err := s.persistTransaction(tx); err != nil {
648 err = s.failTransaction(tx, tx.Targets, nil, err)
649 return RewindResult{OK: false, Error: err.Error()}, err
650 }
651
652 result := RewindResult{TransactionID: tx.ID, Coverage: tx.Coverage, CoverageGaps: append([]CoverageGap(nil), tx.CoverageGaps...)}
653 stages := make([]FileStage, 0, len(tx.Targets))
654 filesDone := 0
655
656 for i := range tx.Targets {
657 t := &tx.Targets[i]
658 st := FileStage{Path: t.Path, Phase: "commit", Action: t.Action}
659 phase := "publish_file"
660 if t.Action == "delete" {
661 phase = "delete_file"
662 }
663 if inject != nil && inject.Phase == phase && filesDone >= inject.AfterFiles {
664 err := fmt.Errorf("injected failure at %s after %d files", inject.Phase, inject.AfterFiles)
665 st.Error = err.Error()
666 stages = append(stages, st)
667 err = s.failTransaction(tx, tx.Targets[:i], stages, err)
668 result.OK = false
669 result.Error = err.Error()
670 result.Files = stages
671 return result, err
672 }
673 // Persist a conservative "may have published" intent before the first
674 // filesystem rename. Recovery can safely compensate even if the crash
675 // happened just before publish.
676 t.Published = true
677 tx.UpdatedAt = time.Now()
678 if err := s.persistTransaction(tx); err != nil {
679 t.Published = false
680 st.Error = err.Error()
681 stages = append(stages, st)
682 err = s.failTransaction(tx, tx.Targets, stages, err)
683 result.Error = err.Error()
684 result.Files = stages
685 return result, err
686 }
687 if err := s.publishTarget(t); err != nil {
688 st.Error = err.Error()
689 stages = append(stages, st)
690 err = s.failTransaction(tx, tx.Targets[:i+1], stages, err)
691 result.OK = false
692 result.Error = err.Error()
693 result.Files = stages
694 return result, err
695 }
696 if inject != nil && inject.Phase == "after_publish_before_progress" && filesDone >= inject.AfterFiles {
697 // Deliberately leave the durable state as committing to simulate a
698 // process crash at the narrowest progress-persistence window.
699 err := fmt.Errorf("injected crash after publish before progress")
700 result.Error = err.Error()
701 result.Files = append(stages, st)
702 return result, err
703 }
704 tx.UpdatedAt = time.Now()
705 if err := s.persistTransaction(tx); err != nil {
706 st.Error = err.Error()
707 stages = append(stages, st)
708 err = s.failTransaction(tx, tx.Targets[:i+1], stages, err)
709 result.Error = err.Error()
710 result.Files = stages
711 return result, err
712 }
713 st.Phase = "done"
714 stages = append(stages, st)
715 filesDone++
716 if t.Action == "write" {
717 result.Written = append(result.Written, t.Path)
718 } else {
719 result.Deleted = append(result.Deleted, t.Path)
720 }
721 }
722 if err := s.persistTransaction(tx); err != nil {
723 err = s.failTransaction(tx, tx.Targets, stages, err)
724 result.Error = err.Error()
725 result.Files = stages
726 return result, err
727 }
728
729 if shouldTruncateConversation(tx) {
730 if inject != nil && inject.Phase == "conversation" {
731 err := fmt.Errorf("injected failure at conversation")
732 err = s.failTransaction(tx, tx.Targets, stages, err)
733 result.OK = false
734 result.Error = err.Error()
735 result.Files = stages
736 return result, err
737 }
738 if applier != nil && tx.HasBoundary {
739 if err := applier.ApplyConversationTruncate(tx.BoundaryIndex, tx.ConversationForward); err != nil {
740 restoreErr := s.restoreTransactionConversation(tx, applier)
741 err = s.failTransactionAfterStateCompensation(tx, tx.Targets, stages, err, restoreErr)
742 result.OK = false
743 result.Error = err.Error()
744 result.Files = stages
745 return result, err
746 }
747 result.ConversationOK = true
748 }
749 if inject != nil && inject.Phase == "truncate" {
750 err := fmt.Errorf("injected failure at truncate")
751 restoreErr := s.restoreTransactionConversation(tx, applier)
752 err = s.failTransactionAfterStateCompensation(tx, tx.Targets, stages, err, restoreErr)
753 result.OK = false
754 result.Error = err.Error()
755 result.Files = stages
756 return result, err
757 }
758 if applier != nil {
759 if err := applier.TruncateCheckpoints(tx.TruncateFrom); err != nil {
760 restoreErr := s.restoreTransactionConversation(tx, applier)
761 err = s.failTransactionAfterStateCompensation(tx, tx.Targets, stages, err, restoreErr)
762 result.OK = false
763 result.Error = err.Error()
764 result.Files = stages
765 return result, err
766 }
767 } else if err := s.TruncateFrom(tx.TruncateFrom); err != nil {
768 restoreErr := s.restoreTransactionConversation(tx, applier)
769 err = s.failTransactionAfterStateCompensation(tx, tx.Targets, stages, err, restoreErr)
770 result.OK = false
771 result.Error = err.Error()
772 result.Files = stages
773 return result, err
774 }
775 }
776
777 if inject != nil && inject.Phase == "finalize" {
778 err := fmt.Errorf("injected failure at finalize")
779 restoreErr := s.restoreTransactionConversation(tx, applier)
780 err = s.failTransactionAfterStateCompensation(tx, tx.Targets, stages, err, restoreErr)
781 result.OK = false
782 result.Error = err.Error()
783 result.Files = stages
784 return result, err
785 }
786 if inject != nil && inject.Phase == "after_conversation_before_finalize" {
787 // Simulate process death after both conversation mutations are durable but
788 // before the transaction can be marked committed. Startup must restore the
789 // forward transcript/checkpoints before compensating files.
790 err := fmt.Errorf("injected crash after conversation before finalize")
791 result.Error = err.Error()
792 result.Files = stages
793 return result, err
794 }
795
796 tx.State = TxCommitted
797 tx.UpdatedAt = time.Now()
798 if err := s.persistTransaction(tx); err != nil {
799 restoreErr := s.restoreTransactionConversation(tx, applier)
800 err = s.failTransactionAfterStateCompensation(tx, tx.Targets, stages, err, restoreErr)
801 result.OK = false
802 result.Error = err.Error()
803 result.Files = stages
804 return result, err
805 }
806 s.mu.Lock()
807 s.lastUndo = tx
808 s.mu.Unlock()
809
810 result.OK = true
811 result.UndoAvailable = true
812 result.Files = stages
813 return result, nil
814 }
815
816 func (s *Store) restoreTransactionConversation(tx *TransactionManifest, applier ConversationApplier) error {
817 if tx == nil {
818 return nil
819 }
820 var restoreErr error
821 if applier != nil {
822 if len(tx.ConversationForward) > 0 {
823 restoreErr = errors.Join(restoreErr, applier.RestoreConversation(tx.ConversationForward))
824 }
825 if len(tx.CheckpointBackup) > 0 {
826 restoreErr = errors.Join(restoreErr, applier.RestoreCheckpoints(tx.CheckpointBackup))
827 }
828 } else if len(tx.CheckpointBackup) > 0 {
829 restoreErr = errors.Join(restoreErr, s.restoreCheckpointBackup(tx.CheckpointBackup))
830 }
831 return restoreErr
832 }
833
834 // SetConversationForward attaches the pre-truncate conversation snapshot to a
835 // prepared transaction before commit. The controller calls this after Prepare.
836 func (s *Store) SetConversationForward(txID string, forward []byte) error {
837 path := s.txManifestPath(txID)
838 var tx TransactionManifest
839 if err := readJSONFile(path, &tx); err != nil {
840 // Also check in-memory last prepare path: store plans don't hold tx yet.
841 // Commit builds tx fresh; controller should pass forward via Commit options.
842 return err
843 }
844 setTransactionConversationForward(&tx, forward)
845 tx.UpdatedAt = time.Now()
846 return s.persistTransaction(&tx)
847 }
848
849 // CommitRewindWithForward is CommitRewind plus conversation forward payload.
850 func (s *Store) CommitRewindWithForward(planID string, forward []byte, applier ConversationApplier, inject *InjectFail) (RewindResult, error) {
851 if s == nil {
852 return RewindResult{}, fmt.Errorf("checkpoints unavailable")
853 }
854 s.mu.Lock()
855 pp, ok := s.plans[planID]
856 if ok {
857 delete(s.plans, planID)
858 }
859 s.mu.Unlock()
860 if !ok {
861 return RewindResult{OK: false, Error: "unknown or expired plan"}, fmt.Errorf("unknown or expired plan %q", planID)
862 }
863 plan := pp.plan
864
865 if plan.Scope == RewindBoth && (!plan.CanFiles || !plan.CanConversation) {
866 return RewindResult{OK: false, Error: plan.DisabledReason, Conflicts: plan.Conflicts}, fmt.Errorf("%s", plan.DisabledReason)
867 }
868 if plan.Scope == RewindCode && !plan.CanFiles {
869 return RewindResult{OK: false, Error: plan.DisabledReason, Conflicts: plan.Conflicts}, fmt.Errorf("%s", plan.DisabledReason)
870 }
871 if (plan.Scope == RewindConversation || plan.Scope == RewindBoth) && !plan.CanConversation {
872 return RewindResult{OK: false, Error: plan.DisabledReason}, fmt.Errorf("%s", plan.DisabledReason)
873 }
874
875 if !s.barrier.TryEnterExclusive() {
876 err := fmt.Errorf("workspace mutation in progress")
877 return RewindResult{OK: false, Error: err.Error(), Conflicts: []RewindConflict{{Reason: ConflictBusyWriter}}}, err
878 }
879 defer s.barrier.ExitExclusive()
880 if conflicts := s.activeWriterConflicts(); len(conflicts) > 0 {
881 err := fmt.Errorf("active background writer")
882 return RewindResult{OK: false, Error: err.Error(), Conflicts: conflicts, Coverage: plan.Coverage}, err
883 }
884
885 if plan.Scope == RewindCode || plan.Scope == RewindBoth {
886 if plan.WorkspaceToken != fmt.Sprintf("%d", s.barrier.Generation()) {
887 conflict := RewindConflict{Reason: ConflictStalePlan}
888 return RewindResult{OK: false, Error: "workspace changed since preview", Conflicts: []RewindConflict{conflict}, Coverage: plan.Coverage}, fmt.Errorf("workspace changed since preview")
889 }
890 if conflicts := s.precheckFiles(plan.Turn); len(conflicts) > 0 {
891 return RewindResult{OK: false, Error: "file conflicts detected", Conflicts: conflicts}, fmt.Errorf("file conflicts detected")
892 }
893 }
894
895 tx, err := s.prepareTransaction(plan, applier)
896 if err != nil {
897 return RewindResult{OK: false, Error: err.Error()}, err
898 }
899 setTransactionConversationForward(tx, forward)
900 if err := s.persistTransaction(tx); err != nil {
901 return RewindResult{OK: false, Error: err.Error()}, err
902 }
903 return s.commitTransaction(tx, applier, inject)
904 }
905
906 func (s *Store) publishTarget(t *TransactionTarget) error {
907 if t.BackupPath == "" {
908 return fmt.Errorf("missing transaction backup path for %s", t.Path)
909 }
910 backupExists, err := securePathExists(s.root, t.BackupPath)
911 if err != nil {
912 return err
913 }
914 if backupExists {
915 return fmt.Errorf("transaction backup already exists for %s", t.Path)
916 }
917 targetExists, err := securePathExists(s.root, t.AbsPath)
918 if err != nil {
919 return err
920 }
921 if targetExists {
922 if err := secureRename(s.root, t.AbsPath, t.BackupPath); err != nil {
923 return fmt.Errorf("backup %s: %w", t.Path, err)
924 }
925 }
926 if t.Action == "delete" {
927 return nil
928 }
929 if t.PublishTmp == "" {
930 return fmt.Errorf("missing publish tmp for %s", t.Path)
931 }
932 if err := secureRename(s.root, t.PublishTmp, t.AbsPath); err != nil {
933 restoreErr := error(nil)
934 if exists, statErr := securePathExists(s.root, t.BackupPath); statErr == nil && exists {
935 restoreErr = secureRename(s.root, t.BackupPath, t.AbsPath)
936 }
937 return errors.Join(fmt.Errorf("publish %s: %w", t.Path, err), restoreErr)
938 }
939 if t.RestoreMode != 0 {
940 if err := secureChmod(s.root, t.AbsPath, os.FileMode(t.RestoreMode)); err != nil {
941 return fmt.Errorf("chmod restored %s: %w", t.Path, err)
942 }
943 }
944 return nil
945 }
946
947 func (s *Store) compensatePublished(targets []TransactionTarget, stages []FileStage) error {
948 var first error
949 for _, v := range slices.Backward(targets) {
950 t := v
951 if !t.Published {
952 if t.PublishTmp != "" {
953 _ = secureRemove(s.root, t.PublishTmp)
954 }
955 continue
956 }
957 // Published is a durable intent. If the target still exactly matches its
958 // forward image, the crash happened before publish and compensation is a
959 // no-op. Any other unrelated state is preserved with a recovery copy.
960 var err error
961 cur, fpErr := FingerprintPath(s.root, t.AbsPath)
962 if fpErr != nil {
963 markCompensationStage(stages, t.Path, fpErr)
964 if first == nil {
965 first = fpErr
966 }
967 continue
968 } else if fingerprintMatches(cur, t.ForwardExisted, t.ForwardSHA, t.ForwardMode) {
969 if t.PublishTmp != "" {
970 _ = secureRemove(s.root, t.PublishTmp)
971 }
972 markCompensationStage(stages, t.Path, nil)
973 continue
974 }
975 // Crash window: publishTarget durably records Published before moving the
976 // target to its backup. A process death after that first rename leaves the
977 // target absent, the forward image in BackupPath, and (for writes) the
978 // publish temp still present. Recognize that owned intermediate state before
979 // classifying the absent target as an external modification.
980 if t.ForwardExisted && !cur.Existed && t.BackupPath != "" {
981 backup, backupErr := FingerprintPath(s.root, t.BackupPath)
982 publishPending := t.Action == "delete"
983 if t.Action == "write" && t.PublishTmp != "" {
984 publishPending, _ = securePathExists(s.root, t.PublishTmp)
985 }
986 if backupErr == nil && publishPending && fingerprintMatches(backup, true, t.ForwardSHA, t.ForwardMode) {
987 err = secureRename(s.root, t.BackupPath, t.AbsPath)
988 if err == nil && t.PublishTmp != "" {
989 if removeErr := secureRemove(s.root, t.PublishTmp); removeErr != nil && !os.IsNotExist(removeErr) {
990 err = removeErr
991 }
992 }
993 markCompensationStage(stages, t.Path, err)
994 if err != nil && first == nil {
995 first = err
996 }
997 continue
998 }
999 }
1000 if t.ForwardExisted {
1001 data, lerr := s.loadBlobOrInline(t.ForwardBlob, t.ForwardInline)
1002 if lerr != nil && t.BackupPath != "" {
1003 data, lerr = secureReadFile(s.root, t.BackupPath)
1004 }
1005 if !fingerprintMatches(cur, t.RestoreExisted, t.RestoreSHA, t.RestoreMode) {
1006 if lerr == nil {
1007 suffix := t.RestoreSHA
1008 if len(suffix) > 8 {
1009 suffix = suffix[:8]
1010 }
1011 if suffix == "" {
1012 suffix = "unknown"
1013 }
1014 recov := t.AbsPath + ".reasonix-recovery-" + suffix
1015 _ = secureWriteNew(s.root, recov, data, os.FileMode(t.ForwardMode))
1016 err = fmt.Errorf("external modification after publish; recovery copy at %s", recov)
1017 } else {
1018 err = lerr
1019 }
1020 } else if backupExists, backupErr := securePathExists(s.root, t.BackupPath); backupErr == nil && backupExists {
1021 if cur.Existed {
1022 err = secureRemove(s.root, t.AbsPath)
1023 }
1024 if err == nil {
1025 err = secureRename(s.root, t.BackupPath, t.AbsPath)
1026 }
1027 } else if lerr != nil {
1028 err = lerr
1029 } else {
1030 mode := os.FileMode(0o644)
1031 if t.ForwardMode != 0 {
1032 mode = os.FileMode(t.ForwardMode)
1033 }
1034 if cur.Existed {
1035 if werr := secureRemove(s.root, t.AbsPath); werr != nil {
1036 err = werr
1037 }
1038 }
1039 tmp, _ := transactionSiblingPaths(t.AbsPath, newID("compensate"), 0)
1040 if werr := s.writePublishTemp(tmp, data, mode); werr != nil {
1041 err = werr
1042 } else if err == nil {
1043 if werr := secureRename(s.root, tmp, t.AbsPath); werr != nil {
1044 err = werr
1045 }
1046 }
1047 }
1048 } else {
1049 // Forward did not exist — remove what we published.
1050 if !fingerprintMatches(cur, t.RestoreExisted, t.RestoreSHA, t.RestoreMode) {
1051 // External rewrite of a file we restored then someone changed —
1052 // for compensate of delete action inverse: leave it.
1053 err = fmt.Errorf("external modification; not removing %s", t.AbsPath)
1054 } else {
1055 err = secureRemove(s.root, t.AbsPath)
1056 if os.IsNotExist(err) {
1057 err = nil
1058 }
1059 }
1060 }
1061 markCompensationStage(stages, t.Path, err)
1062 if err != nil && first == nil {
1063 first = err
1064 }
1065 }
1066 return first
1067 }
1068
1069 func fingerprintMatches(fp Fingerprint, existed bool, sha string, mode uint32) bool {
1070 if fp.Existed != existed {
1071 return false
1072 }
1073 if !existed {
1074 return true
1075 }
1076 if sha != "" && fp.SHA256 != sha {
1077 return false
1078 }
1079 return mode == 0 || fp.Mode == 0 || fp.Mode == mode
1080 }
1081
1082 func markCompensationStage(stages []FileStage, path string, err error) {
1083 for i := range stages {
1084 if stages[i].Path != path {
1085 continue
1086 }
1087 stages[i].Compensated = err == nil
1088 if err != nil {
1089 stages[i].CompError = err.Error()
1090 }
1091 }
1092 }
1093
1094 func (s *Store) failTransaction(tx *TransactionManifest, targets []TransactionTarget, stages []FileStage, cause error) error {
1095 return s.failTransactionAfterStateCompensation(tx, targets, stages, cause, nil)
1096 }
1097
1098 // failTransactionAfterStateCompensation compensates files and records whether
1099 // the conversation/checkpoint side was also restored. Any incomplete side keeps
1100 // the manifest committing so startup can retry the whole compensation.
1101 func (s *Store) failTransactionAfterStateCompensation(tx *TransactionManifest, targets []TransactionTarget, stages []FileStage, cause, stateCompensationErr error) error {
1102 if tx != nil {
1103 targets = tx.Targets
1104 }
1105 compensationErr := s.compensatePublished(targets, stages)
1106 combined := errors.Join(cause, stateCompensationErr)
1107 if compensationErr != nil {
1108 combined = errors.Join(combined, fmt.Errorf("compensation failed: %w", compensationErr))
1109 }
1110 if compensationErr != nil || stateCompensationErr != nil {
1111 // Do not make a failed compensation terminal. Startup recovery retries
1112 // committing manifests; marking this aborted would strand a half-applied
1113 // workspace permanently.
1114 tx.State = TxCommitting
1115 tx.Error = combined.Error()
1116 tx.UpdatedAt = time.Now()
1117 if persistErr := s.persistTransaction(tx); persistErr != nil {
1118 combined = errors.Join(combined, fmt.Errorf("persist pending compensation: %w", persistErr))
1119 }
1120 return combined
1121 }
1122 if abortErr := s.abortTransaction(tx, combined); abortErr != nil {
1123 combined = errors.Join(combined, fmt.Errorf("persist aborted transaction: %w", abortErr))
1124 }
1125 return combined
1126 }
1127
1128 func (s *Store) abortTransaction(tx *TransactionManifest, cause error) error {
1129 tx.State = TxAborted
1130 tx.Error = cause.Error()
1131 tx.UpdatedAt = time.Now()
1132 return s.persistTransaction(tx)
1133 }
1134
1135 func (s *Store) persistTransaction(tx *TransactionManifest) error {
1136 if s.dir == "" {
1137 return nil
1138 }
1139 return writeJSONAtomic(s.txManifestPath(tx.ID), tx)
1140 }
1141
1142 func (s *Store) txDir() string {
1143 if s.dir == "" {
1144 return filepath.Join(os.TempDir(), "reasonix-ckpt-tx")
1145 }
1146 return filepath.Join(s.dir, "transactions")
1147 }
1148
1149 func (s *Store) txManifestPath(id string) string {
1150 return filepath.Join(s.txDir(), id+".json")
1151 }
1152
1153 // RecoverTransactions scans for incomplete file-only transactions. Conversation
1154 // transactions are intentionally deferred until the controller has installed the
1155 // resumed session and can provide a ConversationApplier.
1156 func (s *Store) RecoverTransactions() []string {
1157 return s.recoverTransactions(nil)
1158 }
1159
1160 // RecoverTransactionsWithApplier finishes startup recovery after the resumed
1161 // conversation is live. A committing rewind first restores its forward
1162 // transcript/checkpoints, then compensates files; a committing undo first
1163 // reapplies its parent rewind, then compensates files. The manifest remains
1164 // committing if either side fails so a later startup can retry idempotently.
1165 func (s *Store) RecoverTransactionsWithApplier(applier ConversationApplier) []string {
1166 return s.recoverTransactions(applier)
1167 }
1168
1169 func (s *Store) recoverTransactions(applier ConversationApplier) []string {
1170 if s == nil || s.dir == "" {
1171 return nil
1172 }
1173 dir := s.txDir()
1174 ents, err := os.ReadDir(dir)
1175 if err != nil {
1176 return nil
1177 }
1178 undoneParents := map[string]bool{}
1179 for _, entry := range ents {
1180 if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".json") {
1181 continue
1182 }
1183 var tx TransactionManifest
1184 if readJSONFile(filepath.Join(dir, entry.Name()), &tx) == nil && tx.State == TxCommitted && tx.Kind == "undo" && tx.ParentTransaction != "" {
1185 undoneParents[tx.ParentTransaction] = true
1186 }
1187 }
1188 var notes []string
1189 for _, e := range ents {
1190 if e.IsDir() || !strings.HasSuffix(e.Name(), ".json") {
1191 continue
1192 }
1193 var tx TransactionManifest
1194 if err := readJSONFile(filepath.Join(dir, e.Name()), &tx); err != nil {
1195 continue
1196 }
1197 switch tx.State {
1198 case TxPrepared:
1199 // Never published — safe to discard.
1200 for _, target := range tx.Targets {
1201 if target.PublishTmp != "" {
1202 _ = secureRemove(s.root, target.PublishTmp)
1203 }
1204 }
1205 tx.State = TxAborted
1206 tx.Error = "abandoned prepared transaction on recovery"
1207 tx.UpdatedAt = time.Now()
1208 _ = s.persistTransaction(&tx)
1209 notes = append(notes, fmt.Sprintf("aborted prepared %s", tx.ID))
1210 case TxCommitting:
1211 needsConversation := tx.Scope == RewindConversation || tx.Scope == RewindBoth
1212 if needsConversation && applier == nil {
1213 notes = append(notes, fmt.Sprintf("deferred conversation recovery %s", tx.ID))
1214 continue
1215 }
1216 if needsConversation {
1217 var restoreErr error
1218 if tx.Kind != "undo" && tx.HasBoundary && len(tx.ConversationForward) == 0 {
1219 restoreErr = fmt.Errorf("missing forward conversation payload")
1220 } else if tx.Kind == "undo" {
1221 if tx.HasBoundary {
1222 restoreErr = errors.Join(restoreErr, applier.ApplyConversationTruncate(tx.BoundaryIndex, tx.ConversationForward))
1223 }
1224 restoreErr = errors.Join(restoreErr, applier.TruncateCheckpoints(tx.TruncateFrom))
1225 } else {
1226 restoreErr = s.restoreTransactionConversation(&tx, applier)
1227 }
1228 if restoreErr != nil {
1229 notes = append(notes, fmt.Sprintf("conversation recovery %s pending: %v", tx.ID, restoreErr))
1230 tx.Error = fmt.Sprintf("crash recovery conversation compensation pending: %v", restoreErr)
1231 tx.UpdatedAt = time.Now()
1232 _ = s.persistTransaction(&tx)
1233 continue
1234 }
1235 }
1236 // Compensate published files back to forward images.
1237 stages := make([]FileStage, len(tx.Targets))
1238 for i, t := range tx.Targets {
1239 stages[i] = FileStage{Path: t.Path, Phase: "compensate"}
1240 }
1241 if err := s.compensatePublished(tx.Targets, stages); err != nil {
1242 notes = append(notes, fmt.Sprintf("compensate %s: %v", tx.ID, err))
1243 tx.Error = fmt.Sprintf("crash recovery compensation pending: %v", err)
1244 tx.UpdatedAt = time.Now()
1245 _ = s.persistTransaction(&tx)
1246 } else {
1247 notes = append(notes, fmt.Sprintf("compensated committing %s", tx.ID))
1248 tx.State = TxAborted
1249 tx.Error = "compensated after crash during commit"
1250 tx.UpdatedAt = time.Now()
1251 _ = s.persistTransaction(&tx)
1252 }
1253 case TxCommitted:
1254 if tx.Kind == "undo" || undoneParents[tx.ID] {
1255 continue
1256 }
1257 // Keep as last undo if newer.
1258 s.mu.Lock()
1259 if s.lastUndo == nil || s.lastUndo.UpdatedAt.Before(tx.UpdatedAt) {
1260 cp := tx
1261 s.lastUndo = &cp
1262 }
1263 s.mu.Unlock()
1264 }
1265 }
1266 return notes
1267 }
1268
1269 func (s *Store) precheckFiles(fromTurn int) []RewindConflict {
1270 earliest := s.earliestRevisions(fromTurn)
1271 var conflicts []RewindConflict
1272 for p, rev := range earliest {
1273 abs, err := safePath(s.root, p)
1274 if err != nil {
1275 conflicts = append(conflicts, RewindConflict{Path: p, Reason: ConflictPathUnsafe})
1276 continue
1277 }
1278 if rev.BlobRef == "" && rev.Content == nil && rev.Existed {
1279 if rev.SHA256 != "" && s.blobs != nil && !s.blobs.Has(rev.SHA256) && (rev.BlobRef == "" || !s.blobs.Has(rev.BlobRef)) {
1280 conflicts = append(conflicts, RewindConflict{Path: p, Reason: ConflictMissingPayload, CheckpointSHA: rev.SHA256})
1281 continue
1282 }
1283 if rev.Content == nil && rev.BlobRef == "" {
1284 conflicts = append(conflicts, RewindConflict{Path: p, Reason: ConflictMissingPayload})
1285 continue
1286 }
1287 }
1288 fp, err := FingerprintPath(s.root, abs)
1289 if err != nil {
1290 // unreadable etc.
1291 conflicts = append(conflicts, RewindConflict{Path: p, Reason: ConflictExternalChange, CheckpointSHA: rev.SHA256})
1292 continue
1293 }
1294 if MatchesRestoreImage(fp, rev.SHA256, rev.Existed) {
1295 continue
1296 }
1297 // Prefer after fingerprint for conflict detection.
1298 reason := CompareIdentity(fp, rev.AfterSHA256, rev.AfterExisted, rev.AfterMode)
1299 if reason == ConflictCoverageLegacy {
1300 // Legacy handled at plan level; skip per-file for batch.
1301 continue
1302 }
1303 if reason != "" {
1304 conflicts = append(conflicts, RewindConflict{
1305 Path: p,
1306 Reason: reason,
1307 CheckpointSHA: rev.SHA256,
1308 LastOwnedSHA: rev.AfterSHA256,
1309 CurrentSHA: fp.SHA256,
1310 CheckpointMode: rev.Mode,
1311 CurrentMode: fp.Mode,
1312 CurrentExisted: fp.Existed,
1313 CheckpointExist: rev.Existed,
1314 })
1315 }
1316 }
1317 sort.Slice(conflicts, func(i, j int) bool { return conflicts[i].Path < conflicts[j].Path })
1318 return conflicts
1319 }
1320
1321 func (s *Store) earliestRevisions(fromTurn int) map[string]FileRevision {
1322 s.mu.Lock()
1323 defer s.mu.Unlock()
1324 return s.earliestRevisionsLocked(fromTurn)
1325 }
1326
1327 func (s *Store) earliestRevisionsLocked(fromTurn int) map[string]FileRevision {
1328 earliest := map[string]FileRevision{}
1329 for _, c := range s.all() {
1330 if c.Turn < fromTurn {
1331 continue
1332 }
1333 for _, rev := range c.revisions() {
1334 pathKey := NormalizeRelPath(s.root, rev.Path)
1335 if first, ok := earliest[pathKey]; ok {
1336 // Preserve the earliest preimage, but carry forward the final
1337 // mutation's ownership identity. Missing final identity deliberately
1338 // clears an older proof instead of authorizing an unsafe restore.
1339 first.AfterExisted = rev.AfterExisted
1340 first.AfterSHA256 = rev.AfterSHA256
1341 first.AfterMode = rev.AfterMode
1342 earliest[pathKey] = first
1343 continue
1344 }
1345 rev.Path = pathKey
1346 earliest[pathKey] = rev
1347 }
1348 }
1349 return earliest
1350 }
1351
1352 func (s *Store) filesFromTurnLocked(fromTurn int) []string {
1353 seen := map[string]bool{}
1354 var out []string
1355 for _, c := range s.all() {
1356 if c.Turn < fromTurn {
1357 continue
1358 }
1359 for _, rev := range c.revisions() {
1360 pathKey := NormalizeRelPath(s.root, rev.Path)
1361 if seen[pathKey] {
1362 continue
1363 }
1364 seen[pathKey] = true
1365 out = append(out, pathKey)
1366 }
1367 }
1368 sort.Strings(out)
1369 return out
1370 }
1371
1372 func (s *Store) coverageFromTurnLocked(fromTurn int) (Coverage, []CoverageGap, bool, bool) {
1373 var gaps []CoverageGap
1374 legacy := false
1375 expired := false
1376 hasFiles := false
1377 partial := false
1378 for _, c := range s.all() {
1379 if c.Turn < fromTurn {
1380 continue
1381 }
1382 if c.SchemaVersion < SchemaV2 && c.SchemaVersion != 0 {
1383 legacy = true
1384 }
1385 if c.SchemaVersion == 0 {
1386 // v1 had no schemaVersion field
1387 legacy = true
1388 }
1389 if c.Coverage == CoverageLegacy || c.Legacy {
1390 legacy = true
1391 }
1392 if c.ExpiredFilePayload {
1393 expired = true
1394 }
1395 if c.Coverage == CoveragePartial {
1396 partial = true
1397 }
1398 gaps = append(gaps, c.CoverageGaps...)
1399 if len(c.revisions()) > 0 {
1400 hasFiles = true
1401 }
1402 }
1403 if legacy {
1404 return CoverageLegacy, append(gaps, CoverageGap{Reason: GapLegacyUnverified}), true, expired
1405 }
1406 if expired {
1407 return CoveragePartial, append(gaps, CoverageGap{Reason: GapExpiredPayload}), false, true
1408 }
1409 if !hasFiles {
1410 if len(gaps) > 0 {
1411 return CoverageNone, gaps, false, false
1412 }
1413 return CoverageNone, nil, false, false
1414 }
1415 if partial || len(gaps) > 0 {
1416 return CoveragePartial, gaps, false, false
1417 }
1418 return CoverageComplete, nil, false, false
1419 }
1420
1421 func (s *Store) loadRevisionBytes(rev FileRevision) ([]byte, error) {
1422 if rev.BlobRef != "" && s.blobs != nil {
1423 return s.blobs.Get(rev.BlobRef)
1424 }
1425 if rev.Content != nil {
1426 return []byte(*rev.Content), nil
1427 }
1428 if rev.SHA256 != "" && s.blobs != nil && s.blobs.Has(rev.SHA256) {
1429 return s.blobs.Get(rev.SHA256)
1430 }
1431 return nil, fmt.Errorf("missing payload for %s", rev.Path)
1432 }
1433
1434 func (s *Store) loadBlobOrInline(ref string, inline []byte) ([]byte, error) {
1435 if ref != "" && s.blobs != nil {
1436 return s.blobs.Get(ref)
1437 }
1438 if inline != nil {
1439 return inline, nil
1440 }
1441 return nil, fmt.Errorf("missing blob %q", ref)
1442 }
1443
1444 func (s *Store) backupCheckpointsFrom(fromTurn int) ([]byte, error) {
1445 s.mu.Lock()
1446 defer s.mu.Unlock()
1447 var future []*Checkpoint
1448 for _, c := range s.all() {
1449 if c.Turn >= fromTurn {
1450 cp := *c
1451 future = append(future, &cp)
1452 }
1453 }
1454 return json.Marshal(future)
1455 }
1456
1457 func (s *Store) restoreCheckpointBackup(backup []byte) error {
1458 var future []*Checkpoint
1459 if err := json.Unmarshal(backup, &future); err != nil {
1460 return err
1461 }
1462 s.mu.Lock()
1463 defer s.mu.Unlock()
1464 // Merge future checkpoints back (by turn).
1465 byTurn := map[int]*Checkpoint{}
1466 for _, c := range s.done {
1467 byTurn[c.Turn] = c
1468 }
1469 if s.cur != nil {
1470 byTurn[s.cur.Turn] = s.cur
1471 }
1472 for _, c := range future {
1473 byTurn[c.Turn] = c
1474 if err := s.persist(c); err != nil {
1475 return fmt.Errorf("persist restored checkpoint turn %d: %w", c.Turn, err)
1476 }
1477 counterpart := filepath.Join(s.expiredDir(), fmt.Sprintf("turn-%d.json", c.Turn))
1478 if c.ExpiredFilePayload {
1479 counterpart = filepath.Join(s.dir, fmt.Sprintf("turn-%d.json", c.Turn))
1480 }
1481 if err := os.Remove(counterpart); err != nil && !os.IsNotExist(err) {
1482 return fmt.Errorf("remove stale checkpoint counterpart turn %d: %w", c.Turn, err)
1483 }
1484 }
1485 // Rebuild done/cur: highest turn as cur if it was cur; else all in done.
1486 turns := make([]int, 0, len(byTurn))
1487 for t := range byTurn {
1488 turns = append(turns, t)
1489 }
1490 sort.Ints(turns)
1491 s.done = nil
1492 s.cur = nil
1493 for _, t := range turns {
1494 s.done = append(s.done, byTurn[t])
1495 }
1496 return nil
1497 }
1498
1499 func newID(prefix string) string {
1500 return fmt.Sprintf("%s-%d-%s", prefix, time.Now().UnixNano(), Digest(fmt.Appendf(nil, "%d", time.Now().UnixNano()))[:8])
1501 }
1502
1503 // RestoreCheckpointBackupPublic reloads backed-up checkpoints after an undo.
1504 func (s *Store) RestoreCheckpointBackupPublic(backup []byte) error {
1505 return s.restoreCheckpointBackup(backup)
1506 }
1507
1508 // PrepareFileRevert prepares a single-file restore to the earliest session preimage.
1509 func (s *Store) PrepareFileRevert(path string, sessionRev int64) (RewindPlan, error) {
1510 if s == nil {
1511 return RewindPlan{}, fmt.Errorf("checkpoints unavailable")
1512 }
1513 plan := RewindPlan{
1514 PlanID: newID("plan"),
1515 Scope: RewindCode,
1516 Path: path,
1517 SessionRevision: sessionRev,
1518 CreatedAt: time.Now(),
1519 WorkspaceToken: fmt.Sprintf("%d", s.barrier.Generation()),
1520 Files: []string{path},
1521 FileCount: 1,
1522 }
1523 state, ok := s.FileState(path)
1524 if !ok {
1525 plan.CanFiles = false
1526 plan.DisabledReason = "file is not session-owned"
1527 return plan, nil
1528 }
1529 _ = state
1530 abs, err := safePath(s.root, path)
1531 if err != nil {
1532 plan.CanFiles = false
1533 plan.DisabledReason = "path unsafe"
1534 plan.Conflicts = []RewindConflict{{Path: path, Reason: ConflictPathUnsafe}}
1535 return plan, nil
1536 }
1537 revs := s.earliestRevisions(0)
1538 rev, has := revs[path]
1539 if !has {
1540 for p, r := range revs {
1541 if ap, e := safePath(s.root, p); e == nil && ap == abs {
1542 rev, has = r, true
1543 plan.Path = p
1544 break
1545 }
1546 }
1547 }
1548 if !has {
1549 plan.CanFiles = false
1550 plan.DisabledReason = "file is not session-owned"
1551 return plan, nil
1552 }
1553 if rev.AfterExisted == nil && rev.AfterSHA256 == "" {
1554 // A v1 or incomplete capture has a preimage but no evidence that the
1555 // current file is still the session's last write. Do not turn the
1556 // generic conflict-overwrite affordance into an unsafe legacy restore.
1557 plan.PlanID = ""
1558 plan.CanFiles = false
1559 plan.Legacy = true
1560 plan.Coverage = CoverageLegacy
1561 plan.DisabledReason = "legacy checkpoint cannot verify later manual edits"
1562 return plan, nil
1563 }
1564 fp, fperr := FingerprintPath(s.root, abs)
1565 if fperr == nil {
1566 reason := CompareIdentity(fp, rev.AfterSHA256, rev.AfterExisted, rev.AfterMode)
1567 if reason != "" {
1568 plan.Conflicts = []RewindConflict{{
1569 Path: path, Reason: reason,
1570 CheckpointSHA: rev.SHA256, LastOwnedSHA: rev.AfterSHA256, CurrentSHA: fp.SHA256,
1571 CurrentExisted: fp.Existed, CheckpointExist: rev.Existed,
1572 }}
1573 plan.CanFiles = true
1574 plan.DisabledReason = "conflict requires explicit resolution"
1575 } else {
1576 plan.CanFiles = true
1577 }
1578 } else {
1579 plan.CanFiles = false
1580 plan.DisabledReason = "current file identity unavailable"
1581 plan.Conflicts = []RewindConflict{{Path: path, Reason: ConflictExternalChange}}
1582 }
1583 if rev.BlobRef == "" && rev.Content == nil && rev.Existed {
1584 plan.CanFiles = false
1585 plan.DisabledReason = "missing file payload"
1586 plan.Conflicts = append(plan.Conflicts, RewindConflict{Path: path, Reason: ConflictMissingPayload})
1587 }
1588 s.mu.Lock()
1589 if s.plans == nil {
1590 s.plans = map[string]preparedPlan{}
1591 }
1592 s.plans[plan.PlanID] = preparedPlan{plan: plan, created: time.Now(), previewFingerprint: &fp}
1593 s.mu.Unlock()
1594 return plan, nil
1595 }
1596
1597 // CommitFileRevert commits a single-file restore.
1598 func (s *Store) CommitFileRevert(planID string, resolution ConflictResolution) (RewindResult, error) {
1599 if s == nil {
1600 return RewindResult{}, fmt.Errorf("checkpoints unavailable")
1601 }
1602 s.mu.Lock()
1603 pp, ok := s.plans[planID]
1604 if ok {
1605 delete(s.plans, planID)
1606 }
1607 s.mu.Unlock()
1608 if !ok {
1609 return RewindResult{OK: false, Error: "unknown or expired plan"}, fmt.Errorf("unknown or expired plan")
1610 }
1611 plan := pp.plan
1612 if plan.Path == "" {
1613 return RewindResult{OK: false, Error: "not a file plan"}, fmt.Errorf("not a file plan")
1614 }
1615 if !plan.CanFiles {
1616 return RewindResult{OK: false, Error: plan.DisabledReason, Conflicts: plan.Conflicts}, fmt.Errorf("%s", plan.DisabledReason)
1617 }
1618 if len(plan.Conflicts) > 0 && resolution != ResolveOverwriteCheckpoint {
1619 if resolution == ResolveKeepCurrent {
1620 return RewindResult{OK: true, UndoAvailable: false}, nil
1621 }
1622 return RewindResult{OK: false, Error: "conflict requires explicit resolution", Conflicts: plan.Conflicts}, fmt.Errorf("conflict requires explicit resolution")
1623 }
1624 if !s.barrier.TryEnterExclusive() {
1625 err := fmt.Errorf("workspace mutation in progress")
1626 return RewindResult{OK: false, Error: err.Error(), Conflicts: []RewindConflict{{Path: plan.Path, Reason: ConflictBusyWriter}}}, err
1627 }
1628 defer s.barrier.ExitExclusive()
1629 if conflicts := s.activeWriterConflicts(); len(conflicts) > 0 {
1630 err := fmt.Errorf("active background writer")
1631 return RewindResult{OK: false, Error: err.Error(), Conflicts: conflicts}, err
1632 }
1633 if plan.WorkspaceToken != fmt.Sprintf("%d", s.barrier.Generation()) {
1634 conflict := RewindConflict{Path: plan.Path, Reason: ConflictStalePlan}
1635 return RewindResult{OK: false, Error: "workspace changed since preview", Conflicts: []RewindConflict{conflict}}, fmt.Errorf("workspace changed since preview")
1636 }
1637 absPreview, err := safePath(s.root, plan.Path)
1638 if err != nil {
1639 return RewindResult{OK: false, Error: err.Error()}, err
1640 }
1641 current, err := FingerprintPath(s.root, absPreview)
1642 if err != nil || pp.previewFingerprint == nil || !sameFingerprint(current, *pp.previewFingerprint) {
1643 conflict := RewindConflict{Path: plan.Path, Reason: ConflictStalePlan, CurrentSHA: current.SHA256, CurrentExisted: current.Existed}
1644 return RewindResult{OK: false, Error: "file changed since preview; preview again", Conflicts: []RewindConflict{conflict}}, fmt.Errorf("file changed since preview; preview again")
1645 }
1646
1647 // Restore via restoreCodeLegacy for the single earliest path using turn 0.
1648 // Build synthetic order of one path.
1649 revs := s.earliestRevisions(0)
1650 rev, has := revs[plan.Path]
1651 if !has {
1652 abs, _ := safePath(s.root, plan.Path)
1653 for p, r := range revs {
1654 if ap, e := safePath(s.root, p); e == nil && ap == abs {
1655 rev, has = r, true
1656 plan.Path = p
1657 break
1658 }
1659 }
1660 }
1661 if !has {
1662 return RewindResult{OK: false, Error: "file is not session-owned"}, fmt.Errorf("file is not session-owned")
1663 }
1664
1665 // Find which turn first touched this path for RestoreCode semantics:
1666 // restoring one file = write earliest preimage (not all files from a turn).
1667 abs, err := safePath(s.root, rev.Path)
1668 if err != nil {
1669 return RewindResult{OK: false, Error: err.Error()}, err
1670 }
1671 fwd, _, _ := CapturePath(abs, CaptureOptions{WorkspaceRoot: s.root, ReadContent: true})
1672 tx := &TransactionManifest{
1673 SchemaVersion: SchemaV2,
1674 ID: newID("tx"),
1675 WorkspaceRoot: s.root,
1676 State: TxPrepared,
1677 Kind: "file_revert",
1678 Scope: RewindCode,
1679 Path: rev.Path,
1680 CreatedAt: time.Now(),
1681 UpdatedAt: time.Now(),
1682 }
1683 t := TransactionTarget{
1684 Path: rev.Path, AbsPath: abs,
1685 RestoreExisted: rev.Existed, RestoreMode: rev.Mode, RestoreSHA: rev.SHA256, RestoreBlob: rev.BlobRef, RestoreEncoding: rev.Encoding,
1686 ForwardExisted: fwd.Existed, ForwardMode: fwd.Mode, ForwardSHA: fwd.SHA256,
1687 }
1688 t.PublishTmp, t.BackupPath = transactionSiblingPaths(abs, tx.ID, 0)
1689 if rev.Existed {
1690 t.Action = "write"
1691 data, lerr := s.loadRevisionBytes(rev)
1692 if lerr != nil && rev.Content != nil {
1693 data, lerr = s.encodeRevisionContent(rev, abs)
1694 }
1695 if lerr != nil {
1696 return RewindResult{OK: false, Error: lerr.Error()}, lerr
1697 }
1698 mode := os.FileMode(0o644)
1699 if rev.Mode != 0 {
1700 mode = os.FileMode(rev.Mode)
1701 }
1702 if t.RestoreBlob == "" && s.blobs != nil {
1703 ref, perr := s.blobs.Put(data)
1704 if perr != nil {
1705 return RewindResult{OK: false, Error: perr.Error()}, perr
1706 }
1707 t.RestoreBlob = ref
1708 } else if t.RestoreBlob == "" {
1709 t.RestoreInline = clonePayload(data)
1710 }
1711 if err := s.writePublishTemp(t.PublishTmp, data, mode); err != nil {
1712 return RewindResult{OK: false, Error: err.Error()}, err
1713 }
1714 } else {
1715 t.Action = "delete"
1716 t.PublishTmp = ""
1717 }
1718 if fwd.Existed && s.blobs != nil {
1719 ref, perr := s.blobs.Put(fwd.Content)
1720 if perr != nil {
1721 return RewindResult{OK: false, Error: perr.Error()}, perr
1722 }
1723 t.ForwardBlob = ref
1724 } else if fwd.Existed {
1725 t.ForwardInline = clonePayload(fwd.Content)
1726 }
1727 tx.Targets = []TransactionTarget{t}
1728 if err := s.persistTransaction(tx); err != nil {
1729 return RewindResult{OK: false, Error: err.Error()}, err
1730 }
1731 return s.commitTransaction(tx, nil, nil)
1732 }
1733
1734 func sameFingerprint(a, b Fingerprint) bool {
1735 return a.Existed == b.Existed && a.IsDir == b.IsDir && a.IsSymlink == b.IsSymlink &&
1736 a.Nlink == b.Nlink && a.Mode == b.Mode && a.Size == b.Size && a.SHA256 == b.SHA256
1737 }
1738
1739 func clonePayload(data []byte) []byte {
1740 out := make([]byte, len(data))
1741 copy(out, data)
1742 return out
1743 }
1744
1744 lines GO