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