Skip to main content

phi_telemetry/
storage.rs

1//! Metrics persistence — save, load, and list session metrics from disk.
2
3use std::path::Path;
4
5use anyhow::{Context, Result};
6
7use crate::types::{SessionMetrics, SessionSummary};
8
9/// Metrics file name within a session directory.
10const METRICS_FILE: &str = "session_metrics.json";
11
12/// Incrementally write session metrics to `session_metrics.json`.
13///
14/// Merges with any existing file first: a second writer over the same session
15/// (TUI session switch, a new process resuming it) must extend the per-turn
16/// history rather than clobber it — a plain overwrite once erased a whole
17/// session's stats when a resume started a fresh accumulator over the same
18/// directory. Union is keyed on `(turn_number, started_at)`, so the per-turn
19/// saves from a single process stay idempotent.
20pub fn save_metrics(metrics: &SessionMetrics, session_dir: &Path) -> Result<()> {
21    let path = session_dir.join(METRICS_FILE);
22    let merged = match try_load_metrics(session_dir) {
23        Some(mut existing) => {
24            existing.merge_turns(metrics);
25            // The writer knows its own final state (e.g. `finalize` set
26            // Interrupted right before this save) — let it win.
27            existing.outcome = metrics.outcome.clone();
28            existing
29        }
30        None => metrics.clone(),
31    };
32    let json = serde_json::to_string_pretty(&merged)?;
33    std::fs::write(&path, json)?;
34    tracing::debug!(
35        path = %path.display(),
36        turns = merged.total_turns,
37        "session_metrics saved"
38    );
39    Ok(())
40}
41
42/// Load session metrics from `session_metrics.json` in a session directory.
43pub fn load_metrics(session_dir: &Path) -> Result<SessionMetrics> {
44    let path = session_dir.join(METRICS_FILE);
45    let content = std::fs::read_to_string(&path)
46        .with_context(|| format!("Failed to read metrics file: {}", path.display()))?;
47    let metrics: SessionMetrics = serde_json::from_str(&content)
48        .with_context(|| format!("Failed to parse metrics file: {}", path.display()))?;
49    Ok(metrics)
50}
51
52/// Try to load session metrics, returning `None` if the file doesn't exist.
53pub fn try_load_metrics(session_dir: &Path) -> Option<SessionMetrics> {
54    load_metrics(session_dir).ok()
55}
56
57/// List all session summaries by scanning the sessions directory.
58/// Reads each `session_metrics.json` and extracts summary fields.
59pub fn list_all_metrics(base_dir: &Path) -> Result<Vec<SessionSummary>> {
60    let sessions_dir = base_dir.join("sessions");
61    if !sessions_dir.exists() {
62        return Ok(Vec::new());
63    }
64
65    let mut summaries = Vec::new();
66    for entry in std::fs::read_dir(&sessions_dir)? {
67        let entry = entry?;
68        let path = entry.path();
69        if !path.is_dir() {
70            continue;
71        }
72
73        if let Some(metrics) = try_load_metrics(&path) {
74            let product = metrics
75                .custom
76                .get("product")
77                .and_then(|v| v.as_str())
78                .map(|s| s.to_string());
79
80            summaries.push(SessionSummary {
81                session_id: metrics.session_id,
82                node_id: metrics.node_id,
83                created_at: metrics.created_at,
84                model: metrics.model,
85                total_turns: metrics.total_turns,
86                total_chars: metrics.total_chars,
87                outcome: metrics.outcome,
88                product,
89            });
90        }
91    }
92
93    // Sort by created_at descending (newest first)
94    summaries.sort_by(|a, b| b.created_at.cmp(&a.created_at));
95    Ok(summaries)
96}
97
98#[cfg(test)]
99mod tests {
100    use super::*;
101    use crate::types::{SessionMetrics, SessionOutcome, TurnMetrics, TurnOutcome};
102    use serde_json::Value;
103    use std::fs;
104
105    /// Helper: create a temporary directory for tests.
106    fn temp_dir(name: &str) -> std::path::PathBuf {
107        let dir = std::env::temp_dir().join(format!(
108            "phi_telemetry_test_{}_{}",
109            name,
110            std::time::SystemTime::now()
111                .duration_since(std::time::UNIX_EPOCH)
112                .unwrap()
113                .as_nanos()
114        ));
115        fs::create_dir_all(&dir).unwrap();
116        dir
117    }
118
119    /// Helper: build a sample SessionMetrics with one turn.
120    fn sample_metrics(session_id: &str, model: &str) -> SessionMetrics {
121        let mut m = SessionMetrics::new(
122            session_id.to_string(),
123            "node-1".to_string(),
124            model.to_string(),
125        );
126        let turn = TurnMetrics::new(
127            1,
128            "2026-08-01T12:00:00Z".to_string(),
129            1000,
130            model.to_string(),
131            "test input".to_string(),
132            TurnOutcome::Completed,
133        );
134        m.append_turn(turn);
135        m
136    }
137
138    // ── save_metrics ──
139
140    #[test]
141    fn save_metrics_creates_file() {
142        let dir = temp_dir("save_creates");
143        let m = sample_metrics("s1", "gpt-4o");
144        save_metrics(&m, &dir).unwrap();
145
146        let path = dir.join(METRICS_FILE);
147        assert!(path.exists());
148
149        // Verify content is valid JSON
150        let content = fs::read_to_string(&path).unwrap();
151        let parsed: SessionMetrics = serde_json::from_str(&content).unwrap();
152        assert_eq!(parsed.session_id, "s1");
153
154        let _ = fs::remove_dir_all(&dir);
155    }
156
157    #[test]
158    fn save_metrics_merges_instead_of_clobbering() {
159        let dir = temp_dir("save_merges");
160        // First writer records turn 1.
161        let mut m1 = sample_metrics("s1", "gpt-4o");
162        m1.outcome = SessionOutcome::Cancelled;
163        save_metrics(&m1, &dir).unwrap();
164
165        // A second writer (fresh accumulator, e.g. a resume) records turn 2
166        // over the same directory. Its own outcome is the default Completed.
167        let mut m2 = SessionMetrics::new(
168            "s1".to_string(),
169            "node-1".to_string(),
170            "claude-sonnet".to_string(),
171        );
172        let turn2 = TurnMetrics::new(
173            2,
174            "2026-08-01T12:01:00Z".to_string(),
175            2000,
176            "gpt-4o".to_string(),
177            "second input".to_string(),
178            TurnOutcome::Completed,
179        );
180        m2.append_turn(turn2);
181        save_metrics(&m2, &dir).unwrap();
182
183        let loaded = load_metrics(&dir).unwrap();
184        // Both writers' turns survive — nothing is erased.
185        assert_eq!(loaded.total_turns, 2);
186        assert_eq!(loaded.turns.len(), 2);
187        assert_eq!(loaded.turns[0].turn_number, 1);
188        assert_eq!(loaded.turns[1].turn_number, 2);
189        // Aggregates are rebuilt over the union.
190        assert_eq!(loaded.total_duration_ms, 3000);
191        // The last writer's outcome wins.
192        assert_eq!(loaded.outcome, SessionOutcome::Completed);
193
194        // Re-saving the same data is idempotent (per-turn saves from one
195        // process must not double-count).
196        save_metrics(&m2, &dir).unwrap();
197        let loaded = load_metrics(&dir).unwrap();
198        assert_eq!(loaded.total_turns, 2);
199
200        let _ = fs::remove_dir_all(&dir);
201    }
202
203    // ── load_metrics ──
204
205    #[test]
206    fn load_metrics_roundtrip() {
207        let dir = temp_dir("load_roundtrip");
208        let original = sample_metrics("s1", "gpt-4o");
209        save_metrics(&original, &dir).unwrap();
210
211        let loaded = load_metrics(&dir).unwrap();
212        assert_eq!(loaded.session_id, "s1");
213        assert_eq!(loaded.model, "gpt-4o");
214        assert_eq!(loaded.total_turns, 1);
215        assert_eq!(loaded.turns.len(), 1);
216        assert_eq!(loaded.node_id, "node-1");
217
218        let _ = fs::remove_dir_all(&dir);
219    }
220
221    #[test]
222    fn load_metrics_file_not_found() {
223        let dir = temp_dir("load_not_found");
224        let result = load_metrics(&dir);
225        assert!(result.is_err());
226        assert!(
227            result
228                .unwrap_err()
229                .to_string()
230                .contains("Failed to read metrics file")
231        );
232
233        let _ = fs::remove_dir_all(&dir);
234    }
235
236    #[test]
237    fn load_metrics_invalid_json() {
238        let dir = temp_dir("load_invalid");
239        fs::write(dir.join(METRICS_FILE), "not valid json").unwrap();
240
241        let result = load_metrics(&dir);
242        assert!(result.is_err());
243        assert!(
244            result
245                .unwrap_err()
246                .to_string()
247                .contains("Failed to parse metrics file")
248        );
249
250        let _ = fs::remove_dir_all(&dir);
251    }
252
253    // ── try_load_metrics ──
254
255    #[test]
256    fn try_load_metrics_returns_some_when_exists() {
257        let dir = temp_dir("try_some");
258        let m = sample_metrics("s1", "gpt-4o");
259        save_metrics(&m, &dir).unwrap();
260
261        let result = try_load_metrics(&dir);
262        assert!(result.is_some());
263        assert_eq!(result.unwrap().session_id, "s1");
264
265        let _ = fs::remove_dir_all(&dir);
266    }
267
268    #[test]
269    fn try_load_metrics_returns_none_when_missing() {
270        let dir = temp_dir("try_none");
271        let result = try_load_metrics(&dir);
272        assert!(result.is_none());
273
274        let _ = fs::remove_dir_all(&dir);
275    }
276
277    #[test]
278    fn try_load_metrics_returns_none_on_invalid_json() {
279        let dir = temp_dir("try_invalid");
280        fs::write(dir.join(METRICS_FILE), "not json").unwrap();
281
282        let result = try_load_metrics(&dir);
283        assert!(result.is_none());
284
285        let _ = fs::remove_dir_all(&dir);
286    }
287
288    // ── list_all_metrics ──
289
290    #[test]
291    fn list_all_metrics_empty_when_no_sessions_dir() {
292        let dir = temp_dir("list_no_sessions");
293        let result = list_all_metrics(&dir).unwrap();
294        assert!(result.is_empty());
295
296        let _ = fs::remove_dir_all(&dir);
297    }
298
299    #[test]
300    fn list_all_metrics_returns_summaries_sorted_by_date() {
301        let base = temp_dir("list_sorted");
302        let sessions_dir = base.join("sessions");
303        fs::create_dir_all(&sessions_dir).unwrap();
304
305        // Create two sessions with different created_at
306        let dir1 = sessions_dir.join("s1");
307        fs::create_dir_all(&dir1).unwrap();
308        let mut m1 = sample_metrics("s1", "gpt-4o");
309        m1.created_at = "2026-01-01T00:00:00Z".to_string();
310        save_metrics(&m1, &dir1).unwrap();
311
312        let dir2 = sessions_dir.join("s2");
313        fs::create_dir_all(&dir2).unwrap();
314        let mut m2 = sample_metrics("s2", "claude-sonnet");
315        m2.created_at = "2026-08-01T00:00:00Z".to_string();
316        save_metrics(&m2, &dir2).unwrap();
317
318        let summaries = list_all_metrics(&base).unwrap();
319        assert_eq!(summaries.len(), 2);
320        // Sorted descending by created_at — s2 first
321        assert_eq!(summaries[0].session_id, "s2");
322        assert_eq!(summaries[1].session_id, "s1");
323        assert_eq!(summaries[0].model, "claude-sonnet");
324        assert_eq!(summaries[0].total_turns, 1);
325
326        let _ = fs::remove_dir_all(&base);
327    }
328
329    #[test]
330    fn list_all_metrics_skips_non_session_dirs() {
331        let base = temp_dir("list_skip");
332        let sessions_dir = base.join("sessions");
333        fs::create_dir_all(&sessions_dir).unwrap();
334
335        // A subdirectory without metrics — should be skipped gracefully
336        let empty_dir = sessions_dir.join("empty_session");
337        fs::create_dir_all(&empty_dir).unwrap();
338
339        // A regular file (not a dir) — should be skipped
340        fs::write(sessions_dir.join("not_a_dir.txt"), "data").unwrap();
341
342        // A valid session
343        let valid_dir = sessions_dir.join("valid");
344        fs::create_dir_all(&valid_dir).unwrap();
345        save_metrics(&sample_metrics("valid", "gpt-4o"), &valid_dir).unwrap();
346
347        let summaries = list_all_metrics(&base).unwrap();
348        assert_eq!(summaries.len(), 1);
349        assert_eq!(summaries[0].session_id, "valid");
350
351        let _ = fs::remove_dir_all(&base);
352    }
353
354    #[test]
355    fn list_all_metrics_extracts_product_from_custom() {
356        let base = temp_dir("list_product");
357        let sessions_dir = base.join("sessions");
358        let sdir = sessions_dir.join("s1");
359        fs::create_dir_all(&sdir).unwrap();
360
361        let mut m = sample_metrics("s1", "gpt-4o");
362        if let Value::Object(ref mut map) = m.custom {
363            map.insert("product".to_string(), Value::String("phi-bard".to_string()));
364        }
365        save_metrics(&m, &sdir).unwrap();
366
367        let summaries = list_all_metrics(&base).unwrap();
368        assert_eq!(summaries.len(), 1);
369        assert_eq!(summaries[0].product, Some("phi-bard".to_string()));
370
371        let _ = fs::remove_dir_all(&base);
372    }
373
374    #[test]
375    fn list_all_metrics_no_product_when_custom_empty() {
376        let base = temp_dir("list_no_product");
377        let sessions_dir = base.join("sessions");
378        let sdir = sessions_dir.join("s1");
379        fs::create_dir_all(&sdir).unwrap();
380
381        let m = sample_metrics("s1", "gpt-4o");
382        save_metrics(&m, &sdir).unwrap();
383
384        let summaries = list_all_metrics(&base).unwrap();
385        assert_eq!(summaries[0].product, None);
386
387        let _ = fs::remove_dir_all(&base);
388    }
389
390    #[test]
391    fn list_all_metrics_summary_fields_are_correct() {
392        let base = temp_dir("list_fields");
393        let sessions_dir = base.join("sessions");
394        let sdir = sessions_dir.join("s1");
395        fs::create_dir_all(&sdir).unwrap();
396
397        let mut m = sample_metrics("s1", "gpt-4o");
398        m.node_id = "node-42".to_string();
399        m.created_at = "2026-06-15T10:00:00Z".to_string();
400        m.total_chars = 5000;
401        m.outcome = SessionOutcome::Completed;
402        save_metrics(&m, &sdir).unwrap();
403
404        let summaries = list_all_metrics(&base).unwrap();
405        let s = &summaries[0];
406        assert_eq!(s.session_id, "s1");
407        assert_eq!(s.node_id, "node-42");
408        assert_eq!(s.created_at, "2026-06-15T10:00:00Z");
409        assert_eq!(s.model, "gpt-4o");
410        assert_eq!(s.total_turns, 1);
411        assert_eq!(s.total_chars, 5000);
412        assert_eq!(s.outcome, SessionOutcome::Completed);
413
414        let _ = fs::remove_dir_all(&base);
415    }
416}