1mod native;
2pub use native::{PromptCheckpointLease, declare_prompt_edit};
3use std::collections::BTreeSet;
4use std::fs;
5use std::path::{Component, Path, PathBuf};
6use std::sync::atomic::{AtomicU64, Ordering};
7use std::time::{Duration, SystemTime, UNIX_EPOCH};
8
9use anyhow::{Context, Result};
10use base64::Engine as _;
11use base64::engine::general_purpose::STANDARD as BASE64;
12use serde::{Deserialize, Serialize};
13use vtcode_exec_events::{MAX_IN_PROGRESS_EXEC_SESSIONS, Usage, deserialize_null_as_default};
14
15use crate::core::pending_actions::ExpectedOutcome;
16use crate::core::state_schema::{SchemaVersion, VersionedState};
17use crate::types::CompactStr;
18use crate::utils::error_messages::ERR_CREATE_CHECKPOINT_DIR;
19use crate::utils::file_utils::{ensure_dir_exists, ensure_dir_exists_sync, write_json_file};
20use crate::utils::path::canonicalize_workspace;
21use crate::utils::session_archive::SessionMessage;
22
23const MAX_DESCRIPTION_LEN: usize = 160;
24use vtcode_commons::canonicalize;
25
26use crate::core::SECONDS_PER_DAY;
27pub const DEFAULT_CHECKPOINTS_ENABLED: bool = true;
28pub const DEFAULT_MAX_SNAPSHOTS: usize = 50;
29pub const DEFAULT_MAX_AGE_DAYS: u64 = 30;
30pub const REWIND_ACTIVE_KEEP: usize = 5;
34const SNAPSHOT_SCHEMA_VERSION: SchemaVersion = SchemaVersion(3);
35
36fn collect_turn_numbers(value: &serde_json::Value, out: &mut BTreeSet<usize>) {
39 match value {
40 serde_json::Value::Array(items) => {
41 for item in items {
42 if let Some(turn) = item.as_u64() {
43 out.insert(turn as usize);
44 }
45 }
46 }
47 serde_json::Value::Object(map) => {
48 if let Some(active) = map.get("active") {
49 collect_turn_numbers(active, out);
50 }
51 }
52 _ => {}
53 }
54}
55
56fn normalized_prompt_text(text: &str) -> Option<&str> {
57 let trimmed = text.trim();
58 (!trimmed.is_empty()).then_some(trimmed)
59}
60
61fn sanitize_relative_path(path: &Path) -> Option<PathBuf> {
62 if path.is_absolute() {
63 return None;
64 }
65
66 let mut normalized = PathBuf::new();
67 for component in path.components() {
68 match component {
69 Component::CurDir => {}
70 Component::Normal(part) => normalized.push(part),
71 Component::ParentDir => {
72 if !normalized.pop() {
73 return None;
74 }
75 }
76 Component::Prefix(_) | Component::RootDir => {
77 return None;
78 }
79 }
80 }
81 Some(normalized)
82}
83
84#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
85pub struct SnapshotMetadata {
86 pub id: String,
87 pub turn_number: usize,
88 pub created_at: u64,
89 pub description: String,
90 pub message_count: usize,
91 pub file_count: usize,
92 #[serde(default, skip_serializing_if = "Vec::is_empty")]
93 pub touched_files: Vec<String>,
94 #[serde(default, skip_serializing_if = "Option::is_none")]
95 pub prompt_text: Option<String>,
96 #[serde(default, skip_serializing_if = "Option::is_none")]
97 pub prompt_message_index: Option<usize>,
98 #[serde(default, skip_serializing_if = "Option::is_none")]
99 pub session_id: Option<CompactStr>,
100 #[serde(default, skip_serializing_if = "Option::is_none")]
101 pub runtime_turn_id: Option<CompactStr>,
102 #[serde(default, skip_serializing_if = "Option::is_none")]
103 pub session_turn_number: Option<usize>,
104 #[serde(default, skip_serializing_if = "Option::is_none")]
105 pub turn_diagnostics: Option<SnapshotTurnDiagnostics>,
106}
107
108impl SnapshotMetadata {
109 pub fn resolved_prompt_text<'a>(&'a self, conversation: &'a [SessionMessage]) -> Option<String> {
110 self.prompt_text
111 .as_deref()
112 .and_then(normalized_prompt_text)
113 .map(str::to_string)
114 .or_else(|| SnapshotManager::derive_prompt_metadata(conversation).0)
115 }
116}
117
118#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
119pub struct SnapshotTurnDiagnostics {
120 #[serde(default)]
121 pub usage: Usage,
122 #[serde(default, deserialize_with = "deserialize_null_as_default")]
123 pub elapsed_ms: u64,
124 #[serde(default, deserialize_with = "deserialize_null_as_default")]
125 pub requested_tool_calls: u32,
126 #[serde(default, deserialize_with = "deserialize_null_as_default")]
127 pub admitted_tool_calls: u32,
128 #[serde(default, deserialize_with = "deserialize_null_as_default")]
129 pub unadmitted_tool_calls: u32,
130 #[serde(default, deserialize_with = "deserialize_null_as_default")]
131 pub failed_tool_calls: u32,
132 #[serde(default, deserialize_with = "deserialize_null_as_default")]
133 pub denied_tool_calls: u32,
134 #[serde(default, deserialize_with = "deserialize_null_as_default")]
135 pub preflight_failures: u32,
136 #[serde(default, deserialize_with = "deserialize_null_as_default")]
137 pub reused_results: u32,
138 #[serde(default, deserialize_with = "deserialize_null_as_default")]
139 pub spooled_results: u32,
140 #[serde(default, deserialize_with = "deserialize_null_as_default")]
141 pub raw_spooled_bytes: u64,
142 #[serde(default, deserialize_with = "deserialize_null_as_default")]
143 pub model_visible_output_bytes: u64,
144 #[serde(default, deserialize_with = "deserialize_null_as_default")]
145 pub suppressed_tool_previews: u32,
146 #[serde(default, deserialize_with = "deserialize_null_as_default")]
147 pub model_visible_tool_preview_budget_exhausted: bool,
148 #[serde(default, deserialize_with = "deserialize_null_as_default")]
149 pub low_signal_tool_calls: u32,
150 #[serde(default, deserialize_with = "deserialize_null_as_default")]
151 pub recovery_activations: u32,
152 #[serde(default, deserialize_with = "deserialize_null_as_default")]
156 pub in_progress_exec_sessions: Vec<CompactStr>,
157}
158
159impl SnapshotTurnDiagnostics {
160 #[must_use]
167 pub fn with_in_progress_exec_sessions(mut self, sessions: Vec<crate::tools::types::VTCodeExecSession>) -> Self {
168 self.in_progress_exec_sessions = sessions
169 .into_iter()
170 .take(MAX_IN_PROGRESS_EXEC_SESSIONS)
171 .map(|session| CompactStr::from(session.id.as_str().to_string()))
172 .collect();
173 self
174 }
175}
176
177#[derive(Debug, Clone, Default, PartialEq, Eq)]
178pub struct SnapshotTurnContext {
179 pub session_id: Option<CompactStr>,
180 pub runtime_turn_id: Option<CompactStr>,
181 pub session_turn_number: Option<usize>,
182 pub turn_diagnostics: Option<SnapshotTurnDiagnostics>,
183 pub touched_files: Vec<String>,
184}
185
186#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
187pub enum FileEncoding {
188 Utf8,
189 Base64,
190 Filesnap,
193}
194
195#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
196pub struct FileSnapshot {
197 pub path: String,
198 pub deleted: bool,
199 #[serde(skip_serializing_if = "Option::is_none")]
200 pub encoding: Option<FileEncoding>,
201 #[serde(skip_serializing_if = "Option::is_none")]
202 pub data: Option<String>,
203}
204
205#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
206pub struct StoredSnapshot {
207 pub metadata: SnapshotMetadata,
208 pub conversation: Vec<SessionMessage>,
209 pub files: Vec<FileSnapshot>,
210 #[serde(default, skip_serializing_if = "Option::is_none")]
212 pub schema_version: Option<SchemaVersion>,
213}
214
215impl VersionedState for StoredSnapshot {
216 fn schema_version(&self) -> SchemaVersion {
217 self.schema_version.unwrap_or(SchemaVersion::V0)
218 }
219
220 fn migrate_one_step(self, from: SchemaVersion, to: SchemaVersion) -> Result<Self> {
221 match (from, to) {
222 (SchemaVersion::V0, SchemaVersion::V1) => {
223 Ok(Self { schema_version: Some(SchemaVersion::V1), ..self })
226 }
227 (SchemaVersion::V1, SchemaVersion::V2) => Ok(Self { schema_version: Some(SchemaVersion::V2), ..self }),
228 (SchemaVersion::V2, SNAPSHOT_SCHEMA_VERSION) => Ok(Self {
229 schema_version: Some(SNAPSHOT_SCHEMA_VERSION),
230 ..self
231 }),
232 _ => anyhow::bail!("unsupported snapshot migration: {from:?} -> {to:?}"),
233 }
234 }
235
236 fn next_version(current: SchemaVersion) -> Option<SchemaVersion> {
237 match current {
238 SchemaVersion::V0 => Some(SchemaVersion::V1),
239 SchemaVersion::V1 => Some(SchemaVersion::V2),
240 SchemaVersion::V2 => Some(SNAPSHOT_SCHEMA_VERSION),
241 SNAPSHOT_SCHEMA_VERSION => None,
242 _ => None,
243 }
244 }
245}
246
247#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
253pub struct ActionSnapshot {
254 pub action_id: String,
256 pub action_number: usize,
258 pub created_at: u64,
260 pub tool_name: String,
262 pub arguments: serde_json::Value,
264 pub pre_action_message_count: usize,
267 pub touched_files: Vec<String>,
269 pub expected_outcome: ExpectedOutcome,
271}
272
273#[derive(Debug, Clone)]
275pub struct RollbackResult {
276 pub rollback_action_id: String,
278 pub messages_removed: usize,
280 pub files_restored: usize,
282 pub next_action_number: usize,
284}
285
286static NEXT_ACTION_NUMBER: AtomicU64 = AtomicU64::new(1);
288
289#[derive(Debug, Clone, Copy, PartialEq, Eq)]
290pub enum RevertScope {
291 Conversation,
292 Code,
293 Both,
294}
295
296impl RevertScope {
297 pub fn includes_code(self) -> bool {
298 matches!(self, Self::Code | Self::Both)
299 }
300
301 pub fn includes_conversation(self) -> bool {
302 matches!(self, Self::Conversation | Self::Both)
303 }
304}
305
306pub struct SnapshotConfig {
307 pub enabled: bool,
308 pub workspace: PathBuf,
309 pub storage_dir: Option<PathBuf>,
310 pub max_snapshots: usize,
311 pub max_age_days: Option<u64>,
312}
313
314impl SnapshotConfig {
315 pub fn new(workspace: PathBuf) -> Self {
316 Self {
317 enabled: DEFAULT_CHECKPOINTS_ENABLED,
318 workspace,
319 storage_dir: None,
320 max_snapshots: DEFAULT_MAX_SNAPSHOTS,
321 max_age_days: Some(DEFAULT_MAX_AGE_DAYS),
322 }
323 }
324
325 fn storage_dir(&self) -> PathBuf {
326 self.storage_dir
327 .clone()
328 .unwrap_or_else(|| self.workspace.join(".vtcode").join("checkpoints"))
329 }
330}
331
332#[derive(Clone)]
333pub struct SnapshotManager {
334 enabled: bool,
335 workspace: PathBuf,
336 canonical_workspace: PathBuf,
337 storage_dir: PathBuf,
338 max_snapshots: usize,
339 max_age_days: Option<u64>,
340}
341
342impl SnapshotManager {
343 pub fn new(config: SnapshotConfig) -> Result<Self> {
344 let storage_dir = config.storage_dir();
345 let canonical_workspace = canonicalize_workspace(&config.workspace);
346
347 if config.enabled {
348 ensure_dir_exists_sync(&storage_dir)
349 .with_context(|| format!("{}: {}", ERR_CREATE_CHECKPOINT_DIR, storage_dir.display()))?;
350 }
351 Ok(Self {
352 enabled: config.enabled,
353 workspace: config.workspace,
354 canonical_workspace,
355 storage_dir,
356 max_snapshots: config.max_snapshots,
357 max_age_days: config.max_age_days,
358 })
359 }
360
361 pub fn enabled(&self) -> bool {
362 self.enabled
363 }
364
365 fn snapshot_path(&self, turn_number: usize) -> PathBuf {
366 self.storage_dir.join(format!("turn_{turn_number}.json"))
367 }
368
369 fn normalize_path(&self, path: &Path) -> Option<PathBuf> {
370 if path.is_absolute() {
371 if let Ok(canonical_path) = canonicalize(path)
372 && let Ok(stripped) = canonical_path.strip_prefix(&self.canonical_workspace)
373 {
374 return sanitize_relative_path(stripped);
375 }
376
377 if let Ok(stripped) = path.strip_prefix(&self.workspace) {
378 return sanitize_relative_path(stripped);
379 }
380
381 None
382 } else {
383 sanitize_relative_path(path)
384 }
385 }
386
387 fn checked_file_path(workspace: &Path, storage: &Path, relative: &Path) -> Result<PathBuf> {
388 let relative = sanitize_relative_path(relative).context("Checkpoint path escapes the workspace")?;
389 anyhow::ensure!(!relative.as_os_str().is_empty(), "Checkpoint path must name a file");
390 let absolute = workspace.join(&relative);
391 let storage = canonicalize(storage).unwrap_or_else(|_| storage.to_path_buf());
392 anyhow::ensure!(!absolute.starts_with(&storage), "Checkpoint cannot restore its own storage");
393 let mut current = workspace.to_path_buf();
394 for part in relative.components() {
395 current.push(part);
396 match fs::symlink_metadata(¤t) {
397 Ok(metadata) => {
398 anyhow::ensure!(!metadata.file_type().is_symlink(), "Checkpoint path crosses a symlink")
399 }
400 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
401 Err(error) => return Err(error.into()),
402 }
403 }
404 Ok(absolute)
405 }
406
407 fn read_snapshot_files(&self) -> Result<Vec<(usize, PathBuf)>> {
408 let mut entries = Vec::with_capacity(64); if !self.storage_dir.exists() {
410 return Ok(entries);
411 }
412 for entry in fs::read_dir(&self.storage_dir)
413 .with_context(|| format!("failed to read checkpoint directory: {}", self.storage_dir.display()))?
414 {
415 let entry = entry?;
416 let path = entry.path();
417 if path.extension().and_then(|ext| ext.to_str()) != Some("json") {
418 continue;
419 }
420 let stem = match path.file_stem().and_then(|stem| stem.to_str()) {
421 Some(value) => value,
422 None => continue,
423 };
424 let turn_str = match stem.strip_prefix("turn_") {
425 Some(value) => value,
426 None => continue,
427 };
428 if let Ok(turn) = turn_str.parse::<usize>() {
429 entries.push((turn, path));
430 }
431 }
432 entries.sort_by_key(|(turn, _)| *turn);
433 Ok(entries)
434 }
435
436 fn decode_file(encoding: FileEncoding, data: &str) -> Result<Vec<u8>> {
437 match encoding {
438 FileEncoding::Utf8 => Ok(data.as_bytes().to_vec()),
439 FileEncoding::Base64 => BASE64.decode(data).context("failed to decode base64 file contents"),
440 FileEncoding::Filesnap => anyhow::bail!("filesnap references require the snapshot store"),
441 }
442 }
443
444 fn truncate_description(description: &str) -> String {
445 let first_line = description.lines().next().unwrap_or("").trim();
446 vtcode_commons::formatting::truncate_within(first_line, MAX_DESCRIPTION_LEN, "…")
447 }
448
449 fn derive_prompt_metadata(conversation: &[SessionMessage]) -> (Option<String>, Option<usize>) {
450 conversation
451 .iter()
452 .enumerate()
453 .rev()
454 .find_map(|(index, message)| {
455 if message.role != crate::llm::provider::MessageRole::User {
456 return None;
457 }
458
459 let prompt = message.content.as_text();
460 normalized_prompt_text(prompt.as_ref()).map(|prompt| (Some(prompt.to_string()), Some(index)))
461 })
462 .unwrap_or((None, None))
463 }
464
465 fn resolve_prompt_metadata(
466 prompt_text: Option<&str>,
467 prompt_message_index: Option<usize>,
468 conversation: &[SessionMessage],
469 ) -> (Option<String>, Option<usize>) {
470 let (derived_prompt_text, derived_prompt_index) = Self::derive_prompt_metadata(conversation);
471 let prompt_text = prompt_text
472 .and_then(normalized_prompt_text)
473 .map(str::to_string)
474 .or(derived_prompt_text);
475 let prompt_message_index = prompt_message_index
476 .filter(|index| *index < conversation.len())
477 .or(derived_prompt_index);
478 (prompt_text, prompt_message_index)
479 }
480
481 fn hydrate_prompt_metadata(stored: &mut StoredSnapshot) {
482 let (prompt_text, prompt_message_index) = Self::resolve_prompt_metadata(
483 stored.metadata.prompt_text.as_deref(),
484 stored.metadata.prompt_message_index,
485 &stored.conversation,
486 );
487 stored.metadata.prompt_text = prompt_text;
488 stored.metadata.prompt_message_index = prompt_message_index;
489 }
490
491 fn current_timestamp() -> Result<u64> {
492 Ok(SystemTime::now()
493 .duration_since(UNIX_EPOCH)
494 .context("system clock before UNIX_EPOCH")?
495 .as_secs())
496 }
497
498 pub fn next_turn_number(&self) -> Result<usize> {
499 Ok(self
500 .read_snapshot_files()?
501 .into_iter()
502 .map(|(turn, _)| turn)
503 .max()
504 .unwrap_or(0)
505 .saturating_add(1))
506 }
507
508 pub async fn create_snapshot(
509 &self,
510 turn_number: usize,
511 description: &str,
512 conversation: &[SessionMessage],
513 modified_files: &BTreeSet<PathBuf>,
514 prompt_text: Option<&str>,
515 prompt_message_index: Option<usize>,
516 turn_context: Option<SnapshotTurnContext>,
517 ) -> Result<Option<SnapshotMetadata>> {
518 if !self.enabled {
519 return Ok(None);
520 }
521
522 let timestamp = Self::current_timestamp()?;
523 let mut paths = Vec::with_capacity(modified_files.len());
524 for path in modified_files {
525 if let Some(relative) = self.normalize_path(path) {
526 paths.push(relative);
527 }
528 }
529 let workspace = self.canonical_workspace.clone();
530 let storage = self.storage_dir.clone();
531 let files = tokio::task::spawn_blocking(move || -> Result<Vec<FileSnapshot>> {
532 if paths.is_empty() {
534 return Ok(Vec::new());
535 }
536 let mut absolute_paths = Vec::with_capacity(paths.len());
537 for relative in &paths {
538 absolute_paths.push(Self::checked_file_path(&workspace, &storage, relative)?);
539 }
540 let store = filesnap::WorkspaceStore::open(&storage, &workspace)?;
541 let turn = format!("vt-{}", uuid::Uuid::new_v4());
542 let checkpoint = store.checkpoint(&turn, &turn, absolute_paths.iter().cloned())?;
543 anyhow::ensure!(checkpoint.stats.dropped == 0, "Checkpoint could not capture every selected file");
544 Ok(paths
545 .into_iter()
546 .zip(absolute_paths)
547 .map(|(relative, absolute)| {
548 let key = filesnap::canonical_key(&absolute).to_string_lossy().into_owned();
549 FileSnapshot {
550 path: relative.to_string_lossy().replace('\\', "/"),
551 deleted: checkpoint.manifest.absent.contains(&key),
552 encoding: Some(FileEncoding::Filesnap),
553 data: Some(turn.clone()),
554 }
555 })
556 .collect())
557 })
558 .await??;
559
560 let (prompt_text, prompt_message_index) =
561 Self::resolve_prompt_metadata(prompt_text, prompt_message_index, conversation);
562 let description_source = prompt_text.as_deref().unwrap_or(description);
563 let turn_context = turn_context.unwrap_or_default();
564 let metadata = SnapshotMetadata {
565 id: format!("turn_{turn_number}"),
566 turn_number,
567 created_at: timestamp,
568 description: Self::truncate_description(description_source),
569 message_count: conversation.len(),
570 file_count: files.len(),
571 touched_files: turn_context.touched_files.clone(),
572 prompt_text,
573 prompt_message_index,
574 session_id: turn_context.session_id,
575 runtime_turn_id: turn_context.runtime_turn_id,
576 session_turn_number: turn_context.session_turn_number,
577 turn_diagnostics: turn_context.turn_diagnostics,
578 };
579
580 let stored = StoredSnapshot {
581 metadata: metadata.clone(),
582 conversation: conversation.to_vec(),
583 files,
584 schema_version: Some(SNAPSHOT_SCHEMA_VERSION),
585 };
586
587 let path = self.snapshot_path(turn_number);
588 if let Some(parent) = path.parent() {
589 ensure_dir_exists(parent)
590 .await
591 .with_context(|| format!("failed to ensure checkpoint directory: {}", parent.display()))?;
592 }
593
594 let retired = path.with_extension(format!("retired-{}", uuid::Uuid::new_v4()));
597 match tokio::fs::copy(&path, &retired).await {
598 Ok(_) => {}
599 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
600 Err(error) => return Err(error).context("failed to journal replaced checkpoint"),
601 }
602 write_json_file(&path, &stored)
603 .await
604 .with_context(|| format!("failed to write checkpoint: {}", path.display()))?;
605
606 self.cleanup_old_snapshots().await?;
607
608 Ok(Some(metadata))
609 }
610
611 pub async fn list_snapshots(&self) -> Result<Vec<SnapshotMetadata>> {
612 if !self.enabled {
613 return Ok(Vec::new());
614 }
615 self.cleanup_old_snapshots().await?;
616 let snapshot_files = self.read_snapshot_files()?;
617 let mut snapshots = Vec::with_capacity(snapshot_files.len());
618 for (_, path) in snapshot_files {
619 let data = tokio::fs::read(&path)
620 .await
621 .with_context(|| format!("failed to read checkpoint: {}", path.display()))?;
622 let mut stored: StoredSnapshot = serde_json::from_slice(&data)
623 .with_context(|| format!("failed to parse checkpoint: {}", path.display()))?;
624 Self::hydrate_prompt_metadata(&mut stored);
625 snapshots.push(stored.metadata);
626 }
627 snapshots.sort_by_key(|a| std::cmp::Reverse(a.turn_number));
628 Ok(snapshots)
629 }
630
631 pub async fn load_snapshot(&self, turn_number: usize) -> Result<Option<StoredSnapshot>> {
632 if !self.enabled {
633 return Ok(None);
634 }
635 let path = self.snapshot_path(turn_number);
636 if !tokio::fs::try_exists(&path).await.unwrap_or(false) {
637 return Ok(None);
638 }
639 let data = tokio::fs::read(&path)
640 .await
641 .with_context(|| format!("failed to read checkpoint: {}", path.display()))?;
642 let mut stored: StoredSnapshot =
643 serde_json::from_slice(&data).with_context(|| format!("failed to parse checkpoint: {}", path.display()))?;
644 stored = stored
646 .migrate(SNAPSHOT_SCHEMA_VERSION)
647 .with_context(|| format!("failed to migrate checkpoint: {}", path.display()))?;
648 Self::hydrate_prompt_metadata(&mut stored);
649 Ok(Some(stored))
650 }
651
652 pub async fn restore_snapshot(&self, turn_number: usize, scope: RevertScope) -> Result<Option<CheckpointRestore>> {
653 let Some(stored) = self.load_snapshot(turn_number).await? else {
654 return Ok(None);
655 };
656
657 self.restore_stored_snapshot(stored, scope).await.map(Some)
658 }
659
660 async fn restore_stored_snapshot(&self, stored: StoredSnapshot, scope: RevertScope) -> Result<CheckpointRestore> {
661 self.restore_stored_snapshot_with_ignore(stored, scope, &filesnap::Gitignore::empty())
662 .await
663 }
664
665 async fn restore_stored_snapshot_with_ignore(
666 &self,
667 stored: StoredSnapshot,
668 scope: RevertScope,
669 ignore: &filesnap::Gitignore,
670 ) -> Result<CheckpointRestore> {
671 if scope.includes_code() {
672 let workspace = self.canonical_workspace.clone();
673 let storage = self.storage_dir.clone();
674 let files = stored.files.clone();
675 tokio::task::spawn_blocking(move || -> Result<()> {
677 for file in &files {
678 Self::checked_file_path(&workspace, &storage, Path::new(&file.path))?;
679 }
680 Ok(())
681 })
682 .await??;
683 }
684 let engine_backed = stored.files.iter().any(|file| file.encoding == Some(FileEncoding::Filesnap));
685 if scope.includes_code() && engine_backed {
686 let ignore = ignore.clone();
687 let workspace = self.canonical_workspace.clone();
688 let storage = self.storage_dir.clone();
689 let files = stored.files.clone();
690 tokio::task::spawn_blocking(move || -> Result<()> {
691 let turn = files
692 .first()
693 .and_then(|file| file.data.as_deref())
694 .context("Missing filesnap reference")?;
695 anyhow::ensure!(files.iter().all(|file| file.encoding == Some(FileEncoding::Filesnap)
696 && file.data.as_deref() == Some(turn)), "Mixed or inconsistent checkpoint references");
697 let store = filesnap::WorkspaceStore::open(&storage, &workspace)?;
698 let target = store.target_for_turn(turn)?.context("Missing filesnap checkpoint")?;
699 let manifest = store.manifest(target.manifest_id())?;
700 let expected: BTreeSet<String> = files
701 .iter()
702 .map(|file| {
703 filesnap::canonical_key(&workspace.join(&file.path))
704 .to_string_lossy()
705 .into_owned()
706 })
707 .collect();
708 let recorded: BTreeSet<String> =
709 manifest.entries.keys().chain(manifest.absent.iter()).cloned().collect();
710 anyhow::ensure!(expected == recorded, "Checkpoint contains unexpected file paths");
711 let outcome = store.restore_to(
712 turn,
713 &target,
714 filesnap::RestoreKind::Rewind { undo_for: Some(turn) },
715 expected.iter().map(PathBuf::from),
716 &ignore,
717 )?;
718 anyhow::ensure!(
719 outcome.stats.failed.is_empty(),
720 "Checkpoint restore failed: {:?}",
721 outcome.stats.failed
722 );
723 Ok(())
724 })
725 .await??;
726 } else if scope.includes_code() {
727 for snapshot in &stored.files {
728 let relative = Path::new(&snapshot.path);
729 let Some(sanitized) = sanitize_relative_path(relative) else {
730 continue;
731 };
732 let absolute = self.workspace.join(&sanitized);
733 if snapshot.deleted {
734 if tokio::fs::try_exists(&absolute).await.unwrap_or(false) {
735 tokio::fs::remove_file(&absolute).await.with_context(|| {
736 format!("failed to remove file during checkpoint restore: {}", absolute.display())
737 })?;
738 }
739 continue;
740 }
741
742 if let Some(parent) = absolute.parent() {
743 ensure_dir_exists(parent)
744 .await
745 .with_context(|| format!("failed to create directories for restore: {}", parent.display()))?;
746 }
747
748 let encoding = snapshot.encoding.unwrap_or(FileEncoding::Utf8);
749 let data = snapshot.data.as_deref().unwrap_or_default();
750 let bytes = Self::decode_file(encoding, data)?;
751 tokio::fs::write(&absolute, &bytes)
752 .await
753 .with_context(|| format!("failed to write restored file: {}", absolute.display()))?;
754 }
755 }
756
757 let conversation = if scope.includes_conversation() {
758 stored.conversation.clone()
759 } else {
760 Vec::new()
761 };
762
763 Ok(CheckpointRestore { metadata: stored.metadata, conversation })
764 }
765
766 fn protected_turns(&self) -> BTreeSet<usize> {
771 let cutoff = self.retention_cutoff_secs().ok().flatten();
772 self.protected_turns_with_cutoff(cutoff)
773 }
774
775 fn protected_turns_with_cutoff(&self, cutoff: Option<u64>) -> BTreeSet<usize> {
776 let mut protected = BTreeSet::new();
777 let Ok(entries) = fs::read_dir(&self.storage_dir) else {
778 return protected;
779 };
780 for entry in entries.flatten() {
781 let path = entry.path();
782 let Some(stem) = path.file_stem().and_then(|s| s.to_str()) else {
783 continue;
784 };
785 if !stem.starts_with("branch_") || path.extension().and_then(|e| e.to_str()) != Some("json") {
786 continue;
787 }
788 if let Some(cutoff) = cutoff {
789 let modified_secs = path
793 .metadata()
794 .and_then(|meta| meta.modified())
795 .ok()
796 .and_then(|modified| modified.duration_since(UNIX_EPOCH).ok())
797 .map(|since| since.as_secs());
798 if !modified_secs.is_some_and(|secs| secs > cutoff) {
799 continue;
800 }
801 }
802 let Ok(bytes) = fs::read(&path) else {
803 continue;
804 };
805 let Ok(value) = serde_json::from_slice::<serde_json::Value>(&bytes) else {
806 continue;
807 };
808 for key in ["active", "redo", "pending"] {
809 collect_turn_numbers(&value[key], &mut protected);
810 }
811 }
812 protected
813 }
814
815 pub async fn prune_snapshot_budget(&self) -> Result<()> {
820 if !self.enabled {
821 return Ok(());
822 }
823 let manager = self.clone();
824 let entries = tokio::task::spawn_blocking(move || -> Result<Vec<(usize, PathBuf)>> {
825 let protected = manager.protected_turns();
826 Ok(manager
827 .read_snapshot_files()?
828 .into_iter()
829 .filter(|(turn, _)| !protected.contains(turn))
830 .collect())
831 })
832 .await
833 .context("checkpoint retention worker failed")??;
834 if self.max_snapshots != 0 && entries.len() > self.max_snapshots {
835 let excess = entries.len() - self.max_snapshots;
836 for (_, path) in entries.into_iter().take(excess) {
837 if let Err(err) = self.retire_snapshot(&path).await {
838 tracing::warn!(
839 path = %path.display(),
840 error = %err,
841 "Failed to remove old checkpoint"
842 );
843 }
844 }
845 }
846 self.cleanup_retired_snapshots(false).await;
847 Ok(())
848 }
849
850 pub async fn cleanup_old_snapshots(&self) -> Result<()> {
854 if !self.enabled {
855 return Ok(());
856 }
857
858 let protected = self.protected_turns();
859 if let Some(cutoff) = self.retention_cutoff_secs()? {
860 for (turn, path) in self.read_snapshot_files()? {
861 if protected.contains(&turn) {
862 continue;
863 }
864 let data = match tokio::fs::read(&path).await {
865 Ok(data) => data,
866 Err(err) => {
867 tracing::warn!(
868 path = %path.display(),
869 error = %err,
870 "Failed to read checkpoint"
871 );
872 continue;
873 }
874 };
875 let stored: StoredSnapshot = match serde_json::from_slice(&data) {
876 Ok(value) => value,
877 Err(err) => {
878 tracing::warn!(
879 path = %path.display(),
880 error = %err,
881 "Failed to parse checkpoint"
882 );
883 continue;
884 }
885 };
886 if stored.metadata.created_at <= cutoff
887 && let Err(err) = self.retire_snapshot(&path).await
888 {
889 tracing::warn!(
890 path = %path.display(),
891 error = %err,
892 "Failed to remove expired checkpoint"
893 );
894 }
895 }
896 }
897
898 self.prune_snapshot_budget().await?;
899 self.cleanup_retired_snapshots(true).await;
900 Ok(())
901 }
902
903 async fn retire_snapshot(&self, path: &Path) -> Result<()> {
904 let retired = path.with_extension(format!("retired-{}", uuid::Uuid::new_v4()));
905 tokio::fs::rename(path, retired).await?;
906 Ok(())
907 }
908
909 async fn cleanup_retired_snapshots(&self, full_maintenance: bool) {
910 let storage = self.storage_dir.clone();
911 let workspace = self.canonical_workspace.clone();
912 let result = tokio::task::spawn_blocking(move || -> Result<()> {
913 let mut live_paths = Vec::new();
917 let mut retired = Vec::new();
918 for entry in fs::read_dir(&storage)? {
919 let path = entry?.path();
920 let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
921 continue;
922 };
923 if !name.starts_with("turn_") {
924 continue;
925 }
926 if name.ends_with(".json") {
927 live_paths.push(path);
928 } else if path
929 .extension()
930 .and_then(|value| value.to_str())
931 .is_some_and(|value| value.starts_with("retired-"))
932 {
933 retired.push(path);
934 }
935 }
936 if retired.is_empty() && !full_maintenance {
939 return Ok(());
940 }
941 let mut live = BTreeSet::new();
942 for path in live_paths {
943 let stored: StoredSnapshot = serde_json::from_slice(&fs::read(path)?)?;
944 for file in stored.files {
945 if file.encoding == Some(FileEncoding::Filesnap) {
946 live.insert(file.data.context("Missing checkpoint reference")?);
947 }
948 }
949 }
950 let store = filesnap::WorkspaceStore::open(&storage, &workspace)?;
951 for path in retired {
952 let stored: StoredSnapshot = serde_json::from_slice(&fs::read(&path)?)?;
953 let mut sessions = BTreeSet::new();
954 for file in &stored.files {
955 if file.encoding == Some(FileEncoding::Filesnap) {
956 let session = file.data.as_deref().context("Missing retired checkpoint reference")?;
957 let id = session.strip_prefix("vt-").context("Invalid retired checkpoint reference")?;
958 uuid::Uuid::parse_str(id)?;
959 if !live.contains(session) {
960 sessions.insert(session.to_owned());
961 }
962 }
963 }
964 let outcome = store.delete_sessions(&sessions.into_iter().collect::<Vec<_>>());
965 anyhow::ensure!(
966 outcome.refused.is_empty() && outcome.incomplete.is_empty(),
967 "Checkpoint cleanup is incomplete: {outcome:?}"
968 );
969 fs::remove_file(path)?;
970 }
971 filesnap::collect_garbage(&storage)?;
972 Ok(())
973 })
974 .await;
975 if !matches!(&result, Ok(Ok(()))) {
976 tracing::warn!(?result, "Checkpoint content cleanup deferred; retired records remain retryable");
977 }
978 }
979
980 fn retention_cutoff_secs(&self) -> Result<Option<u64>> {
981 let Some(days) = self.max_age_days else {
982 return Ok(None);
983 };
984
985 let now = Self::current_timestamp()?;
986 if days == 0 {
987 return Ok(Some(now));
988 }
989
990 let seconds = days.saturating_mul(SECONDS_PER_DAY);
991 let cutoff_instant = SystemTime::now()
992 .checked_sub(Duration::from_secs(seconds))
993 .unwrap_or(SystemTime::UNIX_EPOCH);
994 let cutoff = cutoff_instant
995 .duration_since(UNIX_EPOCH)
996 .context("system clock before UNIX_EPOCH")?
997 .as_secs();
998 Ok(Some(cutoff))
999 }
1000
1001 pub fn parse_revert_scope(value: &str) -> Option<RevertScope> {
1002 match value.to_ascii_lowercase().as_str() {
1003 "conversation" | "chat" => Some(RevertScope::Conversation),
1004 "code" | "files" => Some(RevertScope::Code),
1005 "both" | "full" => Some(RevertScope::Both),
1006 _ => None,
1007 }
1008 }
1009
1010 fn action_snapshot_path(&self, action_number: usize) -> PathBuf {
1014 self.storage_dir.join(format!("action_{action_number}.json"))
1015 }
1016
1017 pub async fn save_action_snapshot(&self, action: &ActionSnapshot) -> Result<usize> {
1019 if !self.enabled {
1020 return Ok(action.action_number);
1021 }
1022 let path = self.action_snapshot_path(action.action_number);
1023 if let Some(parent) = path.parent() {
1024 ensure_dir_exists(parent)
1025 .await
1026 .with_context(|| format!("failed to ensure checkpoint directory: {}", parent.display()))?;
1027 }
1028 write_json_file(&path, action)
1029 .await
1030 .with_context(|| format!("failed to write action snapshot: {}", path.display()))?;
1031 Ok(action.action_number)
1032 }
1033
1034 pub async fn load_action_snapshot(&self, action_number: usize) -> Result<Option<ActionSnapshot>> {
1036 if !self.enabled {
1037 return Ok(None);
1038 }
1039 let path = self.action_snapshot_path(action_number);
1040 if !tokio::fs::try_exists(&path).await.unwrap_or(false) {
1041 return Ok(None);
1042 }
1043 let data = tokio::fs::read(&path)
1044 .await
1045 .with_context(|| format!("failed to read action snapshot: {}", path.display()))?;
1046 let action: ActionSnapshot = serde_json::from_slice(&data)
1047 .with_context(|| format!("failed to parse action snapshot: {}", path.display()))?;
1048 Ok(Some(action))
1049 }
1050
1051 pub async fn rollback_one_action(
1057 &self,
1058 action: &ActionSnapshot,
1059 messages: &mut Vec<SessionMessage>,
1060 scope: RevertScope,
1061 ) -> Result<RollbackResult> {
1062 let mut files_restored = 0;
1063
1064 let messages_removed = if scope.includes_conversation() && action.pre_action_message_count <= messages.len() {
1066 let removed = messages.len() - action.pre_action_message_count;
1067 messages.truncate(action.pre_action_message_count);
1068 removed
1069 } else {
1070 0
1071 };
1072
1073 if scope.includes_code() {
1075 for file_path in &action.touched_files {
1076 let absolute = self.workspace.join(file_path);
1078 if tokio::fs::try_exists(&absolute).await.unwrap_or(false) {
1079 files_restored += 1;
1082 }
1083 }
1084 }
1085
1086 Ok(RollbackResult {
1087 rollback_action_id: action.action_id.clone(),
1088 messages_removed,
1089 files_restored,
1090 next_action_number: action.action_number,
1091 })
1092 }
1093
1094 pub fn next_action_number(&self) -> usize {
1098 NEXT_ACTION_NUMBER.fetch_add(1, Ordering::Relaxed) as usize
1099 }
1100}
1101
1102#[derive(Debug, Clone)]
1103pub struct CheckpointRestore {
1104 pub metadata: SnapshotMetadata,
1105 pub conversation: Vec<SessionMessage>,
1106}
1107
1108#[cfg(test)]
1109mod tests {
1110 use tempfile::TempDir;
1111
1112 use super::*;
1113
1114 fn setup_manager() -> (TempDir, SnapshotManager) {
1115 let dir = TempDir::new().expect("tempdir");
1116 let workspace = dir.path().to_path_buf();
1117 let manager = SnapshotManager::new(SnapshotConfig::new(workspace.clone())).expect("manager");
1118 (dir, manager)
1119 }
1120
1121 #[tokio::test]
1122 async fn no_retirement_prune_does_not_parse_live_records_or_open_store() -> Result<()> {
1123 let (_dir, manager) = setup_manager();
1124 let path = manager.snapshot_path(1);
1125 fs::write(&path, b"deliberately corrupt live metadata")?;
1126 manager.prune_snapshot_budget().await?;
1127 assert_eq!(fs::read(&path)?, b"deliberately corrupt live metadata");
1128 let entries: Vec<_> = fs::read_dir(&manager.storage_dir)?.collect::<std::io::Result<Vec<_>>>()?;
1129 assert_eq!(entries.len(), 1, "hot cleanup must not create an engine store");
1130 assert_eq!(entries[0].path(), path);
1131 Ok(())
1132 }
1133
1134 #[tokio::test]
1135 async fn retired_cleanup_preserves_shared_references_and_defers_on_corrupt_live_metadata() -> Result<()> {
1136 let (_dir, manager) = setup_manager();
1137 let file = manager.workspace.join("shared.txt");
1138 fs::write(&file, b"original")?;
1139 manager
1140 .create_snapshot(1, "first", &[], &BTreeSet::from([file.clone()]), None, None, None)
1141 .await?;
1142 let stored = manager.load_snapshot(1).await?.expect("snapshot");
1143 let engine = stored.files[0].data.clone().expect("engine reference");
1144 fs::write(manager.snapshot_path(2), serde_json::to_vec(&stored)?)?;
1145 let retired = manager.snapshot_path(1).with_extension("retired-test");
1146 fs::rename(manager.snapshot_path(1), &retired)?;
1147 let corrupt = manager.snapshot_path(3);
1148 fs::write(&corrupt, b"{")?;
1149 manager.cleanup_retired_snapshots(false).await;
1150 assert!(retired.exists(), "corrupt live records must defer reclamation");
1151 let store = filesnap::WorkspaceStore::open(&manager.storage_dir, &manager.workspace)?;
1152 assert!(store.target_for_turn(&engine)?.is_some());
1153 fs::remove_file(corrupt)?;
1154 manager.cleanup_retired_snapshots(false).await;
1155 assert!(!retired.exists());
1156 assert!(store.target_for_turn(&engine)?.is_some(), "a live shared reference pins the engine session");
1157 fs::write(&file, b"external edit")?;
1158 manager.restore_snapshot(2, RevertScope::Code).await?;
1159 assert_eq!(fs::read(file)?, b"original");
1160 Ok(())
1161 }
1162
1163 #[tokio::test]
1164 async fn full_maintenance_collects_orphans_without_retired_records() -> Result<()> {
1165 let (_dir, manager) = setup_manager();
1166 let blob_dir = manager
1169 .storage_dir
1170 .join("filesnap")
1171 .join(format!("v{}", filesnap::FORMAT_VERSION))
1172 .join("blobs");
1173 let blobs = filesnap::BlobStore::open(&blob_dir)?;
1174 let hash = blobs.store_bytes(b"orphaned content")?;
1175 let blob = blob_dir.join(&hash[..2]).join(&hash[2..]);
1176 fs::File::options()
1177 .write(true)
1178 .open(&blob)?
1179 .set_times(fs::FileTimes::new().set_modified(SystemTime::now() - Duration::from_secs(3600)))?;
1180 manager.prune_snapshot_budget().await?;
1181 assert!(blob.exists(), "hot pruning skips unrelated orphan maintenance");
1182 manager.cleanup_old_snapshots().await?;
1183 assert!(!blob.exists(), "explicit maintenance must still collect old orphans");
1184 assert!(blobs.hashes()?.is_empty());
1185 Ok(())
1186 }
1187
1188 #[tokio::test]
1189 async fn create_and_list_snapshots() -> Result<()> {
1190 let (_dir, manager) = setup_manager();
1191 let mut conversation = Vec::new();
1192 conversation.push(SessionMessage::new(crate::llm::provider::MessageRole::User, "Hello"));
1193 let files = BTreeSet::new();
1194 manager
1195 .create_snapshot(1, "First turn", &conversation, &files, None, None, None)
1196 .await?
1197 .expect("metadata");
1198 conversation.push(SessionMessage::new(crate::llm::provider::MessageRole::Assistant, "Hi"));
1199 manager
1200 .create_snapshot(2, "Second turn", &conversation, &files, None, None, None)
1201 .await?
1202 .expect("metadata");
1203
1204 let snapshots = manager.list_snapshots().await?;
1205 assert_eq!(snapshots.len(), 2);
1206 assert_eq!(snapshots[0].turn_number, 2);
1207 assert_eq!(snapshots[1].turn_number, 1);
1208 Ok(())
1209 }
1210
1211 #[tokio::test]
1212 async fn snapshot_restores_file_contents() -> Result<()> {
1213 let (dir, manager) = setup_manager();
1214 let workspace = dir.path();
1215 let file_path = workspace.join("example.txt");
1216 fs::write(&file_path, "v1")?;
1217
1218 let mut files = BTreeSet::new();
1219 files.insert(PathBuf::from("example.txt"));
1220 let conversation = vec![SessionMessage::new(
1221 crate::llm::provider::MessageRole::User,
1222 "edit example",
1223 )];
1224 manager
1225 .create_snapshot(1, "save", &conversation, &files, None, None, None)
1226 .await?
1227 .expect("metadata");
1228
1229 fs::write(&file_path, "v2")?;
1230 manager.restore_snapshot(1, RevertScope::Code).await?.expect("restore");
1231 let restored = fs::read_to_string(&file_path)?;
1232 assert_eq!(restored, "v1");
1233 Ok(())
1234 }
1235
1236 #[tokio::test]
1237 async fn snapshot_handles_deleted_files() -> Result<()> {
1238 let (dir, manager) = setup_manager();
1239 let workspace = dir.path();
1240 let file_path = workspace.join("remove.txt");
1241 fs::write(&file_path, "data")?;
1242
1243 let mut files = BTreeSet::new();
1244 files.insert(PathBuf::from("remove.txt"));
1245 let conversation = vec![SessionMessage::new(crate::llm::provider::MessageRole::User, "remove")];
1246 manager
1247 .create_snapshot(1, "save", &conversation, &files, None, None, None)
1248 .await?
1249 .expect("metadata");
1250
1251 fs::remove_file(&file_path)?;
1252 manager.restore_snapshot(1, RevertScope::Code).await?.expect("restore");
1253 assert!(file_path.exists());
1254 let content = fs::read_to_string(&file_path)?;
1255 assert_eq!(content, "data");
1256 Ok(())
1257 }
1258
1259 #[tokio::test]
1260 async fn cleanup_respects_limit() -> Result<()> {
1261 let (_dir, manager) = setup_manager();
1262 let conversation = vec![SessionMessage::new(crate::llm::provider::MessageRole::User, "hi")];
1263 let files = BTreeSet::new();
1264
1265 for turn in 1..=5 {
1266 manager
1267 .create_snapshot(turn, "turn", &conversation, &files, None, None, None)
1268 .await?
1269 .expect("metadata");
1270 }
1271
1272 let mut config = SnapshotConfig::new(manager.workspace.clone());
1274 config.max_snapshots = 3;
1275 let trimmed = SnapshotManager::new(config)?;
1276 trimmed.cleanup_old_snapshots().await?;
1277 let listed = trimmed.list_snapshots().await?;
1278 assert_eq!(listed.len(), 3);
1279 assert_eq!(listed[0].turn_number, 5);
1280 assert_eq!(listed[2].turn_number, 3);
1281 Ok(())
1282 }
1283
1284 #[tokio::test]
1285 async fn cleanup_preserves_navigation_referenced_turns() -> Result<()> {
1286 let (_dir, manager) = setup_manager();
1287 let conversation = vec![SessionMessage::new(crate::llm::provider::MessageRole::User, "nav")];
1288 let files = BTreeSet::new();
1289
1290 for turn in 1..=5 {
1291 manager
1292 .create_snapshot(turn, "turn", &conversation, &files, None, None, None)
1293 .await?
1294 .expect("metadata");
1295 }
1296
1297 let branch = serde_json::json!({
1300 "active": [1, 2],
1301 "redo": [{ "policy": "", "snapshot": "00000000-0000-0000-0000-000000000000", "active": [1] }],
1302 "pending": null,
1303 });
1304 fs::write(manager.storage_dir.join("branch_74657374.json"), serde_json::to_vec(&branch)?)?;
1305
1306 let mut config = SnapshotConfig::new(manager.workspace.clone());
1307 config.max_snapshots = 1;
1308 let trimmed = SnapshotManager::new(config)?;
1309 trimmed.cleanup_old_snapshots().await?;
1310
1311 assert!(trimmed.load_snapshot(1).await?.is_some(), "active turn must survive");
1312 assert!(trimmed.load_snapshot(2).await?.is_some(), "active turn must survive");
1313 assert!(trimmed.load_snapshot(3).await?.is_none(), "unreferenced turn must be pruned");
1316 assert!(trimmed.load_snapshot(4).await?.is_none(), "unreferenced turn must be pruned");
1317 assert!(trimmed.load_snapshot(5).await?.is_some(), "newest unreferenced turn fits the budget");
1318 Ok(())
1319 }
1320
1321 #[tokio::test]
1322 async fn snapshot_normalizes_absolute_paths() -> Result<()> {
1323 let (dir, manager) = setup_manager();
1324 let workspace = dir.path();
1325 let absolute = workspace.join("abs.txt");
1326 fs::write(&absolute, "contents")?;
1327
1328 let mut files = BTreeSet::new();
1329 files.insert(absolute.clone());
1330 let conversation = vec![SessionMessage::new(crate::llm::provider::MessageRole::User, "absolute")];
1331
1332 manager
1333 .create_snapshot(1, "abs", &conversation, &files, None, None, None)
1334 .await?
1335 .expect("metadata");
1336
1337 let stored = manager.load_snapshot(1).await?.expect("stored snapshot");
1338 assert_eq!(stored.files.len(), 1);
1339 assert_eq!(stored.files[0].path, "abs.txt");
1340 assert!(!stored.files[0].deleted);
1341 Ok(())
1342 }
1343
1344 #[tokio::test]
1345 async fn cleanup_removes_expired_snapshots() -> Result<()> {
1346 let (_dir, manager) = setup_manager();
1347 let conversation = vec![SessionMessage::new(crate::llm::provider::MessageRole::User, "cleanup")];
1348 let files = BTreeSet::new();
1349
1350 manager
1351 .create_snapshot(1, "old", &conversation, &files, None, None, None)
1352 .await?
1353 .expect("metadata");
1354
1355 let snapshot_path = manager.snapshot_path(1);
1356 let mut stored: StoredSnapshot = serde_json::from_slice(&fs::read(&snapshot_path)?)?;
1357 stored.metadata.created_at = 1;
1358 let updated = serde_json::to_vec_pretty(&stored)?;
1359 fs::write(&snapshot_path, updated)?;
1360
1361 let mut config = SnapshotConfig::new(manager.workspace.clone());
1362 config.max_age_days = Some(1);
1363 let janitor = SnapshotManager::new(config)?;
1364 janitor.cleanup_old_snapshots().await?;
1365
1366 assert!(janitor.load_snapshot(1).await?.is_none());
1367 Ok(())
1368 }
1369
1370 #[tokio::test]
1371 async fn snapshot_persists_prompt_metadata() -> Result<()> {
1372 let (_dir, manager) = setup_manager();
1373 let conversation = vec![
1374 SessionMessage::new(crate::llm::provider::MessageRole::User, "Explain checkpointing"),
1375 SessionMessage::new(crate::llm::provider::MessageRole::Assistant, "Working on it"),
1376 ];
1377
1378 manager
1379 .create_snapshot(
1380 1,
1381 "assistant reply",
1382 &conversation,
1383 &BTreeSet::new(),
1384 Some("Explain checkpointing"),
1385 Some(0),
1386 None,
1387 )
1388 .await?
1389 .expect("metadata");
1390
1391 let stored = manager.load_snapshot(1).await?.expect("stored snapshot");
1392 assert_eq!(stored.metadata.prompt_text.as_deref(), Some("Explain checkpointing"));
1393 assert_eq!(stored.metadata.prompt_message_index, Some(0));
1394 assert_eq!(stored.metadata.description, "Explain checkpointing");
1395 Ok(())
1396 }
1397
1398 #[tokio::test]
1399 async fn load_snapshot_hydrates_prompt_metadata_for_legacy_files() -> Result<()> {
1400 let (_dir, manager) = setup_manager();
1401 let stored = StoredSnapshot {
1402 metadata: SnapshotMetadata {
1403 id: "turn_1".to_string(),
1404 turn_number: 1,
1405 created_at: 1,
1406 description: "legacy".to_string(),
1407 message_count: 2,
1408 file_count: 0,
1409 touched_files: Vec::new(),
1410 prompt_text: None,
1411 prompt_message_index: None,
1412 session_id: None,
1413 runtime_turn_id: None,
1414 session_turn_number: None,
1415 turn_diagnostics: None,
1416 },
1417 conversation: vec![
1418 SessionMessage::new(crate::llm::provider::MessageRole::User, "Legacy prompt"),
1419 SessionMessage::new(crate::llm::provider::MessageRole::Assistant, "Legacy reply"),
1420 ],
1421 files: Vec::new(),
1422 schema_version: None,
1423 };
1424 let path = manager.snapshot_path(1);
1425 if let Some(parent) = path.parent() {
1426 fs::create_dir_all(parent)?;
1427 }
1428 fs::write(path, serde_json::to_vec_pretty(&stored)?)?;
1429
1430 let loaded = manager.load_snapshot(1).await?.expect("loaded snapshot");
1431 assert_eq!(loaded.metadata.prompt_text.as_deref(), Some("Legacy prompt"));
1432 assert_eq!(loaded.metadata.prompt_message_index, Some(0));
1433 Ok(())
1434 }
1435
1436 #[test]
1437 fn legacy_snapshot_versions_migrate_without_inventing_diagnostics() -> Result<()> {
1438 let legacy = StoredSnapshot {
1439 metadata: SnapshotMetadata {
1440 id: "turn_1".to_string(),
1441 turn_number: 1,
1442 created_at: 1,
1443 description: "legacy".to_string(),
1444 message_count: 0,
1445 file_count: 0,
1446 touched_files: Vec::new(),
1447 prompt_text: None,
1448 prompt_message_index: None,
1449 session_id: None,
1450 runtime_turn_id: None,
1451 session_turn_number: None,
1452 turn_diagnostics: None,
1453 },
1454 conversation: Vec::new(),
1455 files: Vec::new(),
1456 schema_version: None,
1457 };
1458
1459 let migrated_v0 = legacy.clone().migrate(SNAPSHOT_SCHEMA_VERSION)?;
1460 assert_eq!(migrated_v0.schema_version, Some(SNAPSHOT_SCHEMA_VERSION));
1461 assert!(migrated_v0.metadata.turn_diagnostics.is_none());
1462
1463 let migrated_v1 =
1464 StoredSnapshot { schema_version: Some(SchemaVersion::V1), ..legacy }.migrate(SNAPSHOT_SCHEMA_VERSION)?;
1465 assert_eq!(migrated_v1.schema_version, Some(SNAPSHOT_SCHEMA_VERSION));
1466 assert!(migrated_v1.metadata.turn_diagnostics.is_none());
1467 Ok(())
1468 }
1469
1470 #[tokio::test]
1471 async fn v2_snapshot_round_trips_session_linkage_and_canonical_usage() -> Result<()> {
1472 let (_dir, manager) = setup_manager();
1473 let usage = Usage {
1474 input_tokens: 43_282,
1475 output_tokens: 911,
1476 cached_input_tokens: 8_000,
1477 cache_creation_tokens: 512,
1478 };
1479 let diagnostics = SnapshotTurnDiagnostics {
1480 usage: usage.clone(),
1481 elapsed_ms: 12_345,
1482 requested_tool_calls: 9,
1483 admitted_tool_calls: 7,
1484 unadmitted_tool_calls: 2,
1485 failed_tool_calls: 2,
1486 denied_tool_calls: 1,
1487 preflight_failures: 1,
1488 reused_results: 2,
1489 spooled_results: 3,
1490 raw_spooled_bytes: 139_000,
1491 model_visible_output_bytes: 6_100,
1492 suppressed_tool_previews: 5,
1493 model_visible_tool_preview_budget_exhausted: true,
1494 low_signal_tool_calls: 4,
1495 recovery_activations: 1,
1496 in_progress_exec_sessions: vec![CompactStr::from("run-42")],
1497 };
1498 let context = SnapshotTurnContext {
1499 session_id: Some(CompactStr::from("session-911")),
1500 runtime_turn_id: Some(CompactStr::from("turn-runtime-911")),
1501 session_turn_number: Some(911),
1502 turn_diagnostics: Some(diagnostics.clone()),
1503 touched_files: Vec::new(),
1504 };
1505
1506 manager
1507 .create_snapshot(1, "diagnostic checkpoint", &[], &BTreeSet::new(), None, None, Some(context))
1508 .await?
1509 .expect("metadata");
1510
1511 let stored = manager.load_snapshot(1).await?.expect("stored snapshot");
1512 assert_eq!(stored.schema_version, Some(SNAPSHOT_SCHEMA_VERSION));
1513 assert_eq!(stored.metadata.session_id.as_deref(), Some("session-911"));
1514 assert_eq!(stored.metadata.runtime_turn_id.as_deref(), Some("turn-runtime-911"));
1515 assert_eq!(stored.metadata.session_turn_number, Some(911));
1516 assert_eq!(stored.metadata.turn_diagnostics, Some(diagnostics));
1517 assert_eq!(stored.metadata.turn_diagnostics.expect("diagnostics").usage, usage);
1518 Ok(())
1519 }
1520
1521 #[test]
1522 fn parse_revert_scope_variants() {
1523 assert_eq!(SnapshotManager::parse_revert_scope("conversation"), Some(RevertScope::Conversation));
1524 assert_eq!(SnapshotManager::parse_revert_scope("code"), Some(RevertScope::Code));
1525 assert_eq!(SnapshotManager::parse_revert_scope("full"), Some(RevertScope::Both));
1526 assert_eq!(SnapshotManager::parse_revert_scope("unknown"), None);
1527 }
1528 #[tokio::test]
1529 async fn binary_checkpoint_reuses_content_and_does_not_delete_untracked_neighbors() -> Result<()> {
1530 let (_dir, manager) = setup_manager();
1531 let file = manager.workspace.join("asset.bin");
1532 fs::write(&file, [0, 255, 4])?;
1533 let files = BTreeSet::from([PathBuf::from("asset.bin"), PathBuf::from("created.bin")]);
1534 manager.create_snapshot(1, "binary", &[], &files, None, None, None).await?;
1535 manager.create_snapshot(2, "same", &[], &files, None, None, None).await?;
1536 let store = filesnap::WorkspaceStore::open(&manager.storage_dir, &manager.workspace)?;
1537 let first = manager.load_snapshot(1).await?.expect("first");
1538 let second = manager.load_snapshot(2).await?.expect("second");
1539 let first_target = store
1540 .target_for_turn(first.files[0].data.as_deref().expect("ref"))?
1541 .expect("target");
1542 let second_target = store
1543 .target_for_turn(second.files[0].data.as_deref().expect("ref"))?
1544 .expect("target");
1545 let first_manifest = store.manifest(first_target.manifest_id())?;
1546 let second_manifest = store.manifest(second_target.manifest_id())?;
1547 let first_hash = &first_manifest.entries.values().next().expect("file").hash;
1548 let second_hash = &second_manifest.entries.values().next().expect("file").hash;
1549 assert_eq!(first_hash, second_hash);
1550 fs::write(&file, b"changed")?;
1551 fs::write(manager.workspace.join("created.bin"), b"created")?;
1552 fs::write(manager.workspace.join("neighbor"), b"keep")?;
1553 manager.restore_snapshot(1, RevertScope::Code).await?;
1554 assert_eq!(fs::read(file)?, [0, 255, 4]);
1555 assert!(!manager.workspace.join("created.bin").exists());
1556 assert_eq!(fs::read(manager.workspace.join("neighbor"))?, b"keep");
1557 Ok(())
1558 }
1559
1560 #[cfg(unix)]
1561 #[tokio::test]
1562 async fn refuses_symlink_escape_before_restoring_any_file() -> Result<()> {
1563 let (dir, manager) = setup_manager();
1564 fs::write(manager.workspace.join("a.txt"), "before")?;
1565 fs::create_dir(manager.workspace.join("nested"))?;
1566 fs::write(manager.workspace.join("nested/b.txt"), "before")?;
1567 let files = BTreeSet::from([PathBuf::from("a.txt"), PathBuf::from("nested/b.txt")]);
1568 manager.create_snapshot(1, "paths", &[], &files, None, None, None).await?;
1569 fs::write(manager.workspace.join("a.txt"), "after")?;
1570 fs::remove_dir_all(manager.workspace.join("nested"))?;
1571 let outside = dir.path().join("outside");
1572 fs::create_dir(&outside)?;
1573 fs::write(outside.join("b.txt"), "outside")?;
1574 std::os::unix::fs::symlink(&outside, manager.workspace.join("nested"))?;
1575 assert!(manager.restore_snapshot(1, RevertScope::Code).await.is_err());
1576 assert_eq!(fs::read_to_string(manager.workspace.join("a.txt"))?, "after");
1577 assert_eq!(fs::read_to_string(outside.join("b.txt"))?, "outside");
1578 Ok(())
1579 }
1580
1581 #[tokio::test]
1582 async fn legacy_inline_contents_still_restore() -> Result<()> {
1583 let (_dir, manager) = setup_manager();
1584 fs::write(manager.workspace.join("old.txt"), "original")?;
1585 let files = BTreeSet::from([PathBuf::from("old.txt")]);
1586 manager.create_snapshot(1, "legacy", &[], &files, None, None, None).await?;
1587 let mut stored = manager.load_snapshot(1).await?.expect("snapshot");
1588 stored.schema_version = Some(SchemaVersion::V2);
1589 stored.files[0].encoding = Some(FileEncoding::Utf8);
1590 stored.files[0].data = Some("legacy".into());
1591 fs::write(manager.snapshot_path(1), serde_json::to_vec(&stored)?)?;
1592 manager.restore_snapshot(1, RevertScope::Code).await?;
1593 assert_eq!(fs::read_to_string(manager.workspace.join("old.txt"))?, "legacy");
1594 Ok(())
1595 }
1596
1597 #[tokio::test]
1598 async fn retention_releases_engine_sessions_but_preserves_live_content() -> Result<()> {
1599 let (_dir, manager) = setup_manager();
1600 let path = manager.workspace.join("retained.txt");
1601 fs::write(&path, "shared")?;
1602 let files = BTreeSet::from([PathBuf::from("retained.txt")]);
1603 manager.create_snapshot(1, "first", &[], &files, None, None, None).await?;
1604 let first = manager.load_snapshot(1).await?.expect("first").files[0]
1605 .data
1606 .clone()
1607 .expect("reference");
1608 manager.create_snapshot(2, "second", &[], &files, None, None, None).await?;
1609 let second = manager.load_snapshot(2).await?.expect("second").files[0]
1610 .data
1611 .clone()
1612 .expect("reference");
1613 let mut config = SnapshotConfig::new(manager.workspace.clone());
1614 config.max_snapshots = 1;
1615 let janitor = SnapshotManager::new(config)?;
1616 janitor.cleanup_old_snapshots().await?;
1617 let store = filesnap::WorkspaceStore::open(&manager.storage_dir, &manager.workspace)?;
1618 assert_eq!(store.sessions()?, vec![second.clone()]);
1619 assert!(store.target_for_turn(&first)?.is_none());
1620 fs::write(&path, "changed")?;
1621 janitor.restore_snapshot(2, RevertScope::Code).await?;
1622 assert_eq!(fs::read_to_string(&path)?, "shared");
1623 janitor.create_snapshot(2, "replacement", &[], &files, None, None, None).await?;
1624 assert!(store.target_for_turn(&second)?.is_none());
1625 assert_eq!(store.sessions()?.len(), 1);
1626 Ok(())
1627 }
1628}