perf(dashboard): 用量统计性能优化 — get_summary/get_day_detail 改异步执行不再阻塞 UI;range_totals 单遍统计(每次调用 7 次全量聚合 → 2 次);扫描缓存改 Arc 共享消除整表深拷贝;wire.jsonl 按 (mtime,len) 逐文件增量解析替换 8s TTL 整表缓存(920MB/6.3万条实测:冷扫描 ~17.7s → ~1.0s,任意间隔切换范围 ~1.5s 且 UI 全程可交互);新增单遍等价性与增量行为测试(含混合增删改 vs 强制全量重扫的全字段指纹比对),聚合口径零变化
This commit is contained in:
1 parent
ac82135dee
commit
d190a55bf2
1 file changed
+668
-67
+668
-67
@@ -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<UsageRecord>, 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<String>,
|
||||
alias2prov: HashMap<String, String>,
|
||||
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<UsageRecord>, 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<String>,
|
||||
alias2prov: &HashMap<String, String>,
|
||||
) -> 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<UsageRecord>,
|
||||
/// 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<PathBuf, FileScanEntry>,
|
||||
/// Last assembled table, shared with callers. When the walk finds nothing
|
||||
/// changed this pointer is handed out instead of rebuilding the table.
|
||||
result: Option<Arc<Vec<UsageRecord>>>,
|
||||
}
|
||||
|
||||
/// 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<Vec<UsageRecord>>, 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<String>| 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<String> = 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<PathBuf> = 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<UsageRecord>, 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<UsageRecord> =
|
||||
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<String>,
|
||||
@@ -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<UsageRecord>,
|
||||
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<UsageRecord>) -> Vec<UsageRecord> {
|
||||
///
|
||||
/// 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<Vec<UsageRecord>>) -> Arc<Vec<UsageRecord>> {
|
||||
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<UsageRecord>,
|
||||
meta: ScanMeta,
|
||||
static SCAN_CACHE: Mutex<Option<ScanCacheState>> = 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<Option<SummaryCacheEntry>> = 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<UsageRecord>, 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<Vec<UsageRecord>>, 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<String, TotalsRow> {
|
||||
const RANGES: [&str; 5] = ["today", "yesterday", "7d", "30d", "all"];
|
||||
let mut starts = [0u64; 5];
|
||||
let mut ends: [Option<u64>; 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<u64> {
|
||||
@@ -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<String>, range: Option<String>, refresh: Option<bool>) -> SummaryResult {
|
||||
let t0 = std::time::Instant::now();
|
||||
let refresh = refresh.unwrap_or(false);
|
||||
@@ -2172,11 +2365,7 @@ pub fn get_summary(home_override: Option<String>, range: Option<String>, 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<DailyRow> {
|
||||
|
||||
/// 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<String>, date: String) -> Option<DailyRow> {
|
||||
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::<Vec<_>>()
|
||||
.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<String> {
|
||||
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");
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user