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
18pub 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#[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
84fn 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
124fn canonicalize_for_strip(path: &Path) -> PathBuf {
127 if let Ok(canonical) = canonicalize(path) {
128 return canonical;
129 }
130 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 "heredoc_redirect" | "comment" => {}
249 _ => return None,
252 }
253 }
254 Some(paths)
255 };
256 collect().unwrap_or_default()
257}
258
259pub 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 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 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 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 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 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 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 match fs::OpenOptions::new().read(true).write(true).open(&lock_path) {
464 Ok(file) => {
465 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 fs::write(dir.path().join("unrelated.txt"), "leave me alone")?;
1010
1011 let restored = manager.navigate_prompt(Some(2), RevertScope::Both, &session, ¤t).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, ¤t).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, ¤t)
1072 .await
1073 .is_err()
1074 );
1075 let restored = manager.navigate_prompt(Some(1), RevertScope::Both, &session, ¤t).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, ¤t).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 let restored = manager.navigate_prompt(Some(2), RevertScope::Both, &session, ¤t).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 assert!(manager.recover_pending_rewind(&session).await.is_err());
1144 let restored = manager.navigate_prompt(Some(1), RevertScope::Both, &session, ¤t).await?;
1145 assert!(manager.recover_pending_rewind(&session).await.is_err());
1146 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 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}