From cd161f4401c2cd2d5c17a2a6a5b6eaacc1945d2a Mon Sep 17 00:00:00 2001 From: Jason Date: Thu, 23 Jul 2026 17:00:03 +0800 Subject: [PATCH] feat(usage): import Grok Build official-mode usage from session logs Grok CLI's official OAuth mode cannot be routed through the local proxy (empty config is the mode switch, so there is no injection point), which left official-mode usage invisible. Add session_usage_grokbuild to import usage from ~/.grok/sessions updates.jsonl: - Only turn_completed events carry usage; each event is the independent per-turn total (accumulated across inference loops within one prompt), so events are imported at face value. Do not reintroduce differencing of adjacent events: counters reset every turn and differencing would massively under-record. - Cost priority: reported costUsdTicks (1 tick = 1e-10 USD) wins when complete, because the backfill only repairs rows with total <= 0 and can never correct a positive mispriced value; local pricing fills the breakdown and raises a drift warning above max(1% of reported, 1e-6). costIsPartial marks the reported value a lower bound: prefer a full local recompute when the model is priced, else record the lower bound. - Idempotency key grok_session:{session}:{prompt_id}:{model} anchors on the upstream per-turn UUID (index fallback only when prompt_id is empty), so rewind truncation cannot shift keys and double count; orphan rows from truncated turns are kept since the tokens were spent. - Anti double-count vs proxy takeover: 10min settle window plus a time-window guard over recent grokbuild proxy activity; guarded skips never mark files as synced. - Seed grok-4.5-build pricing 2/6/0.30, back-derived from exact costUsdTicks samples (cache read bills at 0.30, not the listed 0.50). - Map _grok_session to a friendly provider display name and refresh the takeover-capability comment in services/proxy.rs. --- src-tauri/src/database/schema.rs | 4 + src-tauri/src/services/mod.rs | 1 + src-tauri/src/services/proxy.rs | 7 +- src-tauri/src/services/session_usage.rs | 5 + .../src/services/session_usage_grokbuild.rs | 1151 +++++++++++++++++ src-tauri/src/services/usage_stats.rs | 35 + 6 files changed, 1201 insertions(+), 2 deletions(-) create mode 100644 src-tauri/src/services/session_usage_grokbuild.rs diff --git a/src-tauri/src/database/schema.rs b/src-tauri/src/database/schema.rs index 0dc967bfc..9d25d9dbd 100644 --- a/src-tauri/src/database/schema.rs +++ b/src-tauri/src/database/schema.rs @@ -2246,6 +2246,10 @@ impl Database { ("qwen3-32b", "Qwen3 32B", "0.16", "0.64", "0", "0"), // Grok 系列 (xAI) ("grok-4.5", "Grok 4.5", "2", "6", "0.50", "0"), + // Grok CLI 官方 OAuth 态 modelUsage 上报的内部别名。定价由 + // costUsdTicks(1 tick = 1e-10 USD)双轮实测反推:input/output 与 + // grok-4.5 同为 2/6,cache read 实际按 0.30 计(非 API 挂牌的 0.50) + ("grok-4.5-build", "Grok 4.5 Build", "2", "6", "0.30", "0"), ("grok-4.3", "Grok 4.3", "1.25", "2.50", "0.20", "0"), ( "grok-4.20-0309-reasoning", diff --git a/src-tauri/src/services/mod.rs b/src-tauri/src/services/mod.rs index a796c34e4..1d824cec0 100644 --- a/src-tauri/src/services/mod.rs +++ b/src-tauri/src/services/mod.rs @@ -17,6 +17,7 @@ pub mod s3_sync; pub mod session_usage; pub mod session_usage_codex; pub mod session_usage_gemini; +pub mod session_usage_grokbuild; pub mod session_usage_opencode; pub mod skill; pub mod speedtest; diff --git a/src-tauri/src/services/proxy.rs b/src-tauri/src/services/proxy.rs index 175f0b939..4efee1537 100644 --- a/src-tauri/src/services/proxy.rs +++ b/src-tauri/src/services/proxy.rs @@ -1505,8 +1505,11 @@ impl ProxyService { /// Grok Build live 是否具备可接管的自定义模型表。 /// /// 官方态 live(Grok CLI 自带 OAuth 登录、无 `[model.*]` 表)没有注入 - /// 占位符的落点,且官方供应商本就禁止经代理接管(封号风险),调用方 - /// 应跳过接管或直接报错。 + /// 占位符的落点:Grok CLI 以「config 是否为空」区分官方 OAuth / 自定义 + /// 供应商两种模式,表达不出「官方 OAuth + 自定义 base_url」。官方供应商 + /// 的接管能力门见 `official_provider_supports_proxy_takeover`(按应用逐个 + /// 开,目前仅 Codex),调用方应跳过接管或直接报错。官方态的用量统计由 + /// `session_usage_grokbuild` 从会话日志导入,不依赖代理。 fn grok_live_config_supports_takeover(config: &Value) -> bool { config .get("config") diff --git a/src-tauri/src/services/session_usage.rs b/src-tauri/src/services/session_usage.rs index 37e414964..357588844 100644 --- a/src-tauri/src/services/session_usage.rs +++ b/src-tauri/src/services/session_usage.rs @@ -86,6 +86,11 @@ pub fn sync_all_unlocked(db: &Database) -> SessionSyncResult { "OpenCode", crate::services::session_usage_opencode::sync_opencode_usage(db), ); + merge_sync_step( + &mut result, + "Grok Build", + crate::services::session_usage_grokbuild::sync_grokbuild_usage(db), + ); notify_sync_result(&result); result } diff --git a/src-tauri/src/services/session_usage_grokbuild.rs b/src-tauri/src/services/session_usage_grokbuild.rs new file mode 100644 index 000000000..7f0be931a --- /dev/null +++ b/src-tauri/src/services/session_usage_grokbuild.rs @@ -0,0 +1,1151 @@ +//! Grok Build (Grok CLI) 会话用量追踪 +//! +//! 从 `~/.grok/{sessions,archived_sessions}///updates.jsonl` +//! 的 `turn_completed` 事件中提取用量,写入 proxy_request_logs,实现官方 +//! OAuth 直连态(无代理数据)下的用量统计。 +//! +//! ## 数据流 +//! ```text +//! updates.jsonl(逐轮 turn_completed) → 沉降窗/接管守卫 → 费用计算 → proxy_request_logs +//! ``` +//! +//! ## 事件口径(2026-07-23 单进程双 prompt 实测 + CLI 二进制逆向双重确证) +//! - `sessionUpdate == "turn_completed"` 事件的 usage 是【该 user prompt 一轮 +//! 的独立总量】:轮内跨 inference loop 累加(`modelCalls`/`numTurns` = 本轮 +//! loop 数),下一轮从零起算。【不是】进程或会话累计——进程累计走 CLI 内 +//! 另一条独立通道(`GetSessionUsage`,"since start or last resume"),不落 +//! updates.jsonl。🔴 勿改回相邻事件差分:那是把每轮总量误当累计快照,会把 +//! 第二轮记成两轮之差造成巨量漏记(曾犯,实测单进程双 prompt 证伪)。 +//! - 逐事件按面值入账即为正确的逐轮记录;两轮数值完全相同 = 两笔真实用量, +//! 照常都入账。 +//! - `reasoningTokens` ⊂ `outputTokens`(totalTokens = input + output,且 +//! costUsdTicks 反推 output 未加计 reasoning),不参与计费。 +//! - `costUsdTicks`(1 tick = 1e-10 USD)是 CLI 自报的本轮精确成本,6 个实测 +//! 样本与本地定价 grok-4.5-build 2/6/0.30 分毫不差。**有自报且完整时 +//! total_cost 以自报为准**(回填只补 total<=0 的行、不修正错价,入账后无 +//! 修复路径,所以定价漂移窗口不能押在本地价上);本地定价负责分项成本与 +//! 漂移告警。`costIsPartial` 标记自报为下界:有本地价回退本地全额复算并 +//! 抑制漂移告警,无价才用下界入账(分项记 0)。 +//! - 防接管态双算不用指纹去重:接管态下 CLI 照写 updates.jsonl,但轮事件是 +//! 聚合值(多 loop 求和),与代理逐请求行结构性不相等。改用「沉降窗 + +//! 接管活动时间窗守卫」:只导入足够旧的事件(届时接管态的代理行必已 +//! 落库),插入前按事件时刻查询附近是否存在代理直录行(见 +//! `has_recent_grokbuild_proxy_activity`)。 + +use crate::database::{lock_conn, Database}; +use crate::error::AppError; +use crate::proxy::usage::calculator::CostCalculator; +use crate::proxy::usage::parser::TokenUsage; +use crate::services::session_usage::{ + get_sync_state, metadata_modified_nanos, update_sync_state, SessionSyncResult, +}; +use crate::services::sql_helpers::INPUT_TOKEN_SEMANTICS_TOTAL; +use crate::services::usage_stats::{ + find_model_pricing, has_recent_grokbuild_proxy_activity, SESSION_PROXY_DEDUP_WINDOW_SECONDS, +}; +use rust_decimal::Decimal; +use std::fs; +use std::path::{Path, PathBuf}; +use std::time::SystemTime; + +/// 事件沉降窗:只导入早于「现在 − 窗口」的事件。 +/// +/// 接管态下 CLI 照写 updates.jsonl,同一请求代理已逐请求记账;代理行与 +/// 会话事件几乎同时产生,若导入抢在代理行落库前运行,接管守卫会因查不到 +/// 代理行而放行,双算永久留存。让事件先「沉降」再导入后,守卫查询必然 +/// 能看到已落库的代理行,竞态从源头消除。代价:官方态用量最多延迟约一个 +/// 窗口 + 一次后台同步周期(60s)上屏。 +const SETTLE_WINDOW_SECONDS: i64 = SESSION_PROXY_DEDUP_WINDOW_SECONDS; + +/// 单个模型的本轮用量(从 `modelUsage` 或顶层 usage 提取,均为逐轮口径) +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +struct GrokCounters { + input: u64, + output: u64, + cached: u64, + api_ms: u64, + model_calls: u64, + /// CLI 自报本轮成本,1 tick = 1e-10 USD;0 = 上游未提供 + cost_ticks: u64, + /// 上游标记 cost_ticks 只是部分费用(`costIsPartial`):此时它是下界 + cost_partial: bool, +} + +impl GrokCounters { + fn is_zero(&self) -> bool { + self.input == 0 && self.output == 0 && self.cached == 0 + } + + fn reported_cost_usd(&self) -> Option { + (self.cost_ticks > 0) + .then(|| Decimal::from(self.cost_ticks) / Decimal::from(10_000_000_000u64)) + } +} + +/// 一条 `turn_completed` 用量事件 +#[derive(Debug)] +struct GrokUsageEvent { + created_at: i64, + prompt_id: String, + /// 事件级 `costIsPartial`(顶层 usage 上观测到的位置;对本事件全部模型生效) + cost_is_partial: bool, + per_model: Vec<(String, GrokCounters)>, +} + +/// 同步 Grok Build 使用数据(从 updates.jsonl 会话日志) +pub fn sync_grokbuild_usage(db: &Database) -> Result { + let files = collect_grok_updates_files(); + + let mut result = SessionSyncResult { + files_scanned: files.len() as u32, + ..Default::default() + }; + + for file_path in &files { + match sync_single_grok_file(db, file_path) { + Ok(file_result) => result.merge(file_result), + Err(e) => { + let msg = format!("Grok Build 会话文件解析失败 {}: {e}", file_path.display()); + log::warn!("[GROK-SYNC] {msg}"); + result.errors.push(msg); + } + } + } + + if result.imported > 0 { + log::info!( + "[GROK-SYNC] 同步完成: 导入 {} 条, 跳过 {} 条, 扫描 {} 个文件, 延后 {} 个文件", + result.imported, + result.skipped, + result.files_scanned, + result.deferred_files + ); + } + + Ok(result) +} + +/// 收集所有 Grok 会话的 updates.jsonl(含归档会话,与会话浏览器同根) +fn collect_grok_updates_files() -> Vec { + let mut files = Vec::new(); + for root in crate::session_manager::providers::grokbuild::session_roots() { + collect_files_named(&root, "updates.jsonl", &mut files); + } + files +} + +/// 递归收集目录下指定文件名的文件(容忍布局深度变化,对齐会话浏览器的做法) +fn collect_files_named(root: &Path, name: &str, files: &mut Vec) { + let Ok(entries) = fs::read_dir(root) else { + return; + }; + for entry in entries.flatten() { + let path = entry.path(); + if path.is_dir() { + collect_files_named(&path, name, files); + } else if path.file_name().and_then(|n| n.to_str()) == Some(name) { + files.push(path); + } + } +} + +/// 同步单个 updates.jsonl 文件 +fn sync_single_grok_file(db: &Database, file_path: &Path) -> Result { + let file_path_str = file_path.to_string_lossy().to_string(); + + let metadata = fs::metadata(file_path) + .map_err(|e| AppError::Config(format!("无法读取文件元数据: {e}")))?; + let file_modified = metadata_modified_nanos(&metadata); + + let (last_modified, _last_offset) = get_sync_state(db, &file_path_str)?; + if file_modified <= last_modified { + return Ok(SessionSyncResult::default()); + } + + // 文件变更时全量重读:UPSERT 幂等使重读无害,且沉降窗延后的事件本就 + // 依赖下一轮重读补入。事件已是逐轮独立值,改 offset 增量读在正确性上 + // 可行(无差分基线依赖),但需另行处理延后事件的 offset 回退,收益 + // (活跃会话每周期省一次 O(N) 解析)暂不值得该复杂度。 + let content = fs::read_to_string(file_path) + .map_err(|e| AppError::Config(format!("无法读取文件: {e}")))?; + let events = parse_grok_usage_events(&content); + + // 会话 ID = 会话目录名(与 summary.json 的 info.id 一致)。request_id + // 唯一性押在该 UUIDv7 全局唯一上:同 ID 的归档/活跃副本经 UPSERT 幂等 + // 收敛(有意),不同 下撞 ID 视为不可能。 + let session_id = file_path + .parent() + .and_then(|dir| dir.file_name()) + .and_then(|n| n.to_str()) + .unwrap_or("unknown") + .to_string(); + + let now = SystemTime::now() + .duration_since(SystemTime::UNIX_EPOCH) + .map(|d| d.as_secs() as i64) + .unwrap_or(0); + + let mut result = SessionSyncResult::default(); + let mut deferred = false; + + for (idx, event) in events.iter().enumerate() { + // 沉降窗:事件按 append 顺序时间单调,遇到第一条未沉降的事件即停, + // 后续事件与它一起等下一轮(保持"文件前缀已导入"的简单不变量)。 + // 已知局限:未来时间戳(时钟误设)会让该文件持续延后并整文件重扫, + // 墙钟越过 事件时刻+窗口 后自愈;活跃会话每周期全量重读为设计代价。 + if now.saturating_sub(event.created_at) < SETTLE_WINDOW_SECONDS { + deferred = true; + break; + } + + // 接管守卫按事件时刻判定一次,整条事件的所有模型行同进退; + // 被守卫跳过的 token 已由代理行记账,跳过即终态(同步状态照常 + // 推进)。已知局限:守卫无 session 维度,见 usage_stats.rs 注释。 + let takeover_active = { + let conn = lock_conn!(db.conn); + has_recent_grokbuild_proxy_activity(&conn, event.created_at)? + }; + + for (model, turn) in &event.per_model { + if turn.is_zero() { + continue; + } + if takeover_active { + // 计入 skipped(对齐 gemini 指纹去重跳过的语义:未入账,代理 + // 行权威)。勿改用 suspected_duplicates——codex 对它的语义相反 + // (已入账待查),而 merge() 会把两义直接求和。 + result.skipped += 1; + continue; + } + + // 幂等键锚定上游稳定 ID(prompt_id 是每轮唯一的 UUID),不含文件 + // 内序号:updates.jsonl 前缀被改写(如 rewind 截断)导致事件序号 + // 前移时,幸存轮次仍命中原行不会双算;被移除轮次的行保留—— + // rewind 不退还已消耗的 token,留存即正确记账。若上游对同一 + // prompt_id 写多条 turn_completed(未观测到),UPSERT 取后者, + // 方向是少记不双算。prompt_id 缺失时回退 "idx{N}"(UUID 形态的 + // prompt_id 不可能与之撞名)。 + let turn_key = if event.prompt_id.is_empty() { + format!("idx{idx}") + } else { + event.prompt_id.clone() + }; + let request_id = format!("grok_session:{session_id}:{turn_key}:{model}"); + match insert_grok_session_entry( + db, + &request_id, + turn, + event.cost_is_partial || turn.cost_partial, + model, + &session_id, + event.created_at, + ) { + Ok(true) => result.imported += 1, + Ok(false) => result.skipped += 1, + Err(e) => { + log::warn!("[GROK-SYNC] 插入失败 ({request_id}): {e}"); + result.skipped += 1; + } + } + } + } + + if deferred { + // 不落同步状态:下一轮重读整个文件,把沉降后的事件补入。 + result.deferred_files += 1; + } else { + update_sync_state(db, &file_path_str, file_modified, events.len() as i64)?; + } + + Ok(result) +} + +/// 从 updates.jsonl 内容解析出全部逐轮用量事件(保持文件顺序) +fn parse_grok_usage_events(content: &str) -> Vec { + let mut events = Vec::new(); + + for line in content.lines() { + let line = line.trim(); + if line.is_empty() { + continue; + } + let Ok(record) = serde_json::from_str::(line) else { + continue; + }; + if record.get("method").and_then(|v| v.as_str()) != Some("_x.ai/session/update") { + continue; + } + let update = record.get("params").and_then(|p| p.get("update")); + // 只认 turn_completed(实测全体带 usage 的事件均为此类;判别字段是 + // sessionUpdate,serde internally-tagged)。字段缺失时向后兼容放行, + // 但显式标为其它类型的事件即使带 usage 也不导入——中途快照若与轮末 + // 事件并存,双导会双算。 + let kind = update + .and_then(|u| u.get("sessionUpdate")) + .and_then(|v| v.as_str()); + if kind.is_some() && kind != Some("turn_completed") { + continue; + } + let Some(usage) = update + .and_then(|u| u.get("usage")) + .filter(|u| u.is_object()) + else { + continue; + }; + // 沉降窗与接管守卫都依赖事件时刻,没有时间戳的事件无法安全导入。 + let Some(created_at) = parse_event_timestamp(record.get("timestamp")) else { + continue; + }; + + let prompt_id = update + .and_then(|u| u.get("prompt_id")) + .and_then(|v| v.as_str()) + .unwrap_or("") + .to_string(); + + let mut per_model: Vec<(String, GrokCounters)> = usage + .get("modelUsage") + .and_then(|m| m.as_object()) + .map(|map| { + map.iter() + .map(|(model, counters)| (model.clone(), parse_grok_counters(counters))) + .collect() + }) + .unwrap_or_default(); + if per_model.is_empty() { + // 缺 modelUsage 时退回顶层逐轮值;模型名未知,交由查价层兜底。 + per_model.push(("unknown".to_string(), parse_grok_counters(usage))); + } + // modelUsage 是 JSON object,遍历序不保证稳定;排序保证插入顺序 + // 与日志在多次重扫间确定。 + per_model.sort_by(|a, b| a.0.cmp(&b.0)); + + events.push(GrokUsageEvent { + created_at, + prompt_id, + cost_is_partial: usage + .get("costIsPartial") + .and_then(|v| v.as_bool()) + .unwrap_or(false), + per_model, + }); + } + + events +} + +fn parse_grok_counters(value: &serde_json::Value) -> GrokCounters { + let get = |key: &str| value.get(key).and_then(|v| v.as_u64()).unwrap_or(0); + GrokCounters { + input: get("inputTokens"), + output: get("outputTokens"), + cached: get("cachedReadTokens"), + api_ms: get("apiDurationMs"), + model_calls: get("modelCalls"), + cost_ticks: get("costUsdTicks"), + cost_partial: value + .get("costIsPartial") + .and_then(|v| v.as_bool()) + .unwrap_or(false), + } +} + +/// updates.jsonl 顶层 `timestamp` 实测为数字 epoch 秒(勿与 summary.json 的 +/// RFC3339 字符串混淆);字符串形态仅作防御性兜底。 +fn parse_event_timestamp(value: Option<&serde_json::Value>) -> Option { + let value = value?; + if let Some(n) = value.as_i64() { + // 防未来毫秒形态:超过 1e11 视作毫秒 + return Some(if n > 100_000_000_000 { n / 1000 } else { n }); + } + value + .as_str() + .and_then(|ts| chrono::DateTime::parse_from_rfc3339(ts).ok()) + .map(|dt| dt.timestamp()) +} + +/// 插入单条 Grok 会话记录到 proxy_request_logs +fn insert_grok_session_entry( + db: &Database, + request_id: &str, + turn: &GrokCounters, + cost_is_partial: bool, + model: &str, + session_id: &str, + created_at: i64, +) -> Result { + let conn = lock_conn!(db.conn); + + let clamp = |v: u64| v.min(u32::MAX as u64) as u32; + let usage = TokenUsage { + input_tokens: clamp(turn.input), + output_tokens: clamp(turn.output), + cache_read_tokens: clamp(turn.cached), + cache_creation_tokens: 0, + model: Some(model.to_string()), + message_id: None, + }; + + let pricing = find_model_pricing(&conn, model); + let multiplier = Decimal::from(1); + let reported = turn.reported_cost_usd(); + // 插入成功(changed)后才发,避免重扫时重复刷日志 + let mut deferred_warn: Option = None; + + // total_cost 取值优先级(🔴 回填机制只补 total<=0 的行、从不修正已有正值, + // 见 backfill_missing_usage_costs;本导入器 UPSERT 也不因 cost 单独变化而 + // 更新——所以入账时就必须写对,事后没有修复路径): + // 1. 有自报且完整 → 以自报为准(上游 ground truth,定价漂移窗口内也准确; + // 本地定价负责分项与漂移告警,漂移时分项与 total 允许暂不自洽); + // 2. 自报不完整(costIsPartial)→ 有本地价用本地全额复算(token 数完整), + // 并抑制此时无意义的漂移告警;无价则仍用自报下界(好过记 0); + // 3. 无自报 → 本地复算;彻底无价才整单记 0。 + let (input_cost, output_cost, cache_read_cost, cache_creation_cost, total_cost) = match pricing + { + Some(p) => { + let cost = CostCalculator::calculate_for_app("grokbuild", &usage, &p, multiplier); + let total = match reported { + Some(reported) if !cost_is_partial => { + // 偏差超 1%(微额下限 1e-6)即本地定价漂移——xAI 调价时 + // 最早的可观测信号,提醒更新 seed/repair。 + let tolerance = (reported * Decimal::new(1, 2)).max(Decimal::new(1, 6)); + if (cost.total_cost - reported).abs() > tolerance { + deferred_warn = Some(format!( + "本地定价与 CLI 自报成本偏差超阈值,total 已以自报为准,请更新本地定价: model={model} local={} reported={reported} request_id={request_id}", + cost.total_cost + )); + } + reported + } + _ => cost.total_cost, + }; + ( + cost.input_cost.to_string(), + cost.output_cost.to_string(), + cost.cache_read_cost.to_string(), + cost.cache_creation_cost.to_string(), + total.to_string(), + ) + } + None => { + // 未 seed 的新别名:token 照常入账;有自报成本时直接采用(分项 + // 记 0),彻底无价才整单记 0。xAI 内部别名会周期性变动 + // (grok-4.5-build 即先例),两种情况都要留下可排查的痕迹。 + let total = match reported { + Some(reported) => { + if model != "unknown" { + let partial_note = if cost_is_partial { + "(上游标记为部分费用,实际为下界)" + } else { + "" + }; + deferred_warn = Some(format!( + "模型定价未找到,采用 CLI 自报成本入账{partial_note}: model={model} total={reported} request_id={request_id}" + )); + } + reported.to_string() + } + None => { + if model != "unknown" { + deferred_warn = Some(format!( + "模型定价未找到且无自报成本,成本记 0: model={model} request_id={request_id}" + )); + } + "0".to_string() + } + }; + ( + "0".to_string(), + "0".to_string(), + "0".to_string(), + "0".to_string(), + total, + ) + } + }; + + // UPSERT:重扫幂等;解析口径修正后重扫时更新既有行(token/成本/ + // latency;created_at 保持首插值不动,避免行在沉降窗与 rollup 边界间漂移)。 + // WHERE 的 data_source 守卫是纵深防御:request_id 前缀命名空间已隔离, + // 万一撞上非本导入器的行也绝不改写它。 + // input_token_semantics 显式写 TOTAL——xAI 口径 inputTokens 含 cache read, + // 与代理路径的 grokbuild 行(logger)保持同一语义,勿依赖列默认值。 + conn.execute( + "INSERT INTO proxy_request_logs ( + request_id, provider_id, app_type, model, request_model, + input_tokens, output_tokens, cache_read_tokens, cache_creation_tokens, + input_cost_usd, output_cost_usd, cache_read_cost_usd, cache_creation_cost_usd, total_cost_usd, + latency_ms, first_token_ms, status_code, error_message, session_id, + provider_type, is_streaming, cost_multiplier, created_at, data_source, + input_token_semantics + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19, ?20, ?21, ?22, ?23, ?24, ?25) + ON CONFLICT(request_id) DO UPDATE SET + model = excluded.model, + input_tokens = excluded.input_tokens, + output_tokens = excluded.output_tokens, + cache_read_tokens = excluded.cache_read_tokens, + input_cost_usd = excluded.input_cost_usd, + output_cost_usd = excluded.output_cost_usd, + cache_read_cost_usd = excluded.cache_read_cost_usd, + cache_creation_cost_usd = excluded.cache_creation_cost_usd, + total_cost_usd = excluded.total_cost_usd, + latency_ms = excluded.latency_ms + WHERE data_source = 'grok_session' + AND (input_tokens != excluded.input_tokens + OR output_tokens != excluded.output_tokens + OR cache_read_tokens != excluded.cache_read_tokens + OR latency_ms != excluded.latency_ms + OR model != excluded.model)", + rusqlite::params![ + request_id, + "_grok_session", // provider_id + "grokbuild", // app_type + model, + model, // request_model = model + usage.input_tokens, + usage.output_tokens, + usage.cache_read_tokens, + 0i64, // cache_creation_tokens + input_cost, + output_cost, + cache_read_cost, + cache_creation_cost, + total_cost, + turn.api_ms.min(i64::MAX as u64) as i64, // latency_ms(本轮 API 时长) + Option::::None, // first_token_ms + 200i64, // status_code + Option::::None, // error_message + session_id, + Some("grok_session"), // provider_type + 1i64, // is_streaming + "1.0", // cost_multiplier + created_at, + "grok_session", // data_source + INPUT_TOKEN_SEMANTICS_TOTAL, + ], + ) + .map_err(|e| AppError::Database(format!("插入 Grok Build 会话日志失败: {e}")))?; + + // changes() > 0 表示新插入或已更新,== 0 表示值完全相同(无实际变更) + let changed = conn.changes() > 0; + if changed { + if let Some(msg) = deferred_warn { + log::warn!("[GROK-SYNC] {msg}"); + } + } + Ok(changed) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::io::Write; + use tempfile::tempdir; + + /// 早于沉降窗的固定基准时刻(2023-11-14T22:13:20Z) + const OLD_EPOCH: i64 = 1_700_000_000; + + fn epoch_to_rfc3339(epoch: i64) -> String { + chrono::DateTime::from_timestamp(epoch, 0) + .expect("valid epoch") + .to_rfc3339() + } + + /// 顶层 timestamp 用真实的数字 epoch 秒格式(RFC3339 兜底见 parses 测试) + fn usage_event_line(epoch: i64, prompt_id: &str, model_usage: &str) -> String { + format!( + r#"{{"timestamp":{epoch},"method":"_x.ai/session/update","params":{{"update":{{"sessionUpdate":"turn_completed","prompt_id":"{prompt_id}","stop_reason":"end_turn","usage":{{"modelUsage":{{{model_usage}}}}}}}}}}}"# + ) + } + + /// 带事件级 costIsPartial 标记的变体 + fn usage_event_line_partial(epoch: i64, prompt_id: &str, model_usage: &str) -> String { + format!( + r#"{{"timestamp":{epoch},"method":"_x.ai/session/update","params":{{"update":{{"sessionUpdate":"turn_completed","prompt_id":"{prompt_id}","stop_reason":"end_turn","usage":{{"costIsPartial":true,"modelUsage":{{{model_usage}}}}}}}}}}}"# + ) + } + + fn model_counters(model: &str, input: u64, output: u64, cached: u64, calls: u64) -> String { + model_counters_with_ticks(model, input, output, cached, calls, 0) + } + + fn model_counters_with_ticks( + model: &str, + input: u64, + output: u64, + cached: u64, + calls: u64, + ticks: u64, + ) -> String { + format!( + r#""{model}":{{"inputTokens":{input},"outputTokens":{output},"cachedReadTokens":{cached},"reasoningTokens":0,"modelCalls":{calls},"apiDurationMs":1000,"costUsdTicks":{ticks}}}"# + ) + } + + fn write_session_file(dir: &Path, session_id: &str, lines: &[String]) -> PathBuf { + let session_dir = dir.join("sessions").join("enc-project").join(session_id); + std::fs::create_dir_all(&session_dir).expect("create session dir"); + let path = session_dir.join("updates.jsonl"); + let mut file = std::fs::File::create(&path).expect("create updates.jsonl"); + for line in lines { + writeln!(file, "{line}").expect("write line"); + } + path + } + + /// (request_id, input, output, cache_read, input_token_semantics) + type GrokSessionRow = (String, u32, u32, u32, i64); + + fn query_rows(db: &Database) -> Result, AppError> { + let conn = lock_conn!(db.conn); + let mut stmt = conn + .prepare( + "SELECT request_id, input_tokens, output_tokens, cache_read_tokens, input_token_semantics + FROM proxy_request_logs WHERE data_source = 'grok_session' ORDER BY request_id", + ) + .expect("prepare"); + let rows = stmt + .query_map([], |row| { + Ok(( + row.get(0)?, + row.get(1)?, + row.get(2)?, + row.get(3)?, + row.get(4)?, + )) + }) + .expect("query") + .filter_map(Result::ok) + .collect(); + Ok(rows) + } + + fn query_costs(db: &Database) -> Result, AppError> { + let conn = lock_conn!(db.conn); + let mut stmt = conn + .prepare( + "SELECT request_id, total_cost_usd FROM proxy_request_logs + WHERE data_source = 'grok_session' ORDER BY created_at, request_id", + ) + .expect("prepare"); + let rows = stmt + .query_map([], |row| Ok((row.get(0)?, row.get(1)?))) + .expect("query") + .filter_map(Result::ok) + .collect(); + Ok(rows) + } + + #[test] + fn parses_turn_completed_and_ignores_noise_and_other_kinds() { + let content = concat!( + "{\"timestamp\":\"2026-07-20T13:26:10Z\",\"method\":\"session/update\",\"params\":{\"update\":{\"sessionUpdate\":\"agent_message_chunk\",\"content\":{}}}}\n", + "not json at all\n", + // 显式标为非 turn_completed 却带 usage:防中途快照双算,不得导入 + "{\"timestamp\":\"2026-07-20T13:26:20Z\",\"method\":\"_x.ai/session/update\",\"params\":{\"update\":{\"sessionUpdate\":\"usage_snapshot\",\"prompt_id\":\"px\",\"usage\":{\"inputTokens\":9999,\"outputTokens\":9,\"cachedReadTokens\":0}}}}\n", + "{\"timestamp\":\"2026-07-20T13:26:24Z\",\"method\":\"_x.ai/session/update\",\"params\":{\"update\":{\"sessionUpdate\":\"turn_completed\",\"prompt_id\":\"p1\",\"usage\":{\"inputTokens\":16632,\"outputTokens\":104,\"cachedReadTokens\":0,\"modelUsage\":{\"grok-4.5-build\":{\"inputTokens\":16632,\"outputTokens\":104,\"cachedReadTokens\":0,\"apiDurationMs\":5342,\"costUsdTicks\":338880000}}}}}}\n", + ); + let events = parse_grok_usage_events(content); + assert_eq!(events.len(), 1); + assert_eq!(events[0].prompt_id, "p1"); + assert_eq!(events[0].per_model.len(), 1); + assert_eq!(events[0].per_model[0].0, "grok-4.5-build"); + assert_eq!( + events[0].per_model[0].1, + GrokCounters { + input: 16632, + output: 104, + cached: 0, + api_ms: 5342, + model_calls: 0, + cost_ticks: 338_880_000, + cost_partial: false, + } + ); + } + + #[test] + fn missing_model_usage_falls_back_to_top_level_counters() { + // 同时覆盖:sessionUpdate 字段缺失时向后兼容放行 + let line = format!( + r#"{{"timestamp":"{}","method":"_x.ai/session/update","params":{{"update":{{"prompt_id":"p1","usage":{{"inputTokens":100,"outputTokens":10,"cachedReadTokens":5}}}}}}}}"#, + epoch_to_rfc3339(OLD_EPOCH) + ); + let events = parse_grok_usage_events(&line); + assert_eq!(events.len(), 1); + assert_eq!(events[0].per_model[0].0, "unknown"); + assert_eq!(events[0].per_model[0].1.input, 100); + } + + #[test] + fn two_turns_import_at_face_value_matching_reported_ticks() -> Result<(), AppError> { + use std::str::FromStr; + // 2026-07-23 单进程双 prompt 实测原值:turn_completed 是逐轮独立总量。 + // 若误用相邻差分,第二轮会被记成 53/28/6144 的假增量(曾犯)。 + // 每轮 ticks 同时钉死逐轮口径与 2/6/0.30 定价: + // 轮1 (17294-11136)×2 + 11136×0.30 + 28×6 = 15824.8 µUSD = 158248000 ticks + // 轮2 (17347-17280)×2 + 17280×0.30 + 56×6 = 5654.0 µUSD = 56540000 ticks + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![ + usage_event_line( + OLD_EPOCH, + "p1", + &model_counters_with_ticks("grok-4.5-build", 17294, 28, 11136, 1, 158_248_000), + ), + usage_event_line( + OLD_EPOCH + 60, + "p2", + &model_counters_with_ticks("grok-4.5-build", 17347, 56, 17280, 1, 56_540_000), + ), + ]; + let path = write_session_file(temp.path(), "sess-two-turns", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 2); + assert_eq!(result.deferred_files, 0); + + let rows = query_rows(&db)?; + assert_eq!(rows.len(), 2); + assert_eq!((rows[0].1, rows[0].2, rows[0].3), (17294, 28, 11136)); + assert_eq!((rows[1].1, rows[1].2, rows[1].3), (17347, 56, 17280)); + // 语义列显式为 TOTAL,与代理路径一致 + assert!(rows.iter().all(|r| r.4 == INPUT_TOKEN_SEMANTICS_TOTAL)); + + // 本地定价复算须与 CLI 自报 ticks 分毫不差(漂移告警在此阈值内静默) + let costs = query_costs(&db)?; + let expected1 = Decimal::from(158_248_000u64) / Decimal::from(10_000_000_000u64); + let expected2 = Decimal::from(56_540_000u64) / Decimal::from(10_000_000_000u64); + assert_eq!(Decimal::from_str(&costs[0].1).expect("decimal"), expected1); + assert_eq!(Decimal::from_str(&costs[1].1).expect("decimal"), expected2); + Ok(()) + } + + #[test] + fn second_turn_with_smaller_counters_imports_at_face_value() -> Result<(), AppError> { + // 2026-07-23 跨进程实测原值(进程 A 单轮 27386/74/15360,--resume 的 + // 进程 B 单轮 13793/21/13696)。逐轮口径下"第二轮更小"是常态, + // 与是否跨进程无关,一律按面值入账。 + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![ + usage_event_line( + OLD_EPOCH, + "p1", + &model_counters("grok-4.5-build", 27386, 74, 15360, 2), + ), + usage_event_line( + OLD_EPOCH + 15, + "p2", + &model_counters("grok-4.5-build", 13793, 21, 13696, 1), + ), + ]; + let path = write_session_file(temp.path(), "sess-resume", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 2); + + let rows = query_rows(&db)?; + assert_eq!(rows.len(), 2); + assert_eq!((rows[0].1, rows[0].2, rows[0].3), (27386, 74, 15360)); + assert_eq!((rows[1].1, rows[1].2, rows[1].3), (13793, 21, 13696)); + Ok(()) + } + + #[test] + fn identical_turns_both_import() -> Result<(), AppError> { + // 回归(逐轮口径):两轮数值完全相同 = 两笔真实用量,都必须入账。 + // 差分口径会把第二轮当零增量整轮跳过——那正是被证伪的旧行为。 + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![ + usage_event_line( + OLD_EPOCH, + "p1", + &model_counters("grok-4.5-build", 100, 10, 0, 1), + ), + usage_event_line( + OLD_EPOCH + 60, + "p2", + &model_counters("grok-4.5-build", 100, 10, 0, 1), + ), + ]; + let path = write_session_file(temp.path(), "sess-identical", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 2, "相同数值的两轮都是真实用量"); + assert_eq!(query_rows(&db)?.len(), 2); + Ok(()) + } + + #[test] + fn multi_model_event_produces_row_per_model() -> Result<(), AppError> { + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let both = format!( + "{},{}", + model_counters("grok-4.5-build", 100, 10, 0, 1), + model_counters("grok-4.3", 30, 3, 0, 1) + ); + let lines = vec![usage_event_line(OLD_EPOCH, "p1", &both)]; + let path = write_session_file(temp.path(), "sess-multi", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 2); + let rows = query_rows(&db)?; + assert!(rows[0].0.ends_with(":grok-4.3")); + assert!(rows[1].0.ends_with(":grok-4.5-build")); + Ok(()) + } + + #[test] + fn settle_window_defers_recent_events_without_recording_sync_state() -> Result<(), AppError> { + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let now = SystemTime::now() + .duration_since(SystemTime::UNIX_EPOCH) + .expect("now") + .as_secs() as i64; + let lines = vec![ + usage_event_line( + OLD_EPOCH, + "p1", + &model_counters("grok-4.5-build", 100, 10, 0, 1), + ), + // 未沉降的新事件:本轮延后,且不落同步状态以便下一轮重读 + usage_event_line(now, "p2", &model_counters("grok-4.5-build", 250, 30, 0, 1)), + ]; + let path = write_session_file(temp.path(), "sess-settle", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 1); + assert_eq!(result.deferred_files, 1); + assert_eq!(query_rows(&db)?.len(), 1); + + let (last_modified, _) = get_sync_state(&db, &path.to_string_lossy())?; + assert_eq!(last_modified, 0, "延后时不得记录同步状态"); + + // 下一轮重读:旧事件 UPSERT 无变化,新事件仍未沉降继续延后 + let rerun = sync_single_grok_file(&db, &path)?; + assert_eq!(rerun.imported, 0); + assert_eq!(rerun.skipped, 1); + assert_eq!(rerun.deferred_files, 1); + assert_eq!(query_rows(&db)?.len(), 1); + Ok(()) + } + + #[test] + fn takeover_guard_skips_events_near_proxy_activity() -> Result<(), AppError> { + let db = Database::memory()?; + { + let conn = lock_conn!(db.conn); + conn.execute( + "INSERT INTO proxy_request_logs ( + request_id, provider_id, app_type, model, request_model, + input_tokens, output_tokens, cache_read_tokens, cache_creation_tokens, + total_cost_usd, latency_ms, status_code, created_at, data_source + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + rusqlite::params![ + "grok-proxy-req", + "some-provider", + "grokbuild", + "grok-4.5", + "grok-4.5", + 999, + 88, + 0, + 0, + "0.01", + 100, + 200, + OLD_EPOCH + 30, + "proxy" + ], + )?; + } + let temp = tempdir().expect("tempdir"); + let lines = vec![ + // 事件时刻落在代理行 ±窗口内 → 接管态,跳过(代理行权威) + usage_event_line( + OLD_EPOCH, + "p1", + &model_counters("grok-4.5-build", 100, 10, 0, 1), + ), + // 远离接管窗口的后续事件按面值正常导入 + usage_event_line( + OLD_EPOCH + SESSION_PROXY_DEDUP_WINDOW_SECONDS + 3600, + "p2", + &model_counters("grok-4.5-build", 250, 30, 0, 1), + ), + ]; + let path = write_session_file(temp.path(), "sess-guard", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.skipped, 1, "守卫跳过计入 skipped(未入账)"); + assert_eq!(result.imported, 1); + + let rows = query_rows(&db)?; + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].1, 250, "守卫外事件按本轮面值入账"); + Ok(()) + } + + #[test] + fn rescan_is_idempotent() -> Result<(), AppError> { + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![ + usage_event_line( + OLD_EPOCH, + "p1", + &model_counters("grok-4.5-build", 100, 10, 0, 1), + ), + usage_event_line( + OLD_EPOCH + 60, + "p2", + &model_counters("grok-4.5-build", 250, 30, 50, 1), + ), + ]; + let path = write_session_file(temp.path(), "sess-idem", &lines); + + let first = sync_single_grok_file(&db, &path)?; + assert_eq!(first.imported, 2); + + // mtime 未变 → 短路 + let second = sync_single_grok_file(&db, &path)?; + assert_eq!(second.imported + second.skipped, 0); + + // 强制重读(清同步状态)→ UPSERT 全部无变化 + { + let conn = lock_conn!(db.conn); + conn.execute("DELETE FROM session_log_sync", [])?; + } + let third = sync_single_grok_file(&db, &path)?; + assert_eq!(third.imported, 0); + assert_eq!(third.skipped, 2); + assert_eq!(query_rows(&db)?.len(), 2); + Ok(()) + } + + #[test] + fn rewind_truncation_does_not_double_count() -> Result<(), AppError> { + // 回归(对比评审发现):幂等键若含文件内序号,updates.jsonl 前缀被 + // 改写(rewind 截断)后幸存事件序号前移会生成新 request_id 造成双算。 + // prompt_id 锚定键下:幸存轮命中原行;被移除轮的行保留(rewind 不 + // 退还已消耗 token,留存即正确)。 + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let full = vec![ + usage_event_line( + OLD_EPOCH, + "p1", + &model_counters("grok-4.5-build", 100, 10, 0, 1), + ), + usage_event_line( + OLD_EPOCH + 60, + "p2", + &model_counters("grok-4.5-build", 200, 20, 0, 1), + ), + usage_event_line( + OLD_EPOCH + 120, + "p3", + &model_counters("grok-4.5-build", 300, 30, 0, 1), + ), + ]; + let path = write_session_file(temp.path(), "sess-rewind", &full); + assert_eq!(sync_single_grok_file(&db, &path)?.imported, 3); + + // 模拟 rewind 截掉 p2:p3 从 idx2 前移到 idx1 + let truncated = vec![full[0].clone(), full[2].clone()]; + write_session_file(temp.path(), "sess-rewind", &truncated); + { + let conn = lock_conn!(db.conn); + conn.execute("DELETE FROM session_log_sync", [])?; + } + + let rescan = sync_single_grok_file(&db, &path)?; + assert_eq!(rescan.imported, 0, "幸存轮不得因序号前移重新入账"); + + let rows = query_rows(&db)?; + assert_eq!(rows.len(), 3, "被截掉轮次的行保留(token 已实际消耗)"); + let p3: Vec<_> = rows.iter().filter(|r| r.0.contains(":p3:")).collect(); + assert_eq!(p3.len(), 1); + assert_eq!(p3[0].1, 300); + Ok(()) + } + + #[test] + fn empty_prompt_id_falls_back_to_index_key() -> Result<(), AppError> { + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![usage_event_line( + OLD_EPOCH, + "", + &model_counters("grok-4.5-build", 100, 10, 0, 1), + )]; + let path = write_session_file(temp.path(), "sess-noprompt", &lines); + + assert_eq!(sync_single_grok_file(&db, &path)?.imported, 1); + let rows = query_rows(&db)?; + assert!(rows[0].0.contains(":idx0:"), "空 prompt_id 回退序号键"); + Ok(()) + } + + #[test] + fn cost_matches_cli_reported_ticks_for_seeded_grok45_build() -> Result<(), AppError> { + use std::str::FromStr; + // 真实样本:inputTokens=16632, outputTokens=104, cache=0, + // costUsdTicks=338880000(1 tick = 1e-10 USD)。seed 的 grok-4.5-build + // 定价(2/6)应精确复现 CLI 自报成本。fixture 故意不带 ticks, + // 验证的是本地定价独立复算。 + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![usage_event_line( + OLD_EPOCH, + "p1", + &model_counters("grok-4.5-build", 16632, 104, 0, 1), + )]; + let path = write_session_file(temp.path(), "sess-ticks", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 1); + + let conn = lock_conn!(db.conn); + let total: String = conn.query_row( + "SELECT total_cost_usd FROM proxy_request_logs WHERE data_source = 'grok_session'", + [], + |row| row.get(0), + )?; + let expected = Decimal::from(338_880_000u64) / Decimal::from(10_000_000_000u64); + assert_eq!(Decimal::from_str(&total).expect("decimal"), expected); + Ok(()) + } + + #[test] + fn cost_matches_cli_reported_ticks_with_cache_reads() -> Result<(), AppError> { + use std::str::FromStr; + // 2026-07-23 实测带缓存样本:13793/21/13696,costUsdTicks=44288000。 + // 钉死 cache read 实测单价 0.30:billable_input=(13793-13696)×2/1M + // + 21×6/1M + 13696×0.30/1M = 0.0044288。seed 若改回 0.50 此测试即红。 + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![usage_event_line( + OLD_EPOCH, + "p1", + &model_counters("grok-4.5-build", 13793, 21, 13696, 1), + )]; + let path = write_session_file(temp.path(), "sess-ticks-cache", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 1); + + let conn = lock_conn!(db.conn); + let total: String = conn.query_row( + "SELECT total_cost_usd FROM proxy_request_logs WHERE data_source = 'grok_session'", + [], + |row| row.get(0), + )?; + let expected = Decimal::from(44_288_000u64) / Decimal::from(10_000_000_000u64); + assert_eq!(Decimal::from_str(&total).expect("decimal"), expected); + Ok(()) + } + + #[test] + fn reported_ticks_override_stale_local_pricing() -> Result<(), AppError> { + use std::str::FromStr; + // 定价漂移窗口:CLI 自报为本地复算(338880000 ticks)的两倍,模拟 + // xAI 调价而 seed 未更新。total 必须以自报为准(回填不修正正值行, + // 本地价错就永久错);分项仍按本地价(暂不自洽,有漂移告警提示)。 + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![usage_event_line( + OLD_EPOCH, + "p1", + &model_counters_with_ticks("grok-4.5-build", 16632, 104, 0, 1, 677_760_000), + )]; + let path = write_session_file(temp.path(), "sess-drift", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 1); + + let conn = lock_conn!(db.conn); + let (input_cost, total): (String, String) = conn.query_row( + "SELECT input_cost_usd, total_cost_usd FROM proxy_request_logs + WHERE data_source = 'grok_session'", + [], + |row| Ok((row.get(0)?, row.get(1)?)), + )?; + let expected_total = Decimal::from(677_760_000u64) / Decimal::from(10_000_000_000u64); + assert_eq!( + Decimal::from_str(&total).expect("decimal"), + expected_total, + "total 以自报为准" + ); + assert!( + Decimal::from_str(&input_cost).expect("decimal") > Decimal::ZERO, + "分项仍按本地定价" + ); + Ok(()) + } + + #[test] + fn partial_reported_cost_prefers_local_pricing_when_priced() -> Result<(), AppError> { + use std::str::FromStr; + // costIsPartial=true:自报只是下界,不可作 total。token 数是完整的, + // 有本地价时用本地全额复算(此处应得 338880000 ticks 等值)。 + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![usage_event_line_partial( + OLD_EPOCH, + "p1", + &model_counters_with_ticks("grok-4.5-build", 16632, 104, 0, 1, 1_000), + )]; + let path = write_session_file(temp.path(), "sess-partial", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 1); + + let conn = lock_conn!(db.conn); + let total: String = conn.query_row( + "SELECT total_cost_usd FROM proxy_request_logs WHERE data_source = 'grok_session'", + [], + |row| row.get(0), + )?; + let expected = Decimal::from(338_880_000u64) / Decimal::from(10_000_000_000u64); + assert_eq!(Decimal::from_str(&total).expect("decimal"), expected); + Ok(()) + } + + #[test] + fn unpriced_model_falls_back_to_reported_ticks() -> Result<(), AppError> { + use std::str::FromStr; + // 未 seed 的新别名:total_cost 采用 CLI 自报 ticks(分项记 0), + // 不再整单记 0。 + let db = Database::memory()?; + let temp = tempdir().expect("tempdir"); + let lines = vec![usage_event_line( + OLD_EPOCH, + "p1", + &model_counters_with_ticks("grok-6-future-alias", 1000, 100, 0, 1, 56_540_000), + )]; + let path = write_session_file(temp.path(), "sess-unpriced", &lines); + + let result = sync_single_grok_file(&db, &path)?; + assert_eq!(result.imported, 1); + + let conn = lock_conn!(db.conn); + let (input_cost, total): (String, String) = conn.query_row( + "SELECT input_cost_usd, total_cost_usd FROM proxy_request_logs + WHERE data_source = 'grok_session'", + [], + |row| Ok((row.get(0)?, row.get(1)?)), + )?; + assert_eq!( + Decimal::from_str(&input_cost).expect("decimal"), + Decimal::ZERO + ); + let expected = Decimal::from(56_540_000u64) / Decimal::from(10_000_000_000u64); + assert_eq!(Decimal::from_str(&total).expect("decimal"), expected); + Ok(()) + } +} diff --git a/src-tauri/src/services/usage_stats.rs b/src-tauri/src/services/usage_stats.rs index 26d6f670b..117b221a6 100644 --- a/src-tauri/src/services/usage_stats.rs +++ b/src-tauri/src/services/usage_stats.rs @@ -214,6 +214,7 @@ fn provider_name_coalesce(log_alias: &str, provider_alias: &str) -> String { WHEN '_codex_session' THEN 'Codex (Session)' \ WHEN '_gemini_session' THEN 'Gemini (Session)' \ WHEN '_opencode_session' THEN 'OpenCode (Session)' \ + WHEN '_grok_session' THEN 'Grok Build (Session)' \ ELSE {log_alias}.provider_id END)" ) } @@ -416,6 +417,40 @@ pub(crate) fn has_matching_proxy_usage_log( .map_err(|e| AppError::Database(format!("查询重复代理用量日志失败: {e}"))) } +/// grokbuild 会话导入的接管活动守卫:给定时刻 ±窗口内存在任何 grokbuild +/// 代理直录行,即认为当时处于代理接管态,会话事件应整体跳过——同一请求 +/// 已由代理逐请求记账,会话侧再入账必双算。 +/// +/// 不复用 [`has_matching_proxy_usage_log`] 的指纹匹配:Grok 会话事件是 +/// 逐轮聚合值,与代理逐请求行的 token 值结构性不相等,指纹永不命中。 +/// 这里按"接管态检测"而非"行匹配"设计,故不过滤 status_code——失败的 +/// 代理请求同样证明流量正走代理。 +/// +/// 已知局限(有意取舍,方向保守只漏不双):窗口不含 session 维度,任一 +/// grokbuild 代理行会给 ±窗口内的全部会话事件投下阴影——接管/官方两态在 +/// 十分钟内交替或并行使用时,官方侧轮次会被跳过(漏记而非双算)。 +pub(crate) fn has_recent_grokbuild_proxy_activity( + conn: &Connection, + created_at: i64, +) -> Result { + let l_data_source = data_source_expr("l"); + let sql = format!( + "SELECT EXISTS ( + SELECT 1 + FROM proxy_request_logs l + WHERE {l_data_source} = 'proxy' + AND l.app_type = 'grokbuild' + AND l.created_at BETWEEN ?1 - ?2 AND ?1 + ?2 + )" + ); + conn.query_row( + &sql, + params![created_at, SESSION_PROXY_DEDUP_WINDOW_SECONDS], + |row| row.get::<_, bool>(0), + ) + .map_err(|e| AppError::Database(format!("查询 Grok 接管活动失败: {e}"))) +} + pub(crate) fn has_suspected_codex_session_duplicate( conn: &Connection, request_id: &str,