Skip to main content

vtcode_core/core/
trajectory.rs

1use serde::Serialize;
2use std::fs;
3use std::path::{Path, PathBuf};
4use std::sync::Arc;
5use std::time::{Duration, SystemTime, UNIX_EPOCH};
6
7use crate::telemetry::perf::PerfSpan;
8use crate::utils::async_line_writer::AsyncLineWriter;
9
10const TRAJECTORY_PREFIX: &str = "trajectory-";
11const TRAJECTORY_EXTENSION: &str = "jsonl";
12use super::SECONDS_PER_DAY;
13const BYTES_PER_MB: u64 = 1024 * 1024;
14
15/// Async JSONL logger for agent trajectory records (routes, tool calls).
16#[derive(Clone)]
17pub struct TrajectoryLogger {
18    enabled: bool,
19    writer: Option<Arc<AsyncLineWriter>>,
20}
21
22/// Retention policy for trajectory log files.
23#[derive(Debug, Clone, Copy)]
24pub struct TrajectoryRetention {
25    /// Maximum number of rotated trajectory files to keep.
26    pub max_files: usize,
27    /// Maximum age of trajectory files in days before pruning.
28    pub max_age_days: u64,
29    /// Maximum total size of all trajectory files in bytes.
30    pub max_total_size_bytes: u64,
31}
32
33impl Default for TrajectoryRetention {
34    fn default() -> Self {
35        use vtcode_config::constants::defaults;
36        Self {
37            max_files: defaults::DEFAULT_TRAJECTORY_MAX_FILES,
38            max_age_days: defaults::DEFAULT_TRAJECTORY_MAX_AGE_DAYS,
39            max_total_size_bytes: defaults::DEFAULT_TRAJECTORY_MAX_SIZE_MB.saturating_mul(BYTES_PER_MB),
40        }
41    }
42}
43
44impl TrajectoryLogger {
45    /// Create a trajectory logger in the workspace's `.vtcode/logs/` directory with default retention.
46    pub fn new(workspace: &Path) -> Self {
47        Self::with_retention(workspace, TrajectoryRetention::default())
48    }
49
50    /// Create a trajectory logger without blocking the async executor during
51    /// retention maintenance or file setup.
52    pub async fn new_async(workspace: &Path) -> Self {
53        Self::with_retention_async(workspace, TrajectoryRetention::default()).await
54    }
55
56    /// Create a trajectory logger with a custom retention policy.
57    pub fn with_retention(workspace: &Path, retention: TrajectoryRetention) -> Self {
58        let dir = workspace.join(".vtcode").join("logs");
59        rotate_current_trajectory(&dir);
60        prune_trajectory_logs_best_effort(&dir, retention);
61        let path = dir.join("trajectory.jsonl");
62        let writer = AsyncLineWriter::new(path).ok().map(Arc::new);
63        let enabled = writer.is_some();
64        Self { enabled, writer }
65    }
66
67    /// Create a trajectory logger asynchronously. Synchronous retention work
68    /// is isolated on Tokio's blocking pool, and the writer uses async setup.
69    pub async fn with_retention_async(workspace: &Path, retention: TrajectoryRetention) -> Self {
70        let dir = workspace.join(".vtcode").join("logs");
71        let retention_dir = dir.clone();
72        if let Err(error) = tokio::task::spawn_blocking(move || {
73            rotate_current_trajectory(&retention_dir);
74            prune_trajectory_logs_best_effort(&retention_dir, retention);
75        })
76        .await
77        {
78            tracing::debug!(error = %error, "Failed to prepare trajectory retention asynchronously");
79        }
80
81        let path = dir.join("trajectory.jsonl");
82        let writer = AsyncLineWriter::new_async(path).await.ok().map(Arc::new);
83        let enabled = writer.is_some();
84        Self { enabled, writer }
85    }
86
87    /// Create a disabled trajectory logger that discards all records.
88    pub fn disabled() -> Self {
89        Self { enabled: false, writer: None }
90    }
91
92    /// Write a serializable record to the trajectory log.
93    pub fn log<T: Serialize>(&self, record: &T) {
94        if !self.enabled {
95            return;
96        }
97        let mut perf = PerfSpan::new("vtcode.perf.trajectory_log_ms");
98        perf.tag("mode", "async");
99        if let Ok(line) = serde_json::to_string(record)
100            && let Some(writer) = self.writer.as_ref()
101        {
102            writer.write_line(line);
103        }
104    }
105
106    #[cfg(test)]
107    pub async fn flush(&self) {
108        if let Some(writer) = self.writer.as_ref() {
109            writer.flush().await;
110        }
111    }
112
113    /// Log a model routing decision for a given turn.
114    pub fn log_route(&self, turn: usize, selected_model: &str, class: &str, input_preview: &str) {
115        #[derive(Serialize)]
116        struct RouteRec<'a> {
117            kind: &'static str,
118            turn: usize,
119            selected_model: &'a str,
120            class: &'a str,
121            input_preview: &'a str,
122            ts: i64,
123        }
124        let rec = RouteRec {
125            kind: "route",
126            turn,
127            selected_model,
128            class,
129            input_preview,
130            ts: chrono::Utc::now().timestamp(),
131        };
132        self.log(&rec);
133    }
134
135    /// Log a tool call with its arguments, success status, and agent context.
136    pub fn log_tool_call(
137        &self,
138        turn: usize,
139        name: &str,
140        args: &serde_json::Value,
141        ok: bool,
142        agent_name: Option<&str>,
143        is_subagent: bool,
144    ) {
145        #[derive(Serialize)]
146        struct ToolRec<'a> {
147            kind: &'static str,
148            turn: usize,
149            name: &'a str,
150            args: serde_json::Value,
151            ok: bool,
152            agent_name: Option<&'a str>,
153            is_subagent: bool,
154            ts: i64,
155        }
156        let rec = ToolRec {
157            kind: "tool",
158            turn,
159            name,
160            args: args.clone(),
161            ok,
162            agent_name,
163            is_subagent,
164            ts: chrono::Utc::now().timestamp(),
165        };
166        self.log(&rec);
167    }
168}
169
170fn rotate_current_trajectory(dir: &Path) {
171    let current = dir.join("trajectory.jsonl");
172    if !current.exists() {
173        return;
174    }
175    let metadata = match fs::metadata(&current) {
176        Ok(m) => m,
177        Err(_) => return,
178    };
179    if metadata.len() == 0 {
180        return;
181    }
182    let timestamp = chrono::Utc::now().format("%Y%m%dT%H%M%SZ");
183    let rotated_name = format!("{TRAJECTORY_PREFIX}{timestamp}.{TRAJECTORY_EXTENSION}");
184    let rotated_path = dir.join(rotated_name);
185    if let Err(e) = fs::rename(&current, &rotated_path) {
186        tracing::warn!("Failed to rotate trajectory {} -> {}: {e}", current.display(), rotated_path.display());
187    }
188}
189
190fn is_trajectory_file(path: &Path) -> bool {
191    let name = match path.file_name().and_then(|n| n.to_str()) {
192        Some(n) => n,
193        None => return false,
194    };
195    name.starts_with(TRAJECTORY_PREFIX) && name.ends_with(&format!(".{TRAJECTORY_EXTENSION}"))
196}
197
198struct FileEntry {
199    path: PathBuf,
200    modified: SystemTime,
201    size: u64,
202}
203
204fn prune_trajectory_logs_best_effort(dir: &Path, limits: TrajectoryRetention) {
205    if let Err(err) = prune_trajectory_logs(dir, limits) {
206        tracing::debug!("Failed to prune trajectory logs in {}: {}", dir.display(), err);
207    }
208}
209
210fn prune_trajectory_logs(dir: &Path, limits: TrajectoryRetention) -> anyhow::Result<()> {
211    if !dir.exists() {
212        return Ok(());
213    }
214
215    let mut entries = Vec::new();
216    for entry in fs::read_dir(dir)? {
217        let entry = match entry {
218            Ok(e) => e,
219            Err(_) => continue,
220        };
221        let path = entry.path();
222        if !is_trajectory_file(&path) {
223            continue;
224        }
225        let metadata = match entry.metadata() {
226            Ok(m) if m.is_file() => m,
227            _ => continue,
228        };
229        entries.push(FileEntry {
230            path,
231            modified: metadata.modified().unwrap_or(UNIX_EPOCH),
232            size: metadata.len(),
233        });
234    }
235
236    if entries.is_empty() {
237        return Ok(());
238    }
239
240    let now = SystemTime::now();
241    let age_cutoff = if limits.max_age_days == 0 {
242        now
243    } else {
244        now.checked_sub(Duration::from_secs(limits.max_age_days.saturating_mul(SECONDS_PER_DAY)))
245            .unwrap_or(UNIX_EPOCH)
246    };
247
248    let (expired, mut retained): (Vec<_>, Vec<_>) = entries.into_iter().partition(|entry| entry.modified <= age_cutoff);
249    remove_files(expired);
250
251    retained.sort_by_key(|a| std::cmp::Reverse(a.modified));
252
253    if limits.max_files > 0 && retained.len() > limits.max_files {
254        let overflow = retained.split_off(limits.max_files);
255        remove_files(overflow);
256    }
257
258    if limits.max_total_size_bytes == 0 || retained.is_empty() {
259        return Ok(());
260    }
261
262    let mut total_size = 0u64;
263    let mut size_overflow = Vec::new();
264    let mut keep = Vec::with_capacity(retained.len());
265    for entry in retained {
266        let projected = total_size.saturating_add(entry.size);
267        if keep.is_empty() || projected <= limits.max_total_size_bytes {
268            total_size = projected;
269            keep.push(entry);
270        } else {
271            size_overflow.push(entry);
272        }
273    }
274    remove_files(size_overflow);
275
276    Ok(())
277}
278
279fn remove_files(entries: Vec<FileEntry>) {
280    for entry in entries {
281        if let Err(err) = fs::remove_file(&entry.path) {
282            tracing::debug!("Failed to remove trajectory log {}: {}", entry.path.display(), err);
283        }
284    }
285}
286
287#[cfg(test)]
288mod tests {
289    use super::*;
290    use tempfile::TempDir;
291
292    #[tokio::test]
293    async fn test_trajectory_logger_log_route_integration() {
294        let temp_dir = TempDir::new().unwrap();
295        let logger = TrajectoryLogger::new(temp_dir.path());
296
297        logger.log_route(1, "gemini-3-flash-preview", "standard", "test user input for logging");
298        logger.flush().await;
299
300        let log_path = temp_dir.path().join(".vtcode/logs/trajectory.jsonl");
301        assert!(log_path.exists());
302
303        let content = fs::read_to_string(log_path).unwrap();
304        let lines: Vec<&str> = content.lines().collect();
305        assert_eq!(lines.len(), 1);
306
307        let record: serde_json::Value = serde_json::from_str(lines[0]).unwrap();
308        assert_eq!(record["kind"], "route");
309        assert_eq!(record["turn"], 1);
310        assert_eq!(record["selected_model"], "gemini-3-flash-preview");
311        assert_eq!(record["class"], "standard");
312        assert_eq!(record["input_preview"], "test user input for logging");
313        assert!(record["ts"].is_number());
314    }
315
316    #[tokio::test]
317    async fn test_async_trajectory_logger_log_route_integration() {
318        let temp_dir = TempDir::new().unwrap();
319        let logger = TrajectoryLogger::new_async(temp_dir.path()).await;
320
321        logger.log_route(1, "test-model", "standard", "async setup");
322        logger.flush().await;
323
324        let log_path = temp_dir.path().join(".vtcode/logs/trajectory.jsonl");
325        assert_eq!(fs::read_to_string(log_path).unwrap().lines().count(), 1);
326    }
327
328    #[tokio::test]
329    async fn test_rotation_renames_existing_log() {
330        let temp_dir = TempDir::new().unwrap();
331        let logs_dir = temp_dir.path().join(".vtcode").join("logs");
332        fs::create_dir_all(&logs_dir).unwrap();
333
334        let current = logs_dir.join("trajectory.jsonl");
335        fs::write(&current, r#"{"kind":"route","turn":1}"#).unwrap();
336
337        let _logger = TrajectoryLogger::new(temp_dir.path());
338
339        let rotated: Vec<_> = fs::read_dir(&logs_dir)
340            .unwrap()
341            .filter_map(|e| e.ok())
342            .filter(|e| is_trajectory_file(&e.path()))
343            .collect();
344        assert_eq!(rotated.len(), 1, "Old log should be rotated");
345        assert!(current.exists(), "New current file should be created");
346    }
347
348    #[test]
349    fn test_prune_removes_old_files() {
350        let temp_dir = TempDir::new().unwrap();
351        let logs_dir = temp_dir.path().join(".vtcode").join("logs");
352        fs::create_dir_all(&logs_dir).unwrap();
353
354        for i in 0..5 {
355            let name = format!("trajectory-2024010{i}T000000Z.jsonl");
356            fs::write(logs_dir.join(name), "data").unwrap();
357        }
358
359        let limits = TrajectoryRetention {
360            max_files: 3,
361            max_age_days: 0,
362            max_total_size_bytes: 100 * BYTES_PER_MB,
363        };
364
365        prune_trajectory_logs(&logs_dir, limits).unwrap();
366
367        let remaining: Vec<_> = fs::read_dir(&logs_dir)
368            .unwrap()
369            .filter_map(|e| e.ok())
370            .filter(|e| is_trajectory_file(&e.path()))
371            .collect();
372        assert!(remaining.len() <= 3, "Should keep at most 3 files");
373    }
374
375    #[tokio::test]
376    async fn test_empty_trajectory_not_rotated() {
377        let temp_dir = TempDir::new().unwrap();
378        let logs_dir = temp_dir.path().join(".vtcode").join("logs");
379        fs::create_dir_all(&logs_dir).unwrap();
380
381        let current = logs_dir.join("trajectory.jsonl");
382        fs::write(&current, "").unwrap();
383
384        let _logger = TrajectoryLogger::new(temp_dir.path());
385
386        let rotated: Vec<_> = fs::read_dir(&logs_dir)
387            .unwrap()
388            .filter_map(|e| e.ok())
389            .filter(|e| is_trajectory_file(&e.path()))
390            .collect();
391        assert_eq!(rotated.len(), 0, "Empty file should not be rotated");
392    }
393}