Skip to main content

vtcode_core/core/agent/snapshots/
native.rs

1use super::*;
2use std::collections::HashMap;
3use std::sync::{Arc, Mutex, OnceLock};
4
5#[derive(Clone)]
6struct Active {
7    workspace: PathBuf,
8    storage: PathBuf,
9    record: PathBuf,
10    engine: String,
11    watch: String,
12}
13static ACTIVE: OnceLock<Mutex<HashMap<String, Arc<Mutex<Active>>>>> = OnceLock::new();
14fn active_map() -> &'static Mutex<HashMap<String, Arc<Mutex<Active>>>> {
15    ACTIVE.get_or_init(Mutex::default)
16}
17
18/// Holds exclusive workspace access until the complete agent turn has settled.
19pub struct PromptCheckpointLease {
20    key: String,
21    _lock: fs::File,
22}
23impl Drop for PromptCheckpointLease {
24    fn drop(&mut self) {
25        if let Ok(mut active) = active_map().lock() {
26            active.remove(&self.key);
27        }
28    }
29}
30#[derive(Clone, Serialize, Deserialize)]
31struct Recovery {
32    #[serde(default)]
33    policy: String,
34    snapshot: String,
35    active: Vec<usize>,
36}
37#[derive(Clone, Default, Serialize, Deserialize)]
38struct Navigation {
39    active: Vec<usize>,
40    redo: Vec<Recovery>,
41    pending: Option<Recovery>,
42}
43
44fn atomic_json(path: &Path, value: &impl Serialize) -> Result<()> {
45    use std::io::Write;
46    let temp = path.with_extension(format!("{}.tmp", uuid::Uuid::new_v4()));
47    let mut options = fs::OpenOptions::new();
48    options.create_new(true).write(true);
49    #[cfg(unix)]
50    {
51        use std::os::unix::fs::OpenOptionsExt;
52        options.mode(0o600);
53    }
54    let mut file = options.open(&temp)?;
55    file.write_all(&serde_json::to_vec(value)?)?;
56    file.sync_all()?;
57    fs::rename(temp, path)?;
58    Ok(())
59}
60
61/// Whether a locked fd still refers to the file currently at `lock_path`.
62///
63/// The completion path unlinks `rewind.lock` while holding its lock. A process
64/// that opened the old inode just before that unlink can `try_lock` the ghost
65/// after the holder drops — while another process creates a fresh file — and
66/// both would believe they hold the lock. Comparing inos closes that window.
67#[cfg(unix)]
68fn locked_file_matches_path(file: &fs::File, lock_path: &Path) -> bool {
69    use std::os::unix::fs::MetadataExt;
70    let Ok(locked) = file.metadata() else {
71        return false;
72    };
73    let Ok(current) = fs::metadata(lock_path) else {
74        return false;
75    };
76    locked.ino() == current.ino()
77}
78
79#[cfg(not(unix))]
80fn locked_file_matches_path(_: &fs::File, _: &Path) -> bool {
81    true
82}
83
84/// Acquire the workspace rewind lock, verifying the locked inode is still the
85/// file at `lock_path`. A ghost acquisition (locked inode replaced under us)
86/// is dropped and retried against the current file, so mutual exclusion holds
87/// across the completion path's unlink-while-held cleanup.
88fn acquire_verified_rewind_lock(lock_path: &Path) -> std::io::Result<fs::File> {
89    for _ in 0..3 {
90        let file = fs::OpenOptions::new()
91            .read(true)
92            .write(true)
93            .create(true)
94            .truncate(false)
95            .open(lock_path)?;
96        file.try_lock()?;
97        if locked_file_matches_path(&file, lock_path) {
98            return Ok(file);
99        }
100        drop(file);
101    }
102    Err(std::io::Error::other("rewind.lock kept changing identity; retry the turn"))
103}
104
105fn file_records(manifest: &filesnap::Manifest, workspace: &Path, engine: &str) -> Result<Vec<FileSnapshot>> {
106    manifest
107        .entries
108        .keys()
109        .chain(manifest.absent.iter())
110        .map(|path| {
111            let relative = Path::new(path)
112                .strip_prefix(workspace)
113                .context("Checkpoint escaped the workspace")?;
114            Ok(FileSnapshot {
115                path: relative.to_string_lossy().replace('\\', "/"),
116                deleted: manifest.absent.contains(path),
117                encoding: Some(FileEncoding::Filesnap),
118                data: Some(engine.to_owned()),
119            })
120        })
121        .collect()
122}
123
124// Only literal POSIX-shell redirects whose cwd stays unchanged can provide
125// reliable preimages. Never execute or expand shell text to discover a path.
126fn canonicalize_for_strip(path: &Path) -> PathBuf {
127    if let Ok(canonical) = canonicalize(path) {
128        return canonical;
129    }
130    // The target may not exist yet (creation). Walk up to the nearest
131    // existing ancestor so symlinked workspaces (`/var` vs `/private/var`)
132    // still resolve inside the workspace instead of being skipped.
133    let mut ancestor = path.parent();
134    while let Some(dir) = ancestor {
135        if dir.as_os_str().is_empty() {
136            break;
137        }
138        if let Ok(canonical_parent) = canonicalize(dir)
139            && let Ok(stripped) = path.strip_prefix(dir)
140        {
141            let mut joined = canonical_parent;
142            joined.push(stripped);
143            return joined;
144        }
145        ancestor = dir.parent();
146    }
147    path.to_path_buf()
148}
149
150fn recovery_path(storage: &Path, snapshot: &str) -> PathBuf {
151    storage.join(format!("turn_recovery_{snapshot}.json"))
152}
153
154fn shell_redirect_paths(args: &serde_json::Value) -> BTreeSet<PathBuf> {
155    use crate::command_safety::shell_parser::contains_dynamic_shell_syntax;
156
157    let collect = || -> Option<BTreeSet<PathBuf>> {
158        if cfg!(windows) {
159            return None;
160        }
161        if let Some(shell) = args.get("shell").and_then(serde_json::Value::as_str) {
162            let name = Path::new(shell).file_name()?.to_str()?;
163            if !matches!(name, "bash" | "sh" | "dash" | "zsh" | "ksh") {
164                return None;
165            }
166        }
167        let script = args
168            .get("raw_command")
169            .and_then(serde_json::Value::as_str)
170            .map(str::to_owned)
171            .or_else(|| crate::tools::command_args::raw_command_text(args))?;
172        if !script.contains('>') {
173            return None;
174        }
175        let mut parser = tree_sitter::Parser::new();
176        parser.set_language(&tree_sitter_bash::LANGUAGE.into()).ok()?;
177        let tree = parser.parse(&script, None)?;
178        if tree.root_node().has_error() {
179            return None;
180        }
181        let mut paths = BTreeSet::new();
182        let mut pending = vec![tree.root_node()];
183        while let Some(node) = pending.pop() {
184            match node.kind() {
185                "program" | "list" | "pipeline" | "redirected_statement" => {
186                    let mut cursor = node.walk();
187                    pending.extend(node.named_children(&mut cursor));
188                }
189                "command" => {
190                    let name = node.child_by_field_name("name")?.utf8_text(script.as_bytes()).ok()?;
191                    if contains_dynamic_shell_syntax(name) {
192                        return None;
193                    }
194                    let words = shell_words::split(name).ok()?;
195                    if words.len() != 1
196                        || matches!(
197                            words.first()?.as_str(),
198                            "cd" | "pushd"
199                                | "popd"
200                                | "source"
201                                | "."
202                                | "eval"
203                                | "exec"
204                                | "command"
205                                | "builtin"
206                                | "alias"
207                                | "unalias"
208                                | "trap"
209                                | "enable"
210                                | "shopt"
211                        )
212                    {
213                        return None;
214                    }
215                    let mut cursor = node.walk();
216                    pending.extend(
217                        node.named_children(&mut cursor)
218                            .filter(|child| matches!(child.kind(), "file_redirect" | "heredoc_redirect")),
219                    );
220                }
221                "file_redirect" => {
222                    let text = node.utf8_text(script.as_bytes()).ok()?;
223                    let operator = text.trim_start_matches(|c: char| c.is_ascii_digit());
224                    if !operator.starts_with('>') && !operator.starts_with("&>") {
225                        continue;
226                    }
227                    let mut cursor = node.walk();
228                    for destination in node.children_by_field_name("destination", &mut cursor) {
229                        let raw = destination.utf8_text(script.as_bytes()).ok()?;
230                        if contains_dynamic_shell_syntax(raw) {
231                            return None;
232                        }
233                        let words = shell_words::split(raw).ok()?;
234                        if words.len() != 1 {
235                            return None;
236                        }
237                        let path = words.first()?;
238                        if operator.starts_with(">&") && (path == "-" || path.chars().all(|c| c.is_ascii_digit())) {
239                            continue;
240                        }
241                        if path.is_empty() || path.starts_with('~') {
242                            return None;
243                        }
244                        paths.insert(PathBuf::from(path));
245                    }
246                }
247                // Here-doc contents and comments are data, not shell commands.
248                "heredoc_redirect" | "comment" => {}
249                // Functions, loops, subshells and dynamic cwd changes need a
250                // richer execution model; do not guess their target paths.
251                _ => return None,
252            }
253        }
254        Some(paths)
255    };
256    collect().unwrap_or_default()
257}
258
259/// Capture final local write arguments after host hook rewriting and before tool execution.
260pub async fn declare_prompt_edit(session: String, name: String, args: serde_json::Value) -> Result<()> {
261    let active = active_map()
262        .lock()
263        .map_err(|error| anyhow::anyhow!("Checkpoint lock poisoned: {error}"))?
264        .get(&session)
265        .cloned();
266    let Some(active) = active else {
267        return Ok(());
268    };
269    let is_shell = crate::tools::tool_intent::is_command_run_tool_call(&name, &args);
270    let operation = args.get("action").and_then(serde_json::Value::as_str).unwrap_or(&name);
271    if !is_shell
272        && !["write", "edit", "patch", "delete", "remove", "create", "move", "rename"]
273            .iter()
274            .any(|verb| operation.contains(verb))
275    {
276        return Ok(());
277    }
278    tokio::task::spawn_blocking(move || -> Result<()> {
279        let active = active
280            .lock()
281            .map_err(|error| anyhow::anyhow!("Checkpoint lock poisoned: {error}"))?;
282        let mut paths = BTreeSet::new();
283        if is_shell {
284            let redirects = shell_redirect_paths(&args);
285            if redirects.is_empty() {
286                return Ok(());
287            }
288            let cwd = crate::tools::command_args::working_dir_text(&args)
289                .map_or_else(|| active.workspace.clone(), |path| active.workspace.join(path));
290            // A missing workdir must not abort capture; fall back so the tool
291            // still runs and other paths are still tracked.
292            let cwd = canonicalize(&cwd).unwrap_or(cwd);
293            paths.extend(redirects.into_iter().map(|path| cwd.join(path)));
294        } else {
295            for key in [
296                "path",
297                "file_path",
298                "destination",
299                "destination_path",
300                "new_path",
301                "source",
302            ] {
303                if let Some(path) = args.get(key).and_then(serde_json::Value::as_str) {
304                    paths.insert(PathBuf::from(path));
305                }
306            }
307            for key in ["patch", "input", "patch_text"] {
308                if let Some(patch) = args.get(key).and_then(serde_json::Value::as_str) {
309                    for line in patch.lines() {
310                        for prefix in [
311                            "*** Add File: ",
312                            "*** Update File: ",
313                            "*** Delete File: ",
314                            "*** Move to: ",
315                        ] {
316                            if let Some(path) = line.strip_prefix(prefix) {
317                                paths.insert(PathBuf::from(path));
318                            }
319                        }
320                    }
321                }
322            }
323        }
324        if paths.is_empty() {
325            return Ok(());
326        }
327        let store = filesnap::WorkspaceStore::open(&active.storage, &active.workspace)?;
328        let target = store.target_for_turn(&active.engine)?.context("Missing active checkpoint")?;
329        let before = store.manifest(target.manifest_id())?;
330        let ignore = filesnap::load_ignore(&active.workspace);
331        for path in paths {
332            if path.components().any(|part| part == Component::ParentDir) {
333                continue;
334            }
335            // Canonicalize absolute targets the same way `normalize_path` does
336            // so macOS `/var` vs `/private/var` and symlinked workspaces do not
337            // fail checkpointing. Fall back to the raw path for not-yet-created
338            // files. Anything still outside the workspace is skipped, never fatal.
339            let relative = if path.is_absolute() {
340                let canonical = canonicalize_for_strip(&path);
341                if let Ok(relative) = canonical.strip_prefix(&active.workspace) {
342                    relative.to_path_buf()
343                } else if let Ok(relative) = path.strip_prefix(&active.workspace) {
344                    relative.to_path_buf()
345                } else {
346                    continue;
347                }
348            } else {
349                path
350            };
351            let path = match SnapshotManager::checked_file_path(&active.workspace, &active.storage, &relative) {
352                Ok(path) => path,
353                Err(error) => {
354                    tracing::debug!(%error, "Skipping checkpoint pre-image outside restore authority");
355                    continue;
356                }
357            };
358            if filesnap::is_ignored(&ignore, &path) {
359                continue;
360            }
361            if let Err(error) =
362                store.declare_paths(&active.watch, &active.engine, std::slice::from_ref(&path))
363            {
364                tracing::debug!(%error, "Skipping checkpoint watch declaration");
365                continue;
366            }
367            let key = path.to_string_lossy();
368            // The first preimage owns the turn boundary. A second write must
369            // never replace a recorded absence with the newly created bytes.
370            if before.entries.contains_key(key.as_ref()) || before.absent.contains(key.as_ref()) {
371                continue;
372            }
373            let image = match fs::read(&path) {
374                Ok(bytes) => filesnap::PreEditImage::Existed(bytes),
375                Err(error) if error.kind() == std::io::ErrorKind::NotFound => filesnap::PreEditImage::DidNotExist,
376                Err(error) => {
377                    // Directories and other non-file targets have no pre-image;
378                    // skip them rather than failing the turn.
379                    tracing::debug!(%error, path = %path.display(), "Skipping checkpoint pre-image for non-file target");
380                    continue;
381                }
382            };
383            if let Err(error) = filesnap::declare_edits(
384                &store,
385                &active.engine,
386                &active.engine,
387                &filesnap::TurnScope::at(&active.workspace),
388                vec![(path, image)],
389            ) {
390                tracing::debug!(%error, "Skipping checkpoint pre-image declaration");
391                continue;
392            }
393        }
394        let target = store.target_for_turn(&active.engine)?.context("Missing active checkpoint")?;
395        let manifest = store.manifest(target.manifest_id())?;
396        let mut stored: StoredSnapshot = serde_json::from_slice(&fs::read(&active.record)?)?;
397        stored.files = file_records(&manifest, &active.workspace, &active.engine)?;
398        stored.metadata.file_count = stored.files.len();
399        atomic_json(&active.record, &stored)
400    })
401    .await?
402}
403
404impl SnapshotManager {
405    async fn retire_recovery_record(&self, snapshot: &str) {
406        if uuid::Uuid::parse_str(snapshot).is_err() {
407            return;
408        }
409        let record = recovery_path(&self.storage_dir, snapshot);
410        match self.retire_snapshot(&record).await {
411            Ok(()) => {}
412            Err(error) => {
413                let is_missing = error
414                    .downcast_ref::<std::io::Error>()
415                    .is_some_and(|io| io.kind() == std::io::ErrorKind::NotFound);
416                if !is_missing {
417                    tracing::warn!(%error, "Failed to retire consumed recovery record");
418                }
419            }
420        }
421    }
422
423    fn navigation_path(&self, session: &str) -> Result<PathBuf> {
424        anyhow::ensure!(session.len() <= 80 && !session.is_empty(), "Invalid session ID");
425        let key: String = session.as_bytes().iter().map(|b| format!("{b:02x}")).collect();
426        Ok(self.storage_dir.join(format!("branch_{key}.json")))
427    }
428    fn navigation(&self, session: &str) -> Result<Navigation> {
429        match fs::read(self.navigation_path(session)?) {
430            Ok(bytes) => Ok(serde_json::from_slice(&bytes)?),
431            Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(Navigation::default()),
432            Err(error) => Err(error.into()),
433        }
434    }
435
436    /// Bound a finished session's navigation record so it cannot pin every turn
437    /// snapshot for the full retention window. Keeps the newest
438    /// [`REWIND_ACTIVE_KEEP`] active turns for a possible resume, drops the
439    /// rest (those `turn_*.json` files become prune-eligible), clears the redo
440    /// stack, and releases the workspace rewind lock.
441    pub async fn complete_session_navigation(&self, session: &str) -> Result<()> {
442        if !self.enabled {
443            return Ok(());
444        }
445        let nav_path = self.navigation_path(session)?;
446        let mut state = self.navigation(session)?;
447        let had_record = nav_path.exists();
448        if state.active.len() > REWIND_ACTIVE_KEEP {
449            let drop_count = state.active.len() - REWIND_ACTIVE_KEEP;
450            state.active.drain(..drop_count);
451        }
452        for entry in std::mem::take(&mut state.redo) {
453            self.retire_recovery_record(&entry.snapshot).await;
454        }
455        // Avoid inventing an empty branch file for checkpoint-less sessions.
456        if had_record || !state.active.is_empty() || state.pending.is_some() {
457            atomic_json(&nav_path, &state)?;
458        }
459        let lock_path = self.storage_dir.join("rewind.lock");
460        // Only remove the lock while we hold it. Deleting an flocked file that
461        // another process holds would let later `try_lock` calls create a fresh
462        // inode and break mutual exclusion.
463        match fs::OpenOptions::new().read(true).write(true).open(&lock_path) {
464            Ok(file) => {
465                // Only unlink while holding a lock on the file that is
466                // currently at the path; a ghost lock (path replaced under
467                // us) must not delete the replacement's lock file.
468                if file.try_lock().is_ok() && locked_file_matches_path(&file, &lock_path) {
469                    if let Err(error) = fs::remove_file(&lock_path)
470                        && error.kind() != std::io::ErrorKind::NotFound
471                    {
472                        tracing::debug!(%error, "failed to remove rewind.lock after session completion");
473                    }
474                }
475            }
476            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
477            Err(error) => {
478                tracing::debug!(%error, "failed to open rewind.lock after session completion");
479            }
480        }
481        // Newly unpinned turns can be reclaimed now rather than waiting for the
482        // next prompt's budget prune.
483        if let Err(error) = self.prune_snapshot_budget().await {
484            tracing::debug!(%error, "checkpoint budget prune failed after session completion");
485        }
486        Ok(())
487    }
488
489    fn build_ignore(&self, policy: &str) -> Result<filesnap::Gitignore> {
490        anyhow::ensure!(policy.len() <= 1024 * 1024, "Ignore policy is too large");
491        let mut builder = filesnap::GitignoreBuilder::new(&self.canonical_workspace);
492        for line in policy.lines() {
493            builder.add_line(None, line)?;
494        }
495        Ok(builder.build()?)
496    }
497    /// Save the exact conversation prefix and workspace before sending a prompt.
498    pub async fn begin_prompt(
499        &self,
500        turn: usize,
501        session: &str,
502        prompt: &str,
503        conversation: &[SessionMessage],
504    ) -> Result<PromptCheckpointLease> {
505        self.begin_prompt_with_cancellation(
506            turn,
507            session,
508            prompt,
509            conversation,
510            tokio_util::sync::CancellationToken::new(),
511        )
512        .await?
513        .context("checkpoint preparation cancelled")
514    }
515
516    /// Prepare a checkpoint while allowing the caller to stop waiting. The
517    /// blocking worker retains the rewind lock through atomic publication and
518    /// drops its lease on cancellation; no tool may run without that lease.
519    pub async fn begin_prompt_with_cancellation(
520        &self,
521        turn: usize,
522        session: &str,
523        prompt: &str,
524        conversation: &[SessionMessage],
525        cancellation: tokio_util::sync::CancellationToken,
526    ) -> Result<Option<PromptCheckpointLease>> {
527        let manager = self.clone();
528        let session = session.to_owned();
529        let prompt = prompt.to_owned();
530        let conversation = conversation.to_vec();
531        // Lock acquisition, scanning, serialization and durable publication
532        // all belong on the blocking worker. Its lease retains the same lock
533        // through pruning and the complete agent turn, including on cancellation.
534        let worker_cancel = cancellation.clone();
535        let worker = tokio::task::spawn_blocking(move || -> Result<Option<PromptCheckpointLease>> {
536            if worker_cancel.is_cancelled() {
537                return Ok(None);
538            }
539            let lease = manager.begin_prompt_blocking(turn, &session, &prompt, &conversation)?;
540            if worker_cancel.is_cancelled() {
541                Ok(None)
542            } else {
543                Ok(Some(lease))
544            }
545        });
546        // Cancel the worker's publication ownership even if this future is
547        // dropped by a caller racing a stop request.
548        let cancel_on_drop = cancellation.clone().drop_guard();
549        let lease = tokio::select! {
550            biased;
551            _ = cancellation.cancelled() => return Ok(None),
552            result = worker => result.context("checkpoint preparation worker failed")??,
553        };
554        if lease.is_some() {
555            tokio::select! {
556                biased;
557                _ = cancellation.cancelled() => return Ok(None),
558                result = self.prune_snapshot_budget() => {
559                    if let Err(error) = result { tracing::debug!(%error, "checkpoint budget prune failed"); }
560                }
561            }
562        }
563        let _ = cancel_on_drop.disarm();
564        Ok(lease)
565    }
566
567    fn begin_prompt_blocking(
568        &self,
569        turn: usize,
570        session: &str,
571        prompt: &str,
572        conversation: &[SessionMessage],
573    ) -> Result<PromptCheckpointLease> {
574        let lock = acquire_verified_rewind_lock(&self.storage_dir.join("rewind.lock"))
575            .context("Another turn or rewind is using this workspace")?;
576        let mut state = self.navigation(session)?;
577        anyhow::ensure!(state.pending.is_none(), "Interrupted rewind; run /rewind-recover before continuing");
578        let workspace = self.canonical_workspace.clone();
579        let storage = self.storage_dir.clone();
580        let engine = format!("vt-{}", uuid::Uuid::new_v4());
581        let watch = format!("vt-watch-{session}");
582        let active = Active {
583            workspace: workspace.clone(),
584            storage: storage.clone(),
585            record: self.snapshot_path(turn),
586            engine: engine.clone(),
587            watch: watch.clone(),
588        };
589        let files = {
590            let store = filesnap::WorkspaceStore::open(&storage, &workspace)?;
591            store.note_turn(&watch, &engine)?;
592            let watched: Vec<_> = store
593                .declared_paths(&watch, filesnap::DeclaredWindow::default())?
594                .into_iter()
595                .collect();
596            store.declare_paths(&engine, &engine, &watched)?;
597            let checkpoint = filesnap::capture_turn(&store, &engine, &engine, &filesnap::TurnScope::at(&workspace))?;
598            anyhow::ensure!(checkpoint.stats.dropped == 0, "Checkpoint skipped paths; cannot safely start turn");
599            let files = file_records(&checkpoint.manifest, &workspace, &engine)?;
600            for file in &files {
601                SnapshotManager::checked_file_path(&workspace, &storage, Path::new(&file.path))?;
602            }
603            files
604        };
605        let metadata = SnapshotMetadata {
606            id: format!("turn_{turn}"),
607            turn_number: turn,
608            created_at: Self::current_timestamp()?,
609            description: Self::truncate_description(prompt),
610            message_count: conversation.len(),
611            file_count: files.len(),
612            touched_files: vec![],
613            prompt_text: Some(prompt.into()),
614            prompt_message_index: None,
615            session_id: Some(session.into()),
616            runtime_turn_id: None,
617            session_turn_number: Some(turn),
618            turn_diagnostics: None,
619        };
620        atomic_json(
621            &active.record,
622            &StoredSnapshot {
623                metadata,
624                conversation: conversation.to_vec(),
625                files,
626                schema_version: Some(SNAPSHOT_SCHEMA_VERSION),
627            },
628        )?;
629        state.active.push(turn);
630        for entry in std::mem::take(&mut state.redo) {
631            if uuid::Uuid::parse_str(&entry.snapshot).is_ok() {
632                let record = recovery_path(&self.storage_dir, &entry.snapshot);
633                let retired = record.with_extension(format!("retired-{}", uuid::Uuid::new_v4()));
634                if let Err(error) = fs::rename(&record, retired)
635                    && error.kind() != std::io::ErrorKind::NotFound
636                {
637                    tracing::warn!(%error, "Failed to retire discarded recovery record");
638                }
639            }
640        }
641        atomic_json(&self.navigation_path(session)?, &state)?;
642        active_map()
643            .lock()
644            .map_err(|error| anyhow::anyhow!("Checkpoint lock poisoned: {error}"))?
645            .insert(session.into(), Arc::new(Mutex::new(active)));
646        Ok(PromptCheckpointLease { key: session.into(), _lock: lock })
647    }
648    /// Enumerate only checkpoints on this session's current conversation branch.
649    pub async fn rewind_points(&self, session: &str) -> Result<Vec<SnapshotMetadata>> {
650        let state = self.navigation(session)?;
651        let mut points = Vec::new();
652        for turn in state.active.iter().rev() {
653            if let Some(stored) = self.load_snapshot(*turn).await? {
654                points.push(stored.metadata);
655            }
656        }
657        Ok(points)
658    }
659    /// Restore files and return their matching conversation, with persistent multi-level redo.
660    pub async fn navigate_prompt(
661        &self,
662        turn: Option<usize>,
663        scope: RevertScope,
664        session: &str,
665        conversation: &[SessionMessage],
666    ) -> Result<CheckpointRestore> {
667        let _lock = acquire_verified_rewind_lock(&self.storage_dir.join("rewind.lock"))
668            .context("Wait for the current turn to finish")?;
669        let policy = match fs::read_to_string(self.canonical_workspace.join(".filesnapignore")) {
670            Ok(text) => text,
671            Err(error) if error.kind() == std::io::ErrorKind::NotFound => String::new(),
672            Err(error) => return Err(error.into()),
673        };
674        let ignore = self.build_ignore(&policy)?;
675        let mut state = self.navigation(session)?;
676        let load_recovery = |saved: &Recovery| -> Result<StoredSnapshot> {
677            uuid::Uuid::parse_str(&saved.snapshot)?;
678            Ok(serde_json::from_slice(&fs::read(recovery_path(&self.storage_dir, &saved.snapshot))?)?)
679        };
680        if let Some(pending) = state.pending.clone() {
681            anyhow::ensure!(turn.is_none(), "Interrupted rewind; run /rewind-recover");
682            let restored = self
683                .restore_stored_snapshot_with_ignore(
684                    load_recovery(&pending)?,
685                    RevertScope::Both,
686                    &self.build_ignore(&pending.policy)?,
687                )
688                .await?;
689            state.pending = None;
690            atomic_json(&self.navigation_path(session)?, &state)?;
691            self.retire_recovery_record(&pending.snapshot).await;
692            return Ok(restored);
693        }
694        let (targets, next_active) = if let Some(turn) = turn {
695            let index = state
696                .active
697                .iter()
698                .position(|id| *id == turn)
699                .context("Checkpoint is not on this branch")?;
700            let mut targets = Vec::new();
701            for id in state.active[index..].iter().rev() {
702                targets.push(self.load_snapshot(*id).await?.context("Checkpoint expired")?);
703            }
704            (targets, state.active[..index].to_vec())
705        } else {
706            let saved = state.redo.last().context("Nothing to redo; new prompts clear redo")?;
707            (vec![load_recovery(saved)?], saved.active.clone())
708        };
709        let destination = targets.last().context("No restore target")?;
710        let mut rescue = destination.clone();
711        rescue.conversation = conversation.to_vec();
712        rescue.metadata.prompt_text = None;
713        rescue.metadata.prompt_message_index = None;
714        rescue.metadata.message_count = conversation.len();
715        let paths: BTreeSet<String> = targets
716            .iter()
717            .flat_map(|target| target.files.iter().map(|f| f.path.clone()))
718            .collect();
719        let workspace = self.canonical_workspace.clone();
720        let storage = self.storage_dir.clone();
721        rescue.files = tokio::task::spawn_blocking(move || -> Result<Vec<FileSnapshot>> {
722            let paths = paths
723                .iter()
724                .map(|p| Self::checked_file_path(&workspace, &storage, Path::new(p)))
725                .collect::<Result<Vec<_>>>()?;
726            let engine = format!("vt-{}", uuid::Uuid::new_v4());
727            let store = filesnap::WorkspaceStore::open(&storage, &workspace)?;
728            let checkpoint = store.checkpoint(&engine, &engine, paths)?;
729            anyhow::ensure!(checkpoint.stats.dropped == 0, "Could not save recovery files; rewind cancelled");
730            file_records(&checkpoint.manifest, &workspace, &engine)
731        })
732        .await??;
733        let saved = Recovery {
734            policy,
735            snapshot: uuid::Uuid::new_v4().to_string(),
736            active: state.active.clone(),
737        };
738        atomic_json(&recovery_path(&self.storage_dir, &saved.snapshot), &rescue)?;
739        state.pending = Some(saved.clone());
740        atomic_json(&self.navigation_path(session)?, &state)?;
741        let destination = CheckpointRestore {
742            metadata: destination.metadata.clone(),
743            conversation: destination.conversation.clone(),
744        };
745        for target in targets {
746            if let Err(error) = self.restore_stored_snapshot_with_ignore(target, scope, &ignore).await {
747                self.restore_stored_snapshot_with_ignore(rescue, RevertScope::Both, &ignore)
748                    .await
749                    .context(format!("Rewind failed ({error}); recovery failed; run /rewind-recover"))?;
750                let failed = saved.snapshot.clone();
751                state.pending = None;
752                atomic_json(&self.navigation_path(session)?, &state)?;
753                self.retire_recovery_record(&failed).await;
754                return Err(error.context("Rewind failed; original files recovered"));
755            }
756        }
757        state.pending = None;
758        state.active = next_active;
759        if turn.is_some() {
760            state.redo.push(saved);
761        } else {
762            // Redo writes a fresh rescue record for crash recovery before
763            // restoring. On success neither the consumed redo entry nor the
764            // transient rescue is needed, so retire both instead of leaking.
765            self.retire_recovery_record(&saved.snapshot).await;
766            if let Some(used) = state.redo.pop() {
767                self.retire_recovery_record(&used.snapshot).await;
768            }
769        }
770        atomic_json(&self.navigation_path(session)?, &state)?;
771        Ok(destination)
772    }
773
774    /// Complete an interrupted rewind only. Unlike [`Self::navigate_prompt`],
775    /// this never touches the redo stack, so `/rewind-recover` cannot overwrite
776    /// edits made after a completed rewind when there is nothing to recover.
777    pub async fn recover_pending_rewind(&self, session: &str) -> Result<CheckpointRestore> {
778        let _lock = acquire_verified_rewind_lock(&self.storage_dir.join("rewind.lock"))
779            .context("Wait for the current turn to finish")?;
780        // Fail closed when no recovery is pending; never fall back to redo.
781        let state = self.navigation(session)?;
782        let pending = state.pending.clone().context("No interrupted rewind to recover")?;
783        let stored: StoredSnapshot = {
784            uuid::Uuid::parse_str(&pending.snapshot)?;
785            serde_json::from_slice(&fs::read(recovery_path(&self.storage_dir, &pending.snapshot))?)?
786        };
787        let policy = self.build_ignore(&pending.policy)?;
788        let restored = self
789            .restore_stored_snapshot_with_ignore(stored, RevertScope::Both, &policy)
790            .await?;
791        let mut state = self.navigation(session)?;
792        // Only clear if the same pending is still present; a concurrent rewind
793        // must not lose its recovery record.
794        if state
795            .pending
796            .as_ref()
797            .is_some_and(|current| current.snapshot == pending.snapshot)
798        {
799            state.pending = None;
800            atomic_json(&self.navigation_path(session)?, &state)?;
801            self.retire_recovery_record(&pending.snapshot).await;
802        }
803        Ok(restored)
804    }
805}
806
807#[cfg(test)]
808mod tests {
809    use super::*;
810    use crate::llm::provider::MessageRole;
811    use tempfile::TempDir;
812
813    #[test]
814    #[cfg(unix)]
815    fn rewind_lock_acquire_verifies_and_blocks_concurrent_holders() {
816        let temp = TempDir::new().expect("tempdir");
817        let lock_path = temp.path().join("rewind.lock");
818
819        let first = acquire_verified_rewind_lock(&lock_path).expect("first acquire");
820        assert!(
821            locked_file_matches_path(&first, &lock_path),
822            "a fresh acquisition must cover the current path inode"
823        );
824
825        // A live holder blocks other acquirers (WouldBlock), preserving the
826        // pre-existing mutual-exclusion behavior.
827        let second = acquire_verified_rewind_lock(&lock_path);
828        let error = second.expect_err("held lock must block");
829        assert_eq!(error.kind(), std::io::ErrorKind::WouldBlock);
830
831        drop(first);
832        let third = acquire_verified_rewind_lock(&lock_path).expect("re-acquire after release");
833        assert!(locked_file_matches_path(&third, &lock_path));
834    }
835
836    #[test]
837    #[cfg(unix)]
838    fn rewind_lock_detects_replaced_path_inode() {
839        let temp = TempDir::new().expect("tempdir");
840        let lock_path = temp.path().join("rewind.lock");
841        fs::write(&lock_path, b"").expect("create lock file");
842        let ghost = fs::OpenOptions::new()
843            .read(true)
844            .write(true)
845            .open(&lock_path)
846            .expect("open old inode");
847
848        // The completion path's unlink-while-held replaced the file: the old
849        // fd is now a ghost and must not be treated as covering the path.
850        fs::remove_file(&lock_path).expect("unlink");
851        fs::write(&lock_path, b"").expect("recreate lock file");
852        assert!(
853            !locked_file_matches_path(&ghost, &lock_path),
854            "a ghost fd must not verify against the replaced path"
855        );
856
857        // The acquire helper never accepts a ghost: it retries onto the
858        // current file even while the ghost is still open.
859        let verified = acquire_verified_rewind_lock(&lock_path).expect("acquire onto current inode");
860        assert!(locked_file_matches_path(&verified, &lock_path));
861    }
862
863    #[test]
864    #[cfg(unix)]
865    fn shell_redirects_only_capture_literal_targets_with_a_known_cwd() {
866        let paths = |script: &str| shell_redirect_paths(&serde_json::json!({"cmd": script}));
867        assert_eq!(paths("printf 'hello' > hello.py && cat hello.py"), BTreeSet::from([PathBuf::from("hello.py")]));
868        assert_eq!(
869            paths("cat > 'hello world.py' <<'PY'\nprint('> not-a-path')\nPY"),
870            BTreeSet::from([PathBuf::from("hello world.py")])
871        );
872        assert_eq!(paths("printf hi >> log.txt 2>&1"), BTreeSet::from([PathBuf::from("log.txt")]));
873        for script in [
874            "cat hello.py",
875            "cat < input.txt",
876            "printf '> innocent.txt'",
877            "echo hi 2>&1",
878            "echo hi > $OUTPUT",
879            "echo hi > $(pwd)/hello.py",
880            "echo hi > ~/hello.py",
881            "cd nested && echo hi > hello.py",
882            "command cd nested; echo hi > hello.py",
883            "eval 'cd nested'; echo hi > hello.py",
884            "f() { echo hi > hello.py; }; f",
885            "echo hi >",
886        ] {
887            assert!(paths(script).is_empty(), "must not guess paths for {script}");
888        }
889    }
890
891    #[tokio::test]
892    async fn cancelled_prompt_preparation_never_admits_a_lease_and_a_fresh_prompt_can_start() -> Result<()> {
893        let dir = TempDir::new()?;
894        let manager = SnapshotManager::new(SnapshotConfig::new(dir.path().into()))?;
895        let session = uuid::Uuid::new_v4().to_string();
896        let cancellation = tokio_util::sync::CancellationToken::new();
897        cancellation.cancel();
898        assert!(
899            manager
900                .begin_prompt_with_cancellation(1, &session, "cancelled", &[], cancellation)
901                .await?
902                .is_none()
903        );
904        assert!(!active_map().lock().unwrap().contains_key(&session));
905        let lease = manager.begin_prompt(2, &session, "fresh", &[]).await?;
906        assert!(active_map().lock().unwrap().contains_key(&session));
907        assert!(manager.begin_prompt(3, &session, "must not overlap", &[]).await.is_err());
908        drop(lease);
909        assert!(!active_map().lock().unwrap().contains_key(&session));
910        let _next_lease = manager.begin_prompt(3, &session, "after release", &[]).await?;
911        Ok(())
912    }
913
914    #[tokio::test]
915    #[cfg(unix)]
916    async fn shell_preimages_respect_ignore_and_workspace_boundaries() -> Result<()> {
917        let dir = TempDir::new()?;
918        let outside = TempDir::new()?;
919        fs::write(dir.path().join(".filesnapignore"), "ignored.txt\n")?;
920        let manager = SnapshotManager::new(SnapshotConfig::new(dir.path().into()))?;
921        let session = uuid::Uuid::new_v4().to_string();
922        let lease = manager.begin_prompt(1, &session, "create files", &[]).await?;
923        let outside_file = outside.path().join("outside.txt");
924        let script =
925            format!("echo hi > ignored.txt; echo hi > {}", shell_words::quote(&outside_file.to_string_lossy()));
926        declare_prompt_edit(session.clone(), "exec_command".into(), serde_json::json!({"cmd":script})).await?;
927        fs::write(dir.path().join("ignored.txt"), "ignored")?;
928        fs::write(&outside_file, "outside")?;
929        // No input or display text should be mistaken for a shell write.
930        declare_prompt_edit(session.clone(), "exec_command".into(), serde_json::json!({"cmd":"cat ignored.txt"}))
931            .await?;
932        assert!(manager.load_snapshot(1).await?.expect("checkpoint").files.is_empty());
933        drop(lease);
934        manager.navigate_prompt(Some(1), RevertScope::Both, &session, &[]).await?;
935        assert_eq!(fs::read_to_string(dir.path().join("ignored.txt"))?, "ignored");
936        assert_eq!(fs::read_to_string(outside_file)?, "outside");
937        let _lease = manager.begin_prompt(2, &session, "next prompt", &[]).await?;
938        assert!(manager.load_snapshot(2).await?.expect("next checkpoint").files.is_empty());
939        Ok(())
940    }
941
942    #[tokio::test]
943    #[cfg(unix)]
944    async fn shell_creation_rewind_redo_then_rewind_creation_removes_file() -> Result<()> {
945        let dir = TempDir::new()?;
946        let manager = SnapshotManager::new(SnapshotConfig::new(dir.path().into()))?;
947        let session = uuid::Uuid::new_v4().to_string();
948        let nested = dir.path().join("nested");
949        fs::create_dir(&nested)?;
950        let file = nested.join("hello.py");
951        let neighbor = dir.path().join("hello.py");
952        fs::write(&neighbor, "unrelated existing file")?;
953        let original = vec![SessionMessage::new(MessageRole::User, "pwd")];
954        let created = vec![SessionMessage::new(MessageRole::User, "add a hello world py")];
955        let current = vec![SessionMessage::new(MessageRole::User, "add greetings")];
956        let hello = "print(\"Hello, World!\")\n";
957        let greetings = "print(\"Hello, Ada!\")\n";
958
959        let lease = manager.begin_prompt(1, &session, "add a hello world py", &original).await?;
960        let script = "printf 'print(\"Hello, World!\")\\n' > hello.py && cat hello.py";
961        declare_prompt_edit(
962            session.clone(),
963            "exec_command".into(),
964            serde_json::json!({"cmd":script,"workdir":"nested"}),
965        )
966        .await?;
967        let output = tokio::process::Command::new("sh")
968            .args(["-c", script])
969            .current_dir(&nested)
970            .output()
971            .await?;
972        assert!(output.status.success());
973        assert_eq!(fs::read_to_string(&file)?, hello);
974        // A subsequent write in the same prompt must retain the first absence.
975        declare_prompt_edit(
976            session.clone(),
977            "exec_command".into(),
978            serde_json::json!({"cmd":script,"workdir":"nested"}),
979        )
980        .await?;
981        assert!(
982            manager
983                .load_snapshot(1)
984                .await?
985                .expect("creation checkpoint")
986                .files
987                .iter()
988                .any(|f| f.path == "nested/hello.py" && f.deleted)
989        );
990        drop(lease);
991
992        let lease = manager.begin_prompt(2, &session, "add greetings", &created).await?;
993        let script = "cat > hello.py <<'PY'\nprint(\"Hello, Ada!\")\nPY";
994        declare_prompt_edit(
995            session.clone(),
996            "unified_exec".into(),
997            serde_json::json!({"action":"run","command":script,"working_dir":nested}),
998        )
999        .await?;
1000        let output = tokio::process::Command::new("sh")
1001            .args(["-c", script])
1002            .current_dir(&nested)
1003            .output()
1004            .await?;
1005        assert!(output.status.success());
1006        assert_eq!(fs::read_to_string(&file)?, greetings);
1007        drop(lease);
1008        // A file never declared or captured must not be swept away by rewind.
1009        fs::write(dir.path().join("unrelated.txt"), "leave me alone")?;
1010
1011        let restored = manager.navigate_prompt(Some(2), RevertScope::Both, &session, &current).await?;
1012        assert_eq!(restored.conversation, created);
1013        assert_eq!(fs::read_to_string(&file)?, hello);
1014        let restored = manager
1015            .navigate_prompt(None, RevertScope::Both, &session, &restored.conversation)
1016            .await?;
1017        assert_eq!(restored.conversation, current);
1018        assert_eq!(fs::read_to_string(&file)?, greetings);
1019        let restored = manager
1020            .navigate_prompt(Some(1), RevertScope::Both, &session, &restored.conversation)
1021            .await?;
1022        assert_eq!(restored.conversation, original);
1023        assert!(!file.exists());
1024        assert_eq!(fs::read_to_string(&neighbor)?, "unrelated existing file");
1025        assert_eq!(fs::read_to_string(dir.path().join("unrelated.txt"))?, "leave me alone");
1026        let resumed = SnapshotManager::new(SnapshotConfig::new(dir.path().into()))?;
1027        let restored = resumed
1028            .navigate_prompt(None, RevertScope::Both, &session, &restored.conversation)
1029            .await?;
1030        assert_eq!(restored.conversation, current);
1031        assert_eq!(fs::read_to_string(&file)?, greetings);
1032        Ok(())
1033    }
1034
1035    #[tokio::test]
1036    async fn native_history_restores_prompt_prefix_binary_assets_and_nested_redo() -> Result<()> {
1037        let dir = TempDir::new()?;
1038        let manager = SnapshotManager::new(SnapshotConfig::new(dir.path().into()))?;
1039        let session = uuid::Uuid::new_v4().to_string();
1040        let asset = dir.path().join("asset.bin");
1041        fs::write(&asset, [0, 255, 7])?;
1042        let original = vec![SessionMessage::new(MessageRole::User, "original")];
1043        let lease = manager.begin_prompt(1, &session, "first", &original).await?;
1044        fs::write(&asset, "first result")?;
1045        drop(lease);
1046        let second = vec![SessionMessage::new(MessageRole::User, "second")];
1047        let lease = manager.begin_prompt(2, &session, "second", &second).await?;
1048        declare_prompt_edit(session.clone(), "write_file".into(), serde_json::json!({"path":".created"})).await?;
1049        fs::write(dir.path().join(".created"), "second result")?;
1050        drop(lease);
1051        let current = vec![SessionMessage::new(MessageRole::User, "current")];
1052        let restored = manager.navigate_prompt(Some(2), RevertScope::Both, &session, &current).await?;
1053        assert_eq!(restored.conversation, second);
1054        assert!(!dir.path().join(".created").exists());
1055        let restored = manager
1056            .navigate_prompt(Some(1), RevertScope::Both, &session, &restored.conversation)
1057            .await?;
1058        assert_eq!(restored.conversation, original);
1059        assert_eq!(fs::read(&asset)?, vec![0, 255, 7]);
1060        let restored = manager
1061            .navigate_prompt(None, RevertScope::Both, &session, &restored.conversation)
1062            .await?;
1063        assert_eq!(restored.conversation, second);
1064        let restored = manager
1065            .navigate_prompt(None, RevertScope::Both, &session, &restored.conversation)
1066            .await?;
1067        assert_eq!(restored.conversation, current);
1068        assert_eq!(fs::read_to_string(dir.path().join(".created"))?, "second result");
1069        assert!(
1070            manager
1071                .navigate_prompt(None, RevertScope::Both, &session, &current)
1072                .await
1073                .is_err()
1074        );
1075        // One jump across both prompts must also reverse later file creation.
1076        let restored = manager.navigate_prompt(Some(1), RevertScope::Both, &session, &current).await?;
1077        assert_eq!(restored.conversation, original);
1078        assert!(!dir.path().join(".created").exists());
1079        let resumed = SnapshotManager::new(SnapshotConfig::new(dir.path().into()))?;
1080        let restored = resumed.navigate_prompt(None, RevertScope::Both, &session, &original).await?;
1081        assert_eq!(restored.conversation, current);
1082        Ok(())
1083    }
1084
1085    fn live_recovery_files(workspace: &Path) -> Vec<PathBuf> {
1086        let storage = workspace.join(".vtcode").join("checkpoints");
1087        fs::read_dir(&storage)
1088            .map(|entries| {
1089                entries
1090                    .filter_map(|entry| entry.ok().map(|entry| entry.path()))
1091                    .filter(|path| {
1092                        path.file_name()
1093                            .and_then(|name| name.to_str())
1094                            .is_some_and(|name| name.starts_with("turn_recovery_") && name.ends_with(".json"))
1095                    })
1096                    .collect()
1097            })
1098            .unwrap_or_default()
1099    }
1100
1101    #[tokio::test]
1102    async fn consumed_recovery_records_are_retired_not_leaked() -> Result<()> {
1103        let dir = TempDir::new()?;
1104        let manager = SnapshotManager::new(SnapshotConfig::new(dir.path().into()))?;
1105        let session = uuid::Uuid::new_v4().to_string();
1106        let original = vec![SessionMessage::new(MessageRole::User, "original")];
1107        let lease = manager.begin_prompt(1, &session, "first", &original).await?;
1108        drop(lease);
1109        let second = vec![SessionMessage::new(MessageRole::User, "second")];
1110        let lease = manager.begin_prompt(2, &session, "second", &second).await?;
1111        drop(lease);
1112        let current = vec![SessionMessage::new(MessageRole::User, "current")];
1113
1114        assert!(live_recovery_files(dir.path()).is_empty());
1115        let restored = manager.navigate_prompt(Some(2), RevertScope::Both, &session, &current).await?;
1116        assert_eq!(live_recovery_files(dir.path()).len(), 1);
1117
1118        let restored = manager
1119            .navigate_prompt(None, RevertScope::Both, &session, &restored.conversation)
1120            .await?;
1121        assert_eq!(restored.conversation, current);
1122        assert!(live_recovery_files(dir.path()).is_empty(), "consumed redo must retire its recovery record");
1123
1124        // A new rewind followed by a new prompt must also retire the discarded redo.
1125        let restored = manager.navigate_prompt(Some(2), RevertScope::Both, &session, &current).await?;
1126        assert_eq!(live_recovery_files(dir.path()).len(), 1);
1127        let _lease = manager.begin_prompt(3, &session, "third", &restored.conversation).await?;
1128        assert!(live_recovery_files(dir.path()).is_empty(), "begin_prompt must retire discarded redo records");
1129        Ok(())
1130    }
1131
1132    #[tokio::test]
1133    async fn recover_pending_fails_closed_without_touching_redo() -> Result<()> {
1134        let dir = TempDir::new()?;
1135        let manager = SnapshotManager::new(SnapshotConfig::new(dir.path().into()))?;
1136        let session = uuid::Uuid::new_v4().to_string();
1137        let original = vec![SessionMessage::new(MessageRole::User, "original")];
1138        let lease = manager.begin_prompt(1, &session, "first", &original).await?;
1139        drop(lease);
1140        let current = vec![SessionMessage::new(MessageRole::User, "current")];
1141
1142        // No interrupted rewind: recovery must fail and must not consume redo.
1143        assert!(manager.recover_pending_rewind(&session).await.is_err());
1144        let restored = manager.navigate_prompt(Some(1), RevertScope::Both, &session, &current).await?;
1145        assert!(manager.recover_pending_rewind(&session).await.is_err());
1146        // Redo is still intact for the completed rewind.
1147        let redone = manager
1148            .navigate_prompt(None, RevertScope::Both, &session, &restored.conversation)
1149            .await?;
1150        assert_eq!(redone.conversation, current);
1151        Ok(())
1152    }
1153
1154    #[tokio::test]
1155    async fn complete_session_navigation_trims_active_and_clears_redo() -> Result<()> {
1156        let dir = TempDir::new()?;
1157        let manager = SnapshotManager::new(SnapshotConfig::new(dir.path().into()))?;
1158        let storage = dir.path().join(".vtcode").join("checkpoints");
1159        let session = uuid::Uuid::new_v4().to_string();
1160        let key: String = session.as_bytes().iter().map(|b| format!("{b:02x}")).collect();
1161        let branch_path = storage.join(format!("branch_{key}.json"));
1162        let lock_path = storage.join("rewind.lock");
1163        fs::write(&lock_path, b"")?;
1164
1165        // 20 protected turns plus a redo entry: the completion trim must keep
1166        // only the newest REWIND_ACTIVE_KEEP actives and drop redo entirely.
1167        let active: Vec<usize> = (1435..=1454).collect();
1168        let recovery = Recovery {
1169            policy: String::new(),
1170            snapshot: uuid::Uuid::new_v4().to_string(),
1171            active: vec![1435],
1172        };
1173        let recovery_path = recovery_path(&storage, &recovery.snapshot);
1174        atomic_json(&recovery_path, &serde_json::json!({}))?;
1175        let state = Navigation {
1176            active: active.clone(),
1177            redo: vec![recovery],
1178            pending: None,
1179        };
1180        atomic_json(&branch_path, &state)?;
1181
1182        manager.complete_session_navigation(&session).await?;
1183
1184        let saved: Navigation = serde_json::from_slice(&fs::read(&branch_path)?)?;
1185        assert_eq!(
1186            saved.active,
1187            active[active.len() - REWIND_ACTIVE_KEEP..].to_vec(),
1188            "only the newest rewind window stays pinned"
1189        );
1190        assert!(saved.redo.is_empty(), "completion must clear redo");
1191        assert!(saved.pending.is_none());
1192        assert!(!lock_path.exists(), "completion must release the workspace rewind lock");
1193        assert!(!recovery_path.exists(), "completion must retire redo recovery records");
1194        Ok(())
1195    }
1196}