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#[derive(Clone)]
17pub struct TrajectoryLogger {
18 enabled: bool,
19 writer: Option<Arc<AsyncLineWriter>>,
20}
21
22#[derive(Debug, Clone, Copy)]
24pub struct TrajectoryRetention {
25 pub max_files: usize,
27 pub max_age_days: u64,
29 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 pub fn new(workspace: &Path) -> Self {
47 Self::with_retention(workspace, TrajectoryRetention::default())
48 }
49
50 pub async fn new_async(workspace: &Path) -> Self {
53 Self::with_retention_async(workspace, TrajectoryRetention::default()).await
54 }
55
56 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 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 pub fn disabled() -> Self {
89 Self { enabled: false, writer: None }
90 }
91
92 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 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 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(¤t) {
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(¤t, &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(¤t, 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(¤t, "").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}