Skip to main content

vtcode_core/core/agent/
snapshots.rs

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;
30/// How many newest active turns a finished session keeps pinned for rewind.
31/// Older entries are dropped on thread completion so their snapshots become
32/// prune-eligible instead of sitting protected for the full age window.
33pub const REWIND_ACTIVE_KEEP: usize = 5;
34const SNAPSHOT_SCHEMA_VERSION: SchemaVersion = SchemaVersion(3);
35
36/// Collect turn numbers from a navigation value: either a bare array of turns
37/// (`active`) or a recovery object with an `active` array (`redo`/`pending`).
38fn 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    /// Exec sessions still running when the turn ended (bounded, newest
153    /// first). Downstream consumers use this to correlate cross-turn resume
154    /// hints; empty when every command settled within the turn.
155    #[serde(default, deserialize_with = "deserialize_null_as_default")]
156    pub in_progress_exec_sessions: Vec<CompactStr>,
157}
158
159impl SnapshotTurnDiagnostics {
160    /// Attach the in-progress exec session ids captured at turn end.
161    ///
162    /// Kept as a builder step because the ids come from the live exec-session
163    /// registry (async), not from the synchronous turn-state counters.
164    /// Bound shares `vtcode_exec_events::MAX_IN_PROGRESS_EXEC_SESSIONS` with
165    /// `TurnCompletedEvent` so checkpoint and event streams cannot drift.
166    #[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    /// A filesnap turn reference, not inline file contents. Old readers reject
191    /// this variant rather than treating an absent payload as an empty file.
192    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    /// Schema version for forward/backward migration. `None` means legacy (v0).
211    #[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                // v0 -> v1: set the explicit schema version; no structural changes yet.
224                // Future versions add per-message metadata here.
225                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/// A snapshot of the agent state taken before executing a single action (tool call).
248///
249/// Unlike turn-level snapshots that capture the full conversation and file state,
250/// action snapshots are lightweight records capturing only the state delta before
251/// a specific tool call. This enables incremental undo of the last N actions.
252#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
253pub struct ActionSnapshot {
254    /// Matches the tool_call_id for this action.
255    pub action_id: String,
256    /// Sequential counter — strictly increasing per session.
257    pub action_number: usize,
258    /// Unix timestamp (seconds) when the snapshot was created.
259    pub created_at: u64,
260    /// Name of the tool being invoked.
261    pub tool_name: String,
262    /// JSON arguments passed to the tool.
263    pub arguments: serde_json::Value,
264    /// The number of messages in the conversation *before* this action was executed.
265    /// Used to truncate the messages vector on rollback.
266    pub pre_action_message_count: usize,
267    /// Files we know were touched by this action (from modified_files).
268    pub touched_files: Vec<String>,
269    /// Expected outcome category, used to determine rollback strategy.
270    pub expected_outcome: ExpectedOutcome,
271}
272
273/// Result of a rollback operation — describes what was undone.
274#[derive(Debug, Clone)]
275pub struct RollbackResult {
276    /// The action_id of the action that was rolled back.
277    pub rollback_action_id: String,
278    /// Number of messages removed from the conversation.
279    pub messages_removed: usize,
280    /// Number of files restored to their pre-action state.
281    pub files_restored: usize,
282    /// The action_number after which this rollback takes effect.
283    pub next_action_number: usize,
284}
285
286/// Global monotonic counter for action snapshots across sessions.
287static 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(&current) {
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); // Typical directory has ~20-50 snapshot files
409        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            // Conversation-only checkpoints need no unreferenced engine session.
533            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        // Preserve an overwritten turn as a cleanup journal until the new JSON
595        // is published. It is invisible to checkpoint enumeration.
596        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        // Migrate legacy snapshots to the current schema version
645        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            // Validate the complete path set before any write, including legacy records.
676            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    /// Turn numbers still referenced by any *recent* session navigation
767    /// (active branch, redo stack, or pending recovery). Rewind/redo must
768    /// never lose these. Branch files older than the retention window are
769    /// ignored so a completed session cannot pin its turns forever.
770    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                // Stale navigation (completed/abandoned session past retention)
790                // must not keep its turn files alive. `cutoff` is a UNIX
791                // timestamp: branch files last modified before it are stale.
792                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    /// Cheap retention prune: enforce the snapshot count budget without
816    /// reading checkpoint JSON bodies. Safe to call on the per-turn hot path
817    /// (`native::begin_prompt`); age-based expiry stays in
818    /// [`Self::cleanup_old_snapshots`], which needs each file's `created_at`.
819    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    /// Full retention: age-expire checkpoints past `max_age_days`, then apply
851    /// the count budget. Age scan reads each candidate's `created_at`, so prefer
852    /// [`Self::prune_snapshot_budget`] on hot paths.
853    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            // A replaced record can still be live if publication failed. Read
914            // every live record before deleting anything; corrupt metadata defers
915            // cleanup rather than guessing which content can be discarded.
916            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            // The per-prompt path does not need live JSON or a content store
937            // when no records were retired. Explicit maintenance still runs GC.
938            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    // ─── Action-level Snapshots (Incremental Rollback) ────────────────────
1011
1012    /// File path for an action-level snapshot.
1013    fn action_snapshot_path(&self, action_number: usize) -> PathBuf {
1014        self.storage_dir.join(format!("action_{action_number}.json"))
1015    }
1016
1017    /// Save an action-level snapshot to disk. Returns the action number.
1018    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    /// Load an action-level snapshot from disk.
1035    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    /// Rollback a single action by restoring conversation state and files.
1052    ///
1053    /// This truncates the messages vector to `pre_action_message_count` and
1054    /// restores file contents from their pre-action state (requires the workspace
1055    /// manager to have tracked file versions).
1056    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        // Restore conversation: truncate to pre-action length
1065        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        // Restore files that were touched by this action
1074        if scope.includes_code() {
1075            for file_path in &action.touched_files {
1076                // Read the current file content as a snapshot if possible
1077                let absolute = self.workspace.join(file_path);
1078                if tokio::fs::try_exists(&absolute).await.unwrap_or(false) {
1079                    // In a full implementation, we'd restore from a file version store.
1080                    // For now, we track that the file was touched and would need restoration.
1081                    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    /// Generate the next monotonic action number.
1095    ///
1096    /// Uses an atomic counter to avoid conflicts between concurrent sessions.
1097    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        // Use the versioned engine store, rather than a standalone BlobStore
1167        // that the workspace collector cannot discover.
1168        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        // Manager default limit is 50, shrink artificially for test
1273        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        // Simulate a live session navigation that still needs turns 1 and 2
1298        // for rewind/redo. Unreferenced turns 3-5 are fair game.
1299        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        // Unreferenced turns 3-5 compete for a single budget slot; the oldest
1314        // two are retired and the newest (5) remains.
1315        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}