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 = 2;
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 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
105/// Debounced flush: persists at most once per FLUSH_INTERVAL_SECS.
106/// Call from post_dispatch on every tool call to bound data loss
107/// on abrupt MCP termination to at most 2 seconds of events.
108pub 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
123/// Adjust saved tokens after post-processing (terse, hints) changed the output size.
124/// Positive delta = savings were over-reported, negative = under-reported.
125///
126/// Persist immediately: this adjustment follows a durable `record()` and must not
127/// be lost when a short-lived MCP process exits before the periodic CEP flush.
128pub 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
194/// Records a tool result against an observed provider turn. `turn == 0` keeps
195/// daemon-free callers honest: first injection is known, re-read count is not.
196pub 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    // Tool-call accounting is durable per event, matching the savings ledger.
209    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            // Stamp the running version so a mid-day update attributes the day
244            // to the release in use for its latest activity (#307).
245            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    // MCP stdio servers are routinely terminated without a graceful shutdown.
272    // Delaying this write made the append-only ledger survive while aggregate
273    // stats disappeared, so persist every completed accounting event.
274    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/// Fold one CEP snapshot into `cep`. Pure (no globals, no I/O) so the
351/// delta/aggregation rules are unit-testable in isolation.
352///
353/// `cache_hits`, `cache_reads`, `tokens_original` and `tokens_compressed` arrive
354/// as **cumulative per-process** counters. For repeated snapshots within the same
355/// PID only the delta since the previous snapshot is added, so the lifetime
356/// totals keep tracking cache activity instead of freezing at the first
357/// checkpoint's value (#361). A new PID starts a fresh session and seeds the
358/// cumulative baselines.
359#[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    /// #706: a corrupt `stats.json` must never be silently replaced by an
450    /// empty store. The loader quarantines the bytes to `stats.json.corrupt`
451    /// (recoverable), and an existing quarantine is never overwritten — the
452    /// older copy is the one closest to the lost history.
453    #[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        // A second corruption must NOT clobber the first quarantine.
474        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        // And a healthy file still loads normally.
484        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        // The common case: only the primary dir has stats. Aggregation must
492        // return the primary untouched so non-split users see no change (#500).
493        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        // A data-dir split (#408/#414/#500): the CLI's primary dir is empty but
503        // the MCP server wrote its savings into a sibling tree. The displayed
504        // total must reflect both so `gain` no longer reports a false `0`.
505        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, &current, &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, &current, &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, &current, &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, &current, &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(&current, &baseline).is_none());
622
623        std::fs::remove_file(lock_path).unwrap();
624        let merged = io::merge_and_save(&current, &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, &current, &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, &current_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        // #307: the most recent known version is carried into the merge.
692        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        // #361: repeated snapshots within one process must keep counting cache
711        // hits/reads (cumulative counters → add the delta), not freeze at the
712        // first checkpoint's value while only tokens advanced.
713        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        // Same PID, cumulative counters grew: hits 2→9, reads 4→20.
717        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}