diff --git a/src-tauri/src/services/session_usage_codex.rs b/src-tauri/src/services/session_usage_codex.rs index f23be19c3..5757f9ded 100644 --- a/src-tauri/src/services/session_usage_codex.rs +++ b/src-tauri/src/services/session_usage_codex.rs @@ -9,7 +9,7 @@ //! ``` //! //! ## 解析的事件类型 -//! - `session_meta` → 提取 session_id +//! - `session_meta` → 提取唯一 thread_id(子代理的 session_id 指向父线程) //! - `turn_context` → 提取当前 model //! - `event_msg` (type=token_count) → 提取累计 token 用量,计算 delta @@ -28,6 +28,8 @@ use std::io::{BufRead, BufReader}; use std::path::{Path, PathBuf}; use std::time::SystemTime; +const CODEX_THREAD_REQUEST_ID_PREFIX: &str = "codex_session:thread-v1"; + /// 累计 token 用量(跟踪 total_token_usage 字段) #[derive(Debug, Clone, Default)] struct CumulativeTokens { @@ -52,10 +54,162 @@ impl DeltaTokens { /// 单文件解析时的运行状态 struct FileParseState { - session_id: Option, + thread_id: Option, current_model: String, prev_total: Option, event_index: u32, + history_replay_boundary: Option, +} + +/// Codex 子代理日志中的 `id` 是当前线程的唯一 ID,`session_id` 则指向父线程。 +#[derive(Debug, Clone, PartialEq, Eq)] +struct CodexSessionIdentity { + thread_id: String, + carries_history_snapshot: bool, +} + +fn parse_codex_session_identity(payload: &serde_json::Value) -> Option { + let thread_id = payload + .get("id") + .or_else(|| payload.get("thread_id")) + .or_else(|| payload.get("threadId")) + .or_else(|| payload.get("session_id")) + .or_else(|| payload.get("sessionId")) + .and_then(|value| value.as_str())? + .to_string(); + let session_id = payload + .get("session_id") + .or_else(|| payload.get("sessionId")) + .and_then(|value| value.as_str()); + let carries_history_snapshot = payload + .get("forked_from_id") + .and_then(|value| value.as_str()) + .is_some_and(|value| !value.is_empty()) + || payload + .get("source") + .and_then(|source| source.get("subagent")) + .is_some() + || session_id.is_some_and(|session_id| session_id != thread_id); + + Some(CodexSessionIdentity { + thread_id, + carries_history_snapshot, + }) +} + +fn read_codex_session_identity(file_path: &Path) -> Option { + let file = fs::File::open(file_path).ok()?; + + for line in BufReader::new(file).lines() { + let Ok(line) = line else { + continue; + }; + if !line.contains("\"session_meta\"") { + continue; + } + let Ok(value) = serde_json::from_str::(&line) else { + continue; + }; + if value.get("type").and_then(|value| value.as_str()) != Some("session_meta") { + continue; + } + if let Some(identity) = value.get("payload").and_then(parse_codex_session_identity) { + return Some(identity); + } + } + + None +} + +/// fork/子代理日志会先重放父线程历史,再以接管事件开始当前线程。 +/// 返回接管事件所在行;此前的 token_count 只用于恢复累计值基线。 +fn codex_history_replay_boundary( + file_path: &Path, + identity: Option<&CodexSessionIdentity>, +) -> Option { + if !identity.is_some_and(|identity| identity.carries_history_snapshot) { + return None; + } + + let file = fs::File::open(file_path).ok()?; + for (index, line) in BufReader::new(file).lines().enumerate() { + let Ok(line) = line else { + continue; + }; + if !line.contains("\"thread_settings_applied\"") + && !line.contains("\"inter_agent_communication") + { + continue; + } + let Ok(value) = serde_json::from_str::(&line) else { + continue; + }; + let Some(event_type) = value.get("type").and_then(|value| value.as_str()) else { + continue; + }; + let is_replay_boundary = event_type.starts_with("inter_agent_communication") + || (event_type == "event_msg" + && value + .get("payload") + .and_then(|payload| payload.get("type")) + .and_then(|value| value.as_str()) + == Some("thread_settings_applied")); + if is_replay_boundary { + return Some(index as i64 + 1); + } + } + + None +} + +fn is_history_snapshot_event(state: &FileParseState, line_offset: i64) -> bool { + state + .history_replay_boundary + .is_some_and(|boundary| line_offset < boundary) +} + +fn get_codex_sync_state(db: &Database, file_path: &Path) -> Result<(i64, i64), AppError> { + let file_path_str = file_path.to_string_lossy().to_string(); + let state = get_sync_state(db, &file_path_str)?; + if state != (0, 0) + || file_path + .parent() + .and_then(Path::file_name) + .and_then(|name| name.to_str()) + != Some("archived_sessions") + { + return Ok(state); + } + + let Some(file_name) = file_path.file_name().and_then(|name| name.to_str()) else { + return Ok(state); + }; + let slash_suffix = format!("/{file_name}"); + let backslash_suffix = format!("\\{file_name}"); + let conn = lock_conn!(db.conn); + let inherited = conn.query_row( + "SELECT last_modified, last_line_offset + FROM session_log_sync + WHERE file_path <> ?1 + AND (substr(file_path, -length(?2)) = ?2 + OR substr(file_path, -length(?3)) = ?3) + ORDER BY last_line_offset DESC, last_modified DESC + LIMIT 1", + rusqlite::params![file_path_str, slash_suffix, backslash_suffix], + |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)), + ); + drop(conn); + + match inherited { + Ok(inherited) => { + update_sync_state(db, &file_path_str, inherited.0, inherited.1)?; + Ok(inherited) + } + Err(rusqlite::Error::QueryReturnedNoRows) => Ok(state), + Err(error) => Err(AppError::Database(format!( + "查询 Codex 归档文件同步状态失败: {error}" + ))), + } } /// 归一化 Codex 模型名 @@ -238,23 +392,27 @@ fn sync_single_codex_file(db: &Database, file_path: &Path) -> Result<(u32, u32), let file_modified = metadata_modified_nanos(&metadata); // 检查同步状态 - let (last_modified, last_offset) = get_sync_state(db, &file_path_str)?; + let (last_modified, last_offset) = get_codex_sync_state(db, file_path)?; // 文件未变化则跳过 if file_modified <= last_modified { return Ok((0, 0)); } + let identity = read_codex_session_identity(file_path); + // 打开文件逐行解析 let file = fs::File::open(file_path).map_err(|e| AppError::Config(format!("无法打开文件: {e}")))?; let reader = BufReader::new(file); + let history_replay_boundary = codex_history_replay_boundary(file_path, identity.as_ref()); let mut state = FileParseState { - session_id: None, + thread_id: identity.map(|identity| identity.thread_id), current_model: "unknown".to_string(), prev_total: None, event_index: 0, + history_replay_boundary, }; let mut line_offset: i64 = 0; @@ -296,16 +454,11 @@ fn sync_single_codex_file(db: &Database, file_path: &Path) -> Result<(u32, u32), }; match event_type { - "session_meta" if state.session_id.is_none() => { - let payload = value.get("payload"); - state.session_id = payload - .and_then(|p| { - p.get("session_id") - .or_else(|| p.get("sessionId")) - .or_else(|| p.get("id")) - }) - .and_then(|v| v.as_str()) - .map(|s| s.to_string()); + "session_meta" if state.thread_id.is_none() => { + state.thread_id = value + .get("payload") + .and_then(parse_codex_session_identity) + .map(|identity| identity.thread_id); } "turn_context" => { if let Some(payload) = value.get("payload") { @@ -383,16 +536,28 @@ fn sync_single_codex_file(db: &Database, file_path: &Path) -> Result<(u32, u32), continue; // 跳过 task 边界的零 delta 事件 } + // 所有非零事件都占据稳定序号,包括已同步事件与 replay 快照。 state.event_index += 1; + // replay 快照更新了 prev_total,但不是当前线程的新用量。 + if is_history_snapshot_event(&state, line_offset) { + if line_offset > last_offset { + skipped += 1; + } + continue; + } + // 跳过已处理的行(但仍需解析以恢复状态) if line_offset <= last_offset { continue; } // 生成唯一 request_id - let session_id_str = state.session_id.as_deref().unwrap_or("unknown"); - let request_id = format!("codex_session:{}:{}", session_id_str, state.event_index); + let thread_id = state.thread_id.as_deref().unwrap_or("unknown"); + let request_id = format!( + "{CODEX_THREAD_REQUEST_ID_PREFIX}:{thread_id}:{}", + state.event_index + ); // 提取时间戳 let timestamp = value @@ -405,7 +570,7 @@ fn sync_single_codex_file(db: &Database, file_path: &Path) -> Result<(u32, u32), &request_id, &delta, &state.current_model, - state.session_id.as_deref(), + state.thread_id.as_deref(), timestamp.as_deref(), ) { Ok(true) => imported += 1, @@ -549,6 +714,56 @@ fn find_codex_pricing(conn: &rusqlite::Connection, model_id: &str) -> Option>() + .join("\n") + + "\n"; + fs::write(path, contents).unwrap(); + } + + fn session_meta(thread_id: &str, session_id: &str) -> serde_json::Value { + serde_json::json!({ + "timestamp": "2026-07-10T03:00:00Z", + "type": "session_meta", + "payload": { + "id": thread_id, + "session_id": session_id, + "source": if thread_id == session_id { + serde_json::Value::String("cli".to_string()) + } else { + serde_json::json!({ "subagent": {} }) + } + } + }) + } + + fn turn_context() -> serde_json::Value { + serde_json::json!({ + "timestamp": "2026-07-10T03:00:01Z", + "type": "turn_context", + "payload": { "model": "gpt-5.6-sol" } + }) + } + + fn token_count(input: u64, cached: u64, output: u64) -> serde_json::Value { + serde_json::json!({ + "timestamp": "2026-07-10T03:00:02Z", + "type": "event_msg", + "payload": { + "type": "token_count", + "info": { "total_token_usage": { + "input_tokens": input, + "cached_input_tokens": cached, + "output_tokens": output + }} + } + }) + } #[test] fn test_delta_first_event() { @@ -659,6 +874,159 @@ mod tests { assert!(files.is_empty()); } + #[test] + fn test_subagent_identity_prefers_unique_thread_id() { + let identity = + parse_codex_session_identity(session_meta("child", "parent").get("payload").unwrap()) + .unwrap(); + + assert_eq!(identity.thread_id, "child"); + assert!(identity.carries_history_snapshot); + } + + #[test] + fn test_subagent_replay_only_establishes_token_baseline() -> Result<(), AppError> { + let db = Database::memory()?; + let temp = tempdir().unwrap(); + let child = temp.path().join("child.jsonl"); + write_jsonl( + &child, + &[ + session_meta("child", "parent"), + turn_context(), + token_count(1_000, 900, 100), + token_count(1_200, 1_000, 120), + serde_json::json!({ + "timestamp": "2026-07-10T03:00:03Z", + "type": "event_msg", + "payload": { "type": "thread_settings_applied" } + }), + token_count(1_300, 1_050, 150), + ], + ); + + assert_eq!(sync_single_codex_file(&db, &child)?, (1, 2)); + + let conn = lock_conn!(db.conn); + let usage: (i64, i64, i64) = conn.query_row( + "SELECT input_tokens, cache_read_tokens, output_tokens + FROM proxy_request_logs + WHERE request_id = 'codex_session:thread-v1:child:3'", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + )?; + assert_eq!(usage, (100, 50, 30)); + + Ok(()) + } + + #[test] + fn test_subagents_under_same_parent_use_distinct_request_ids() -> Result<(), AppError> { + let db = Database::memory()?; + let temp = tempdir().unwrap(); + let child_a = temp.path().join("child-a.jsonl"); + let child_b = temp.path().join("child-b.jsonl"); + write_jsonl( + &child_a, + &[ + session_meta("child-a", "parent"), + turn_context(), + token_count(100, 50, 10), + ], + ); + write_jsonl( + &child_b, + &[ + session_meta("child-b", "parent"), + turn_context(), + token_count(200, 100, 20), + ], + ); + + assert_eq!(sync_single_codex_file(&db, &child_a)?, (1, 0)); + assert_eq!(sync_single_codex_file(&db, &child_b)?, (1, 0)); + + let conn = lock_conn!(db.conn); + let request_ids = conn + .prepare( + "SELECT request_id FROM proxy_request_logs + WHERE data_source = 'codex_session' ORDER BY request_id", + )? + .query_map([], |row| row.get::<_, String>(0))? + .collect::, _>>()?; + assert_eq!( + request_ids, + vec![ + "codex_session:thread-v1:child-a:1", + "codex_session:thread-v1:child-b:1" + ] + ); + + Ok(()) + } + + #[test] + fn test_archived_log_inherits_cursor_and_only_imports_appended_usage() -> Result<(), AppError> { + let db = Database::memory()?; + let temp = tempdir().unwrap(); + let sessions = temp.path().join("sessions"); + let archived = temp.path().join("archived_sessions"); + fs::create_dir_all(&sessions).unwrap(); + fs::create_dir_all(&archived).unwrap(); + let source = sessions.join("rollout-parent.jsonl"); + let archived_file = archived.join("rollout-parent.jsonl"); + write_jsonl( + &archived_file, + &[ + session_meta("parent", "parent"), + turn_context(), + token_count(100, 50, 10), + token_count(200, 100, 20), + ], + ); + + { + 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, + total_cost_usd, latency_ms, status_code, session_id, + created_at, data_source + ) VALUES ('codex_session:parent:2', '_codex_session', 'codex', + 'gpt-5.6-sol', 'gpt-5.6-sol', 999, 99, 0, '0', 0, + 200, 'parent', 1, 'codex_session')", + [], + )?; + } + let source_path = source.to_string_lossy().to_string(); + update_sync_state(&db, &source_path, 1, 3)?; + + assert_eq!(sync_single_codex_file(&db, &archived_file)?, (1, 0)); + assert_eq!(sync_single_codex_file(&db, &archived_file)?, (0, 0)); + + let conn = lock_conn!(db.conn); + let old_row_count: i64 = conn.query_row( + "SELECT COUNT(*) FROM proxy_request_logs + WHERE request_id = 'codex_session:parent:2'", + [], + |row| row.get(0), + )?; + assert_eq!(old_row_count, 1); + let usage: (i64, i64, i64) = conn.query_row( + "SELECT input_tokens, cache_read_tokens, output_tokens + FROM proxy_request_logs + WHERE request_id = 'codex_session:thread-v1:parent:2'", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + )?; + assert_eq!(usage, (100, 50, 10)); + drop(conn); + assert_eq!(get_sync_state(&db, &archived_file.to_string_lossy())?.1, 4); + + Ok(()) + } + #[test] fn test_insert_codex_session_skips_matching_proxy_log() -> Result<(), AppError> { let db = Database::memory()?;