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
35pub fn save(store: &StatsStore) {
36    io::locked_write(store);
37}
38
39fn maybe_flush(store: &mut StatsStore, baseline: &mut StatsStore, last_flush: &mut Instant) {
40    if last_flush.elapsed().as_secs() >= FLUSH_INTERVAL_SECS {
41        let merged = io::merge_and_save(store, baseline);
42        *store = merged.clone();
43        *baseline = merged;
44        *last_flush = Instant::now();
45    }
46}
47
48pub fn flush() {
49    let mut guard = STATS_BUFFER
50        .lock()
51        .unwrap_or_else(std::sync::PoisonError::into_inner);
52    if let Some((ref mut store, ref mut baseline, ref mut last_flush)) = *guard {
53        let merged = io::merge_and_save(store, baseline);
54        *store = merged.clone();
55        *baseline = merged;
56        *last_flush = Instant::now();
57    }
58}
59
60/// Adjust saved tokens after post-processing (terse, hints) changed the output size.
61/// Positive delta = savings were over-reported, negative = under-reported.
62pub fn adjust_savings(command: &str, over_report_delta: i64) {
63    let mut guard = STATS_BUFFER
64        .lock()
65        .unwrap_or_else(std::sync::PoisonError::into_inner);
66    let Some((store, _, _)) = guard.as_mut() else {
67        return;
68    };
69    if over_report_delta > 0 {
70        let adj = over_report_delta as u64;
71        store.total_output_tokens = store.total_output_tokens.saturating_add(adj);
72        if let Some(cmd) = store.commands.get_mut(command) {
73            cmd.output_tokens = cmd.output_tokens.saturating_add(adj);
74        }
75    } else {
76        let adj = over_report_delta.unsigned_abs();
77        store.total_output_tokens = store.total_output_tokens.saturating_sub(adj);
78        if let Some(cmd) = store.commands.get_mut(command) {
79            cmd.output_tokens = cmd.output_tokens.saturating_sub(adj);
80        }
81    }
82}
83
84pub fn record(command: &str, input_tokens: usize, output_tokens: usize) {
85    let mut guard = STATS_BUFFER
86        .lock()
87        .unwrap_or_else(std::sync::PoisonError::into_inner);
88    if guard.is_none() {
89        let disk = io::load_from_disk();
90        *guard = Some((disk.clone(), disk, Instant::now()));
91    }
92    let Some((store, baseline, last_flush)) = guard.as_mut() else {
93        return;
94    };
95
96    let is_first_command = store.total_commands == baseline.total_commands;
97    let now = chrono::Local::now();
98    let today = now.format("%Y-%m-%d").to_string();
99    let timestamp = now.to_rfc3339();
100
101    store.total_commands = store.total_commands.saturating_add(1);
102    store.total_input_tokens = store.total_input_tokens.saturating_add(input_tokens as u64);
103    store.total_output_tokens = store
104        .total_output_tokens
105        .saturating_add(output_tokens as u64);
106
107    if store.first_use.is_none() {
108        store.first_use = Some(timestamp.clone());
109    }
110    store.last_use = Some(timestamp);
111
112    let cmd_key = format::normalize_command(command);
113    let entry = store.commands.entry(cmd_key).or_default();
114    entry.count = entry.count.saturating_add(1);
115    entry.input_tokens = entry.input_tokens.saturating_add(input_tokens as u64);
116    entry.output_tokens = entry.output_tokens.saturating_add(output_tokens as u64);
117
118    let current_version = env!("CARGO_PKG_VERSION").to_string();
119    if let Some(day) = store.daily.last_mut() {
120        if day.date == today {
121            day.commands = day.commands.saturating_add(1);
122            day.input_tokens = day.input_tokens.saturating_add(input_tokens as u64);
123            day.output_tokens = day.output_tokens.saturating_add(output_tokens as u64);
124            // Stamp the running version so a mid-day update attributes the day
125            // to the release in use for its latest activity (#307).
126            day.version = current_version;
127        } else {
128            store.daily.push(DayStats {
129                date: today,
130                commands: 1,
131                input_tokens: input_tokens as u64,
132                output_tokens: output_tokens as u64,
133                version: current_version,
134            });
135        }
136    } else {
137        store.daily.push(DayStats {
138            date: today,
139            commands: 1,
140            input_tokens: input_tokens as u64,
141            output_tokens: output_tokens as u64,
142            version: current_version,
143        });
144    }
145
146    if store.daily.len() > MAX_DAILY_HISTORY_DAYS {
147        store
148            .daily
149            .drain(..store.daily.len() - MAX_DAILY_HISTORY_DAYS);
150    }
151
152    if is_first_command {
153        let merged = io::merge_and_save(store, baseline);
154        *store = merged.clone();
155        *baseline = merged;
156        *last_flush = Instant::now();
157    } else {
158        maybe_flush(store, baseline, last_flush);
159    }
160}
161
162pub fn reset_cep() {
163    let mut guard = STATS_BUFFER
164        .lock()
165        .unwrap_or_else(std::sync::PoisonError::into_inner);
166    let mut store = io::load_from_disk();
167    store.cep = CepStats::default();
168    io::locked_write(&store);
169    *guard = Some((store.clone(), store, Instant::now()));
170}
171
172pub fn reset_all() {
173    let mut guard = STATS_BUFFER
174        .lock()
175        .unwrap_or_else(std::sync::PoisonError::into_inner);
176    let store = StatsStore::default();
177    io::locked_write(&store);
178    *guard = Some((store.clone(), store, Instant::now()));
179    crate::core::heatmap::reset();
180}
181
182pub fn load_stats() -> GainSummary {
183    let store = load();
184    let input_saved = store
185        .total_input_tokens
186        .saturating_sub(store.total_output_tokens);
187    GainSummary {
188        total_saved: input_saved,
189        total_calls: store.total_commands,
190    }
191}
192
193#[allow(clippy::too_many_arguments)]
194pub fn record_cep_session(
195    score: u32,
196    cache_hits: u64,
197    cache_reads: u64,
198    tokens_original: u64,
199    tokens_compressed: u64,
200    modes: &HashMap<String, u64>,
201    tool_calls: u64,
202    complexity: &str,
203) {
204    let mut guard = STATS_BUFFER
205        .lock()
206        .unwrap_or_else(std::sync::PoisonError::into_inner);
207    if guard.is_none() {
208        let disk = io::load_from_disk();
209        *guard = Some((disk.clone(), disk, Instant::now()));
210    }
211    let Some((store, baseline, last_flush)) = guard.as_mut() else {
212        return;
213    };
214
215    apply_cep_snapshot(
216        &mut store.cep,
217        std::process::id(),
218        score,
219        cache_hits,
220        cache_reads,
221        tokens_original,
222        tokens_compressed,
223        modes,
224        tool_calls,
225        complexity,
226    );
227
228    maybe_flush(store, baseline, last_flush);
229}
230
231/// Fold one CEP snapshot into `cep`. Pure (no globals, no I/O) so the
232/// delta/aggregation rules are unit-testable in isolation.
233///
234/// `cache_hits`, `cache_reads`, `tokens_original` and `tokens_compressed` arrive
235/// as **cumulative per-process** counters. For repeated snapshots within the same
236/// PID only the delta since the previous snapshot is added, so the lifetime
237/// totals keep tracking cache activity instead of freezing at the first
238/// checkpoint's value (#361). A new PID starts a fresh session and seeds the
239/// cumulative baselines.
240#[allow(clippy::too_many_arguments)]
241fn apply_cep_snapshot(
242    cep: &mut CepStats,
243    pid: u32,
244    score: u32,
245    cache_hits: u64,
246    cache_reads: u64,
247    tokens_original: u64,
248    tokens_compressed: u64,
249    modes: &HashMap<String, u64>,
250    tool_calls: u64,
251    complexity: &str,
252) {
253    let prev_original = cep.last_session_original.unwrap_or(0);
254    let prev_compressed = cep.last_session_compressed.unwrap_or(0);
255    let prev_cache_hits = cep.last_session_cache_hits.unwrap_or(0);
256    let prev_cache_reads = cep.last_session_cache_reads.unwrap_or(0);
257    let is_same_session = cep.last_session_pid == Some(pid);
258
259    if is_same_session {
260        cep.total_tokens_original += tokens_original.saturating_sub(prev_original);
261        cep.total_tokens_compressed += tokens_compressed.saturating_sub(prev_compressed);
262        cep.total_cache_hits += cache_hits.saturating_sub(prev_cache_hits);
263        cep.total_cache_reads += cache_reads.saturating_sub(prev_cache_reads);
264    } else {
265        cep.sessions += 1;
266        cep.total_cache_hits += cache_hits;
267        cep.total_cache_reads += cache_reads;
268        cep.total_tokens_original += tokens_original;
269        cep.total_tokens_compressed += tokens_compressed;
270
271        for (mode, count) in modes {
272            *cep.modes.entry(mode.clone()).or_insert(0) += count;
273        }
274    }
275
276    cep.last_session_pid = Some(pid);
277    cep.last_session_original = Some(tokens_original);
278    cep.last_session_compressed = Some(tokens_compressed);
279    cep.last_session_cache_hits = Some(cache_hits);
280    cep.last_session_cache_reads = Some(cache_reads);
281
282    let cache_hit_rate = if cache_reads > 0 {
283        (cache_hits as f64 / cache_reads as f64 * 100.0).round() as u32
284    } else {
285        0
286    };
287
288    let compression_rate = if tokens_original > 0 {
289        ((tokens_original - tokens_compressed) as f64 / tokens_original as f64 * 100.0).round()
290            as u32
291    } else {
292        0
293    };
294
295    let total_modes = 6u32;
296    let mode_diversity =
297        ((modes.len() as f64 / total_modes as f64).min(1.0) * 100.0).round() as u32;
298
299    let tokens_saved = tokens_original.saturating_sub(tokens_compressed);
300
301    cep.scores.push(CepSessionSnapshot {
302        timestamp: chrono::Local::now().to_rfc3339(),
303        score,
304        cache_hit_rate,
305        mode_diversity,
306        compression_rate,
307        tool_calls,
308        tokens_saved,
309        complexity: complexity.to_string(),
310    });
311
312    if cep.scores.len() > 100 {
313        cep.scores.drain(..cep.scores.len() - 100);
314    }
315}
316
317#[cfg(test)]
318mod tests {
319    use super::*;
320
321    fn make_store(commands: u64, input: u64, output: u64) -> StatsStore {
322        StatsStore {
323            total_commands: commands,
324            total_input_tokens: input,
325            total_output_tokens: output,
326            ..Default::default()
327        }
328    }
329
330    #[test]
331    fn apply_deltas_merges_mcp_and_shell() {
332        let baseline = make_store(0, 0, 0);
333        let mut current = make_store(0, 0, 0);
334        current.total_commands = 5;
335        current.total_input_tokens = 1000;
336        current.total_output_tokens = 200;
337        current.commands.insert(
338            "ctx_read".to_string(),
339            CommandStats {
340                count: 5,
341                input_tokens: 1000,
342                output_tokens: 200,
343            },
344        );
345
346        let mut disk = make_store(20, 500, 490);
347        disk.commands.insert(
348            "echo".to_string(),
349            CommandStats {
350                count: 20,
351                input_tokens: 500,
352                output_tokens: 490,
353            },
354        );
355
356        let merged = io::apply_deltas(&disk, &current, &baseline);
357
358        assert_eq!(merged.total_commands, 25);
359        assert_eq!(merged.total_input_tokens, 1500);
360        assert_eq!(merged.total_output_tokens, 690);
361        assert_eq!(merged.commands["ctx_read"].count, 5);
362        assert_eq!(merged.commands["echo"].count, 20);
363    }
364
365    #[test]
366    fn apply_deltas_incremental_flush() {
367        let baseline = make_store(10, 200, 100);
368        let current = make_store(15, 700, 300);
369
370        let disk = make_store(30, 600, 500);
371
372        let merged = io::apply_deltas(&disk, &current, &baseline);
373
374        assert_eq!(merged.total_commands, 35);
375        assert_eq!(merged.total_input_tokens, 1100);
376        assert_eq!(merged.total_output_tokens, 700);
377    }
378
379    #[test]
380    fn apply_deltas_preserves_disk_commands() {
381        let baseline = make_store(0, 0, 0);
382        let mut current = make_store(2, 100, 50);
383        current.commands.insert(
384            "ctx_read".to_string(),
385            CommandStats {
386                count: 2,
387                input_tokens: 100,
388                output_tokens: 50,
389            },
390        );
391
392        let mut disk = make_store(10, 300, 280);
393        disk.commands.insert(
394            "echo".to_string(),
395            CommandStats {
396                count: 8,
397                input_tokens: 200,
398                output_tokens: 200,
399            },
400        );
401        disk.commands.insert(
402            "ctx_read".to_string(),
403            CommandStats {
404                count: 3,
405                input_tokens: 150,
406                output_tokens: 80,
407            },
408        );
409
410        let merged = io::apply_deltas(&disk, &current, &baseline);
411
412        assert_eq!(merged.commands["echo"].count, 8);
413        assert_eq!(merged.commands["ctx_read"].count, 5);
414        assert_eq!(merged.commands["ctx_read"].input_tokens, 250);
415    }
416
417    #[test]
418    fn merge_daily_combines_same_date() {
419        let baseline_daily = vec![];
420        let current_daily = vec![DayStats {
421            date: "2026-04-18".to_string(),
422            commands: 5,
423            input_tokens: 1000,
424            output_tokens: 200,
425            version: "3.7.0".to_string(),
426        }];
427        let mut merged_daily = vec![DayStats {
428            date: "2026-04-18".to_string(),
429            commands: 20,
430            input_tokens: 500,
431            output_tokens: 490,
432            version: String::new(),
433        }];
434
435        io::merge_daily(&mut merged_daily, &current_daily, &baseline_daily);
436
437        assert_eq!(merged_daily.len(), 1);
438        assert_eq!(merged_daily[0].commands, 25);
439        assert_eq!(merged_daily[0].input_tokens, 1500);
440        // #307: the most recent known version is carried into the merge.
441        assert_eq!(merged_daily[0].version, "3.7.0");
442    }
443
444    #[test]
445    fn cep_snapshot_seeds_new_session() {
446        let mut cep = CepStats::default();
447        let modes = HashMap::from([("full".to_string(), 3)]);
448        apply_cep_snapshot(&mut cep, 100, 80, 5, 10, 1000, 200, &modes, 4, "Medium");
449        assert_eq!(cep.sessions, 1);
450        assert_eq!(cep.total_cache_hits, 5);
451        assert_eq!(cep.total_cache_reads, 10);
452        assert_eq!(cep.total_tokens_original, 1000);
453        assert_eq!(cep.total_tokens_compressed, 200);
454        assert_eq!(cep.scores.len(), 1);
455    }
456
457    #[test]
458    fn cep_snapshot_same_pid_accumulates_cache_delta() {
459        // #361: repeated snapshots within one process must keep counting cache
460        // hits/reads (cumulative counters → add the delta), not freeze at the
461        // first checkpoint's value while only tokens advanced.
462        let mut cep = CepStats::default();
463        let modes = HashMap::new();
464        apply_cep_snapshot(&mut cep, 100, 80, 2, 4, 500, 100, &modes, 2, "Low");
465        // Same PID, cumulative counters grew: hits 2→9, reads 4→20.
466        apply_cep_snapshot(&mut cep, 100, 85, 9, 20, 1500, 300, &modes, 6, "Low");
467
468        assert_eq!(cep.sessions, 1, "same PID must not start a new session");
469        assert_eq!(cep.total_cache_hits, 9, "2 + delta(9-2)");
470        assert_eq!(cep.total_cache_reads, 20, "4 + delta(20-4)");
471        assert_eq!(cep.total_tokens_original, 1500);
472        assert_eq!(cep.total_tokens_compressed, 300);
473        assert_eq!(cep.scores.len(), 2);
474    }
475
476    #[test]
477    fn cep_snapshot_new_pid_starts_fresh_session() {
478        let mut cep = CepStats::default();
479        let modes = HashMap::new();
480        apply_cep_snapshot(&mut cep, 100, 80, 5, 10, 1000, 200, &modes, 4, "Medium");
481        apply_cep_snapshot(&mut cep, 200, 80, 3, 6, 800, 150, &modes, 4, "Medium");
482        assert_eq!(cep.sessions, 2);
483        assert_eq!(
484            cep.total_cache_hits, 8,
485            "5 (session 1) + 3 (session 2, fresh)"
486        );
487        assert_eq!(cep.total_cache_reads, 16);
488    }
489}