Skip to main content

lean_ctx/core/stats/
mod.rs

1mod format;
2mod io;
3mod model;
4
5pub use format::*;
6pub use model::*;
7
8use std::collections::HashMap;
9use std::sync::Mutex;
10use std::time::Instant;
11
12/// (current_state, baseline_from_disk, last_flush_time)
13static STATS_BUFFER: Mutex<Option<(StatsStore, StatsStore, Instant)>> = Mutex::new(None);
14
15const FLUSH_INTERVAL_SECS: u64 = 30;
16
17/// Daily savings history retained on disk (~10 years of active days). This is a
18/// storage/sync safety bound, NOT a display limit: all-time token totals are
19/// unbounded, and the cumulative-savings chart baselines any pre-window savings
20/// so it always reaches the true all-time total regardless of this window.
21pub(super) const MAX_DAILY_HISTORY_DAYS: usize = 3650;
22
23pub fn load() -> StatsStore {
24    let guard = STATS_BUFFER
25        .lock()
26        .unwrap_or_else(std::sync::PoisonError::into_inner);
27    if let Some((ref current, ref baseline, _)) = *guard {
28        let disk = io::load_from_disk();
29        return io::apply_deltas(&disk, current, baseline);
30    }
31    drop(guard);
32    io::load_from_disk()
33}
34
35/// Loads stats for **display**, summing across every auto-resolved data dir that
36/// holds a `stats.json` (#500). When the MCP server process and the CLI resolve
37/// different XDG dirs — the documented #408/#414 split, common when an agent host
38/// (e.g. a containerised Hermes) launches the MCP server with a different `HOME`
39/// or `XDG_*` than the user's shell — the bulk of the savings can land in a
40/// sibling tree. Reading only the primary dir then makes `gain` report `0` while
41/// the real data sits one directory over. Folding the siblings in keeps the
42/// headline honest regardless of which process wrote where.
43///
44/// Safe by construction:
45/// - **No-op without a split** — when only the primary dir has stats (the
46///   overwhelmingly common case) the result equals [`load`].
47/// - **Respects an explicit pin** — when `LEAN_CTX_DATA_DIR` is set the user has
48///   chosen exactly one dir, so nothing is auto-merged.
49/// - **Read-only** — never writes back; recording still targets the primary dir.
50pub fn load_for_display() -> StatsStore {
51    let primary = load();
52    // An explicit override means "use exactly this dir" — never auto-merge.
53    if std::env::var_os("LEAN_CTX_DATA_DIR").is_some() {
54        return primary;
55    }
56    let primary_dir = crate::core::data_dir::lean_ctx_data_dir()
57        .ok()
58        .and_then(|p| std::fs::canonicalize(&p).ok());
59    let siblings: Vec<StatsStore> = crate::core::data_dir::all_data_dirs_with_stats()
60        .into_iter()
61        .filter(|d| std::fs::canonicalize(d).ok() != primary_dir)
62        .map(|d| io::load_from_dir(&d))
63        .collect();
64    aggregate_for_display(primary, &siblings)
65}
66
67/// Folds sibling-dir stores into the primary for display. Pure (no I/O, no
68/// globals) so the cross-dir summation is unit-testable. Reuses
69/// [`io::apply_deltas`] with a zero baseline so each sibling store is added in
70/// full (delta-from-empty == the whole store).
71fn aggregate_for_display(primary: StatsStore, siblings: &[StatsStore]) -> StatsStore {
72    let zero = StatsStore::default();
73    siblings
74        .iter()
75        .fold(primary, |acc, other| io::apply_deltas(&acc, other, &zero))
76}
77
78pub fn save(store: &StatsStore) {
79    io::locked_write(store);
80}
81
82fn maybe_flush(store: &mut StatsStore, baseline: &mut StatsStore, last_flush: &mut Instant) {
83    if last_flush.elapsed().as_secs() >= FLUSH_INTERVAL_SECS {
84        let merged = io::merge_and_save(store, baseline);
85        *store = merged.clone();
86        *baseline = merged;
87        *last_flush = Instant::now();
88    }
89}
90
91pub fn flush() {
92    let mut guard = STATS_BUFFER
93        .lock()
94        .unwrap_or_else(std::sync::PoisonError::into_inner);
95    if let Some((ref mut store, ref mut baseline, ref mut last_flush)) = *guard {
96        let merged = io::merge_and_save(store, baseline);
97        *store = merged.clone();
98        *baseline = merged;
99        *last_flush = Instant::now();
100    }
101}
102
103/// Adjust saved tokens after post-processing (terse, hints) changed the output size.
104/// Positive delta = savings were over-reported, negative = under-reported.
105pub fn adjust_savings(command: &str, over_report_delta: i64) {
106    let mut guard = STATS_BUFFER
107        .lock()
108        .unwrap_or_else(std::sync::PoisonError::into_inner);
109    let Some((store, _, _)) = guard.as_mut() else {
110        return;
111    };
112    if over_report_delta > 0 {
113        let adj = over_report_delta as u64;
114        store.total_output_tokens = store.total_output_tokens.saturating_add(adj);
115        if let Some(cmd) = store.commands.get_mut(command) {
116            cmd.output_tokens = cmd.output_tokens.saturating_add(adj);
117        }
118    } else {
119        let adj = over_report_delta.unsigned_abs();
120        store.total_output_tokens = store.total_output_tokens.saturating_sub(adj);
121        if let Some(cmd) = store.commands.get_mut(command) {
122            cmd.output_tokens = cmd.output_tokens.saturating_sub(adj);
123        }
124    }
125}
126
127pub fn record(command: &str, input_tokens: usize, output_tokens: usize) {
128    let mut guard = STATS_BUFFER
129        .lock()
130        .unwrap_or_else(std::sync::PoisonError::into_inner);
131    if guard.is_none() {
132        let disk = io::load_from_disk();
133        *guard = Some((disk.clone(), disk, Instant::now()));
134    }
135    let Some((store, baseline, last_flush)) = guard.as_mut() else {
136        return;
137    };
138
139    let is_first_command = store.total_commands == baseline.total_commands;
140    let now = chrono::Local::now();
141    let today = now.format("%Y-%m-%d").to_string();
142    let timestamp = now.to_rfc3339();
143
144    store.total_commands = store.total_commands.saturating_add(1);
145    store.total_input_tokens = store.total_input_tokens.saturating_add(input_tokens as u64);
146    store.total_output_tokens = store
147        .total_output_tokens
148        .saturating_add(output_tokens as u64);
149
150    if store.first_use.is_none() {
151        store.first_use = Some(timestamp.clone());
152    }
153    store.last_use = Some(timestamp);
154
155    let cmd_key = format::normalize_command(command);
156    let entry = store.commands.entry(cmd_key).or_default();
157    entry.count = entry.count.saturating_add(1);
158    entry.input_tokens = entry.input_tokens.saturating_add(input_tokens as u64);
159    entry.output_tokens = entry.output_tokens.saturating_add(output_tokens as u64);
160
161    let current_version = env!("CARGO_PKG_VERSION").to_string();
162    if let Some(day) = store.daily.last_mut() {
163        if day.date == today {
164            day.commands = day.commands.saturating_add(1);
165            day.input_tokens = day.input_tokens.saturating_add(input_tokens as u64);
166            day.output_tokens = day.output_tokens.saturating_add(output_tokens as u64);
167            // Stamp the running version so a mid-day update attributes the day
168            // to the release in use for its latest activity (#307).
169            day.version = current_version;
170        } else {
171            store.daily.push(DayStats {
172                date: today,
173                commands: 1,
174                input_tokens: input_tokens as u64,
175                output_tokens: output_tokens as u64,
176                version: current_version,
177            });
178        }
179    } else {
180        store.daily.push(DayStats {
181            date: today,
182            commands: 1,
183            input_tokens: input_tokens as u64,
184            output_tokens: output_tokens as u64,
185            version: current_version,
186        });
187    }
188
189    if store.daily.len() > MAX_DAILY_HISTORY_DAYS {
190        store
191            .daily
192            .drain(..store.daily.len() - MAX_DAILY_HISTORY_DAYS);
193    }
194
195    if is_first_command {
196        let merged = io::merge_and_save(store, baseline);
197        *store = merged.clone();
198        *baseline = merged;
199        *last_flush = Instant::now();
200    } else {
201        maybe_flush(store, baseline, last_flush);
202    }
203}
204
205pub fn reset_cep() {
206    let mut guard = STATS_BUFFER
207        .lock()
208        .unwrap_or_else(std::sync::PoisonError::into_inner);
209    let mut store = io::load_from_disk();
210    store.cep = CepStats::default();
211    io::locked_write(&store);
212    *guard = Some((store.clone(), store, Instant::now()));
213}
214
215pub fn reset_all() {
216    let mut guard = STATS_BUFFER
217        .lock()
218        .unwrap_or_else(std::sync::PoisonError::into_inner);
219    let store = StatsStore::default();
220    io::locked_write(&store);
221    *guard = Some((store.clone(), store, Instant::now()));
222    crate::core::heatmap::reset();
223}
224
225pub fn load_stats() -> GainSummary {
226    let store = load();
227    let input_saved = store
228        .total_input_tokens
229        .saturating_sub(store.total_output_tokens);
230    GainSummary {
231        total_saved: input_saved,
232        total_calls: store.total_commands,
233    }
234}
235
236#[allow(clippy::too_many_arguments)]
237pub fn record_cep_session(
238    score: u32,
239    cache_hits: u64,
240    cache_reads: u64,
241    tokens_original: u64,
242    tokens_compressed: u64,
243    modes: &HashMap<String, u64>,
244    tool_calls: u64,
245    complexity: &str,
246) {
247    let mut guard = STATS_BUFFER
248        .lock()
249        .unwrap_or_else(std::sync::PoisonError::into_inner);
250    if guard.is_none() {
251        let disk = io::load_from_disk();
252        *guard = Some((disk.clone(), disk, Instant::now()));
253    }
254    let Some((store, baseline, last_flush)) = guard.as_mut() else {
255        return;
256    };
257
258    apply_cep_snapshot(
259        &mut store.cep,
260        std::process::id(),
261        score,
262        cache_hits,
263        cache_reads,
264        tokens_original,
265        tokens_compressed,
266        modes,
267        tool_calls,
268        complexity,
269    );
270
271    maybe_flush(store, baseline, last_flush);
272}
273
274/// Fold one CEP snapshot into `cep`. Pure (no globals, no I/O) so the
275/// delta/aggregation rules are unit-testable in isolation.
276///
277/// `cache_hits`, `cache_reads`, `tokens_original` and `tokens_compressed` arrive
278/// as **cumulative per-process** counters. For repeated snapshots within the same
279/// PID only the delta since the previous snapshot is added, so the lifetime
280/// totals keep tracking cache activity instead of freezing at the first
281/// checkpoint's value (#361). A new PID starts a fresh session and seeds the
282/// cumulative baselines.
283#[allow(clippy::too_many_arguments)]
284fn apply_cep_snapshot(
285    cep: &mut CepStats,
286    pid: u32,
287    score: u32,
288    cache_hits: u64,
289    cache_reads: u64,
290    tokens_original: u64,
291    tokens_compressed: u64,
292    modes: &HashMap<String, u64>,
293    tool_calls: u64,
294    complexity: &str,
295) {
296    let prev_original = cep.last_session_original.unwrap_or(0);
297    let prev_compressed = cep.last_session_compressed.unwrap_or(0);
298    let prev_cache_hits = cep.last_session_cache_hits.unwrap_or(0);
299    let prev_cache_reads = cep.last_session_cache_reads.unwrap_or(0);
300    let is_same_session = cep.last_session_pid == Some(pid);
301
302    if is_same_session {
303        cep.total_tokens_original += tokens_original.saturating_sub(prev_original);
304        cep.total_tokens_compressed += tokens_compressed.saturating_sub(prev_compressed);
305        cep.total_cache_hits += cache_hits.saturating_sub(prev_cache_hits);
306        cep.total_cache_reads += cache_reads.saturating_sub(prev_cache_reads);
307    } else {
308        cep.sessions += 1;
309        cep.total_cache_hits += cache_hits;
310        cep.total_cache_reads += cache_reads;
311        cep.total_tokens_original += tokens_original;
312        cep.total_tokens_compressed += tokens_compressed;
313
314        for (mode, count) in modes {
315            *cep.modes.entry(mode.clone()).or_insert(0) += count;
316        }
317    }
318
319    cep.last_session_pid = Some(pid);
320    cep.last_session_original = Some(tokens_original);
321    cep.last_session_compressed = Some(tokens_compressed);
322    cep.last_session_cache_hits = Some(cache_hits);
323    cep.last_session_cache_reads = Some(cache_reads);
324
325    let cache_hit_rate = if cache_reads > 0 {
326        (cache_hits as f64 / cache_reads as f64 * 100.0).round() as u32
327    } else {
328        0
329    };
330
331    let compression_rate = if tokens_original > 0 {
332        ((tokens_original - tokens_compressed) as f64 / tokens_original as f64 * 100.0).round()
333            as u32
334    } else {
335        0
336    };
337
338    let total_modes = 6u32;
339    let mode_diversity =
340        ((modes.len() as f64 / total_modes as f64).min(1.0) * 100.0).round() as u32;
341
342    let tokens_saved = tokens_original.saturating_sub(tokens_compressed);
343
344    cep.scores.push(CepSessionSnapshot {
345        timestamp: chrono::Local::now().to_rfc3339(),
346        score,
347        cache_hit_rate,
348        mode_diversity,
349        compression_rate,
350        tool_calls,
351        tokens_saved,
352        complexity: complexity.to_string(),
353    });
354
355    if cep.scores.len() > 100 {
356        cep.scores.drain(..cep.scores.len() - 100);
357    }
358}
359
360#[cfg(test)]
361mod tests {
362    use super::*;
363
364    fn make_store(commands: u64, input: u64, output: u64) -> StatsStore {
365        StatsStore {
366            total_commands: commands,
367            total_input_tokens: input,
368            total_output_tokens: output,
369            ..Default::default()
370        }
371    }
372
373    /// #706: a corrupt `stats.json` must never be silently replaced by an
374    /// empty store. The loader quarantines the bytes to `stats.json.corrupt`
375    /// (recoverable), and an existing quarantine is never overwritten — the
376    /// older copy is the one closest to the lost history.
377    #[test]
378    fn corrupt_stats_file_is_quarantined_not_silently_reset() {
379        let dir = crate::core::data_dir::isolated_data_dir();
380        let stats_path = dir.path().join("stats.json");
381        let quarantine = dir.path().join("stats.json.corrupt");
382        let truncated = r#"{"total_commands": 15677, "total_input_tok"#;
383        std::fs::write(&stats_path, truncated).unwrap();
384
385        let loaded = io::load_from_disk();
386        assert_eq!(loaded.total_commands, 0, "fresh store after corruption");
387        assert!(
388            !stats_path.exists(),
389            "corrupt file must be moved aside, not left to be overwritten"
390        );
391        assert_eq!(
392            std::fs::read_to_string(&quarantine).unwrap(),
393            truncated,
394            "quarantine preserves the corrupt bytes verbatim for recovery"
395        );
396
397        // A second corruption must NOT clobber the first quarantine.
398        std::fs::write(&stats_path, "{ newer corruption").unwrap();
399        let loaded = io::load_from_disk();
400        assert_eq!(loaded.total_commands, 0);
401        assert_eq!(
402            std::fs::read_to_string(&quarantine).unwrap(),
403            truncated,
404            "the OLDER quarantine wins — it is closest to the lost history"
405        );
406
407        // And a healthy file still loads normally.
408        let healthy = make_store(42, 9000, 1000);
409        std::fs::write(&stats_path, serde_json::to_string(&healthy).unwrap()).unwrap();
410        assert_eq!(io::load_from_disk().total_commands, 42);
411    }
412
413    #[test]
414    fn aggregate_for_display_is_noop_without_siblings() {
415        // The common case: only the primary dir has stats. Aggregation must
416        // return the primary untouched so non-split users see no change (#500).
417        let primary = make_store(7, 1000, 250);
418        let agg = aggregate_for_display(primary.clone(), &[]);
419        assert_eq!(agg.total_commands, 7);
420        assert_eq!(agg.total_input_tokens, 1000);
421        assert_eq!(agg.total_output_tokens, 250);
422    }
423
424    #[test]
425    fn aggregate_for_display_sums_split_dirs() {
426        // A data-dir split (#408/#414/#500): the CLI's primary dir is empty but
427        // the MCP server wrote its savings into a sibling tree. The displayed
428        // total must reflect both so `gain` no longer reports a false `0`.
429        let primary = make_store(0, 0, 0);
430        let mcp_dir = make_store(12, 8000, 1200);
431        let legacy_dir = make_store(3, 500, 100);
432
433        let agg = aggregate_for_display(primary, &[mcp_dir, legacy_dir]);
434
435        assert_eq!(agg.total_commands, 15, "12 (mcp) + 3 (legacy)");
436        assert_eq!(agg.total_input_tokens, 8500);
437        assert_eq!(agg.total_output_tokens, 1300);
438        let saved = agg
439            .total_input_tokens
440            .saturating_sub(agg.total_output_tokens);
441        assert_eq!(saved, 7200, "savings surface despite an empty primary dir");
442    }
443
444    #[test]
445    fn apply_deltas_merges_mcp_and_shell() {
446        let baseline = make_store(0, 0, 0);
447        let mut current = make_store(0, 0, 0);
448        current.total_commands = 5;
449        current.total_input_tokens = 1000;
450        current.total_output_tokens = 200;
451        current.commands.insert(
452            "ctx_read".to_string(),
453            CommandStats {
454                count: 5,
455                input_tokens: 1000,
456                output_tokens: 200,
457            },
458        );
459
460        let mut disk = make_store(20, 500, 490);
461        disk.commands.insert(
462            "echo".to_string(),
463            CommandStats {
464                count: 20,
465                input_tokens: 500,
466                output_tokens: 490,
467            },
468        );
469
470        let merged = io::apply_deltas(&disk, &current, &baseline);
471
472        assert_eq!(merged.total_commands, 25);
473        assert_eq!(merged.total_input_tokens, 1500);
474        assert_eq!(merged.total_output_tokens, 690);
475        assert_eq!(merged.commands["ctx_read"].count, 5);
476        assert_eq!(merged.commands["echo"].count, 20);
477    }
478
479    #[test]
480    fn apply_deltas_incremental_flush() {
481        let baseline = make_store(10, 200, 100);
482        let current = make_store(15, 700, 300);
483
484        let disk = make_store(30, 600, 500);
485
486        let merged = io::apply_deltas(&disk, &current, &baseline);
487
488        assert_eq!(merged.total_commands, 35);
489        assert_eq!(merged.total_input_tokens, 1100);
490        assert_eq!(merged.total_output_tokens, 700);
491    }
492
493    #[test]
494    fn apply_deltas_preserves_disk_commands() {
495        let baseline = make_store(0, 0, 0);
496        let mut current = make_store(2, 100, 50);
497        current.commands.insert(
498            "ctx_read".to_string(),
499            CommandStats {
500                count: 2,
501                input_tokens: 100,
502                output_tokens: 50,
503            },
504        );
505
506        let mut disk = make_store(10, 300, 280);
507        disk.commands.insert(
508            "echo".to_string(),
509            CommandStats {
510                count: 8,
511                input_tokens: 200,
512                output_tokens: 200,
513            },
514        );
515        disk.commands.insert(
516            "ctx_read".to_string(),
517            CommandStats {
518                count: 3,
519                input_tokens: 150,
520                output_tokens: 80,
521            },
522        );
523
524        let merged = io::apply_deltas(&disk, &current, &baseline);
525
526        assert_eq!(merged.commands["echo"].count, 8);
527        assert_eq!(merged.commands["ctx_read"].count, 5);
528        assert_eq!(merged.commands["ctx_read"].input_tokens, 250);
529    }
530
531    #[test]
532    fn merge_daily_combines_same_date() {
533        let baseline_daily = vec![];
534        let current_daily = vec![DayStats {
535            date: "2026-04-18".to_string(),
536            commands: 5,
537            input_tokens: 1000,
538            output_tokens: 200,
539            version: "3.7.0".to_string(),
540        }];
541        let mut merged_daily = vec![DayStats {
542            date: "2026-04-18".to_string(),
543            commands: 20,
544            input_tokens: 500,
545            output_tokens: 490,
546            version: String::new(),
547        }];
548
549        io::merge_daily(&mut merged_daily, &current_daily, &baseline_daily);
550
551        assert_eq!(merged_daily.len(), 1);
552        assert_eq!(merged_daily[0].commands, 25);
553        assert_eq!(merged_daily[0].input_tokens, 1500);
554        // #307: the most recent known version is carried into the merge.
555        assert_eq!(merged_daily[0].version, "3.7.0");
556    }
557
558    #[test]
559    fn cep_snapshot_seeds_new_session() {
560        let mut cep = CepStats::default();
561        let modes = HashMap::from([("full".to_string(), 3)]);
562        apply_cep_snapshot(&mut cep, 100, 80, 5, 10, 1000, 200, &modes, 4, "Medium");
563        assert_eq!(cep.sessions, 1);
564        assert_eq!(cep.total_cache_hits, 5);
565        assert_eq!(cep.total_cache_reads, 10);
566        assert_eq!(cep.total_tokens_original, 1000);
567        assert_eq!(cep.total_tokens_compressed, 200);
568        assert_eq!(cep.scores.len(), 1);
569    }
570
571    #[test]
572    fn cep_snapshot_same_pid_accumulates_cache_delta() {
573        // #361: repeated snapshots within one process must keep counting cache
574        // hits/reads (cumulative counters → add the delta), not freeze at the
575        // first checkpoint's value while only tokens advanced.
576        let mut cep = CepStats::default();
577        let modes = HashMap::new();
578        apply_cep_snapshot(&mut cep, 100, 80, 2, 4, 500, 100, &modes, 2, "Low");
579        // Same PID, cumulative counters grew: hits 2→9, reads 4→20.
580        apply_cep_snapshot(&mut cep, 100, 85, 9, 20, 1500, 300, &modes, 6, "Low");
581
582        assert_eq!(cep.sessions, 1, "same PID must not start a new session");
583        assert_eq!(cep.total_cache_hits, 9, "2 + delta(9-2)");
584        assert_eq!(cep.total_cache_reads, 20, "4 + delta(20-4)");
585        assert_eq!(cep.total_tokens_original, 1500);
586        assert_eq!(cep.total_tokens_compressed, 300);
587        assert_eq!(cep.scores.len(), 2);
588    }
589
590    #[test]
591    fn cep_snapshot_new_pid_starts_fresh_session() {
592        let mut cep = CepStats::default();
593        let modes = HashMap::new();
594        apply_cep_snapshot(&mut cep, 100, 80, 5, 10, 1000, 200, &modes, 4, "Medium");
595        apply_cep_snapshot(&mut cep, 200, 80, 3, 6, 800, 150, &modes, 4, "Medium");
596        assert_eq!(cep.sessions, 2);
597        assert_eq!(
598            cep.total_cache_hits, 8,
599            "5 (session 1) + 3 (session 2, fresh)"
600        );
601        assert_eq!(cep.total_cache_reads, 16);
602    }
603}