mirror of
https://github.com/farion1231/cc-switch.git
synced 2026-08-04 03:32:25 +08:00
56fb46c093
* perf(codex): cache parent rollout timelines * fix(codex): harden parent timeline cache * fix(codex): tighten replay cache invalidation --------- Co-authored-by: Ayanami <ay@nami.ltd> Co-authored-by: SaladDay <92240037+SaladDay@users.noreply.github.com>
2328 lines
79 KiB
Rust
2328 lines
79 KiB
Rust
//! Codex 会话日志使用追踪
|
||
//!
|
||
//! 从 ~/.codex/sessions/ 下的 JSONL 会话文件中提取精确 token 使用数据,
|
||
//! 替代原有的 state_5.sqlite 估算方案。
|
||
//!
|
||
//! ## 数据流
|
||
//! ```text
|
||
//! ~/.codex/sessions/YYYY/MM/DD/*.jsonl → 增量解析 → delta 计算 → 费用计算 → proxy_request_logs 表
|
||
//! ```
|
||
//!
|
||
//! ## 解析的事件类型
|
||
//! - `session_meta` → 提取唯一 thread_id(子代理的 session_id 指向父线程)
|
||
//! - `turn_context` → 提取当前 model
|
||
//! - `event_msg` (type=token_count) → 提取累计 token 用量,计算 delta
|
||
|
||
use crate::codex_config::get_codex_config_dir;
|
||
use crate::database::{lock_conn, Database};
|
||
use crate::error::AppError;
|
||
use crate::proxy::usage::calculator::{CostCalculator, ModelPricing};
|
||
use crate::proxy::usage::parser::TokenUsage;
|
||
use crate::services::session_usage::{
|
||
get_sync_state, metadata_modified_nanos, update_sync_state, SessionSyncResult,
|
||
};
|
||
use crate::services::usage_stats::{
|
||
find_model_pricing, has_suspected_codex_session_duplicate, should_skip_session_insert, DedupKey,
|
||
};
|
||
use chrono::{DateTime, Utc};
|
||
use rust_decimal::Decimal;
|
||
use std::collections::HashMap;
|
||
use std::fs;
|
||
use std::io::{BufRead, BufReader};
|
||
#[cfg(unix)]
|
||
use std::os::unix::fs::MetadataExt;
|
||
#[cfg(windows)]
|
||
use std::os::windows::io::AsRawHandle;
|
||
use std::path::{Path, PathBuf};
|
||
use std::sync::{Arc, Mutex, OnceLock};
|
||
use std::time::SystemTime;
|
||
#[cfg(windows)]
|
||
use windows_sys::Win32::Storage::FileSystem::{
|
||
FileIdInfo, GetFileInformationByHandleEx, FILE_ID_INFO,
|
||
};
|
||
|
||
const CODEX_THREAD_REQUEST_ID_PREFIX: &str = "codex_session:thread-v1";
|
||
|
||
/// 累计 token 用量(跟踪 total_token_usage 字段)
|
||
#[derive(Debug, Clone, Default)]
|
||
struct CumulativeTokens {
|
||
input: u64,
|
||
cached_input: u64,
|
||
output: u64,
|
||
}
|
||
|
||
/// 单次 API 调用的 token 增量
|
||
#[derive(Debug)]
|
||
struct DeltaTokens {
|
||
input: u32,
|
||
cached_input: u32,
|
||
output: u32,
|
||
}
|
||
|
||
impl DeltaTokens {
|
||
fn is_zero(&self) -> bool {
|
||
self.input == 0 && self.cached_input == 0 && self.output == 0
|
||
}
|
||
}
|
||
|
||
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
||
struct TokenCountersSignature {
|
||
input: Option<u64>,
|
||
cached_input: Option<u64>,
|
||
output: Option<u64>,
|
||
reasoning_output: Option<u64>,
|
||
total: Option<u64>,
|
||
}
|
||
|
||
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
||
struct TokenUsageSignature {
|
||
total: Option<TokenCountersSignature>,
|
||
last: Option<TokenCountersSignature>,
|
||
}
|
||
|
||
#[derive(Debug)]
|
||
struct TimestampedTokenSignature {
|
||
timestamp: DateTime<Utc>,
|
||
signature: TokenUsageSignature,
|
||
}
|
||
|
||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||
struct ParentFileStamp {
|
||
modified_nanos: i64,
|
||
size: u64,
|
||
#[cfg(unix)]
|
||
device: u64,
|
||
#[cfg(unix)]
|
||
inode: u64,
|
||
#[cfg(windows)]
|
||
volume_serial: u64,
|
||
#[cfg(windows)]
|
||
file_id: [u8; 16],
|
||
}
|
||
|
||
impl ParentFileStamp {
|
||
fn from_file(file: &fs::File) -> Option<Self> {
|
||
let metadata = file.metadata().ok()?;
|
||
#[cfg(windows)]
|
||
let (volume_serial, file_id) = windows_file_identity(file)?;
|
||
Some(Self {
|
||
modified_nanos: metadata_modified_nanos(&metadata),
|
||
size: metadata.len(),
|
||
#[cfg(unix)]
|
||
device: metadata.dev(),
|
||
#[cfg(unix)]
|
||
inode: metadata.ino(),
|
||
#[cfg(windows)]
|
||
volume_serial,
|
||
#[cfg(windows)]
|
||
file_id,
|
||
})
|
||
}
|
||
}
|
||
|
||
#[cfg(windows)]
|
||
fn windows_file_identity(file: &fs::File) -> Option<(u64, [u8; 16])> {
|
||
let mut information = FILE_ID_INFO::default();
|
||
// SAFETY: `file` owns a live handle for this call, and `information` is a
|
||
// valid writable FILE_ID_INFO buffer of the size passed to Windows.
|
||
let succeeded = unsafe {
|
||
GetFileInformationByHandleEx(
|
||
file.as_raw_handle(),
|
||
FileIdInfo,
|
||
std::ptr::addr_of_mut!(information).cast(),
|
||
std::mem::size_of::<FILE_ID_INFO>() as u32,
|
||
)
|
||
} != 0;
|
||
succeeded.then_some((
|
||
information.VolumeSerialNumber,
|
||
information.FileId.Identifier,
|
||
))
|
||
}
|
||
|
||
#[derive(Debug)]
|
||
struct ParentTokenTimeline {
|
||
events: Vec<TimestampedTokenSignature>,
|
||
max_timestamp: Option<DateTime<Utc>>,
|
||
has_token_without_timestamp: bool,
|
||
}
|
||
|
||
impl ParentTokenTimeline {
|
||
fn signatures_before(
|
||
&self,
|
||
parent_path: &Path,
|
||
cutoff: DateTime<Utc>,
|
||
) -> Result<Vec<TokenUsageSignature>, String> {
|
||
if self.has_token_without_timestamp {
|
||
return Err(format!(
|
||
"父 rollout {} 的 token_count 缺少有效 timestamp",
|
||
parent_path.display()
|
||
));
|
||
}
|
||
if self
|
||
.max_timestamp
|
||
.is_none_or(|timestamp| timestamp < cutoff)
|
||
{
|
||
return Err(format!(
|
||
"父 rollout {} 尚未写到 child fork 时刻",
|
||
parent_path.display()
|
||
));
|
||
}
|
||
Ok(self
|
||
.events
|
||
.iter()
|
||
.filter(|event| event.timestamp <= cutoff)
|
||
.map(|event| event.signature.clone())
|
||
.collect())
|
||
}
|
||
}
|
||
|
||
#[derive(Debug)]
|
||
struct CachedParentTimeline {
|
||
stamp: ParentFileStamp,
|
||
timeline: Arc<ParentTokenTimeline>,
|
||
}
|
||
|
||
#[derive(Debug)]
|
||
struct CachedReplayPrefix {
|
||
modified: i64,
|
||
size: u64,
|
||
prefix: usize,
|
||
}
|
||
|
||
#[derive(Debug)]
|
||
struct ParsedTokenEvent {
|
||
line_offset: i64,
|
||
signature: TokenUsageSignature,
|
||
delta: DeltaTokens,
|
||
event_index: Option<u32>,
|
||
model: String,
|
||
timestamp: Option<String>,
|
||
}
|
||
|
||
#[derive(Debug)]
|
||
enum ParentResolution {
|
||
None,
|
||
Parent(String),
|
||
Deferred(String),
|
||
}
|
||
|
||
#[derive(Debug)]
|
||
struct ParsedCodexFile {
|
||
root_thread_id: Option<String>,
|
||
root_meta_seen: bool,
|
||
root_timestamp: Option<DateTime<Utc>>,
|
||
parent: ParentResolution,
|
||
token_events: Vec<ParsedTokenEvent>,
|
||
line_offset: i64,
|
||
has_billable_tokens: bool,
|
||
}
|
||
|
||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||
enum PendingReason {
|
||
MissingParent(String),
|
||
Stable(String),
|
||
Retryable(String),
|
||
}
|
||
|
||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||
struct PendingEntry {
|
||
modified: i64,
|
||
size: u64,
|
||
reason: PendingReason,
|
||
}
|
||
|
||
#[derive(Debug, Default)]
|
||
struct CodexReplayCaches {
|
||
parent_timelines: HashMap<PathBuf, CachedParentTimeline>,
|
||
replay_prefixes: HashMap<PathBuf, CachedReplayPrefix>,
|
||
pending: HashMap<PathBuf, PendingEntry>,
|
||
}
|
||
|
||
static CODEX_REPLAY_CACHES: OnceLock<Mutex<CodexReplayCaches>> = OnceLock::new();
|
||
|
||
fn replay_caches() -> &'static Mutex<CodexReplayCaches> {
|
||
CODEX_REPLAY_CACHES.get_or_init(|| Mutex::new(CodexReplayCaches::default()))
|
||
}
|
||
|
||
pub(crate) fn clear_codex_replay_caches() {
|
||
if let Ok(mut caches) = replay_caches().lock() {
|
||
*caches = CodexReplayCaches::default();
|
||
}
|
||
}
|
||
|
||
fn is_rollout_filename(file_name: &str) -> bool {
|
||
if !file_name.starts_with("rollout-") || !file_name.ends_with(".jsonl") {
|
||
return false;
|
||
}
|
||
let stem = file_name.trim_end_matches(".jsonl");
|
||
stem.get(stem.len().saturating_sub(36)..)
|
||
.is_some_and(|candidate| uuid::Uuid::parse_str(candidate).is_ok())
|
||
}
|
||
|
||
fn is_codex_cursor_path(file_path: &str, codex_dir: &Path) -> bool {
|
||
let path = Path::new(file_path);
|
||
let file_name = file_path.rsplit(['/', '\\']).next().unwrap_or_default();
|
||
if !is_rollout_filename(file_name) {
|
||
return false;
|
||
}
|
||
|
||
if path.starts_with(codex_dir.join("sessions"))
|
||
|| path.starts_with(codex_dir.join("archived_sessions"))
|
||
{
|
||
return true;
|
||
}
|
||
|
||
// 兼容用户改过 CODEX_HOME 后遗留、且源文件已不存在的 cursor。只接受
|
||
// 明确目录段 + Codex rollout UUID 文件名,避免宽 codex_dir 误删其他 importer。
|
||
file_path
|
||
.replace('\\', "/")
|
||
.split('/')
|
||
.any(|segment| matches!(segment, "sessions" | "archived_sessions"))
|
||
}
|
||
|
||
fn sqlite_table_exists(conn: &rusqlite::Connection, table: &str) -> Result<bool, AppError> {
|
||
conn.query_row(
|
||
"SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?1)",
|
||
[table],
|
||
|row| row.get(0),
|
||
)
|
||
.map_err(|error| AppError::Database(format!("查询表 {table} 失败: {error}")))
|
||
}
|
||
|
||
fn sqlite_column_exists(
|
||
conn: &rusqlite::Connection,
|
||
table: &str,
|
||
column: &str,
|
||
) -> Result<bool, AppError> {
|
||
conn.query_row(
|
||
"SELECT EXISTS(SELECT 1 FROM pragma_table_info(?1) WHERE name = ?2)",
|
||
rusqlite::params![table, column],
|
||
|row| row.get(0),
|
||
)
|
||
.map_err(|error| AppError::Database(format!("查询列 {table}.{column} 失败: {error}")))
|
||
}
|
||
|
||
pub(crate) fn reset_codex_usage_on_conn(
|
||
conn: &rusqlite::Connection,
|
||
codex_dir: &Path,
|
||
) -> Result<(), AppError> {
|
||
if sqlite_table_exists(conn, "proxy_request_logs")?
|
||
&& sqlite_column_exists(conn, "proxy_request_logs", "data_source")?
|
||
{
|
||
conn.execute(
|
||
"DELETE FROM proxy_request_logs WHERE data_source = 'codex_session'",
|
||
[],
|
||
)
|
||
.map_err(|error| AppError::Database(format!("清理 Codex 会话明细失败: {error}")))?;
|
||
}
|
||
if sqlite_table_exists(conn, "usage_daily_rollups")?
|
||
&& sqlite_column_exists(conn, "usage_daily_rollups", "provider_id")?
|
||
{
|
||
conn.execute(
|
||
"DELETE FROM usage_daily_rollups WHERE provider_id = '_codex_session'",
|
||
[],
|
||
)
|
||
.map_err(|error| AppError::Database(format!("清理 Codex 用量汇总失败: {error}")))?;
|
||
}
|
||
if sqlite_table_exists(conn, "session_log_sync")?
|
||
&& sqlite_column_exists(conn, "session_log_sync", "file_path")?
|
||
{
|
||
let paths = {
|
||
let mut statement = conn
|
||
.prepare("SELECT file_path FROM session_log_sync")
|
||
.map_err(|error| {
|
||
AppError::Database(format!("读取会话同步 cursor 失败: {error}"))
|
||
})?;
|
||
let paths = statement
|
||
.query_map([], |row| row.get::<_, String>(0))
|
||
.map_err(|error| AppError::Database(format!("查询会话同步 cursor 失败: {error}")))?
|
||
.collect::<Result<Vec<_>, _>>()
|
||
.map_err(|error| {
|
||
AppError::Database(format!("解析会话同步 cursor 失败: {error}"))
|
||
})?;
|
||
paths
|
||
};
|
||
for file_path in paths
|
||
.into_iter()
|
||
.filter(|path| is_codex_cursor_path(path, codex_dir))
|
||
{
|
||
conn.execute(
|
||
"DELETE FROM session_log_sync WHERE file_path = ?1",
|
||
[file_path],
|
||
)
|
||
.map_err(|error| AppError::Database(format!("清理 Codex 同步 cursor 失败: {error}")))?;
|
||
}
|
||
}
|
||
Ok(())
|
||
}
|
||
|
||
impl Database {
|
||
pub(crate) fn reset_codex_usage(&self) -> Result<(), AppError> {
|
||
let codex_dir = get_codex_config_dir();
|
||
let conn = lock_conn!(self.conn);
|
||
conn.execute("SAVEPOINT reset_codex_usage", [])
|
||
.map_err(|error| AppError::Database(format!("开启 Codex 重建事务失败: {error}")))?;
|
||
let result = reset_codex_usage_on_conn(&conn, &codex_dir);
|
||
match result {
|
||
Ok(()) => {
|
||
conn.execute("RELEASE reset_codex_usage", [])
|
||
.map_err(|error| {
|
||
AppError::Database(format!("提交 Codex 重建事务失败: {error}"))
|
||
})?;
|
||
drop(conn);
|
||
clear_codex_replay_caches();
|
||
Ok(())
|
||
}
|
||
Err(error) => {
|
||
conn.execute("ROLLBACK TO reset_codex_usage", []).ok();
|
||
conn.execute("RELEASE reset_codex_usage", []).ok();
|
||
Err(error)
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
fn non_empty_string(value: Option<&serde_json::Value>) -> Option<String> {
|
||
value
|
||
.and_then(serde_json::Value::as_str)
|
||
.filter(|value| !value.is_empty())
|
||
.map(str::to_owned)
|
||
}
|
||
|
||
fn thread_id_from_filename(path: &Path) -> Option<String> {
|
||
let stem = path.file_stem()?.to_str()?;
|
||
let candidate = stem.get(stem.len().checked_sub(36)?..)?;
|
||
uuid::Uuid::parse_str(candidate)
|
||
.ok()
|
||
.map(|value| value.hyphenated().to_string())
|
||
}
|
||
|
||
fn explicit_parent_from_meta(payload: &serde_json::Value) -> ParentResolution {
|
||
let forked_from = non_empty_string(payload.get("forked_from_id"));
|
||
let spawned_from = payload
|
||
.get("source")
|
||
.and_then(|source| source.get("subagent"))
|
||
.and_then(|subagent| subagent.get("thread_spawn"))
|
||
.and_then(|spawn| non_empty_string(spawn.get("parent_thread_id")));
|
||
|
||
match (forked_from, spawned_from) {
|
||
(None, None) => ParentResolution::None,
|
||
(Some(parent), None) | (None, Some(parent)) => ParentResolution::Parent(parent),
|
||
(Some(forked), Some(spawned)) if forked == spawned => ParentResolution::Parent(forked),
|
||
(Some(forked), Some(spawned)) => ParentResolution::Deferred(format!(
|
||
"forked_from_id ({forked}) 与 thread_spawn.parent_thread_id ({spawned}) 不一致"
|
||
)),
|
||
}
|
||
}
|
||
|
||
fn parse_timestamp(value: Option<&serde_json::Value>) -> Option<DateTime<Utc>> {
|
||
value
|
||
.and_then(serde_json::Value::as_str)
|
||
.and_then(|value| DateTime::parse_from_rfc3339(value).ok())
|
||
.map(|value| value.with_timezone(&Utc))
|
||
}
|
||
|
||
fn parse_signature_counters(value: Option<&serde_json::Value>) -> Option<TokenCountersSignature> {
|
||
let value = value?.as_object()?;
|
||
Some(TokenCountersSignature {
|
||
input: value
|
||
.get("input_tokens")
|
||
.and_then(serde_json::Value::as_u64),
|
||
cached_input: value
|
||
.get("cached_input_tokens")
|
||
.or_else(|| value.get("cache_read_input_tokens"))
|
||
.and_then(serde_json::Value::as_u64),
|
||
output: value
|
||
.get("output_tokens")
|
||
.and_then(serde_json::Value::as_u64),
|
||
reasoning_output: value
|
||
.get("reasoning_output_tokens")
|
||
.and_then(serde_json::Value::as_u64),
|
||
total: value
|
||
.get("total_tokens")
|
||
.and_then(serde_json::Value::as_u64),
|
||
})
|
||
}
|
||
|
||
fn parse_token_signature(info: &serde_json::Value) -> Option<TokenUsageSignature> {
|
||
let total = parse_signature_counters(info.get("total_token_usage"));
|
||
let last = parse_signature_counters(info.get("last_token_usage"));
|
||
(total.is_some() || last.is_some()).then_some(TokenUsageSignature { total, last })
|
||
}
|
||
|
||
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 模型名
|
||
///
|
||
/// 处理规则(按顺序):
|
||
/// 1. 转小写:`GLM-4.6` → `glm-4.6`
|
||
/// 2. 剥离 provider 前缀:`openai/gpt-5.4` → `gpt-5.4`
|
||
/// 3. 剥离 ISO 日期后缀:`gpt-5.4-2026-03-05` → `gpt-5.4`
|
||
/// 4. 剥离紧凑日期后缀:`gpt-5.4-20260305` → `gpt-5.4`
|
||
fn normalize_codex_model(raw: &str) -> String {
|
||
// Step 1: 小写
|
||
let mut name = raw.to_lowercase();
|
||
|
||
// Step 2: 剥离 "provider/" 前缀(如 openai/, azure/)
|
||
if let Some(pos) = name.rfind('/') {
|
||
name = name[pos + 1..].to_string();
|
||
}
|
||
|
||
// Step 3: 剥离 ISO 日期后缀 -YYYY-MM-DD(正好 11 字符)
|
||
if name.len() > 11 && name.is_char_boundary(name.len() - 11) {
|
||
let suffix = &name[name.len() - 11..];
|
||
if suffix.is_ascii()
|
||
&& suffix.as_bytes()[0] == b'-'
|
||
&& suffix[1..5].chars().all(|c| c.is_ascii_digit())
|
||
&& suffix.as_bytes()[5] == b'-'
|
||
&& suffix[6..8].chars().all(|c| c.is_ascii_digit())
|
||
&& suffix.as_bytes()[8] == b'-'
|
||
&& suffix[9..11].chars().all(|c| c.is_ascii_digit())
|
||
{
|
||
name.truncate(name.len() - 11);
|
||
}
|
||
}
|
||
|
||
// Step 4: 剥离紧凑日期后缀 -YYYYMMDD(正好 9 字符)
|
||
if name.len() > 9 {
|
||
let parts: Vec<&str> = name.rsplitn(2, '-').collect();
|
||
if parts.len() == 2 {
|
||
if let Some(suffix) = parts.first() {
|
||
if suffix.len() == 8 && suffix.chars().all(|c| c.is_ascii_digit()) {
|
||
name = parts[1].to_string();
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
name
|
||
}
|
||
|
||
/// 计算两次累计值之间的 delta
|
||
fn compute_delta(prev: &Option<CumulativeTokens>, current: &CumulativeTokens) -> DeltaTokens {
|
||
match prev {
|
||
None => DeltaTokens {
|
||
input: current.input as u32,
|
||
cached_input: current.cached_input as u32,
|
||
output: current.output as u32,
|
||
},
|
||
Some(p) => DeltaTokens {
|
||
input: current.input.saturating_sub(p.input) as u32,
|
||
cached_input: current.cached_input.saturating_sub(p.cached_input) as u32,
|
||
output: current.output.saturating_sub(p.output) as u32,
|
||
},
|
||
}
|
||
}
|
||
|
||
/// 从 JSON Value 中提取累计 token 用量
|
||
fn parse_cumulative_tokens(total_usage: &serde_json::Value) -> Option<CumulativeTokens> {
|
||
if total_usage.is_null() || !total_usage.is_object() {
|
||
return None;
|
||
}
|
||
Some(CumulativeTokens {
|
||
input: total_usage
|
||
.get("input_tokens")
|
||
.and_then(|v| v.as_u64())
|
||
.unwrap_or(0),
|
||
cached_input: total_usage
|
||
.get("cached_input_tokens")
|
||
.or_else(|| total_usage.get("cache_read_input_tokens"))
|
||
.and_then(|v| v.as_u64())
|
||
.unwrap_or(0),
|
||
output: total_usage
|
||
.get("output_tokens")
|
||
.and_then(|v| v.as_u64())
|
||
.unwrap_or(0),
|
||
})
|
||
}
|
||
|
||
type RolloutIndex = HashMap<String, Vec<PathBuf>>;
|
||
|
||
#[derive(Debug, Default)]
|
||
struct CodexFileSyncResult {
|
||
imported: u32,
|
||
skipped: u32,
|
||
suspected_duplicates: u32,
|
||
deferred: bool,
|
||
}
|
||
|
||
/// 同步 Codex 使用数据(从 JSONL 会话日志)
|
||
pub fn sync_codex_usage(db: &Database) -> Result<SessionSyncResult, AppError> {
|
||
let codex_dir = get_codex_config_dir();
|
||
let files = collect_codex_session_files(&codex_dir);
|
||
let rollout_index = build_rollout_index(&files);
|
||
|
||
let mut result = SessionSyncResult {
|
||
imported: 0,
|
||
skipped: 0,
|
||
files_scanned: files.len() as u32,
|
||
suspected_duplicates: 0,
|
||
deferred_files: 0,
|
||
errors: vec![],
|
||
};
|
||
|
||
for file_path in &files {
|
||
match sync_single_codex_file(db, file_path, &rollout_index) {
|
||
Ok(file_result) => {
|
||
result.imported = result.imported.saturating_add(file_result.imported);
|
||
result.skipped = result.skipped.saturating_add(file_result.skipped);
|
||
result.suspected_duplicates = result
|
||
.suspected_duplicates
|
||
.saturating_add(file_result.suspected_duplicates);
|
||
if file_result.deferred {
|
||
result.deferred_files = result.deferred_files.saturating_add(1);
|
||
}
|
||
}
|
||
Err(e) => {
|
||
let msg = format!("Codex 会话文件解析失败 {}: {e}", file_path.display());
|
||
log::warn!("[CODEX-SYNC] {msg}");
|
||
result.errors.push(msg);
|
||
}
|
||
}
|
||
}
|
||
|
||
if result.imported > 0 || result.deferred_files > 0 {
|
||
log::info!(
|
||
"[CODEX-SYNC] 同步完成: 导入 {} 条, 跳过 {} 条, deferred {} 个, 扫描 {} 个文件",
|
||
result.imported,
|
||
result.skipped,
|
||
result.deferred_files,
|
||
result.files_scanned
|
||
);
|
||
}
|
||
|
||
Ok(result)
|
||
}
|
||
|
||
/// 收集所有 Codex 会话 JSONL 文件
|
||
fn collect_codex_session_files(codex_dir: &Path) -> Vec<PathBuf> {
|
||
let mut files = Vec::new();
|
||
|
||
// 1. 扫描 sessions/YYYY/MM/DD/*.jsonl(日期分区目录)
|
||
let sessions_dir = codex_dir.join("sessions");
|
||
if sessions_dir.is_dir() {
|
||
collect_jsonl_recursive(&sessions_dir, &mut files, 0, 3);
|
||
}
|
||
|
||
// 2. 扫描 archived_sessions/*.jsonl(扁平归档目录)
|
||
let archived_dir = codex_dir.join("archived_sessions");
|
||
if archived_dir.is_dir() {
|
||
if let Ok(entries) = fs::read_dir(&archived_dir) {
|
||
for entry in entries.flatten() {
|
||
let path = entry.path();
|
||
if path.extension().and_then(|e| e.to_str()) == Some("jsonl") {
|
||
files.push(path);
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
files.sort();
|
||
files
|
||
}
|
||
|
||
fn build_rollout_index(files: &[PathBuf]) -> RolloutIndex {
|
||
let mut index = RolloutIndex::new();
|
||
for path in files {
|
||
if let Some(thread_id) = thread_id_from_filename(path) {
|
||
index.entry(thread_id).or_default().push(path.clone());
|
||
}
|
||
}
|
||
for paths in index.values_mut() {
|
||
paths.sort();
|
||
}
|
||
index
|
||
}
|
||
|
||
/// 递归扫描目录下的 .jsonl 文件(限制最大深度)
|
||
fn collect_jsonl_recursive(dir: &Path, files: &mut Vec<PathBuf>, depth: u32, max_depth: u32) {
|
||
let entries = match fs::read_dir(dir) {
|
||
Ok(e) => e,
|
||
Err(_) => return,
|
||
};
|
||
|
||
for entry in entries.flatten() {
|
||
let path = entry.path();
|
||
if path.is_dir() && depth < max_depth {
|
||
collect_jsonl_recursive(&path, files, depth + 1, max_depth);
|
||
} else if path.extension().and_then(|e| e.to_str()) == Some("jsonl") {
|
||
files.push(path);
|
||
}
|
||
}
|
||
}
|
||
|
||
fn parse_codex_file(
|
||
file_path: &Path,
|
||
root_thread_id: Option<String>,
|
||
) -> Result<ParsedCodexFile, AppError> {
|
||
let file =
|
||
fs::File::open(file_path).map_err(|e| AppError::Config(format!("无法打开文件: {e}")))?;
|
||
let reader = BufReader::new(file);
|
||
let mut root_meta_seen = false;
|
||
let mut root_timestamp = None;
|
||
let mut parent = ParentResolution::None;
|
||
let mut current_model = "unknown".to_string();
|
||
let mut prev_total: Option<CumulativeTokens> = None;
|
||
let mut event_index = 0u32;
|
||
let mut token_events = Vec::new();
|
||
let mut line_offset = 0i64;
|
||
let mut has_billable_tokens = false;
|
||
|
||
for line_result in reader.lines() {
|
||
line_offset += 1;
|
||
let line = match line_result {
|
||
Ok(line) => line,
|
||
Err(_) => continue,
|
||
};
|
||
if line.trim().is_empty() {
|
||
continue;
|
||
}
|
||
|
||
let is_event_msg = line.contains("\"event_msg\"");
|
||
let is_turn_context = line.contains("\"turn_context\"");
|
||
let is_session_meta = line.contains("\"session_meta\"");
|
||
if !is_event_msg && !is_turn_context && !is_session_meta {
|
||
continue;
|
||
}
|
||
if is_event_msg && !line.contains("\"token_count\"") {
|
||
continue;
|
||
}
|
||
|
||
let value: serde_json::Value = match serde_json::from_str(&line) {
|
||
Ok(value) => value,
|
||
Err(_) => continue,
|
||
};
|
||
let Some(event_type) = value.get("type").and_then(serde_json::Value::as_str) else {
|
||
continue;
|
||
};
|
||
|
||
match event_type {
|
||
"session_meta" if !root_meta_seen => {
|
||
root_meta_seen = true;
|
||
root_timestamp = parse_timestamp(value.get("timestamp"));
|
||
let payload = value.get("payload").unwrap_or(&serde_json::Value::Null);
|
||
parent = explicit_parent_from_meta(payload);
|
||
|
||
let meta_thread_id = non_empty_string(
|
||
payload
|
||
.get("id")
|
||
.or_else(|| payload.get("thread_id"))
|
||
.or_else(|| payload.get("threadId")),
|
||
);
|
||
if let (Some(filename_id), Some(meta_id)) = (&root_thread_id, meta_thread_id) {
|
||
if filename_id != &meta_id {
|
||
parent = ParentResolution::Deferred(format!(
|
||
"文件名线程 ID ({filename_id}) 与 root meta ID ({meta_id}) 不一致"
|
||
));
|
||
}
|
||
}
|
||
|
||
if let ParentResolution::Parent(parent_id) = &mut parent {
|
||
match uuid::Uuid::parse_str(parent_id) {
|
||
Ok(value) => *parent_id = value.hyphenated().to_string(),
|
||
Err(_) => {
|
||
parent = ParentResolution::Deferred(format!(
|
||
"显式 parent_thread_id 不是有效 UUID: {parent_id}"
|
||
));
|
||
}
|
||
}
|
||
}
|
||
if matches!((&root_thread_id, &parent), (Some(root), ParentResolution::Parent(parent_id)) if root == parent_id)
|
||
{
|
||
parent = ParentResolution::Deferred(
|
||
"parent_thread_id 与 root_thread_id 相同".to_string(),
|
||
);
|
||
}
|
||
}
|
||
"turn_context" => {
|
||
if let Some(payload) = value.get("payload") {
|
||
if let Some(model) = payload
|
||
.get("model")
|
||
.or_else(|| payload.get("info").and_then(|info| info.get("model")))
|
||
.and_then(serde_json::Value::as_str)
|
||
{
|
||
current_model = normalize_codex_model(model);
|
||
}
|
||
}
|
||
}
|
||
"event_msg" => {
|
||
let Some(payload) = value.get("payload") else {
|
||
continue;
|
||
};
|
||
if payload.get("type").and_then(serde_json::Value::as_str) != Some("token_count") {
|
||
continue;
|
||
}
|
||
let Some(info) = payload.get("info").filter(|info| !info.is_null()) else {
|
||
continue;
|
||
};
|
||
let Some(signature) = parse_token_signature(info) else {
|
||
continue;
|
||
};
|
||
|
||
if let Some(model) = info
|
||
.get("model")
|
||
.or_else(|| info.get("model_name"))
|
||
.or_else(|| payload.get("model"))
|
||
.and_then(serde_json::Value::as_str)
|
||
{
|
||
current_model = normalize_codex_model(model);
|
||
}
|
||
|
||
let (cumulative, is_total) = if let Some(total) = info.get("total_token_usage") {
|
||
(parse_cumulative_tokens(total), true)
|
||
} else if let Some(last) = info.get("last_token_usage") {
|
||
(parse_cumulative_tokens(last), false)
|
||
} else {
|
||
continue;
|
||
};
|
||
let Some(cumulative) = cumulative else {
|
||
continue;
|
||
};
|
||
let delta = if is_total {
|
||
let delta = compute_delta(&prev_total, &cumulative);
|
||
prev_total = Some(cumulative);
|
||
delta
|
||
} else {
|
||
DeltaTokens {
|
||
input: cumulative.input as u32,
|
||
cached_input: cumulative.cached_input as u32,
|
||
output: cumulative.output as u32,
|
||
}
|
||
};
|
||
let delta = DeltaTokens {
|
||
cached_input: delta.cached_input.min(delta.input),
|
||
..delta
|
||
};
|
||
let nonzero_index = if delta.is_zero() {
|
||
None
|
||
} else {
|
||
has_billable_tokens = true;
|
||
event_index = event_index.saturating_add(1);
|
||
Some(event_index)
|
||
};
|
||
|
||
token_events.push(ParsedTokenEvent {
|
||
line_offset,
|
||
signature,
|
||
delta,
|
||
event_index: nonzero_index,
|
||
model: current_model.clone(),
|
||
timestamp: value
|
||
.get("timestamp")
|
||
.and_then(serde_json::Value::as_str)
|
||
.map(str::to_owned),
|
||
});
|
||
}
|
||
_ => {}
|
||
}
|
||
}
|
||
|
||
Ok(ParsedCodexFile {
|
||
root_thread_id,
|
||
root_meta_seen,
|
||
root_timestamp,
|
||
parent,
|
||
token_events,
|
||
line_offset,
|
||
has_billable_tokens,
|
||
})
|
||
}
|
||
|
||
fn parent_signatures_before(
|
||
parent_path: &Path,
|
||
cutoff: DateTime<Utc>,
|
||
) -> Result<Vec<TokenUsageSignature>, String> {
|
||
let file = fs::File::open(parent_path)
|
||
.map_err(|error| format!("无法打开父 rollout {}: {error}", parent_path.display()))?;
|
||
let stamp = ParentFileStamp::from_file(&file);
|
||
let cached_timeline = stamp.and_then(|stamp| {
|
||
replay_caches().lock().ok().and_then(|caches| {
|
||
caches
|
||
.parent_timelines
|
||
.get(parent_path)
|
||
.filter(|entry| entry.stamp == stamp)
|
||
.map(|entry| Arc::clone(&entry.timeline))
|
||
})
|
||
});
|
||
if let Some(timeline) = cached_timeline {
|
||
return timeline.signatures_before(parent_path, cutoff);
|
||
}
|
||
|
||
let mut events = Vec::new();
|
||
let mut max_timestamp: Option<DateTime<Utc>> = None;
|
||
let mut has_token_without_timestamp = false;
|
||
|
||
// 必须扫描完整父文件,不能在首个未来时间戳处 break:rollout 写入顺序
|
||
// 不承诺时间戳严格单调。缓存完整时间线后,不同 child cutoff 只需内存过滤。
|
||
for line in BufReader::new(file).lines() {
|
||
let Ok(line) = line else {
|
||
continue;
|
||
};
|
||
let Ok(value) = serde_json::from_str::<serde_json::Value>(&line) else {
|
||
continue;
|
||
};
|
||
let timestamp = parse_timestamp(value.get("timestamp"));
|
||
if let Some(timestamp) = timestamp {
|
||
max_timestamp = Some(max_timestamp.map_or(timestamp, |current| current.max(timestamp)));
|
||
}
|
||
if value.get("type").and_then(serde_json::Value::as_str) != Some("event_msg")
|
||
|| value
|
||
.get("payload")
|
||
.and_then(|payload| payload.get("type"))
|
||
.and_then(serde_json::Value::as_str)
|
||
!= Some("token_count")
|
||
{
|
||
continue;
|
||
}
|
||
let Some(info) = value
|
||
.get("payload")
|
||
.and_then(|payload| payload.get("info"))
|
||
.filter(|info| !info.is_null())
|
||
else {
|
||
continue;
|
||
};
|
||
let Some(signature) = parse_token_signature(info) else {
|
||
continue;
|
||
};
|
||
let Some(timestamp) = timestamp else {
|
||
has_token_without_timestamp = true;
|
||
continue;
|
||
};
|
||
events.push(TimestampedTokenSignature {
|
||
timestamp,
|
||
signature,
|
||
});
|
||
}
|
||
|
||
let timeline = Arc::new(ParentTokenTimeline {
|
||
events,
|
||
max_timestamp,
|
||
has_token_without_timestamp,
|
||
});
|
||
let result = timeline.signatures_before(parent_path, cutoff);
|
||
if let (Some(stamp), Ok(mut caches)) = (stamp, replay_caches().lock()) {
|
||
caches.parent_timelines.insert(
|
||
parent_path.to_path_buf(),
|
||
CachedParentTimeline {
|
||
stamp,
|
||
timeline: Arc::clone(&timeline),
|
||
},
|
||
);
|
||
}
|
||
result
|
||
}
|
||
|
||
fn resolve_parent_signatures(
|
||
parent_id: &str,
|
||
cutoff: DateTime<Utc>,
|
||
rollout_index: &RolloutIndex,
|
||
) -> Result<Vec<TokenUsageSignature>, String> {
|
||
let Some(candidates) = rollout_index.get(parent_id) else {
|
||
return Err(format!("找不到父 rollout: {parent_id}"));
|
||
};
|
||
|
||
let mut snapshots = Vec::with_capacity(candidates.len());
|
||
for candidate in candidates {
|
||
snapshots.push(parent_signatures_before(candidate, cutoff)?);
|
||
}
|
||
let Some(first) = snapshots.first() else {
|
||
return Err(format!("找不到父 rollout: {parent_id}"));
|
||
};
|
||
if snapshots.iter().skip(1).any(|snapshot| snapshot != first) {
|
||
return Err(format!(
|
||
"父 rollout UUID {parent_id} 对应多个内容不一致的文件"
|
||
));
|
||
}
|
||
Ok(first.clone())
|
||
}
|
||
|
||
fn matching_replay_prefix(child: &[ParsedTokenEvent], parent: &[TokenUsageSignature]) -> usize {
|
||
let mut parent_offset = 0usize;
|
||
let mut matched = 0usize;
|
||
for event in child {
|
||
let Some(relative_match) = parent[parent_offset..]
|
||
.iter()
|
||
.position(|signature| signature == &event.signature)
|
||
else {
|
||
break;
|
||
};
|
||
parent_offset += relative_match + 1;
|
||
matched += 1;
|
||
}
|
||
matched
|
||
}
|
||
|
||
fn mark_deferred(
|
||
file_path: &Path,
|
||
modified: i64,
|
||
size: u64,
|
||
reason: PendingReason,
|
||
) -> CodexFileSyncResult {
|
||
let entry = PendingEntry {
|
||
modified,
|
||
size,
|
||
reason,
|
||
};
|
||
let should_warn = replay_caches()
|
||
.lock()
|
||
.ok()
|
||
.and_then(|mut caches| {
|
||
caches
|
||
.pending
|
||
.insert(file_path.to_path_buf(), entry.clone())
|
||
})
|
||
.as_ref()
|
||
!= Some(&entry);
|
||
if should_warn {
|
||
let reason = match &entry.reason {
|
||
PendingReason::MissingParent(parent) => format!("找不到父 rollout {parent}"),
|
||
PendingReason::Stable(reason) | PendingReason::Retryable(reason) => reason.clone(),
|
||
};
|
||
log::warn!("[CODEX-SYNC] deferred {}: {reason}", file_path.display());
|
||
}
|
||
CodexFileSyncResult {
|
||
deferred: true,
|
||
..CodexFileSyncResult::default()
|
||
}
|
||
}
|
||
|
||
/// 同步单个 Codex JSONL 文件。
|
||
fn sync_single_codex_file(
|
||
db: &Database,
|
||
file_path: &Path,
|
||
rollout_index: &RolloutIndex,
|
||
) -> Result<CodexFileSyncResult, AppError> {
|
||
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 file_size = metadata.len();
|
||
|
||
// 检查同步状态
|
||
let (last_modified, last_offset) = get_codex_sync_state(db, file_path)?;
|
||
|
||
// 文件未变化则跳过
|
||
if file_modified <= last_modified {
|
||
return Ok(CodexFileSyncResult::default());
|
||
}
|
||
|
||
if let Ok(mut caches) = replay_caches().lock() {
|
||
if let Some(pending) = caches.pending.get(file_path).cloned() {
|
||
if pending.modified == file_modified && pending.size == file_size {
|
||
match &pending.reason {
|
||
PendingReason::MissingParent(parent) if !rollout_index.contains_key(parent) => {
|
||
return Ok(CodexFileSyncResult {
|
||
deferred: true,
|
||
..CodexFileSyncResult::default()
|
||
});
|
||
}
|
||
PendingReason::Stable(_) => {
|
||
return Ok(CodexFileSyncResult {
|
||
deferred: true,
|
||
..CodexFileSyncResult::default()
|
||
});
|
||
}
|
||
PendingReason::Retryable(_) => {
|
||
caches.pending.remove(file_path);
|
||
}
|
||
_ => {
|
||
caches.pending.remove(file_path);
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
let parsed = parse_codex_file(file_path, thread_id_from_filename(file_path))?;
|
||
if !parsed.has_billable_tokens {
|
||
update_sync_state(db, &file_path_str, file_modified, parsed.line_offset)?;
|
||
return Ok(CodexFileSyncResult::default());
|
||
}
|
||
let Some(root_thread_id) = parsed.root_thread_id.as_deref() else {
|
||
return Ok(mark_deferred(
|
||
file_path,
|
||
file_modified,
|
||
file_size,
|
||
PendingReason::Stable("文件名缺少有效的尾部 UUID".to_string()),
|
||
));
|
||
};
|
||
if !parsed.root_meta_seen {
|
||
return Ok(mark_deferred(
|
||
file_path,
|
||
file_modified,
|
||
file_size,
|
||
PendingReason::Stable("含计费 token 但尚无 session_meta".to_string()),
|
||
));
|
||
}
|
||
|
||
let replay_prefix = match &parsed.parent {
|
||
ParentResolution::None => 0,
|
||
ParentResolution::Deferred(reason) => {
|
||
return Ok(mark_deferred(
|
||
file_path,
|
||
file_modified,
|
||
file_size,
|
||
PendingReason::Stable(reason.clone()),
|
||
));
|
||
}
|
||
ParentResolution::Parent(parent_id) => {
|
||
let Some(cutoff) = parsed.root_timestamp else {
|
||
return Ok(mark_deferred(
|
||
file_path,
|
||
file_modified,
|
||
file_size,
|
||
PendingReason::Stable(
|
||
"parented rollout 的 root meta 缺少有效 timestamp".to_string(),
|
||
),
|
||
));
|
||
};
|
||
if let Ok(caches) = replay_caches().lock() {
|
||
if let Some(prefix) = caches
|
||
.replay_prefixes
|
||
.get(file_path)
|
||
.filter(|cached| cached.modified == file_modified && cached.size == file_size)
|
||
.map(|cached| cached.prefix)
|
||
{
|
||
prefix
|
||
} else {
|
||
drop(caches);
|
||
let parent_signatures =
|
||
match resolve_parent_signatures(parent_id, cutoff, rollout_index) {
|
||
Ok(signatures) => signatures,
|
||
Err(reason) => {
|
||
let pending_reason = if rollout_index.contains_key(parent_id) {
|
||
PendingReason::Retryable(reason)
|
||
} else {
|
||
PendingReason::MissingParent(parent_id.clone())
|
||
};
|
||
return Ok(mark_deferred(
|
||
file_path,
|
||
file_modified,
|
||
file_size,
|
||
pending_reason,
|
||
));
|
||
}
|
||
};
|
||
let prefix = matching_replay_prefix(&parsed.token_events, &parent_signatures);
|
||
if let Ok(mut caches) = replay_caches().lock() {
|
||
caches.replay_prefixes.insert(
|
||
file_path.to_path_buf(),
|
||
CachedReplayPrefix {
|
||
modified: file_modified,
|
||
size: file_size,
|
||
prefix,
|
||
},
|
||
);
|
||
}
|
||
prefix
|
||
}
|
||
} else {
|
||
let parent_signatures = resolve_parent_signatures(parent_id, cutoff, rollout_index)
|
||
.map_err(AppError::Config)?;
|
||
matching_replay_prefix(&parsed.token_events, &parent_signatures)
|
||
}
|
||
}
|
||
};
|
||
|
||
if let Ok(mut caches) = replay_caches().lock() {
|
||
caches.pending.remove(file_path);
|
||
}
|
||
|
||
let mut result = CodexFileSyncResult::default();
|
||
for (token_offset, event) in parsed.token_events.iter().enumerate() {
|
||
let Some(event_index) = event.event_index else {
|
||
continue;
|
||
};
|
||
if token_offset < replay_prefix {
|
||
if event.line_offset > last_offset {
|
||
result.skipped = result.skipped.saturating_add(1);
|
||
}
|
||
continue;
|
||
}
|
||
if event.line_offset <= last_offset {
|
||
continue;
|
||
}
|
||
|
||
let request_id = format!("{CODEX_THREAD_REQUEST_ID_PREFIX}:{root_thread_id}:{event_index}");
|
||
match insert_codex_session_entry(
|
||
db,
|
||
&request_id,
|
||
&event.delta,
|
||
&event.model,
|
||
Some(root_thread_id),
|
||
event.timestamp.as_deref(),
|
||
&mut result.suspected_duplicates,
|
||
) {
|
||
Ok(true) => result.imported = result.imported.saturating_add(1),
|
||
Ok(false) => result.skipped = result.skipped.saturating_add(1),
|
||
Err(e) => {
|
||
log::warn!("[CODEX-SYNC] 插入失败 ({request_id}): {e}");
|
||
result.skipped = result.skipped.saturating_add(1);
|
||
}
|
||
}
|
||
}
|
||
|
||
update_sync_state(db, &file_path_str, file_modified, parsed.line_offset)?;
|
||
Ok(result)
|
||
}
|
||
|
||
/// 插入单条 Codex 会话记录到 proxy_request_logs
|
||
fn insert_codex_session_entry(
|
||
db: &Database,
|
||
request_id: &str,
|
||
delta: &DeltaTokens,
|
||
model: &str,
|
||
session_id: Option<&str>,
|
||
timestamp: Option<&str>,
|
||
suspected_duplicates: &mut u32,
|
||
) -> Result<bool, AppError> {
|
||
let conn = lock_conn!(db.conn);
|
||
|
||
let created_at = timestamp
|
||
.and_then(|ts| {
|
||
chrono::DateTime::parse_from_rfc3339(ts)
|
||
.ok()
|
||
.map(|dt| dt.timestamp())
|
||
})
|
||
.unwrap_or_else(|| {
|
||
SystemTime::now()
|
||
.duration_since(SystemTime::UNIX_EPOCH)
|
||
.map(|d| d.as_secs() as i64)
|
||
.unwrap_or(0)
|
||
});
|
||
|
||
let dedup_key = DedupKey {
|
||
app_type: "codex",
|
||
model,
|
||
input_tokens: delta.input,
|
||
output_tokens: delta.output,
|
||
cache_read_tokens: delta.cached_input,
|
||
cache_creation_tokens: 0,
|
||
created_at,
|
||
};
|
||
if should_skip_session_insert(&conn, request_id, &dedup_key)? {
|
||
return Ok(false);
|
||
}
|
||
if has_suspected_codex_session_duplicate(&conn, request_id, &dedup_key)? {
|
||
*suspected_duplicates = suspected_duplicates.saturating_add(1);
|
||
log::warn!(
|
||
"[CODEX-SYNC] 疑似重复会话用量: request_id={request_id}, model={model}, input={}, output={}, cache_read={}",
|
||
delta.input,
|
||
delta.output,
|
||
delta.cached_input
|
||
);
|
||
}
|
||
|
||
// 计算费用
|
||
let usage = TokenUsage {
|
||
input_tokens: delta.input,
|
||
output_tokens: delta.output,
|
||
cache_read_tokens: delta.cached_input,
|
||
cache_creation_tokens: 0,
|
||
model: Some(model.to_string()),
|
||
message_id: None,
|
||
};
|
||
|
||
let pricing = find_codex_pricing(&conn, model);
|
||
let multiplier = Decimal::from(1);
|
||
let (input_cost, output_cost, cache_read_cost, cache_creation_cost, total_cost) = match pricing
|
||
{
|
||
Some(p) => {
|
||
let cost = CostCalculator::calculate_for_app("codex", &usage, &p, multiplier);
|
||
(
|
||
cost.input_cost.to_string(),
|
||
cost.output_cost.to_string(),
|
||
cost.cache_read_cost.to_string(),
|
||
cost.cache_creation_cost.to_string(),
|
||
cost.total_cost.to_string(),
|
||
)
|
||
}
|
||
None => (
|
||
"0".to_string(),
|
||
"0".to_string(),
|
||
"0".to_string(),
|
||
"0".to_string(),
|
||
"0".to_string(),
|
||
),
|
||
};
|
||
|
||
let inserted_rows = conn
|
||
.execute(
|
||
"INSERT OR IGNORE 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
|
||
) 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)",
|
||
rusqlite::params![
|
||
request_id,
|
||
"_codex_session", // provider_id
|
||
"codex", // app_type
|
||
model,
|
||
model, // request_model = model
|
||
delta.input,
|
||
delta.output,
|
||
delta.cached_input,
|
||
0i64, // cache_creation_tokens: Codex 日志无此数据
|
||
input_cost,
|
||
output_cost,
|
||
cache_read_cost,
|
||
cache_creation_cost,
|
||
total_cost,
|
||
0i64, // latency_ms
|
||
Option::<i64>::None, // first_token_ms
|
||
200i64, // status_code
|
||
Option::<String>::None, // error_message
|
||
session_id.map(|s| s.to_string()),
|
||
Some("codex_session"), // provider_type
|
||
1i64, // is_streaming
|
||
"1.0", // cost_multiplier
|
||
created_at,
|
||
"codex_session", // data_source
|
||
],
|
||
)
|
||
.map_err(|e| AppError::Database(format!("插入 Codex 会话日志失败: {e}")))?;
|
||
|
||
Ok(inserted_rows > 0)
|
||
}
|
||
|
||
/// 查找 Codex 模型定价(带归一化)
|
||
fn find_codex_pricing(conn: &rusqlite::Connection, model_id: &str) -> Option<ModelPricing> {
|
||
find_model_pricing(conn, &normalize_codex_model(model_id))
|
||
}
|
||
|
||
#[cfg(test)]
|
||
mod tests {
|
||
use super::*;
|
||
use tempfile::tempdir;
|
||
|
||
const PARENT_ID: &str = "00000000-0000-4000-8000-000000000001";
|
||
const CHILD_A_ID: &str = "00000000-0000-4000-8000-000000000002";
|
||
const CHILD_B_ID: &str = "00000000-0000-4000-8000-000000000003";
|
||
|
||
fn write_jsonl(path: &Path, values: &[serde_json::Value]) {
|
||
let contents = values
|
||
.iter()
|
||
.map(serde_json::Value::to_string)
|
||
.collect::<Vec<_>>()
|
||
.join("\n")
|
||
+ "\n";
|
||
fs::write(path, contents).unwrap();
|
||
}
|
||
|
||
fn rollout_path(dir: &Path, thread_id: &str) -> PathBuf {
|
||
dir.join(format!("rollout-2026-07-10T03-00-00-{thread_id}.jsonl"))
|
||
}
|
||
|
||
fn session_meta_at(
|
||
thread_id: &str,
|
||
forked_from_id: Option<&str>,
|
||
spawned_from_id: Option<&str>,
|
||
timestamp: &str,
|
||
) -> serde_json::Value {
|
||
let source = spawned_from_id.map_or_else(
|
||
|| serde_json::Value::String("cli".to_string()),
|
||
|parent| {
|
||
serde_json::json!({
|
||
"subagent": {
|
||
"thread_spawn": { "parent_thread_id": parent }
|
||
}
|
||
})
|
||
},
|
||
);
|
||
serde_json::json!({
|
||
"timestamp": timestamp,
|
||
"type": "session_meta",
|
||
"payload": {
|
||
"id": thread_id,
|
||
"forked_from_id": forked_from_id,
|
||
"source": source
|
||
}
|
||
})
|
||
}
|
||
|
||
fn session_meta(thread_id: &str) -> serde_json::Value {
|
||
session_meta_at(thread_id, None, None, "2026-07-10T03:00:00Z")
|
||
}
|
||
|
||
fn turn_context_at(timestamp: &str) -> serde_json::Value {
|
||
serde_json::json!({
|
||
"timestamp": timestamp,
|
||
"type": "turn_context",
|
||
"payload": { "model": "gpt-5.6-sol" }
|
||
})
|
||
}
|
||
|
||
fn turn_context() -> serde_json::Value {
|
||
turn_context_at("2026-07-10T03:00:01Z")
|
||
}
|
||
|
||
fn token_count_at(input: u64, cached: u64, output: u64, timestamp: &str) -> serde_json::Value {
|
||
serde_json::json!({
|
||
"timestamp": timestamp,
|
||
"type": "event_msg",
|
||
"payload": {
|
||
"type": "token_count",
|
||
"info": { "total_token_usage": {
|
||
"input_tokens": input,
|
||
"cached_input_tokens": cached,
|
||
"output_tokens": output,
|
||
"reasoning_output_tokens": 0,
|
||
"total_tokens": input + output
|
||
}}
|
||
}
|
||
})
|
||
}
|
||
|
||
fn token_count(input: u64, cached: u64, output: u64) -> serde_json::Value {
|
||
token_count_at(input, cached, output, "2026-07-10T03:00:02Z")
|
||
}
|
||
|
||
fn token_count_without_timestamp(input: u64, cached: u64, output: u64) -> serde_json::Value {
|
||
let mut value = token_count(input, cached, output);
|
||
value
|
||
.as_object_mut()
|
||
.expect("token_count must be an object")
|
||
.remove("timestamp");
|
||
value
|
||
}
|
||
|
||
fn sync_test_file(
|
||
db: &Database,
|
||
file: &Path,
|
||
all_files: &[&Path],
|
||
) -> Result<CodexFileSyncResult, AppError> {
|
||
let files = all_files
|
||
.iter()
|
||
.map(|path| path.to_path_buf())
|
||
.collect::<Vec<_>>();
|
||
sync_single_codex_file(db, file, &build_rollout_index(&files))
|
||
}
|
||
|
||
#[test]
|
||
fn test_delta_first_event() {
|
||
let prev = None;
|
||
let current = CumulativeTokens {
|
||
input: 17934,
|
||
cached_input: 9600,
|
||
output: 454,
|
||
};
|
||
let delta = compute_delta(&prev, ¤t);
|
||
assert_eq!(delta.input, 17934);
|
||
assert_eq!(delta.cached_input, 9600);
|
||
assert_eq!(delta.output, 454);
|
||
assert!(!delta.is_zero());
|
||
}
|
||
|
||
#[test]
|
||
fn test_delta_subsequent_event() {
|
||
let prev = Some(CumulativeTokens {
|
||
input: 17934,
|
||
cached_input: 9600,
|
||
output: 454,
|
||
});
|
||
let current = CumulativeTokens {
|
||
input: 36722,
|
||
cached_input: 27904,
|
||
output: 804,
|
||
};
|
||
let delta = compute_delta(&prev, ¤t);
|
||
assert_eq!(delta.input, 36722 - 17934);
|
||
assert_eq!(delta.cached_input, 27904 - 9600);
|
||
assert_eq!(delta.output, 804 - 454);
|
||
}
|
||
|
||
#[test]
|
||
fn test_delta_zero_at_task_boundary() {
|
||
let prev = Some(CumulativeTokens {
|
||
input: 58346,
|
||
cached_input: 46976,
|
||
output: 1045,
|
||
});
|
||
// task 边界:相同的累计值
|
||
let current = CumulativeTokens {
|
||
input: 58346,
|
||
cached_input: 46976,
|
||
output: 1045,
|
||
};
|
||
let delta = compute_delta(&prev, ¤t);
|
||
assert!(delta.is_zero());
|
||
}
|
||
|
||
#[test]
|
||
fn test_delta_saturating_sub() {
|
||
// 异常情况:当前值小于前值(不应发生,但需防护)
|
||
let prev = Some(CumulativeTokens {
|
||
input: 100,
|
||
cached_input: 50,
|
||
output: 30,
|
||
});
|
||
let current = CumulativeTokens {
|
||
input: 80,
|
||
cached_input: 40,
|
||
output: 20,
|
||
};
|
||
let delta = compute_delta(&prev, ¤t);
|
||
assert_eq!(delta.input, 0);
|
||
assert_eq!(delta.cached_input, 0);
|
||
assert_eq!(delta.output, 0);
|
||
assert!(delta.is_zero());
|
||
}
|
||
|
||
#[test]
|
||
fn test_parse_cumulative_tokens_valid() {
|
||
let json: serde_json::Value = serde_json::json!({
|
||
"input_tokens": 17934,
|
||
"cached_input_tokens": 9600,
|
||
"output_tokens": 454,
|
||
"reasoning_output_tokens": 233,
|
||
"total_tokens": 18388
|
||
});
|
||
let tokens = parse_cumulative_tokens(&json).unwrap();
|
||
assert_eq!(tokens.input, 17934);
|
||
assert_eq!(tokens.cached_input, 9600);
|
||
assert_eq!(tokens.output, 454);
|
||
}
|
||
|
||
#[test]
|
||
fn test_parse_cumulative_tokens_null() {
|
||
let json = serde_json::Value::Null;
|
||
assert!(parse_cumulative_tokens(&json).is_none());
|
||
}
|
||
|
||
#[test]
|
||
fn test_parse_cumulative_tokens_alt_field_names() {
|
||
// 某些版本可能使用 cache_read_input_tokens 而非 cached_input_tokens
|
||
let json: serde_json::Value = serde_json::json!({
|
||
"input_tokens": 1000,
|
||
"cache_read_input_tokens": 500,
|
||
"output_tokens": 200
|
||
});
|
||
let tokens = parse_cumulative_tokens(&json).unwrap();
|
||
assert_eq!(tokens.cached_input, 500);
|
||
}
|
||
|
||
#[test]
|
||
fn test_collect_codex_session_files_nonexistent() {
|
||
let files = collect_codex_session_files(Path::new("/nonexistent/path"));
|
||
assert!(files.is_empty());
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_thread_spawn_parent_strips_replay_and_keeps_live_usage() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
let child = rollout_path(temp.path(), CHILD_A_ID);
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_at(1_000, 900, 100, "2026-07-10T03:00:01Z"),
|
||
turn_context_at("2026-07-10T03:00:10Z"),
|
||
],
|
||
);
|
||
write_jsonl(
|
||
&child,
|
||
&[
|
||
session_meta_at(CHILD_A_ID, None, Some(PARENT_ID), "2026-07-10T03:00:05Z"),
|
||
turn_context(),
|
||
token_count_at(1_000, 900, 100, "2026-07-10T03:00:06Z"),
|
||
token_count_at(1_300, 1_050, 150, "2026-07-10T03:00:07Z"),
|
||
],
|
||
);
|
||
|
||
let result = sync_test_file(&db, &child, &[&parent, &child])?;
|
||
assert_eq!(
|
||
(result.imported, result.skipped, result.deferred),
|
||
(1, 1, false)
|
||
);
|
||
|
||
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 = ?1",
|
||
[format!("{CODEX_THREAD_REQUEST_ID_PREFIX}:{CHILD_A_ID}:2")],
|
||
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
|
||
)?;
|
||
assert_eq!(usage, (300, 150, 50));
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_filtered_parent_events_use_subsequence_prefix_alignment() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
let child = rollout_path(temp.path(), CHILD_A_ID);
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
|
||
token_count_at(200, 100, 20, "2026-07-10T03:00:02Z"),
|
||
token_count_at(300, 150, 30, "2026-07-10T03:00:03Z"),
|
||
turn_context_at("2026-07-10T03:00:10Z"),
|
||
],
|
||
);
|
||
write_jsonl(
|
||
&child,
|
||
&[
|
||
session_meta_at(CHILD_A_ID, Some(PARENT_ID), None, "2026-07-10T03:00:05Z"),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:06Z"),
|
||
token_count_at(300, 150, 30, "2026-07-10T03:00:07Z"),
|
||
token_count_at(450, 220, 45, "2026-07-10T03:00:08Z"),
|
||
],
|
||
);
|
||
|
||
let result = sync_test_file(&db, &child, &[&parent, &child])?;
|
||
assert_eq!((result.imported, result.skipped), (1, 2));
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_parent_rollout_is_cached_once_across_fork_cutoffs() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
|
||
token_count_at(200, 100, 20, "2026-07-10T03:00:10Z"),
|
||
turn_context_at("2026-07-10T03:00:20Z"),
|
||
],
|
||
);
|
||
|
||
let early = "2026-07-10T03:00:05Z".parse::<DateTime<Utc>>().unwrap();
|
||
let late = "2026-07-10T03:00:15Z".parse::<DateTime<Utc>>().unwrap();
|
||
assert_eq!(parent_signatures_before(&parent, early).unwrap().len(), 1);
|
||
let first_timeline =
|
||
Arc::clone(&replay_caches().lock().unwrap().parent_timelines[&parent].timeline);
|
||
assert_eq!(parent_signatures_before(&parent, late).unwrap().len(), 2);
|
||
|
||
let caches = replay_caches().lock().unwrap();
|
||
assert_eq!(caches.parent_timelines.len(), 1);
|
||
assert!(Arc::ptr_eq(
|
||
&first_timeline,
|
||
&caches.parent_timelines[&parent].timeline
|
||
));
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_parent_rollout_cache_invalidates_after_append() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
let cutoff = "2026-07-10T03:00:15Z".parse::<DateTime<Utc>>().unwrap();
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
|
||
turn_context_at("2026-07-10T03:00:20Z"),
|
||
],
|
||
);
|
||
assert_eq!(parent_signatures_before(&parent, cutoff).unwrap().len(), 1);
|
||
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
|
||
token_count_at(200, 100, 20, "2026-07-10T03:00:10Z"),
|
||
turn_context_at("2026-07-10T03:00:20Z"),
|
||
],
|
||
);
|
||
assert_eq!(parent_signatures_before(&parent, cutoff).unwrap().len(), 2);
|
||
|
||
let caches = replay_caches().lock().unwrap();
|
||
assert_eq!(caches.parent_timelines.len(), 1);
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_parent_rollout_content_error_cache_preserves_open_errors() {
|
||
clear_codex_replay_caches();
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
let cutoff = "2026-07-10T03:00:05Z".parse::<DateTime<Utc>>().unwrap();
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_without_timestamp(100, 50, 10),
|
||
turn_context_at("2026-07-10T03:00:20Z"),
|
||
],
|
||
);
|
||
|
||
let first_error = parent_signatures_before(&parent, cutoff).unwrap_err();
|
||
assert!(first_error.contains("token_count 缺少有效 timestamp"));
|
||
let cached_timeline =
|
||
|| Arc::clone(&replay_caches().lock().unwrap().parent_timelines[&parent].timeline);
|
||
let first_timeline = cached_timeline();
|
||
|
||
let second_error = parent_signatures_before(&parent, cutoff).unwrap_err();
|
||
assert_eq!(second_error, first_error);
|
||
assert!(Arc::ptr_eq(&first_timeline, &cached_timeline()));
|
||
|
||
fs::remove_file(&parent).unwrap();
|
||
let open_error = parent_signatures_before(&parent, cutoff).unwrap_err();
|
||
assert!(open_error.contains("无法打开父 rollout"));
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_parent_rollout_nanosecond_cutoffs_are_exact() {
|
||
clear_codex_replay_caches();
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:00.000000500Z"),
|
||
turn_context_at("2026-07-10T03:00:00.000000900Z"),
|
||
],
|
||
);
|
||
|
||
let before = "2026-07-10T03:00:00.000000300Z"
|
||
.parse::<DateTime<Utc>>()
|
||
.unwrap();
|
||
let after = "2026-07-10T03:00:00.000000700Z"
|
||
.parse::<DateTime<Utc>>()
|
||
.unwrap();
|
||
assert!(parent_signatures_before(&parent, before)
|
||
.unwrap()
|
||
.is_empty());
|
||
assert_eq!(parent_signatures_before(&parent, after).unwrap().len(), 1);
|
||
assert_eq!(replay_caches().lock().unwrap().parent_timelines.len(), 1);
|
||
}
|
||
|
||
#[cfg(any(unix, windows))]
|
||
#[test]
|
||
fn test_parent_file_stamp_distinguishes_same_size_same_mtime_files() {
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
let replacement = temp.path().join("replacement.jsonl");
|
||
let values = [session_meta(PARENT_ID), token_count(100, 50, 10)];
|
||
write_jsonl(&parent, &values);
|
||
write_jsonl(&replacement, &values);
|
||
let original_file = fs::File::open(&parent).unwrap();
|
||
let original_metadata = original_file.metadata().unwrap();
|
||
let replacement_file = fs::OpenOptions::new()
|
||
.write(true)
|
||
.open(&replacement)
|
||
.unwrap();
|
||
replacement_file
|
||
.set_times(fs::FileTimes::new().set_modified(original_metadata.modified().unwrap()))
|
||
.unwrap();
|
||
let original_stamp = ParentFileStamp::from_file(&original_file).unwrap();
|
||
let replacement_stamp = ParentFileStamp::from_file(&replacement_file).unwrap();
|
||
assert_eq!(
|
||
(original_stamp.size, original_stamp.modified_nanos),
|
||
(replacement_stamp.size, replacement_stamp.modified_nanos)
|
||
);
|
||
assert_ne!(original_stamp, replacement_stamp);
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_empty_fork_imports_no_parent_usage() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
let child = rollout_path(temp.path(), CHILD_A_ID);
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
|
||
token_count_at(200, 100, 20, "2026-07-10T03:00:02Z"),
|
||
turn_context_at("2026-07-10T03:00:10Z"),
|
||
],
|
||
);
|
||
write_jsonl(
|
||
&child,
|
||
&[
|
||
session_meta_at(CHILD_A_ID, Some(PARENT_ID), None, "2026-07-10T03:00:05Z"),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:06Z"),
|
||
token_count_at(200, 100, 20, "2026-07-10T03:00:07Z"),
|
||
serde_json::json!({
|
||
"timestamp": "2026-07-10T03:00:08Z",
|
||
"type": "event_msg",
|
||
"payload": { "type": "thread_settings_applied" }
|
||
}),
|
||
],
|
||
);
|
||
|
||
let result = sync_test_file(&db, &child, &[&parent, &child])?;
|
||
assert_eq!(
|
||
(result.imported, result.skipped, result.deferred),
|
||
(0, 2, false)
|
||
);
|
||
let conn = lock_conn!(db.conn);
|
||
let count: i64 = conn.query_row(
|
||
"SELECT COUNT(*) FROM proxy_request_logs WHERE data_source = 'codex_session'",
|
||
[],
|
||
|row| row.get(0),
|
||
)?;
|
||
assert_eq!(count, 0);
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_conflicting_explicit_parents_are_deferred() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let child = rollout_path(temp.path(), CHILD_A_ID);
|
||
write_jsonl(
|
||
&child,
|
||
&[
|
||
session_meta_at(
|
||
CHILD_A_ID,
|
||
Some(PARENT_ID),
|
||
Some(CHILD_B_ID),
|
||
"2026-07-10T03:00:05Z",
|
||
),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:06Z"),
|
||
],
|
||
);
|
||
|
||
let result = sync_test_file(&db, &child, &[&child])?;
|
||
assert!(result.deferred);
|
||
assert_eq!(get_sync_state(&db, &child.to_string_lossy())?, (0, 0));
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_parent_future_signature_cannot_extend_replay_prefix() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
let child = rollout_path(temp.path(), CHILD_A_ID);
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
|
||
token_count_at(200, 100, 20, "2026-07-10T03:00:06Z"),
|
||
],
|
||
);
|
||
write_jsonl(
|
||
&child,
|
||
&[
|
||
session_meta_at(CHILD_A_ID, Some(PARENT_ID), None, "2026-07-10T03:00:05Z"),
|
||
token_count_at(200, 100, 20, "2026-07-10T03:00:07Z"),
|
||
],
|
||
);
|
||
|
||
let result = sync_test_file(&db, &child, &[&parent, &child])?;
|
||
assert_eq!(
|
||
(result.imported, result.skipped, result.deferred),
|
||
(1, 0, false)
|
||
);
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_missing_parent_is_deferred_and_recovered_without_child_change() -> Result<(), AppError>
|
||
{
|
||
clear_codex_replay_caches();
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let parent = rollout_path(temp.path(), PARENT_ID);
|
||
let child = rollout_path(temp.path(), CHILD_A_ID);
|
||
write_jsonl(
|
||
&child,
|
||
&[
|
||
session_meta_at(CHILD_A_ID, None, Some(PARENT_ID), "2026-07-10T03:00:05Z"),
|
||
token_count_at(900, 400, 90, "2026-07-10T03:00:06Z"),
|
||
],
|
||
);
|
||
|
||
let deferred = sync_test_file(&db, &child, &[&child])?;
|
||
assert!(deferred.deferred);
|
||
assert_eq!(get_sync_state(&db, &child.to_string_lossy())?, (0, 0));
|
||
|
||
write_jsonl(
|
||
&parent,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
token_count_at(100, 50, 10, "2026-07-10T03:00:01Z"),
|
||
turn_context_at("2026-07-10T03:00:10Z"),
|
||
],
|
||
);
|
||
let recovered = sync_test_file(&db, &child, &[&parent, &child])?;
|
||
assert_eq!((recovered.imported, recovered.deferred), (1, false));
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_billable_file_without_meta_is_deferred_without_cursor() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let child = rollout_path(temp.path(), CHILD_A_ID);
|
||
write_jsonl(&child, &[turn_context(), token_count(100, 50, 10)]);
|
||
|
||
let result = sync_test_file(&db, &child, &[&child])?;
|
||
assert!(result.deferred);
|
||
assert_eq!(get_sync_state(&db, &child.to_string_lossy())?, (0, 0));
|
||
|
||
std::thread::sleep(std::time::Duration::from_millis(2));
|
||
write_jsonl(
|
||
&child,
|
||
&[
|
||
turn_context(),
|
||
token_count(100, 50, 10),
|
||
session_meta_at(CHILD_A_ID, None, None, "2026-07-10T03:00:03Z"),
|
||
],
|
||
);
|
||
let recovered = sync_test_file(&db, &child, &[&child])?;
|
||
assert_eq!((recovered.imported, recovered.deferred), (1, false));
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_non_billable_file_without_meta_advances_cursor() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let child = rollout_path(temp.path(), CHILD_A_ID);
|
||
write_jsonl(
|
||
&child,
|
||
&[
|
||
turn_context(),
|
||
token_count_at(0, 0, 0, "2026-07-10T03:00:02Z"),
|
||
],
|
||
);
|
||
|
||
let result = sync_test_file(&db, &child, &[&child])?;
|
||
assert!(!result.deferred);
|
||
assert_eq!(get_sync_state(&db, &child.to_string_lossy())?.1, 2);
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
#[serial_test::serial]
|
||
fn test_subagents_use_filename_thread_ids() -> Result<(), AppError> {
|
||
clear_codex_replay_caches();
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let child_a = rollout_path(temp.path(), CHILD_A_ID);
|
||
let child_b = rollout_path(temp.path(), CHILD_B_ID);
|
||
write_jsonl(
|
||
&child_a,
|
||
&[
|
||
session_meta(CHILD_A_ID),
|
||
turn_context(),
|
||
token_count(100, 50, 10),
|
||
],
|
||
);
|
||
write_jsonl(
|
||
&child_b,
|
||
&[
|
||
session_meta(CHILD_B_ID),
|
||
turn_context(),
|
||
token_count(200, 100, 20),
|
||
],
|
||
);
|
||
|
||
assert_eq!(
|
||
sync_test_file(&db, &child_a, &[&child_a, &child_b])?.imported,
|
||
1
|
||
);
|
||
assert_eq!(
|
||
sync_test_file(&db, &child_b, &[&child_a, &child_b])?.imported,
|
||
1
|
||
);
|
||
|
||
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::<Result<Vec<_>, _>>()?;
|
||
assert_eq!(
|
||
request_ids,
|
||
vec![
|
||
format!("{CODEX_THREAD_REQUEST_ID_PREFIX}:{CHILD_A_ID}:1"),
|
||
format!("{CODEX_THREAD_REQUEST_ID_PREFIX}:{CHILD_B_ID}: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 = rollout_path(&sessions, PARENT_ID);
|
||
let archived_file = rollout_path(&archived, PARENT_ID);
|
||
write_jsonl(
|
||
&archived_file,
|
||
&[
|
||
session_meta(PARENT_ID),
|
||
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_test_file(&db, &archived_file, &[&archived_file])?.imported,
|
||
1
|
||
);
|
||
assert_eq!(
|
||
sync_test_file(&db, &archived_file, &[&archived_file])?.imported,
|
||
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 = ?1",
|
||
[format!("{CODEX_THREAD_REQUEST_ID_PREFIX}:{PARENT_ID}: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()?;
|
||
{
|
||
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![
|
||
"codex-proxy",
|
||
"openai",
|
||
"codex",
|
||
"gpt-5.4",
|
||
"gpt-5.4",
|
||
10,
|
||
2,
|
||
1,
|
||
7,
|
||
"0.01",
|
||
100,
|
||
200,
|
||
1000,
|
||
"proxy"
|
||
],
|
||
)?;
|
||
}
|
||
|
||
let delta = DeltaTokens {
|
||
input: 10,
|
||
cached_input: 1,
|
||
output: 2,
|
||
};
|
||
let mut suspected_duplicates = 0;
|
||
let inserted = insert_codex_session_entry(
|
||
&db,
|
||
"codex-session-dup",
|
||
&delta,
|
||
"gpt-5.4",
|
||
Some("session-1"),
|
||
Some("1970-01-01T00:16:45Z"),
|
||
&mut suspected_duplicates,
|
||
)?;
|
||
assert!(!inserted);
|
||
|
||
let conn = lock_conn!(db.conn);
|
||
let count: i64 = conn.query_row("SELECT COUNT(*) FROM proxy_request_logs", [], |row| {
|
||
row.get(0)
|
||
})?;
|
||
assert_eq!(count, 1);
|
||
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
fn test_codex_session_duplicate_is_observed_but_still_inserted() -> Result<(), AppError> {
|
||
let db = Database::memory()?;
|
||
let delta = DeltaTokens {
|
||
input: 10,
|
||
cached_input: 1,
|
||
output: 2,
|
||
};
|
||
let mut suspected_duplicates = 0;
|
||
assert!(insert_codex_session_entry(
|
||
&db,
|
||
"codex-session-a",
|
||
&delta,
|
||
"gpt-5.4",
|
||
Some("session-a"),
|
||
Some("1970-01-01T00:16:40Z"),
|
||
&mut suspected_duplicates,
|
||
)?);
|
||
assert!(insert_codex_session_entry(
|
||
&db,
|
||
"codex-session-b",
|
||
&delta,
|
||
"gpt-5.4",
|
||
Some("session-b"),
|
||
Some("1970-01-01T00:16:45Z"),
|
||
&mut suspected_duplicates,
|
||
)?);
|
||
assert_eq!(suspected_duplicates, 1);
|
||
|
||
let conn = lock_conn!(db.conn);
|
||
let count: i64 = conn.query_row(
|
||
"SELECT COUNT(*) FROM proxy_request_logs WHERE data_source = 'codex_session'",
|
||
[],
|
||
|row| row.get(0),
|
||
)?;
|
||
assert_eq!(count, 2);
|
||
Ok(())
|
||
}
|
||
|
||
#[test]
|
||
fn reset_codex_usage_only_removes_codex_rows_and_structural_cursors() -> Result<(), AppError> {
|
||
let db = Database::memory()?;
|
||
let temp = tempdir().unwrap();
|
||
let wide_dir = temp.path();
|
||
let current_codex = rollout_path(&wide_dir.join("sessions"), CHILD_A_ID);
|
||
let legacy_codex =
|
||
format!("C:\\old-codex\\archived_sessions\\rollout-old-{CHILD_B_ID}.jsonl");
|
||
let gemini_cursor = wide_dir.join("gemini/sessions/session-123.json");
|
||
let claude_cursor = wide_dir.join(format!("projects/rollout-{PARENT_ID}.jsonl"));
|
||
|
||
{
|
||
let conn = lock_conn!(db.conn);
|
||
conn.execute_batch(
|
||
"INSERT INTO proxy_request_logs (
|
||
request_id, provider_id, app_type, model, input_tokens,
|
||
output_tokens, cache_read_tokens, latency_ms, status_code,
|
||
created_at, data_source
|
||
) VALUES
|
||
('codex-row', '_codex_session', 'codex', 'gpt', 1, 1, 0, 0, 200, 1, 'codex_session'),
|
||
('gemini-row', '_gemini_session', 'gemini', 'gemini', 1, 1, 0, 0, 200, 1, 'gemini_session');
|
||
INSERT INTO usage_daily_rollups (date, app_type, provider_id, model)
|
||
VALUES
|
||
('2026-07-10', 'codex', '_codex_session', 'gpt'),
|
||
('2026-07-10', 'gemini', '_gemini_session', 'gemini');",
|
||
)?;
|
||
for path in [
|
||
current_codex.to_string_lossy().to_string(),
|
||
legacy_codex,
|
||
gemini_cursor.to_string_lossy().to_string(),
|
||
claude_cursor.to_string_lossy().to_string(),
|
||
] {
|
||
conn.execute(
|
||
"INSERT INTO session_log_sync
|
||
(file_path, last_modified, last_line_offset, last_synced_at)
|
||
VALUES (?1, 1, 1, 1)",
|
||
[path],
|
||
)?;
|
||
}
|
||
|
||
reset_codex_usage_on_conn(&conn, wide_dir)?;
|
||
let codex_rows: i64 = conn.query_row(
|
||
"SELECT COUNT(*) FROM proxy_request_logs WHERE data_source = 'codex_session'",
|
||
[],
|
||
|row| row.get(0),
|
||
)?;
|
||
let gemini_rows: i64 = conn.query_row(
|
||
"SELECT COUNT(*) FROM proxy_request_logs WHERE data_source = 'gemini_session'",
|
||
[],
|
||
|row| row.get(0),
|
||
)?;
|
||
let codex_rollups: i64 = conn.query_row(
|
||
"SELECT COUNT(*) FROM usage_daily_rollups WHERE provider_id = '_codex_session'",
|
||
[],
|
||
|row| row.get(0),
|
||
)?;
|
||
let remaining_cursors: i64 =
|
||
conn.query_row("SELECT COUNT(*) FROM session_log_sync", [], |row| {
|
||
row.get(0)
|
||
})?;
|
||
assert_eq!((codex_rows, gemini_rows, codex_rollups), (0, 1, 0));
|
||
assert_eq!(remaining_cursors, 2);
|
||
}
|
||
Ok(())
|
||
}
|
||
|
||
// ── 模型名归一化测试 ──
|
||
|
||
#[test]
|
||
fn test_normalize_codex_model_lowercase() {
|
||
assert_eq!(normalize_codex_model("GLM-4.6"), "glm-4.6");
|
||
assert_eq!(normalize_codex_model("DeepSeek-Chat"), "deepseek-chat");
|
||
assert_eq!(normalize_codex_model("GPT-5.4"), "gpt-5.4");
|
||
}
|
||
|
||
#[test]
|
||
fn test_normalize_codex_model_strip_prefix() {
|
||
assert_eq!(normalize_codex_model("openai/gpt-5.4"), "gpt-5.4");
|
||
assert_eq!(
|
||
normalize_codex_model("azure/gpt-5.2-codex"),
|
||
"gpt-5.2-codex"
|
||
);
|
||
assert_eq!(normalize_codex_model("OPENAI/GPT-5.4"), "gpt-5.4");
|
||
}
|
||
|
||
#[test]
|
||
fn test_normalize_codex_model_strip_iso_date() {
|
||
assert_eq!(normalize_codex_model("gpt-5.4-2026-03-05"), "gpt-5.4");
|
||
assert_eq!(
|
||
normalize_codex_model("gpt-5.4-pro-2026-03-05"),
|
||
"gpt-5.4-pro"
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn test_normalize_codex_model_strip_compact_date() {
|
||
assert_eq!(normalize_codex_model("gpt-5.4-20260305"), "gpt-5.4");
|
||
assert_eq!(
|
||
normalize_codex_model("claude-opus-4-6-20260206"),
|
||
"claude-opus-4-6"
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn test_normalize_codex_model_no_change() {
|
||
assert_eq!(normalize_codex_model("gpt-5.4"), "gpt-5.4");
|
||
assert_eq!(normalize_codex_model("gpt-5.2-codex"), "gpt-5.2-codex");
|
||
assert_eq!(normalize_codex_model("o3"), "o3");
|
||
assert_eq!(normalize_codex_model("deepseek-chat"), "deepseek-chat");
|
||
}
|
||
|
||
#[test]
|
||
fn test_normalize_codex_model_combined() {
|
||
// prefix + uppercase + ISO date
|
||
assert_eq!(
|
||
normalize_codex_model("openai/GPT-5.4-2026-03-05"),
|
||
"gpt-5.4"
|
||
);
|
||
// prefix + compact date
|
||
assert_eq!(normalize_codex_model("openai/gpt-5.4-20260305"), "gpt-5.4");
|
||
}
|
||
|
||
#[test]
|
||
fn test_cached_clamped_to_input() {
|
||
// cached > input 的异常场景应被 min() 钳制
|
||
let prev = Some(CumulativeTokens {
|
||
input: 100,
|
||
cached_input: 0,
|
||
output: 50,
|
||
});
|
||
let current = CumulativeTokens {
|
||
input: 110, // delta = 10
|
||
cached_input: 80, // delta = 80(异常:大于 input delta)
|
||
output: 60,
|
||
};
|
||
let delta = compute_delta(&prev, ¤t);
|
||
// 钳制前:cached_input = 80, input = 10
|
||
assert_eq!(delta.cached_input, 80);
|
||
assert_eq!(delta.input, 10);
|
||
// 实际钳制在调用侧:delta.cached_input.min(delta.input)
|
||
let clamped = delta.cached_input.min(delta.input);
|
||
assert_eq!(clamped, 10);
|
||
}
|
||
}
|