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
12static STATS_BUFFER: Mutex<Option<(StatsStore, StatsStore, Instant)>> = Mutex::new(None);
14
15const FLUSH_INTERVAL_SECS: u64 = 2;
16
17pub(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 load_for_display() -> StatsStore {
51 let primary = load();
52 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
67fn 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 Some(merged) = io::merge_and_save(store, baseline)
85 {
86 *store = merged.clone();
87 *baseline = merged;
88 *last_flush = Instant::now();
89 }
90}
91
92pub fn flush() {
93 let mut guard = STATS_BUFFER
94 .lock()
95 .unwrap_or_else(std::sync::PoisonError::into_inner);
96 if let Some((ref mut store, ref mut baseline, ref mut last_flush)) = *guard
97 && let Some(merged) = io::merge_and_save(store, baseline)
98 {
99 *store = merged.clone();
100 *baseline = merged;
101 *last_flush = Instant::now();
102 }
103}
104
105pub fn flush_if_due() {
109 let due = {
110 let guard = STATS_BUFFER
111 .lock()
112 .unwrap_or_else(std::sync::PoisonError::into_inner);
113 match guard.as_ref() {
114 Some((_, _, last_flush)) => last_flush.elapsed().as_secs() >= FLUSH_INTERVAL_SECS,
115 None => false,
116 }
117 };
118 if due {
119 flush();
120 }
121}
122
123pub fn adjust_savings(command: &str, over_report_delta: i64) {
129 let mut guard = STATS_BUFFER
130 .lock()
131 .unwrap_or_else(std::sync::PoisonError::into_inner);
132 let Some((store, baseline, last_flush)) = guard.as_mut() else {
133 return;
134 };
135 let cmd_key = format::normalize_command(command);
136 let stream_tracked = classify_command(&cmd_key) == TrafficClass::Compressible;
137 if over_report_delta > 0 {
138 let adj = over_report_delta as u64;
139 store.total_output_tokens = store.total_output_tokens.saturating_add(adj);
140 if let Some(cmd) = store.commands.get_mut(&cmd_key) {
141 cmd.output_tokens = cmd.output_tokens.saturating_add(adj);
142 }
143 if stream_tracked {
144 store.first_inject_tokens_saved = store.first_inject_tokens_saved.saturating_sub(adj);
145 store.active_tool_result_tokens_saved =
146 store.active_tool_result_tokens_saved.saturating_sub(adj);
147 }
148 } else {
149 let adj = over_report_delta.unsigned_abs();
150 store.total_output_tokens = store.total_output_tokens.saturating_sub(adj);
151 if let Some(cmd) = store.commands.get_mut(&cmd_key) {
152 cmd.output_tokens = cmd.output_tokens.saturating_sub(adj);
153 }
154 if stream_tracked {
155 store.first_inject_tokens_saved = store.first_inject_tokens_saved.saturating_add(adj);
156 if store.last_tool_result_turn > 0 {
157 store.active_tool_result_tokens_saved =
158 store.active_tool_result_tokens_saved.saturating_add(adj);
159 }
160 }
161 }
162 if let Some(merged) = io::merge_and_save(store, baseline) {
163 *store = merged.clone();
164 *baseline = merged;
165 *last_flush = Instant::now();
166 }
167}
168
169pub fn record(command: &str, input_tokens: usize, output_tokens: usize) {
170 record_at_turn(command, input_tokens, output_tokens, 0);
171}
172
173pub fn record_reread(tokens_saved: usize) {
174 let mut guard = STATS_BUFFER
175 .lock()
176 .unwrap_or_else(std::sync::PoisonError::into_inner);
177 if guard.is_none() {
178 let disk = io::load_from_disk();
179 *guard = Some((disk.clone(), disk, Instant::now()));
180 }
181 let Some((store, baseline, last_flush)) = guard.as_mut() else {
182 return;
183 };
184 store.reread_tokens_saved = store
185 .reread_tokens_saved
186 .saturating_add(tokens_saved as u64);
187 if let Some(merged) = io::merge_and_save(store, baseline) {
188 *store = merged.clone();
189 *baseline = merged;
190 *last_flush = Instant::now();
191 }
192}
193
194pub fn record_at_turn(command: &str, input_tokens: usize, output_tokens: usize, turn: u64) {
197 let mut guard = STATS_BUFFER
198 .lock()
199 .unwrap_or_else(std::sync::PoisonError::into_inner);
200 if guard.is_none() {
201 let disk = io::load_from_disk();
202 *guard = Some((disk.clone(), disk, Instant::now()));
203 }
204 let Some((store, baseline, last_flush)) = guard.as_mut() else {
205 return;
206 };
207
208 let now = chrono::Local::now();
210 let today = now.format("%Y-%m-%d").to_string();
211 let timestamp = now.to_rfc3339();
212
213 store.total_commands = store.total_commands.saturating_add(1);
214 store.total_input_tokens = store.total_input_tokens.saturating_add(input_tokens as u64);
215 store.total_output_tokens = store
216 .total_output_tokens
217 .saturating_add(output_tokens as u64);
218
219 if store.first_use.is_none() {
220 store.first_use = Some(timestamp.clone());
221 }
222 store.last_use = Some(timestamp);
223
224 let cmd_key = format::normalize_command(command);
225 if classify_command(&cmd_key) == TrafficClass::Compressible {
226 let saved = input_tokens.saturating_sub(output_tokens) as u64;
227 store.record_tool_result_savings(saved, turn);
228 }
229 store
230 .command_classes
231 .insert(cmd_key.clone(), classify_command(&cmd_key));
232 let entry = store.commands.entry(cmd_key).or_default();
233 entry.count = entry.count.saturating_add(1);
234 entry.input_tokens = entry.input_tokens.saturating_add(input_tokens as u64);
235 entry.output_tokens = entry.output_tokens.saturating_add(output_tokens as u64);
236
237 let current_version = env!("CARGO_PKG_VERSION").to_string();
238 if let Some(day) = store.daily.last_mut() {
239 if day.date == today {
240 day.commands = day.commands.saturating_add(1);
241 day.input_tokens = day.input_tokens.saturating_add(input_tokens as u64);
242 day.output_tokens = day.output_tokens.saturating_add(output_tokens as u64);
243 day.version = current_version;
246 } else {
247 store.daily.push(DayStats {
248 date: today,
249 commands: 1,
250 input_tokens: input_tokens as u64,
251 output_tokens: output_tokens as u64,
252 version: current_version,
253 });
254 }
255 } else {
256 store.daily.push(DayStats {
257 date: today,
258 commands: 1,
259 input_tokens: input_tokens as u64,
260 output_tokens: output_tokens as u64,
261 version: current_version,
262 });
263 }
264
265 if store.daily.len() > MAX_DAILY_HISTORY_DAYS {
266 store
267 .daily
268 .drain(..store.daily.len() - MAX_DAILY_HISTORY_DAYS);
269 }
270
271 if let Some(merged) = io::merge_and_save(store, baseline) {
275 *store = merged.clone();
276 *baseline = merged;
277 *last_flush = Instant::now();
278 }
279}
280
281pub fn reset_cep() {
282 let mut guard = STATS_BUFFER
283 .lock()
284 .unwrap_or_else(std::sync::PoisonError::into_inner);
285 let mut store = io::load_from_disk();
286 store.cep = CepStats::default();
287 io::locked_write(&store);
288 *guard = Some((store.clone(), store, Instant::now()));
289}
290
291pub fn reset_all() {
292 let mut guard = STATS_BUFFER
293 .lock()
294 .unwrap_or_else(std::sync::PoisonError::into_inner);
295 let store = StatsStore::default();
296 io::locked_write(&store);
297 *guard = Some((store.clone(), store, Instant::now()));
298 crate::core::heatmap::reset();
299}
300
301pub fn load_stats() -> GainSummary {
302 let store = load();
303 let input_saved = store
304 .total_input_tokens
305 .saturating_sub(store.total_output_tokens);
306 GainSummary {
307 total_saved: input_saved,
308 total_calls: store.total_commands,
309 }
310}
311
312#[allow(clippy::too_many_arguments)]
313pub fn record_cep_session(
314 score: u32,
315 cache_hits: u64,
316 cache_reads: u64,
317 tokens_original: u64,
318 tokens_compressed: u64,
319 modes: &HashMap<String, u64>,
320 tool_calls: u64,
321 complexity: &str,
322) {
323 let mut guard = STATS_BUFFER
324 .lock()
325 .unwrap_or_else(std::sync::PoisonError::into_inner);
326 if guard.is_none() {
327 let disk = io::load_from_disk();
328 *guard = Some((disk.clone(), disk, Instant::now()));
329 }
330 let Some((store, baseline, last_flush)) = guard.as_mut() else {
331 return;
332 };
333
334 apply_cep_snapshot(
335 &mut store.cep,
336 std::process::id(),
337 score,
338 cache_hits,
339 cache_reads,
340 tokens_original,
341 tokens_compressed,
342 modes,
343 tool_calls,
344 complexity,
345 );
346
347 maybe_flush(store, baseline, last_flush);
348}
349
350#[allow(clippy::too_many_arguments)]
360fn apply_cep_snapshot(
361 cep: &mut CepStats,
362 pid: u32,
363 score: u32,
364 cache_hits: u64,
365 cache_reads: u64,
366 tokens_original: u64,
367 tokens_compressed: u64,
368 modes: &HashMap<String, u64>,
369 tool_calls: u64,
370 complexity: &str,
371) {
372 let prev_original = cep.last_session_original.unwrap_or(0);
373 let prev_compressed = cep.last_session_compressed.unwrap_or(0);
374 let prev_cache_hits = cep.last_session_cache_hits.unwrap_or(0);
375 let prev_cache_reads = cep.last_session_cache_reads.unwrap_or(0);
376 let is_same_session = cep.last_session_pid == Some(pid);
377
378 if is_same_session {
379 cep.total_tokens_original += tokens_original.saturating_sub(prev_original);
380 cep.total_tokens_compressed += tokens_compressed.saturating_sub(prev_compressed);
381 cep.total_cache_hits += cache_hits.saturating_sub(prev_cache_hits);
382 cep.total_cache_reads += cache_reads.saturating_sub(prev_cache_reads);
383 } else {
384 cep.sessions += 1;
385 cep.total_cache_hits += cache_hits;
386 cep.total_cache_reads += cache_reads;
387 cep.total_tokens_original += tokens_original;
388 cep.total_tokens_compressed += tokens_compressed;
389
390 for (mode, count) in modes {
391 *cep.modes.entry(mode.clone()).or_insert(0) += count;
392 }
393 }
394
395 cep.last_session_pid = Some(pid);
396 cep.last_session_original = Some(tokens_original);
397 cep.last_session_compressed = Some(tokens_compressed);
398 cep.last_session_cache_hits = Some(cache_hits);
399 cep.last_session_cache_reads = Some(cache_reads);
400
401 let cache_hit_rate = if cache_reads > 0 {
402 (cache_hits as f64 / cache_reads as f64 * 100.0).round() as u32
403 } else {
404 0
405 };
406
407 let compression_rate = if tokens_original > 0 {
408 ((tokens_original - tokens_compressed) as f64 / tokens_original as f64 * 100.0).round()
409 as u32
410 } else {
411 0
412 };
413
414 let total_modes = 6u32;
415 let mode_diversity =
416 ((modes.len() as f64 / total_modes as f64).min(1.0) * 100.0).round() as u32;
417
418 let tokens_saved = tokens_original.saturating_sub(tokens_compressed);
419
420 cep.scores.push(CepSessionSnapshot {
421 timestamp: chrono::Local::now().to_rfc3339(),
422 score,
423 cache_hit_rate,
424 mode_diversity,
425 compression_rate,
426 tool_calls,
427 tokens_saved,
428 complexity: complexity.to_string(),
429 });
430
431 if cep.scores.len() > 100 {
432 cep.scores.drain(..cep.scores.len() - 100);
433 }
434}
435
436#[cfg(test)]
437mod tests {
438 use super::*;
439
440 fn make_store(commands: u64, input: u64, output: u64) -> StatsStore {
441 StatsStore {
442 total_commands: commands,
443 total_input_tokens: input,
444 total_output_tokens: output,
445 ..Default::default()
446 }
447 }
448
449 #[test]
454 fn corrupt_stats_file_is_quarantined_not_silently_reset() {
455 let dir = crate::core::data_dir::isolated_data_dir();
456 let stats_path = dir.path().join("stats.json");
457 let quarantine = dir.path().join("stats.json.corrupt");
458 let truncated = r#"{"total_commands": 15677, "total_input_tok"#;
459 std::fs::write(&stats_path, truncated).unwrap();
460
461 let loaded = io::load_from_disk();
462 assert_eq!(loaded.total_commands, 0, "fresh store after corruption");
463 assert!(
464 !stats_path.exists(),
465 "corrupt file must be moved aside, not left to be overwritten"
466 );
467 assert_eq!(
468 std::fs::read_to_string(&quarantine).unwrap(),
469 truncated,
470 "quarantine preserves the corrupt bytes verbatim for recovery"
471 );
472
473 std::fs::write(&stats_path, "{ newer corruption").unwrap();
475 let loaded = io::load_from_disk();
476 assert_eq!(loaded.total_commands, 0);
477 assert_eq!(
478 std::fs::read_to_string(&quarantine).unwrap(),
479 truncated,
480 "the OLDER quarantine wins — it is closest to the lost history"
481 );
482
483 let healthy = make_store(42, 9000, 1000);
485 std::fs::write(&stats_path, serde_json::to_string(&healthy).unwrap()).unwrap();
486 assert_eq!(io::load_from_disk().total_commands, 42);
487 }
488
489 #[test]
490 fn aggregate_for_display_is_noop_without_siblings() {
491 let primary = make_store(7, 1000, 250);
494 let agg = aggregate_for_display(primary.clone(), &[]);
495 assert_eq!(agg.total_commands, 7);
496 assert_eq!(agg.total_input_tokens, 1000);
497 assert_eq!(agg.total_output_tokens, 250);
498 }
499
500 #[test]
501 fn aggregate_for_display_sums_split_dirs() {
502 let primary = make_store(0, 0, 0);
506 let mcp_dir = make_store(12, 8000, 1200);
507 let legacy_dir = make_store(3, 500, 100);
508
509 let agg = aggregate_for_display(primary, &[mcp_dir, legacy_dir]);
510
511 assert_eq!(agg.total_commands, 15, "12 (mcp) + 3 (legacy)");
512 assert_eq!(agg.total_input_tokens, 8500);
513 assert_eq!(agg.total_output_tokens, 1300);
514 let saved = agg
515 .total_input_tokens
516 .saturating_sub(agg.total_output_tokens);
517 assert_eq!(saved, 7200, "savings surface despite an empty primary dir");
518 }
519
520 #[test]
521 fn apply_deltas_merges_mcp_and_shell() {
522 let baseline = make_store(0, 0, 0);
523 let mut current = make_store(0, 0, 0);
524 current.total_commands = 5;
525 current.total_input_tokens = 1000;
526 current.total_output_tokens = 200;
527 current.commands.insert(
528 "ctx_read".to_string(),
529 CommandStats {
530 count: 5,
531 input_tokens: 1000,
532 output_tokens: 200,
533 },
534 );
535 current
536 .command_classes
537 .insert("ctx_read".into(), TrafficClass::Compressible);
538
539 let mut disk = make_store(20, 500, 490);
540 disk.commands.insert(
541 "echo".to_string(),
542 CommandStats {
543 count: 20,
544 input_tokens: 500,
545 output_tokens: 490,
546 },
547 );
548
549 let merged = io::apply_deltas(&disk, ¤t, &baseline);
550
551 assert_eq!(merged.total_commands, 25);
552 assert_eq!(merged.total_input_tokens, 1500);
553 assert_eq!(merged.total_output_tokens, 690);
554 assert_eq!(merged.commands["ctx_read"].count, 5);
555 assert_eq!(merged.commands["echo"].count, 20);
556 assert_eq!(
557 merged.command_classes["ctx_read"],
558 TrafficClass::Compressible
559 );
560 }
561
562 #[test]
563 fn apply_deltas_incremental_flush() {
564 let baseline = make_store(10, 200, 100);
565 let current = make_store(15, 700, 300);
566
567 let disk = make_store(30, 600, 500);
568
569 let merged = io::apply_deltas(&disk, ¤t, &baseline);
570
571 assert_eq!(merged.total_commands, 35);
572 assert_eq!(merged.total_input_tokens, 1100);
573 assert_eq!(merged.total_output_tokens, 700);
574 }
575
576 #[test]
577 fn apply_deltas_merges_stream_counters_without_replaying_baseline() {
578 let mut baseline = StatsStore {
579 first_inject_tokens_saved: 1_000,
580 reread_tokens_saved: 2_000,
581 active_tool_result_tokens_saved: 500,
582 last_tool_result_turn: 8,
583 stream_tracked_results: 2,
584 ..StatsStore::default()
585 };
586 let mut current = baseline.clone();
587 current.first_inject_tokens_saved = 1_400;
588 current.reread_tokens_saved = 2_900;
589 current.active_tool_result_tokens_saved = 700;
590 current.last_tool_result_turn = 10;
591 current.stream_tracked_results = 3;
592 let disk = StatsStore {
593 first_inject_tokens_saved: 5_000,
594 reread_tokens_saved: 7_000,
595 active_tool_result_tokens_saved: 1_000,
596 last_tool_result_turn: 9,
597 stream_tracked_results: 10,
598 ..StatsStore::default()
599 };
600
601 let merged = io::apply_deltas(&disk, ¤t, &baseline);
602 assert_eq!(merged.first_inject_tokens_saved, 5_400);
603 assert_eq!(merged.reread_tokens_saved, 7_900);
604 assert_eq!(merged.active_tool_result_tokens_saved, 1_200);
605 assert_eq!(merged.last_tool_result_turn, 10);
606 assert_eq!(merged.stream_tracked_results, 11);
607
608 baseline.first_inject_tokens_saved = current.first_inject_tokens_saved;
609 let no_replay = io::apply_deltas(&merged, ¤t, &baseline);
610 assert_eq!(no_replay.first_inject_tokens_saved, 5_400);
611 }
612
613 #[test]
614 fn merge_and_save_keeps_delta_when_lock_is_busy() {
615 let dir = crate::core::data_dir::isolated_data_dir();
616 let baseline = make_store(10, 200, 100);
617 let current = make_store(11, 300, 120);
618 let lock_path = dir.path().join(".stats.lock");
619 std::fs::write(&lock_path, "busy").unwrap();
620
621 assert!(io::merge_and_save(¤t, &baseline).is_none());
622
623 std::fs::remove_file(lock_path).unwrap();
624 let merged = io::merge_and_save(¤t, &baseline).unwrap();
625 assert_eq!(merged.total_commands, 1);
626 assert_eq!(merged.total_input_tokens, 100);
627 assert_eq!(merged.total_output_tokens, 20);
628 }
629
630 #[test]
631 fn apply_deltas_preserves_disk_commands() {
632 let baseline = make_store(0, 0, 0);
633 let mut current = make_store(2, 100, 50);
634 current.commands.insert(
635 "ctx_read".to_string(),
636 CommandStats {
637 count: 2,
638 input_tokens: 100,
639 output_tokens: 50,
640 },
641 );
642
643 let mut disk = make_store(10, 300, 280);
644 disk.commands.insert(
645 "echo".to_string(),
646 CommandStats {
647 count: 8,
648 input_tokens: 200,
649 output_tokens: 200,
650 },
651 );
652 disk.commands.insert(
653 "ctx_read".to_string(),
654 CommandStats {
655 count: 3,
656 input_tokens: 150,
657 output_tokens: 80,
658 },
659 );
660
661 let merged = io::apply_deltas(&disk, ¤t, &baseline);
662
663 assert_eq!(merged.commands["echo"].count, 8);
664 assert_eq!(merged.commands["ctx_read"].count, 5);
665 assert_eq!(merged.commands["ctx_read"].input_tokens, 250);
666 }
667
668 #[test]
669 fn merge_daily_combines_same_date() {
670 let baseline_daily = vec![];
671 let current_daily = vec![DayStats {
672 date: "2026-04-18".to_string(),
673 commands: 5,
674 input_tokens: 1000,
675 output_tokens: 200,
676 version: "3.7.0".to_string(),
677 }];
678 let mut merged_daily = vec![DayStats {
679 date: "2026-04-18".to_string(),
680 commands: 20,
681 input_tokens: 500,
682 output_tokens: 490,
683 version: String::new(),
684 }];
685
686 io::merge_daily(&mut merged_daily, ¤t_daily, &baseline_daily);
687
688 assert_eq!(merged_daily.len(), 1);
689 assert_eq!(merged_daily[0].commands, 25);
690 assert_eq!(merged_daily[0].input_tokens, 1500);
691 assert_eq!(merged_daily[0].version, "3.7.0");
693 }
694
695 #[test]
696 fn cep_snapshot_seeds_new_session() {
697 let mut cep = CepStats::default();
698 let modes = HashMap::from([("full".to_string(), 3)]);
699 apply_cep_snapshot(&mut cep, 100, 80, 5, 10, 1000, 200, &modes, 4, "Medium");
700 assert_eq!(cep.sessions, 1);
701 assert_eq!(cep.total_cache_hits, 5);
702 assert_eq!(cep.total_cache_reads, 10);
703 assert_eq!(cep.total_tokens_original, 1000);
704 assert_eq!(cep.total_tokens_compressed, 200);
705 assert_eq!(cep.scores.len(), 1);
706 }
707
708 #[test]
709 fn cep_snapshot_same_pid_accumulates_cache_delta() {
710 let mut cep = CepStats::default();
714 let modes = HashMap::new();
715 apply_cep_snapshot(&mut cep, 100, 80, 2, 4, 500, 100, &modes, 2, "Low");
716 apply_cep_snapshot(&mut cep, 100, 85, 9, 20, 1500, 300, &modes, 6, "Low");
718
719 assert_eq!(cep.sessions, 1, "same PID must not start a new session");
720 assert_eq!(cep.total_cache_hits, 9, "2 + delta(9-2)");
721 assert_eq!(cep.total_cache_reads, 20, "4 + delta(20-4)");
722 assert_eq!(cep.total_tokens_original, 1500);
723 assert_eq!(cep.total_tokens_compressed, 300);
724 assert_eq!(cep.scores.len(), 2);
725 }
726
727 #[test]
728 fn cep_snapshot_new_pid_starts_fresh_session() {
729 let mut cep = CepStats::default();
730 let modes = HashMap::new();
731 apply_cep_snapshot(&mut cep, 100, 80, 5, 10, 1000, 200, &modes, 4, "Medium");
732 apply_cep_snapshot(&mut cep, 200, 80, 3, 6, 800, 150, &modes, 4, "Medium");
733 assert_eq!(cep.sessions, 2);
734 assert_eq!(
735 cep.total_cache_hits, 8,
736 "5 (session 1) + 3 (session 2, fresh)"
737 );
738 assert_eq!(cep.total_cache_reads, 16);
739 }
740}