From d190a55bf2e2583c03e6a1450c728cbc0406bc8e Mon Sep 17 00:00:00 2001 From: KimiSwitch Dev Date: Wed, 7 Oct 2026 23:35:17 +0800 Subject: [PATCH] =?UTF-8?q?perf(dashboard):=20=E7=94=A8=E9=87=8F=E7=BB=9F?= =?UTF-8?q?=E8=AE=A1=E6=80=A7=E8=83=BD=E4=BC=98=E5=8C=96=20=E2=80=94=20get?= =?UTF-8?q?=5Fsummary/get=5Fday=5Fdetail=20=E6=94=B9=E5=BC=82=E6=AD=A5?= =?UTF-8?q?=E6=89=A7=E8=A1=8C=E4=B8=8D=E5=86=8D=E9=98=BB=E5=A1=9E=20UI?= =?UTF-8?q?=EF=BC=9Brange=5Ftotals=20=E5=8D=95=E9=81=8D=E7=BB=9F=E8=AE=A1?= =?UTF-8?q?=EF=BC=88=E6=AF=8F=E6=AC=A1=E8=B0=83=E7=94=A8=207=20=E6=AC=A1?= =?UTF-8?q?=E5=85=A8=E9=87=8F=E8=81=9A=E5=90=88=20=E2=86=92=202=20?= =?UTF-8?q?=E6=AC=A1=EF=BC=89=EF=BC=9B=E6=89=AB=E6=8F=8F=E7=BC=93=E5=AD=98?= =?UTF-8?q?=E6=94=B9=20Arc=20=E5=85=B1=E4=BA=AB=E6=B6=88=E9=99=A4=E6=95=B4?= =?UTF-8?q?=E8=A1=A8=E6=B7=B1=E6=8B=B7=E8=B4=9D=EF=BC=9Bwire.jsonl=20?= =?UTF-8?q?=E6=8C=89=20(mtime,len)=20=E9=80=90=E6=96=87=E4=BB=B6=E5=A2=9E?= =?UTF-8?q?=E9=87=8F=E8=A7=A3=E6=9E=90=E6=9B=BF=E6=8D=A2=208s=20TTL=20?= =?UTF-8?q?=E6=95=B4=E8=A1=A8=E7=BC=93=E5=AD=98=EF=BC=88920MB/6.3=E4=B8=87?= =?UTF-8?q?=E6=9D=A1=E5=AE=9E=E6=B5=8B=EF=BC=9A=E5=86=B7=E6=89=AB=E6=8F=8F?= =?UTF-8?q?=20~17.7s=20=E2=86=92=20~1.0s=EF=BC=8C=E4=BB=BB=E6=84=8F?= =?UTF-8?q?=E9=97=B4=E9=9A=94=E5=88=87=E6=8D=A2=E8=8C=83=E5=9B=B4=20~1.5s?= =?UTF-8?q?=20=E4=B8=94=20UI=20=E5=85=A8=E7=A8=8B=E5=8F=AF=E4=BA=A4?= =?UTF-8?q?=E4=BA=92=EF=BC=89=EF=BC=9B=E6=96=B0=E5=A2=9E=E5=8D=95=E9=81=8D?= =?UTF-8?q?=E7=AD=89=E4=BB=B7=E6=80=A7=E4=B8=8E=E5=A2=9E=E9=87=8F=E8=A1=8C?= =?UTF-8?q?=E4=B8=BA=E6=B5=8B=E8=AF=95=EF=BC=88=E5=90=AB=E6=B7=B7=E5=90=88?= =?UTF-8?q?=E5=A2=9E=E5=88=A0=E6=94=B9=20vs=20=E5=BC=BA=E5=88=B6=E5=85=A8?= =?UTF-8?q?=E9=87=8F=E9=87=8D=E6=89=AB=E7=9A=84=E5=85=A8=E5=AD=97=E6=AE=B5?= =?UTF-8?q?=E6=8C=87=E7=BA=B9=E6=AF=94=E5=AF=B9=EF=BC=89=EF=BC=8C=E8=81=9A?= =?UTF-8?q?=E5=90=88=E5=8F=A3=E5=BE=84=E9=9B=B6=E5=8F=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src-tauri/src/dashboard.rs | 735 +++++++++++++++++++++++++++++++++---- 1 file changed, 668 insertions(+), 67 deletions(-) diff --git a/src-tauri/src/dashboard.rs b/src-tauri/src/dashboard.rs index d0c830f..54fef46 100644 --- a/src-tauri/src/dashboard.rs +++ b/src-tauri/src/dashboard.rs @@ -2,7 +2,7 @@ use chrono::{DateTime, Datelike, Local, TimeZone}; use crate::db::{ArchivedSessionSnapshot, DayModelStat, DayStat}; use regex::Regex; use serde::{Deserialize, Serialize}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::fs; use std::io::{BufRead, BufReader}; use std::path::{Path, PathBuf}; @@ -833,8 +833,23 @@ fn cost_for_usage(input_other: u64, output: u64, cache_read: u64, cache_create: /// even after the secondary config is removed or changed. const SECONDARY_MODEL_CACHE_KEY: &str = "dashboard.secondary_model_cache"; -fn scan_usage(home: &Path) -> (Vec, ScanMeta) { - let root = sessions_root(home); +/// The two external inputs the wire parser prices a record with, resolved +/// once per scan: the configured secondary-model alias and the legacy +/// alias→provider map. `fingerprint` changes whenever either of them does, +/// which is what invalidates the per-file parse cache. +struct ScanInputs { + secondary_alias: Option, + alias2prov: HashMap, + fingerprint: String, +} + +/// Resolve the scan inputs for `home`. +/// +/// `read_secondary_model_alias` and `build_alias_provider_map` used to run +/// inside the scanner on every pass; they move up here so the incremental +/// cache can compare them before walking the tree, and so a call that reuses +/// every cached file never repeats the config reads. +fn resolve_scan_inputs(home: &Path) -> ScanInputs { let alias2prov = build_alias_provider_map(home); // Resolve the `__secondary__` marker (subagent requests bound to the // configured secondary model) to the real model alias once per scan. @@ -851,20 +866,113 @@ fn scan_usage(home: &Path) -> (Vec, ScanMeta) { if let Some(alias) = ¤t_secondary { let _ = crate::db::set_setting_pub(SECONDARY_MODEL_CACHE_KEY, alias); } - let mut records = Vec::new(); - let mut files_scanned = 0; - let mut lines_seen = 0; - let errors = Vec::new(); + let fingerprint = scan_inputs_fingerprint(&secondary_alias, &alias2prov); + ScanInputs { secondary_alias, alias2prov, fingerprint } +} + +/// Change key for `ScanInputs`: the resolved secondary alias plus the +/// sorted alias→provider pairs, so any edit either input picked up from +/// config.toml or its backups flips it. +fn scan_inputs_fingerprint( + secondary_alias: &Option, + alias2prov: &HashMap, +) -> String { + let mut pairs: Vec<(&str, &str)> = alias2prov + .iter() + .map(|(k, v)| (k.as_str(), v.as_str())) + .collect(); + pairs.sort_unstable(); + let mut out = String::new(); + out.push_str(secondary_alias.as_deref().unwrap_or("")); + out.push('\u{1}'); + for (k, v) in pairs { + out.push_str(k); + out.push('\u{1}'); + out.push_str(v); + out.push('\u{1}'); + } + out +} + +/// Cached parse of one `wire.jsonl`. Both halves of the stat key are kept: a +/// change in either one invalidates this entry. +struct FileScanEntry { + mtime_ms: u64, + len: u64, + records: Vec, + /// Total lines of the file. `ScanMeta.lines_seen` counts every line, not + /// only the usage ones, so the count has to survive with the records. + lines_seen: usize, +} + +/// Per-file incremental cache for the sessions walk. +/// +/// Replaces the old 8s TTL: a file keeps its parsed records while its +/// (mtime, len) is unchanged, so a refresh only re-parses what actually +/// moved. `parse_ctx` holds the `ScanInputs` fingerprint — when the secondary +/// alias or the alias→provider map moves, every cached parse is stale. +struct ScanCacheState { + home: String, + parse_ctx: String, + files: HashMap, + /// Last assembled table, shared with callers. When the walk finds nothing + /// changed this pointer is handed out instead of rebuilding the table. + result: Option>>, +} + +/// Modification time of an already-statted file in epoch ms (0 when the +/// platform reports nothing usable). +fn file_mtime_ms(meta: &fs::Metadata) -> u64 { + meta.modified() + .ok() + .and_then(|t| t.duration_since(SystemTime::UNIX_EPOCH).ok()) + .map(|d| d.as_millis() as u64) + .unwrap_or(0) +} + +/// Walk the sessions tree and assemble the usage table, reusing the cached +/// parse of every `wire.jsonl` whose (mtime, len) is unchanged. +/// +/// The record set and its newest-first ordering are identical whether a file +/// was reused or re-parsed: reuse is keyed on the file stat alone, and a file +/// that cannot be stat'ed or read is re-parsed (and left uncached) rather +/// than trusted. +fn scan_usage_incremental( + state: &mut ScanCacheState, + home: &Path, + inputs: &ScanInputs, +) -> (Arc>, ScanMeta) { + let root = sessions_root(home); + let home_s = home.to_string_lossy().to_string(); + let root_s = root.to_string_lossy().to_string(); + let meta = |files_scanned: usize, lines_seen: usize, record_count: usize, errors: Vec| ScanMeta { + files_scanned, + lines_seen, + record_count, + home: home_s.clone(), + sessions_root: root_s.clone(), + errors: errors.into_iter().take(20).collect(), + }; if !root.exists() { - return (records, ScanMeta { - files_scanned: 0, lines_seen: 0, record_count: 0, - home: home.to_string_lossy().to_string(), - sessions_root: root.to_string_lossy().to_string(), - errors: vec!["sessions directory not found".into()], - }); + state.files.clear(); + state.result = None; + return ( + Arc::new(Vec::new()), + meta(0, 0, 0, vec!["sessions directory not found".into()]), + ); } + let mut files_scanned = 0usize; + let mut lines_seen = 0usize; + let errors: Vec = Vec::new(); + // Every path this walk produced; anything else left in `state.files` is + // gone from disk and gets dropped below. + let mut walked: HashSet = HashSet::new(); + // Stays true only when the walk reproduced the cached file set and stat + // exactly, which is what allows the assembled table to be reused. + let mut unchanged = true; + for entry in WalkDir::new(&root).into_iter().filter_map(|e| e.ok()) { if entry.file_name() != "wire.jsonl" { continue; } if !entry.file_type().is_file() { continue; } @@ -874,27 +982,67 @@ fn scan_usage(home: &Path) -> (Vec, ScanMeta) { continue; } files_scanned += 1; - let content = match fs::read_to_string(entry.path()) { + let path = p.to_path_buf(); + walked.insert(path.clone()); + // Only a successful stat yields a key, so a file that cannot be + // stat'ed is always re-parsed and never enters the cache. + let stat = fs::metadata(&path).ok().map(|m| (file_mtime_ms(&m), m.len())); + if let (Some((mtime_ms, len)), Some(cached)) = (stat, state.files.get(&path)) { + if cached.mtime_ms == mtime_ms && cached.len == len { + lines_seen += cached.lines_seen; + continue; + } + } + let content = match fs::read_to_string(&path) { Ok(c) => c, - Err(_) => continue, + // Unreadable now: drop any stale parse rather than counting it. + Err(_) => { + state.files.remove(&path); + unchanged = false; + continue; + } }; - lines_seen += content.lines().count(); - records.extend(parse_wire_usage_records(&content, &secondary_alias, &alias2prov)); + let file_lines = content.lines().count(); + let records = parse_wire_usage_records(&content, &inputs.secondary_alias, &inputs.alias2prov); + lines_seen += file_lines; + unchanged = false; + if let Some((mtime_ms, len)) = stat { + state.files.insert(path, FileScanEntry { mtime_ms, len, records, lines_seen: file_lines }); + } } + let before = state.files.len(); + state.files.retain(|p, _| walked.contains(p)); + if state.files.len() != before { + unchanged = false; + } + + if unchanged { + if let Some(result) = &state.result { + return (Arc::clone(result), meta(files_scanned, lines_seen, result.len(), errors)); + } + } + + // Assemble the table from the reused plus freshly parsed records. Files + // are visited in path order so records sharing a timestamp keep a stable + // order regardless of hash or directory iteration order. + let mut paths: Vec<&PathBuf> = state.files.keys().collect(); + paths.sort_unstable(); + let mut records: Vec = + Vec::with_capacity(state.files.values().map(|e| e.records.len()).sum()); + for p in paths { + records.extend(state.files[p].records.iter().cloned()); + } records.sort_by(|a, b| b.time.cmp(&a.time)); - let count = records.len(); - (records, ScanMeta { - files_scanned, lines_seen, record_count: count, - home: home.to_string_lossy().to_string(), - sessions_root: root.to_string_lossy().to_string(), - errors: errors.into_iter().take(20).collect(), - }) + let result = Arc::new(records); + state.result = Some(Arc::clone(&result)); + (Arc::clone(&result), meta(files_scanned, lines_seen, result.len(), errors)) } -/// Parse one `wire.jsonl` body into usage records. Extracted from -/// `scan_usage` so a single session can be snapshotted without walking the -/// whole sessions tree; the record filter chain is unchanged. +/// Parse one `wire.jsonl` body into usage records. Extracted from the +/// sessions walk so a single session can be snapshotted without walking the +/// whole sessions tree; the record filter chain is unchanged, and the +/// incremental cache reuses this exact parser. fn parse_wire_usage_records( content: &str, secondary_alias: &Option, @@ -1045,7 +1193,7 @@ fn synthesize_snapshot_records(rows: &[ArchivedSessionSnapshot], home: &Path) -> } /// Append the synthesized records for `rows` to a live scan result and -/// restore the newest-first ordering `scan_usage` produces. +/// restore the newest-first ordering the sessions walk produces. fn merge_snapshot_records( mut records: Vec, rows: &[ArchivedSessionSnapshot], @@ -1061,48 +1209,55 @@ fn merge_snapshot_records( /// Live scan records plus the stored usage of archived sessions whose files /// have been deleted, so dashboard stats survive session deletion. -fn records_with_archive_merge(home: &Path, records: Vec) -> Vec { +/// +/// With nothing archived the scanned `Arc` is handed straight back — the +/// merged table is only materialized (one copy) when there are snapshots to +/// append, which is the rare path. +fn records_with_archive_merge(home: &Path, records: Arc>) -> Arc> { let rows = crate::db::list_archived_sessions().unwrap_or_default(); - merge_snapshot_records(records, &rows, home) + if rows.is_empty() { return records; } + Arc::new(merge_snapshot_records((*records).clone(), &rows, home)) } // --------------------------------------------------------------------------- // Scan cache // --------------------------------------------------------------------------- -/// Short-TTL cache for the expensive `scan_usage` walk (reads every -/// wire.jsonl under the sessions root and prices every record). Tab switches -/// call get_summary repeatedly; re-walking the whole tree on every switch is -/// the main dashboard lag. Keyed by home, expires after 8s; `refresh=true` -/// bypasses it. -struct SummaryCacheEntry { - home: String, - scanned_at: std::time::Instant, - records: Vec, - meta: ScanMeta, +static SCAN_CACHE: Mutex> = Mutex::new(None); + +/// Drop everything the incremental cache holds for tests. Production never +/// calls this — a stale cache is corrected by the file stats, not by a purge. +#[cfg(test)] +fn clear_scan_cache() { + *SCAN_CACHE.lock().unwrap() = None; } -static SUMMARY_CACHE: Mutex> = Mutex::new(None); -const SUMMARY_CACHE_TTL: std::time::Duration = std::time::Duration::from_secs(8); - -fn scan_usage_cached(home: &Path, refresh: bool) -> (Vec, ScanMeta) { +/// The record table is shared behind an `Arc` so the unchanged path hands out +/// a pointer bump instead of a deep copy of the whole scan (tens of thousands +/// of records on a used install). +/// +/// `refresh=true` forces every file to be re-parsed (the cache for this home +/// is dropped first), matching the old "bypass the cache" semantics; an +/// ordinary call only re-parses files whose (mtime, len) moved. +fn scan_usage_cached(home: &Path, refresh: bool) -> (Arc>, ScanMeta) { let home_s = home.to_string_lossy().to_string(); - let mut cache = SUMMARY_CACHE.lock().unwrap(); - let hit = cache + let inputs = resolve_scan_inputs(home); + let mut cache = SCAN_CACHE.lock().unwrap(); + // A different home or a moved parsing input invalidates every cached + // parse, so the owner is replaced wholesale rather than patched. + let reusable = cache .as_ref() - .is_some_and(|e| !refresh && e.home == home_s && e.scanned_at.elapsed() < SUMMARY_CACHE_TTL); - if hit { - let e = cache.as_ref().unwrap(); - return (e.records.clone(), e.meta.clone()); + .is_some_and(|s| !refresh && s.home == home_s && s.parse_ctx == inputs.fingerprint); + if !reusable { + *cache = Some(ScanCacheState { + home: home_s, + parse_ctx: inputs.fingerprint.clone(), + files: HashMap::new(), + result: None, + }); } - let (records, meta) = scan_usage(home); - *cache = Some(SummaryCacheEntry { - home: home_s, - scanned_at: std::time::Instant::now(), - records: records.clone(), - meta: meta.clone(), - }); - (records, meta) + let state = cache.as_mut().unwrap(); + scan_usage_incremental(state, home, &inputs) } // --------------------------------------------------------------------------- @@ -1294,6 +1449,44 @@ fn filter_by_range<'a>(records: &'a [UsageRecord], range: &str, now_ms: u64) -> } } +/// Totals for every dashboard range in one pass over `records`, instead of one +/// full scan per range. Membership follows `filter_by_range` exactly, and the +/// accumulation/finalization mirror `aggregate`'s so the numbers are identical. +fn range_totals_single_pass(records: &[UsageRecord], now_ms: u64) -> HashMap { + const RANGES: [&str; 5] = ["today", "yesterday", "7d", "30d", "all"]; + let mut starts = [0u64; 5]; + let mut ends: [Option; 5] = [None; 5]; + for (i, r) in RANGES.iter().enumerate() { + starts[i] = range_start(r, now_ms); + ends[i] = range_end(r, now_ms); + } + + let mut totals: [TotalsRow; 5] = std::array::from_fn(|_| TotalsRow::default()); + for rec in records { + for i in 0..RANGES.len() { + let in_range = if RANGES[i] == "all" { + true + } else { + match ends[i] { + Some(end) => rec.time >= starts[i] && rec.time < end, + None => rec.time >= starts[i], + } + }; + if in_range { totals[i].add(rec); } + } + } + + RANGES + .iter() + .zip(totals) + .map(|(range, mut t)| { + let total_input = t.input_other + t.input_cache_read + t.input_cache_creation; + t.cache_hit_rate = if total_input > 0 { t.input_cache_read as f64 / total_input as f64 } else { 0.0 }; + (range.to_string(), t) + }) + .collect() +} + /// Upper bound (exclusive) for bounded ranges. Only "yesterday" has one — /// it must not bleed into today. fn range_end(range: &str, now_ms: u64) -> Option { @@ -2147,7 +2340,7 @@ pub fn get_prices() -> PricesResult { PricesResult { prices: list_prices() } } -#[tauri::command] +#[tauri::command(async)] pub fn get_summary(home_override: Option, range: Option, refresh: Option) -> SummaryResult { let t0 = std::time::Instant::now(); let refresh = refresh.unwrap_or(false); @@ -2172,11 +2365,7 @@ pub fn get_summary(home_override: Option, range: Option, refresh }).collect(); let all_model_count = all_models.len(); - let mut range_totals = HashMap::new(); - for r_k in ["today", "yesterday", "7d", "30d", "all"] { - let s = aggregate(&records, r_k, now_ms); - range_totals.insert(r_k.to_string(), s.totals); - } + let range_totals = range_totals_single_pass(&records, now_ms); let default_model = None; let env_model = std::env::var("KIMI_MODEL_NAME").ok().map(|name| EnvModelInfo { @@ -2261,7 +2450,7 @@ fn build_day_detail(date: &str, records: &[UsageRecord]) -> Option { /// Lazy per-day detail for the heatmap double-click modal (fetched on demand, /// so `get_summary` stays lean). Returns null when the date has no records. -#[tauri::command] +#[tauri::command(async)] pub fn get_day_detail(home_override: Option, date: String) -> Option { let home = resolve_kimi_home(home_override); let (records, _) = scan_usage_cached(&home, false); @@ -2722,6 +2911,82 @@ mod range_and_model_totals_tests { let total: u64 = stats.models_by_name.iter().map(|m| m.total_tokens).sum(); assert_eq!(total, 3 * 200); } + + fn rec_amounts(time: u64, tokens: (u64, u64, u64, u64), cost: f64) -> UsageRecord { + UsageRecord { + time, + model: "openai/gpt-x".to_string(), + input_other: tokens.0, + output: tokens.1, + input_cache_read: tokens.2, + input_cache_creation: tokens.3, + cost_usd: cost, + cost_estimated: false, + price_id: String::new(), + model_resolved: "gpt-x".to_string(), + model_display: "openai/gpt-x".to_string(), + provider: Some("openai".to_string()), + from_env: false, + is_secondary: false, + } + } + + fn assert_totals_eq(got: &TotalsRow, want: &TotalsRow, range: &str) { + assert_eq!(got.requests, want.requests, "{range}: requests"); + assert_eq!(got.input_other, want.input_other, "{range}: input_other"); + assert_eq!(got.output, want.output, "{range}: output"); + assert_eq!(got.input_cache_read, want.input_cache_read, "{range}: input_cache_read"); + assert_eq!(got.input_cache_creation, want.input_cache_creation, "{range}: input_cache_creation"); + assert_eq!(got.cost_usd, want.cost_usd, "{range}: cost_usd"); + assert_eq!(got.total_tokens, want.total_tokens, "{range}: total_tokens"); + assert_eq!(got.cache_hit_rate, want.cache_hit_rate, "{range}: cache_hit_rate"); + } + + #[test] + fn single_pass_range_totals_match_per_range_aggregate() { + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH).unwrap().as_millis() as u64; + let today_start = local_midnight_ms(now, 0); + let yesterday_start = local_midnight_ms(now, 1); + let day = 24 * 3600 * 1000u64; + + let records = vec![ + // Boundary hits: exactly at (and just inside) each day edge. + rec_amounts(today_start, (10, 1, 1, 0), 0.1), + rec_amounts(today_start + 1, (20, 2, 2, 1), 0.2), + rec_amounts(yesterday_start, (30, 3, 3, 0), 0.3), + rec_amounts(yesterday_start + 1, (40, 4, 4, 2), 0.4), + // Just before yesterday's midnight → outside yesterday and today. + rec_amounts(yesterday_start.saturating_sub(1), (50, 5, 5, 5), 0.5), + // Inside 7d / outside today+yesterday. + rec_amounts(now - 3 * day, (60, 6, 6, 3), 0.6), + // Exactly at the 7d cutoff → inside 7d. + rec_amounts(now - 7 * day, (70, 7, 7, 7), 0.7), + // 20d ago → only 30d and all. + rec_amounts(now - 20 * day, (80, 8, 8, 4), 0.8), + // 45d ago → only all. + rec_amounts(now - 45 * day, (90, 9, 9, 9), 0.9), + // Zero-token record: pins the cache_hit_rate 0/0 → 0.0 path. + rec_amounts(now - 2 * day, (0, 0, 0, 0), 0.0), + ]; + + let single = range_totals_single_pass(&records, now); + assert_eq!(single.len(), 5, "every dashboard range is present"); + for range in ["today", "yesterday", "7d", "30d", "all"] { + let want = aggregate(&records, range, now).totals; + let got = single.get(range).expect("range present in single pass"); + assert_totals_eq(got, &want, range); + } + + // Sanity: the ranges really do differ, so the comparison above is not + // trivially satisfied by all-equal numbers. + let all = single.get("all").unwrap(); + let today = single.get("today").unwrap(); + assert_eq!(all.requests, records.len()); + assert_eq!(today.requests, 2, "only the two today-boundary records"); + assert!(all.requests > today.requests); + assert!(all.cache_hit_rate > 0.0, "cache_hit_rate is finalized, not left at 0"); + } } #[cfg(test)] @@ -2979,15 +3244,351 @@ mod archive_snapshot_tests { assert_eq!(row.day_stats[0].models[0].input_other, 100); // While the files are still there, the live scan owns the usage. - let merged = records_with_archive_merge(&home, vec![]); + let merged = records_with_archive_merge(&home, Arc::new(vec![])); assert!(merged.is_empty(), "no synthesized records while the dir exists"); // Once the archived directory is deleted, stats come from the snapshot. fs::remove_dir_all(&old_dir).unwrap(); - let merged = records_with_archive_merge(&home, vec![]); + let merged = records_with_archive_merge(&home, Arc::new(vec![])); assert_eq!(merged.len(), 1); assert_eq!(merged[0].input_other, 100); assert_eq!(merged[0].output, 50); assert_eq!(day_key(merged[0].time), day_key(old_ms)); } } + +#[cfg(test)] +mod scan_cache_tests { + use super::*; + + const WID: &str = "wd_proj_a1b2c3d4e5f6"; + const SID: &str = "session_11111111-2222-3333-4444-555555555555"; + + /// `scan_usage_cached` keeps one `static` cache, so the tests that want to + /// observe reuse (rather than just correctness) take this lock to keep a + /// second test from replacing the state mid-run. Results stay correct + /// either way — a foreign home or a purged cache only forces a full rescan + /// — but ptr-equality assertions need a quiet cache. + fn scan_lock() -> std::sync::MutexGuard<'static, ()> { + static LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + let guard = LOCK.lock().unwrap_or_else(|e| e.into_inner()); + clear_scan_cache(); + guard + } + + /// A fresh temp home with one session directory and its `wire.jsonl`. + fn setup_home_with_wire(lines: &[serde_json::Value]) -> (tempfile::TempDir, PathBuf, PathBuf) { + let td = tempfile::tempdir().unwrap(); + let home = td.path().join("home"); + let dir = home.join("sessions").join(WID).join(SID); + fs::create_dir_all(&dir).unwrap(); + let wire = dir.join("wire.jsonl"); + fs::write(&wire, wire_body(lines)).unwrap(); + (td, home, wire) + } + + fn wire_body(lines: &[serde_json::Value]) -> String { + lines + .iter() + .map(|l| format!("{l}\n")) + .collect::>() + .concat() + } + + fn usage_line(time: u64, model: &str) -> serde_json::Value { + serde_json::json!({ + "type": "usage.record", + "usageScope": "turn", + "time": time, + "model": model, + "usage": { + "inputOther": 100, + "output": 50, + "inputCacheRead": 10, + "inputCacheCreation": 5, + }, + }) + } + + #[test] + fn first_scan_parses_the_wire_file_and_reuses_the_table_when_nothing_moved() { + let _guard = scan_lock(); + let (_td, home, _wire) = setup_home_with_wire(&[usage_line(1_700_000_000_000, "openai/gpt-x")]); + + let (records, meta) = scan_usage_cached(&home, false); + assert_eq!(records.len(), 1); + assert_eq!(meta.files_scanned, 1); + assert_eq!(meta.lines_seen, 1); + assert_eq!(meta.record_count, 1); + assert!(meta.errors.is_empty(), "{:?}", meta.errors); + assert_eq!(records[0].model, "openai/gpt-x"); + assert_eq!(records[0].provider.as_deref(), Some("openai")); + assert_eq!(records[0].time, 1_700_000_000_000); + + // Nothing changed on disk: the assembled table itself is handed back. + let (again, meta2) = scan_usage_cached(&home, false); + assert!(Arc::ptr_eq(&records, &again), "unchanged scan must reuse the cached table"); + assert_eq!(meta2.record_count, 1); + assert_eq!(meta2.lines_seen, 1); + } + + #[test] + fn appended_record_is_reparsed_and_extends_the_table() { + let _guard = scan_lock(); + let (_td, home, wire) = setup_home_with_wire(&[usage_line(1_700_000_000_000, "openai/gpt-x")]); + let (first, _) = scan_usage_cached(&home, false); + assert_eq!(first.len(), 1); + + // Appending changes the file length, so the cached parse is stale. + let mut body = fs::read_to_string(&wire).unwrap(); + body.push_str(&format!("{}\n", usage_line(1_700_000_100_000, "openai/gpt-y"))); + fs::write(&wire, body).unwrap(); + + let (records, meta) = scan_usage_cached(&home, false); + assert_eq!(records.len(), 2, "the new record must show up"); + assert_eq!(meta.record_count, 2); + assert_eq!(meta.lines_seen, 2); + assert_eq!(meta.files_scanned, 1); + assert!(!Arc::ptr_eq(&first, &records), "a changed file must rebuild the table"); + // Newest first, matching the full re-scan ordering. + assert_eq!(records[0].model, "openai/gpt-y"); + assert_eq!(records[1].model, "openai/gpt-x"); + } + + #[test] + fn deleted_file_drops_its_records() { + let _guard = scan_lock(); + let (_td, home, wire) = setup_home_with_wire(&[usage_line(1_700_000_000_000, "openai/gpt-x")]); + assert_eq!(scan_usage_cached(&home, false).0.len(), 1); + + fs::remove_file(&wire).unwrap(); + + let (records, meta) = scan_usage_cached(&home, false); + assert!(records.is_empty(), "a deleted file must stop contributing records"); + assert_eq!(meta.files_scanned, 0); + assert_eq!(meta.lines_seen, 0); + assert_eq!(meta.record_count, 0); + } + + #[test] + fn a_second_file_adds_its_records_on_top_of_the_cached_ones() { + let _guard = scan_lock(); + let (_td, home, _wire) = setup_home_with_wire(&[usage_line(1_700_000_000_000, "openai/gpt-x")]); + assert_eq!(scan_usage_cached(&home, false).0.len(), 1); + + let other = home.join("sessions").join(WID) + .join("session_22222222-2222-3333-4444-555555555555"); + fs::create_dir_all(&other).unwrap(); + fs::write( + other.join("wire.jsonl"), + wire_body(&[usage_line(1_700_000_200_000, "moonshotai/kimi-k3")]), + ) + .unwrap(); + + let (records, meta) = scan_usage_cached(&home, false); + assert_eq!(records.len(), 2); + assert_eq!(meta.files_scanned, 2); + assert_eq!(meta.lines_seen, 2); + assert_eq!(records[0].model, "moonshotai/kimi-k3"); + assert_eq!(records[1].model, "openai/gpt-x"); + } + + #[test] + fn blob_and_task_wire_files_stay_ignored() { + let _guard = scan_lock(); + let (_td, home, _wire) = setup_home_with_wire(&[usage_line(1_700_000_000_000, "openai/gpt-x")]); + // Records reachable only through a blobs/ or tasks/ directory. + for sub in ["blobs", "tasks"] { + let dir = home.join("sessions").join(WID).join(SID).join(sub); + fs::create_dir_all(&dir).unwrap(); + fs::write( + dir.join("wire.jsonl"), + wire_body(&[usage_line(1_700_000_300_000, "openai/skipped")]), + ) + .unwrap(); + } + + let (records, meta) = scan_usage_cached(&home, false); + assert_eq!(records.len(), 1, "blobs/tasks directories stay filtered out"); + assert_eq!(meta.files_scanned, 1); + assert!(records.iter().all(|r| r.model != "openai/skipped")); + } + + #[test] + fn changed_alias_provider_map_invalidates_cached_parses() { + let _guard = scan_lock(); + // A bare alias with no `provider/` prefix only resolves through the + // legacy alias→provider map, so a config change flips its provider. + let line = usage_line(1_700_000_000_000, "legacy-alias"); + let (td, home, _wire) = setup_home_with_wire(&[line]); + + let (before, _) = scan_usage_cached(&home, false); + assert_eq!(before.len(), 1); + assert_eq!(before[0].provider, None, "no config yet leaves the provider unresolved"); + + // Adding the alias to config.toml moves the parse context, so the + // cached parse of an untouched wire.jsonl must be thrown away. + fs::write( + td.path().join("home").join("config.toml"), + "[models.legacy-alias]\nprovider = \"CodingPlanSite\"\n", + ) + .unwrap(); + + let (after, meta) = scan_usage_cached(&home, false); + assert_eq!(after.len(), 1); + assert_eq!(meta.files_scanned, 1); + assert_eq!(meta.lines_seen, 1); + assert_eq!( + after[0].provider.as_deref(), + Some("CodingPlan.site"), + "the new alias→provider map must be applied" + ); + } + + #[test] + fn refresh_forces_a_reparse_of_unchanged_files() { + let _guard = scan_lock(); + let (_td, home, _wire) = setup_home_with_wire(&[usage_line(1_700_000_000_000, "openai/gpt-x")]); + + let (cached, _) = scan_usage_cached(&home, false); + let (refreshed, meta) = scan_usage_cached(&home, true); + + assert!(!Arc::ptr_eq(&cached, &refreshed), "refresh=true must not hand back the cached table"); + assert_eq!(refreshed.len(), 1); + assert_eq!(meta.files_scanned, 1); + assert_eq!(meta.record_count, 1); + // Same numbers as the incremental path — only the work differs. + assert_eq!(refreshed[0].model, cached[0].model); + assert_eq!(refreshed[0].cost_usd, cached[0].cost_usd); + } + + #[test] + fn missing_sessions_root_reports_the_error_and_empty_table() { + let _guard = scan_lock(); + let td = tempfile::tempdir().unwrap(); + let home = td.path().join("no-such-home"); + + let (records, meta) = scan_usage_cached(&home, false); + assert!(records.is_empty()); + assert_eq!(meta.files_scanned, 0); + assert_eq!(meta.record_count, 0); + assert_eq!(meta.errors, vec!["sessions directory not found".to_string()]); + } + + /// Every field of every record, in order — the strongest equality check + /// available (`UsageRecord` implements neither `PartialEq` nor + /// `Serialize`). + fn records_fingerprint(records: &[UsageRecord]) -> Vec { + records + .iter() + .map(|r| { + format!( + "{}|{}|{}|{}|{}|{}|{}|{}|{}|{}|{}|{:?}|{}|{}|{:?}", + r.time, + r.model, + r.model_resolved, + r.model_display, + r.input_other, + r.output, + r.input_cache_read, + r.input_cache_creation, + r.cost_usd, + r.cost_estimated, + r.price_id, + r.provider, + r.from_env, + r.is_secondary, + r.cost_usd.is_nan(), + ) + }) + .collect() + } + + #[test] + fn mixed_edits_match_a_forced_full_rescan_exactly() { + let _guard = scan_lock(); + let (_td, home, wire_a) = setup_home_with_wire(&[ + usage_line(1_700_000_000_000, "openai/gpt-x"), + usage_line(1_700_000_600_000, "kimi/k3"), + ]); + let sessions = home.join("sessions").join(WID); + let wire_b = sessions.join("session_22222222-2222-3333-4444-555555555555").join("wire.jsonl"); + let wire_c = sessions.join("session_33333333-2222-3333-4444-555555555555").join("wire.jsonl"); + for w in [&wire_b, &wire_c] { + fs::create_dir_all(w.parent().unwrap()).unwrap(); + fs::write(w, wire_body(&[usage_line(1_700_000_300_000, "openai/gpt-y")])).unwrap(); + } + // A blobs copy that must never be counted by either scan. + let blob = sessions.join(SID).join("blobs").join("wire.jsonl"); + fs::create_dir_all(blob.parent().unwrap()).unwrap(); + fs::write(&blob, wire_body(&[usage_line(1_700_000_900_000, "openai/ignored")])).unwrap(); + + let (warm, warm_meta) = scan_usage_cached(&home, false); + assert_eq!(warm.len(), 4); + assert_eq!(warm_meta.files_scanned, 3); + + // Touch every file in a different way between the two scans. + fs::write(&wire_b, wire_body(&[ + usage_line(1_700_000_300_000, "__secondary__"), + usage_line(1_700_001_000_000, "zhipuai/glm-4.6"), + ])).unwrap(); + fs::remove_file(&wire_c).unwrap(); + let mut a = fs::read_to_string(&wire_a).unwrap(); + a.push_str(&format!("{}\n", usage_line(1_700_001_200_000, "legacy-alias"))); + fs::write(&wire_a, a).unwrap(); + + let (incremental, inc_meta) = scan_usage_cached(&home, false); + // The same tree read from scratch, with the cache dropped. + let (full, full_meta) = scan_usage_cached(&home, true); + + assert_eq!( + records_fingerprint(&incremental), + records_fingerprint(&full), + "incremental scan must equal a full rescan, field for field" + ); + assert_eq!(inc_meta.files_scanned, full_meta.files_scanned); + assert_eq!(inc_meta.lines_seen, full_meta.lines_seen); + assert_eq!(inc_meta.record_count, full_meta.record_count); + assert_eq!(inc_meta.record_count, incremental.len()); + assert_eq!(inc_meta.errors, full_meta.errors); + + // Sanity on the contents the two scans agreed on: A grew to 3 lines, + // B was rewritten with 2, C was deleted, and the blob copy never counts. + assert_eq!(incremental.len(), 5); + assert_eq!(incremental[0].model, "legacy-alias"); + assert_eq!(incremental[1].model, "zhipuai/glm-4.6"); + assert_eq!(incremental[2].model, "kimi/k3"); + let secondary = incremental.iter().find(|r| r.model == "__secondary__").expect("secondary record"); + assert!(secondary.is_secondary, "the marker stays a subagent record"); + assert!(incremental.iter().all(|r| r.model != "openai/ignored")); + assert!(incremental.iter().all(|r| r.model != "openai/gpt-y"), "the deleted file is gone"); + // Newest-first throughout. + assert!(incremental.windows(2).all(|w| w[0].time >= w[1].time)); + } + + #[test] + fn secondary_alias_change_invalidates_cached_parses() { + let _guard = scan_lock(); + let (td, home, _wire) = setup_home_with_wire(&[usage_line(1_700_000_000_000, "__secondary__")]); + + let (before, _) = scan_usage_cached(&home, false); + assert_eq!(before.len(), 1); + assert!(before[0].is_secondary); + let first_cost = before[0].cost_usd; + + // Point [secondary_model] at a much pricier model; the wire file itself + // is untouched, so only the parse context can catch this. + fs::write( + td.path().join("home").join("config.toml"), + "[secondary_model]\nmodel = \"openai/gpt-5\"\n", + ) + .unwrap(); + + let (after, meta) = scan_usage_cached(&home, false); + assert_eq!(after.len(), 1); + assert_eq!(meta.files_scanned, 1); + assert_eq!(meta.lines_seen, 1, "the reused count still comes from the file"); + assert_eq!(after[0].model, "__secondary__", "the record keeps its stable marker"); + assert!(after[0].cost_usd > first_cost, "the new secondary model re-prices the record"); + } +}