diff --git a/src-tauri/src/shared/local_usage_core.rs b/src-tauri/src/shared/local_usage_core.rs index ca2f40871c..c73969ee08 100644 --- a/src-tauri/src/shared/local_usage_core.rs +++ b/src-tauri/src/shared/local_usage_core.rs @@ -21,13 +21,26 @@ struct DailyTotals { agent_runs: i64, } -#[derive(Default, Clone, Copy)] +#[derive(Default, Clone, Copy, PartialEq, Eq, Hash)] struct UsageTotals { input: i64, cached: i64, output: i64, } +#[derive(PartialEq, Eq, Hash)] +enum UsageIdentity { + Request(String), + Legacy(String, i64, UsageTotals), +} + +struct UsageEntry { + timestamp_ms: i64, + usage: UsageTotals, + model: String, + identity: UsageIdentity, +} + const MAX_ACTIVITY_GAP_MS: i64 = 2 * 60 * 1000; pub(crate) async fn local_usage_snapshot_core( @@ -72,6 +85,7 @@ fn scan_local_usage( .map(|key| (key.clone(), DailyTotals::default())) .collect(); let mut model_totals: HashMap = HashMap::new(); + let mut seen_usage = HashSet::new(); if sessions_roots.is_empty() { return Ok(build_snapshot(updated_at, day_keys, daily, HashMap::new())); @@ -92,7 +106,13 @@ fn scan_local_usage( if path.extension().and_then(|ext| ext.to_str()) != Some("jsonl") { continue; } - scan_file(&path, &mut daily, &mut model_totals, workspace_path)?; + scan_file( + &path, + &mut daily, + &mut model_totals, + workspace_path, + &mut seen_usage, + )?; } } } @@ -184,6 +204,7 @@ fn scan_file( daily: &mut HashMap, model_totals: &mut HashMap, workspace_path: Option<&Path>, + seen_usage: &mut HashSet, ) -> Result<(), String> { let file = match File::open(path) { Ok(file) => file, @@ -193,6 +214,9 @@ fn scan_file( }; let reader = BufReader::new(file); let mut previous_totals: Option = None; + let mut session_id = path.to_string_lossy().into_owned(); + let mut legacy_usage = Vec::new(); + let mut has_request_usage = false; let mut current_model: Option = None; let mut last_activity_ms: Option = None; let mut seen_runs: HashSet = HashSet::new(); @@ -237,6 +261,9 @@ fn scan_file( } if entry_type == "session_meta" { + if let Some(id) = value.pointer("/payload/id").and_then(Value::as_str) { + session_id = id.to_string(); + } continue; } @@ -251,6 +278,40 @@ fn scan_file( continue; } + if entry_type == "token_usage_record" { + if let (Some(usage), Some(timestamp_ms)) = ( + value.pointer("/payload/usage").and_then(Value::as_object), + read_timestamp_ms(&value), + ) { + has_request_usage = true; + legacy_usage.clear(); + let usage = read_usage_totals(usage); + let identity = value + .pointer("/payload/response_id") + .and_then(Value::as_str) + .filter(|id| !id.is_empty()) + .map(|id| UsageIdentity::Request(id.to_string())) + .unwrap_or_else(|| { + UsageIdentity::Legacy(session_id.clone(), timestamp_ms, usage) + }); + record_usage( + UsageEntry { + timestamp_ms, + usage, + model: current_model + .clone() + .unwrap_or_else(|| "unknown".to_string()), + identity, + }, + daily, + model_totals, + seen_usage, + ); + track_activity(daily, &mut last_activity_ms, timestamp_ms); + } + continue; + } + if entry_type == "event_msg" || entry_type.is_empty() { let payload = value.get("payload").and_then(|value| value.as_object()); let payload_type = payload @@ -282,95 +343,29 @@ fn scan_file( continue; } - let info = payload - .and_then(|payload| payload.get("info")) - .and_then(|v| v.as_object()); - let (input, cached, output, used_total) = if let Some(info) = info { - if let Some(total) = find_usage_map(info, &["total_token_usage", "totalTokenUsage"]) - { - ( - read_i64(total, &["input_tokens", "inputTokens"]), - read_i64( - total, - &[ - "cached_input_tokens", - "cache_read_input_tokens", - "cachedInputTokens", - "cacheReadInputTokens", - ], - ), - read_i64(total, &["output_tokens", "outputTokens"]), - true, - ) - } else if let Some(last) = - find_usage_map(info, &["last_token_usage", "lastTokenUsage"]) - { - ( - read_i64(last, &["input_tokens", "inputTokens"]), - read_i64( - last, - &[ - "cached_input_tokens", - "cache_read_input_tokens", - "cachedInputTokens", - "cacheReadInputTokens", - ], - ), - read_i64(last, &["output_tokens", "outputTokens"]), - false, - ) - } else { - continue; - } - } else { - continue; - }; - - let mut delta = UsageTotals { - input, - cached, - output, - }; - - if used_total { - let prev = previous_totals.unwrap_or_default(); - delta = UsageTotals { - input: (input - prev.input).max(0), - cached: (cached - prev.cached).max(0), - output: (output - prev.output).max(0), - }; - previous_totals = Some(UsageTotals { - input, - cached, - output, - }); - } else { - // Some streams emit `last_token_usage` deltas between `total_token_usage` snapshots. - // Treat those as already-counted to avoid double-counting when the next total arrives. - let mut next = previous_totals.unwrap_or_default(); - next.input += delta.input; - next.cached += delta.cached; - next.output += delta.output; - previous_totals = Some(next); - } - - if delta.input == 0 && delta.cached == 0 && delta.output == 0 { - continue; - } - let timestamp_ms = read_timestamp_ms(&value); - if let Some(day_key) = timestamp_ms.and_then(|ms| day_key_for_timestamp_ms(ms)) { - if let Some(entry) = daily.get_mut(&day_key) { - let cached = delta.cached.min(delta.input); - entry.input += delta.input; - entry.cached += cached; - entry.output += delta.output; - - let model = current_model - .clone() - .or_else(|| extract_model_from_token_count(&value)) - .unwrap_or_else(|| "unknown".to_string()); - *model_totals.entry(model).or_insert(0) += delta.input + delta.output; + if !has_request_usage { + if let Some(info) = payload + .and_then(|payload| payload.get("info")) + .and_then(Value::as_object) + { + if let (Some(usage), Some(timestamp_ms)) = + (legacy_usage_delta(info, &mut previous_totals), timestamp_ms) + { + legacy_usage.push(UsageEntry { + timestamp_ms, + usage, + model: current_model + .clone() + .or_else(|| extract_model_from_token_count(&value)) + .unwrap_or_else(|| "unknown".to_string()), + identity: UsageIdentity::Legacy( + session_id.clone(), + timestamp_ms, + usage, + ), + }); + } } } @@ -409,9 +404,91 @@ fn scan_file( } } + // Request records are authoritative for modern files. Buffer legacy counters so + // display events preceding (or delayed after) a request cannot count it twice. + if !has_request_usage { + for entry in legacy_usage { + record_usage(entry, daily, model_totals, seen_usage); + } + } + Ok(()) } +fn read_usage_totals(map: &serde_json::Map) -> UsageTotals { + UsageTotals { + input: read_i64(map, &["input_tokens", "inputTokens"]).max(0), + cached: read_i64( + map, + &[ + "cached_input_tokens", + "cache_read_input_tokens", + "cachedInputTokens", + "cacheReadInputTokens", + ], + ) + .max(0), + output: read_i64(map, &["output_tokens", "outputTokens"]).max(0), + } +} + +fn legacy_usage_delta( + info: &serde_json::Map, + previous_totals: &mut Option, +) -> Option { + let total = + find_usage_map(info, &["total_token_usage", "totalTokenUsage"]).map(read_usage_totals); + let last = find_usage_map(info, &["last_token_usage", "lastTokenUsage"]).map(read_usage_totals); + if let Some(total) = total { + let previous = previous_totals.replace(total); + if let Some(previous) = previous { + // Repeated display snapshots and counter resets are not new requests. + if total == previous + || total.input < previous.input + || total.cached < previous.cached + || total.output < previous.output + { + return None; + } + return last.or(Some(UsageTotals { + input: total.input - previous.input, + cached: total.cached - previous.cached, + output: total.output - previous.output, + })); + } + // A continuation can inherit arbitrary history; only its per-request + // usage is countable. Cumulative-only first records establish a baseline. + return last; + } + if let (Some(last), Some(previous)) = (last, previous_totals.as_mut()) { + previous.input += last.input; + previous.cached += last.cached; + previous.output += last.output; + } + last +} + +fn record_usage( + entry: UsageEntry, + daily: &mut HashMap, + model_totals: &mut HashMap, + seen_usage: &mut HashSet, +) { + let Some(day_key) = day_key_for_timestamp_ms(entry.timestamp_ms) else { + return; + }; + let Some(totals) = daily.get_mut(&day_key) else { + return; + }; + if !seen_usage.insert(entry.identity) { + return; + } + totals.input += entry.usage.input; + totals.cached += entry.usage.cached.min(entry.usage.input); + totals.output += entry.usage.output; + *model_totals.entry(entry.model).or_insert(0) += entry.usage.input + entry.usage.output; +} + fn extract_model_from_turn_context(value: &Value) -> Option { let payload = value.get("payload").and_then(|value| value.as_object())?; if let Some(model) = payload.get("model").and_then(|value| value.as_str()) { @@ -645,7 +722,14 @@ mod tests { let mut daily: HashMap = HashMap::new(); daily.insert(day_key.to_string(), DailyTotals::default()); let mut model_totals: HashMap = HashMap::new(); - scan_file(&path, &mut daily, &mut model_totals, None).expect("scan file"); + scan_file( + &path, + &mut daily, + &mut model_totals, + None, + &mut HashSet::new(), + ) + .expect("scan file"); let totals = daily.get(day_key).copied().unwrap_or_default(); assert_eq!(totals.input, 10); @@ -653,7 +737,7 @@ mod tests { } #[test] - fn scan_file_counts_last_deltas_before_total_snapshot_once() { + fn scan_file_treats_first_total_after_last_as_baseline() { let day_key = "2026-01-19"; let path = write_temp_jsonl(&[ r#"{"timestamp":"2026-01-19T12:00:00.000Z","payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":10,"cached_input_tokens":0,"output_tokens":5}}}}"#, @@ -663,11 +747,18 @@ mod tests { let mut daily: HashMap = HashMap::new(); daily.insert(day_key.to_string(), DailyTotals::default()); let mut model_totals: HashMap = HashMap::new(); - scan_file(&path, &mut daily, &mut model_totals, None).expect("scan file"); + scan_file( + &path, + &mut daily, + &mut model_totals, + None, + &mut HashSet::new(), + ) + .expect("scan file"); let totals = daily.get(day_key).copied().unwrap_or_default(); - assert_eq!(totals.input, 20); - assert_eq!(totals.output, 10); + assert_eq!(totals.input, 10); + assert_eq!(totals.output, 5); } #[test] @@ -682,11 +773,18 @@ mod tests { let mut daily: HashMap = HashMap::new(); daily.insert(day_key.to_string(), DailyTotals::default()); let mut model_totals: HashMap = HashMap::new(); - scan_file(&path, &mut daily, &mut model_totals, None).expect("scan file"); + scan_file( + &path, + &mut daily, + &mut model_totals, + None, + &mut HashSet::new(), + ) + .expect("scan file"); let totals = daily.get(day_key).copied().unwrap_or_default(); - assert_eq!(totals.input, 12); - assert_eq!(totals.output, 6); + assert_eq!(totals.input, 2); + assert_eq!(totals.output, 1); } #[test] @@ -700,7 +798,14 @@ mod tests { let mut daily: HashMap = HashMap::new(); daily.insert(day_key.to_string(), DailyTotals::default()); let mut model_totals: HashMap = HashMap::new(); - scan_file(&path, &mut daily, &mut model_totals, None).expect("scan file"); + scan_file( + &path, + &mut daily, + &mut model_totals, + None, + &mut HashSet::new(), + ) + .expect("scan file"); let totals = daily.get(day_key).copied().unwrap_or_default(); assert_eq!(totals.agent_ms, 5_000); @@ -717,7 +822,14 @@ mod tests { let mut daily: HashMap = HashMap::new(); daily.insert(day_key.to_string(), DailyTotals::default()); let mut model_totals: HashMap = HashMap::new(); - scan_file(&path, &mut daily, &mut model_totals, None).expect("scan file"); + scan_file( + &path, + &mut daily, + &mut model_totals, + None, + &mut HashSet::new(), + ) + .expect("scan file"); let totals = daily.get(day_key).copied().unwrap_or_default(); assert_eq!(totals.agent_runs, 2); @@ -735,7 +847,14 @@ mod tests { let mut daily: HashMap = HashMap::new(); daily.insert(day_key.to_string(), DailyTotals::default()); let mut model_totals: HashMap = HashMap::new(); - scan_file(&path, &mut daily, &mut model_totals, None).expect("scan file"); + scan_file( + &path, + &mut daily, + &mut model_totals, + None, + &mut HashSet::new(), + ) + .expect("scan file"); let totals = daily.get(day_key).copied().unwrap_or_default(); assert_eq!(totals.agent_ms, 10_000); @@ -758,6 +877,7 @@ mod tests { &mut daily, &mut model_totals, Some(Path::new("/tmp/other-project")), + &mut HashSet::new(), ) .expect("scan file"); @@ -786,10 +906,10 @@ mod tests { let root_b = make_temp_sessions_root(); let line_a = format!( - r#"{{"timestamp":{timestamp_ms},"payload":{{"type":"token_count","info":{{"total_token_usage":{{"input_tokens":5,"cached_input_tokens":0,"output_tokens":2}}}}}}}}"# + r#"{{"timestamp":{timestamp_ms},"payload":{{"type":"token_count","info":{{"last_token_usage":{{"input_tokens":5,"cached_input_tokens":0,"output_tokens":2}}}}}}}}"# ); let line_b = format!( - r#"{{"timestamp":{timestamp_ms},"payload":{{"type":"token_count","info":{{"total_token_usage":{{"input_tokens":3,"cached_input_tokens":0,"output_tokens":1}}}}}}}}"# + r#"{{"timestamp":{timestamp_ms},"payload":{{"type":"token_count","info":{{"last_token_usage":{{"input_tokens":3,"cached_input_tokens":0,"output_tokens":1}}}}}}}}"# ); write_session_file(&root_a, &day_key, &[line_a]); @@ -807,6 +927,187 @@ mod tests { assert_eq!(snapshot.totals.last30_days_tokens, 11); } + fn fixture_timestamp() -> i64 { + Local::now().timestamp_millis() + } + + fn legacy_record(timestamp: i64, total: i64, last: Option) -> Value { + let mut info = serde_json::json!({ + "total_token_usage": {"input_tokens": total, "output_tokens": 0} + }); + if let Some(last) = last { + info["last_token_usage"] = + serde_json::json!({"input_tokens": last, "output_tokens": 0}); + } + serde_json::json!({ + "timestamp": timestamp, "type": "event_msg", + "payload": {"type": "token_count", "info": info} + }) + } + + fn request_record(timestamp: i64, response_id: &str, usage: Value) -> Value { + serde_json::json!({ + "timestamp": timestamp, "type": "token_usage_record", + "payload": {"response_id": response_id, "usage": usage} + }) + } + + fn scan_usage_fixtures(fixtures: Vec>, filter: Option<&Path>) -> LocalUsageSnapshot { + let day_keys = make_day_keys(1); + let mut daily = HashMap::from([(day_keys[0].clone(), DailyTotals::default())]); + let mut models = HashMap::new(); + let mut seen = HashSet::new(); + for fixture in fixtures { + let mut lines = vec![ + serde_json::json!({"type":"session_meta","payload":{"id":"session-a","cwd":"/tmp/project-a"}}).to_string(), + serde_json::json!({"type":"turn_context","payload":{"model":"test-model"}}).to_string(), + ]; + lines.extend(fixture.into_iter().map(|value| value.to_string())); + let refs: Vec<&str> = lines.iter().map(String::as_str).collect(); + let path = write_temp_jsonl(&refs); + scan_file(&path, &mut daily, &mut models, filter, &mut seen).expect("scan fixture"); + fs::remove_file(path).expect("remove fixture"); + } + build_snapshot(0, day_keys, daily, models) + } + + #[test] + fn resumed_legacy_file_counts_last_usage_and_deduplicates_copied_events() { + let timestamp = fixture_timestamp(); + let record = legacy_record(timestamp, 1_000_010, Some(10)); + let snapshot = scan_usage_fixtures( + vec![vec![record.clone(), record.clone()], vec![record]], + None, + ); + assert_eq!(snapshot.totals.last30_days_tokens, 10); + assert_eq!(snapshot.top_models[0].tokens, 10); + } + + #[test] + fn cumulative_only_records_establish_baselines_after_start_and_reset() { + let timestamp = fixture_timestamp(); + let snapshot = scan_usage_fixtures( + vec![vec![ + legacy_record(timestamp, 1000, None), + legacy_record(timestamp + 1, 1020, None), + legacy_record(timestamp + 2, 3, None), + legacy_record(timestamp + 3, 8, None), + ]], + None, + ); + assert_eq!(snapshot.days[0].input_tokens, 25); + } + + #[test] + fn request_records_override_display_counters_regardless_of_event_order() { + let timestamp = fixture_timestamp(); + let snapshot = scan_usage_fixtures( + vec![vec![ + legacy_record(timestamp, 1_000_010, Some(10)), + request_record( + timestamp + 1, + "response-a", + serde_json::json!({"input_tokens":10,"output_tokens":5}), + ), + legacy_record(timestamp + 5000, 1_000_010, Some(10)), + ]], + None, + ); + assert_eq!(snapshot.days[0].input_tokens, 10); + assert_eq!(snapshot.days[0].output_tokens, 5); + assert_eq!(snapshot.top_models[0].tokens, 15); + } + + #[test] + fn request_totals_include_cache_once_and_do_not_add_reasoning_or_stale_totals() { + let timestamp = fixture_timestamp(); + let snapshot = scan_usage_fixtures( + vec![vec![ + request_record( + timestamp, + "response-a", + serde_json::json!({ + "input_tokens":100,"cached_input_tokens":80,"output_tokens":20, + "reasoning_output_tokens":15,"total_tokens":999999 + }), + ), + request_record( + timestamp + 1, + "response-b", + serde_json::json!({ + "input_tokens":10,"cached_input_tokens":50,"output_tokens":2 + }), + ), + ]], + None, + ); + assert_eq!(snapshot.days[0].input_tokens, 110); + assert_eq!(snapshot.days[0].cached_input_tokens, 90); + assert_eq!(snapshot.days[0].output_tokens, 22); + assert_eq!(snapshot.days[0].total_tokens, 132); + } + + #[test] + fn request_deduplication_spans_session_roots_without_dropping_distinct_requests() { + let timestamp = fixture_timestamp(); + let day_key = day_key_for_timestamp_ms(timestamp).expect("day key"); + let roots = [make_temp_sessions_root(), make_temp_sessions_root()]; + let record = request_record( + timestamp, + "response-a", + serde_json::json!({"input_tokens":10,"output_tokens":5}), + ); + write_session_file(&roots[0], &day_key, &[record.to_string()]); + write_session_file( + &roots[1], + &day_key, + &[ + record.to_string(), + request_record( + timestamp, + "response-b", + serde_json::json!({"input_tokens":10,"output_tokens":5}), + ) + .to_string(), + ], + ); + let snapshot = scan_local_usage(1, None, &roots).expect("scan roots"); + assert_eq!(snapshot.totals.last30_days_tokens, 30); + for root in roots { + fs::remove_dir_all(root).expect("remove root"); + } + } + + #[test] + fn workspace_filter_also_applies_to_request_records() { + let snapshot = scan_usage_fixtures( + vec![vec![request_record( + fixture_timestamp(), + "response-a", + serde_json::json!({"input_tokens":10,"output_tokens":5}), + )]], + Some(Path::new("/tmp/other-project")), + ); + assert_eq!(snapshot.totals.last30_days_tokens, 0); + assert!(snapshot.top_models.is_empty()); + } + + #[test] + fn legacy_usage_accepts_camel_case_counters() { + let snapshot = scan_usage_fixtures( + vec![vec![serde_json::json!({ + "timestamp":fixture_timestamp(),"payload":{"type":"token_count","info":{ + "totalTokenUsage":{"inputTokens":1000,"outputTokens":500}, + "lastTokenUsage":{"inputTokens":10,"cachedInputTokens":8,"outputTokens":5} + }} + })]], + None, + ); + assert_eq!(snapshot.days[0].input_tokens, 10); + assert_eq!(snapshot.days[0].cached_input_tokens, 8); + assert_eq!(snapshot.days[0].output_tokens, 5); + } + #[test] fn resolve_sessions_roots_uses_single_default_root() { let mut workspaces = HashMap::new();