feat(db): rebuild Codex usage on upgrade and via maintenance action

Schema v16 wipes codex_session detail rows, _codex_session rollups and
Codex rollout cursors inside the migration savepoint (cursor deletion
uses pure shape matching so CODEX_HOME drift cannot orphan cursors);
the next session sync re-imports history from source JSONL under the
corrected importer. Fresh installs traverse the same branch as a no-op.

Add a manual "Rebuild Codex usage" maintenance action (single-flight,
hard-fail backup before reset, unconditional refresh notification even
when reimport is empty or fails after reset) with a destructive confirm
dialog, result toast and four-locale strings. Historical proxy-side
duplicate rows are intentionally left untouched; history whose source
JSONL was already deleted cannot be reconstructed.
This commit is contained in:
Jason
2026-07-20 12:18:47 +08:00
parent c9ac6efd69
commit eff1e0ccfc
11 changed files with 258 additions and 1 deletions
+59
View File
@@ -289,6 +289,35 @@ pub async fn sync_session_usage(
.map_err(|error| AppError::Message(format!("会话用量同步任务失败: {error}")))
}
/// Codex reset 成功后,无论重导是否导入新行或返回错误,都必须通知前端刷新。
/// 调用方应只在 reset 成功后调用,避免把未发生的数据变更误报为重建完成。
fn finish_codex_rebuild(
result: Result<crate::services::session_usage::SessionSyncResult, AppError>,
) -> Result<crate::services::session_usage::SessionSyncResult, AppError> {
crate::usage_events::notify_log_recorded();
result
}
/// 备份数据库后,仅重建 Codex session 用量。锁覆盖 backup → reset → import
/// 整个序列,避免后台同步在清理和重导之间插入数据。
#[tauri::command]
pub async fn rebuild_codex_usage(
state: State<'_, AppState>,
) -> Result<crate::services::session_usage::SessionSyncResult, AppError> {
let db = state.db.clone();
let _guard = crate::services::session_usage::session_sync_mutex()
.lock()
.await;
tauri::async_runtime::spawn_blocking(move || {
db.backup_database_file()?;
db.reset_codex_usage()?;
let result = crate::services::session_usage_codex::sync_codex_usage(&db);
finish_codex_rebuild(result)
})
.await
.map_err(|error| AppError::Message(format!("Codex 用量重建任务失败: {error}")))?
}
/// 获取数据来源分布
#[tauri::command]
pub fn get_usage_data_sources(
@@ -308,3 +337,33 @@ pub struct ModelPricingInfo {
pub cache_read_cost_per_million: String,
pub cache_creation_cost_per_million: String,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn codex_rebuild_notifies_when_reimport_is_empty() {
crate::usage_events::take_test_notify_count();
let result = finish_codex_rebuild(Ok(
crate::services::session_usage::SessionSyncResult::default(),
))
.expect("空重导应成功");
assert_eq!(result.imported, 0);
assert_eq!(crate::usage_events::take_test_notify_count(), 1);
}
#[test]
fn codex_rebuild_notifies_when_reimport_fails_after_reset() {
crate::usage_events::take_test_notify_count();
let result = finish_codex_rebuild(Err(AppError::Message(
"synthetic reimport failure".to_string(),
)));
assert!(result.is_err());
assert_eq!(crate::usage_events::take_test_notify_count(), 1);
}
}
+1 -1
View File
@@ -52,7 +52,7 @@ use std::sync::Mutex;
/// 当前 Schema 版本号
/// 每次修改表结构时递增,并在 schema.rs 中添加相应的迁移逻辑
pub(crate) const SCHEMA_VERSION: i32 = 15;
pub(crate) const SCHEMA_VERSION: i32 = 16;
/// 安全地序列化 JSON,避免 unwrap panic
pub(crate) fn to_json_string<T: Serialize>(value: &T) -> Result<String, AppError> {
+53
View File
@@ -506,6 +506,11 @@ impl Database {
Self::migrate_v14_to_v15(conn)?;
Self::set_user_version(conn, 15)?;
}
15 => {
log::info!("迁移数据库从 v15 到 v16(重建 Codex 会话用量)");
Self::migrate_v15_to_v16(conn)?;
Self::set_user_version(conn, 16)?;
}
_ => {
return Err(AppError::Database(format!(
"未知的数据库版本 {version},无法迁移到 {SCHEMA_VERSION}"
@@ -1510,6 +1515,14 @@ impl Database {
Ok(())
}
/// v15 -> v16: remove Codex session rows and cursors so startup sync can
/// rebuild them with fork-history alignment. Must stay connection-level:
/// schema migration already owns the Database connection mutex.
fn migrate_v15_to_v16(conn: &Connection) -> Result<(), AppError> {
let codex_dir = crate::codex_config::get_codex_config_dir();
crate::services::session_usage_codex::reset_codex_usage_on_conn(conn, &codex_dir)
}
/// 插入默认模型定价数据
/// 格式: (model_id, display_name, input, output, cache_read, cache_creation)
/// 注意: model_id 使用短横线格式(如 claude-haiku-4-5),与 API 返回的模型名称标准化后一致
@@ -3025,4 +3038,44 @@ mod tests {
Ok(())
}
#[test]
fn migrate_v15_to_v16_resets_only_codex_session_usage() -> Result<(), AppError> {
let conn = Connection::open_in_memory()?;
Database::create_tables_on_conn(&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');
INSERT INTO session_log_sync
(file_path, last_modified, last_line_offset, last_synced_at)
VALUES
('/old/sessions/rollout-old-00000000-0000-4000-8000-000000000001.jsonl', 1, 1, 1),
('/gemini/tmp/session-123.json', 1, 1, 1);",
)?;
Database::set_user_version(&conn, 15)?;
Database::apply_schema_migrations_on_conn(&conn)?;
assert_eq!(Database::get_user_version(&conn)?, 16);
let counts: (i64, i64, i64, i64) = conn.query_row(
"SELECT
(SELECT COUNT(*) FROM proxy_request_logs WHERE data_source = 'codex_session'),
(SELECT COUNT(*) FROM proxy_request_logs WHERE data_source = 'gemini_session'),
(SELECT COUNT(*) FROM usage_daily_rollups WHERE provider_id = '_codex_session'),
(SELECT COUNT(*) FROM session_log_sync)",
[],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
)?;
assert_eq!(counts, (0, 1, 0, 1));
Ok(())
}
}
+1
View File
@@ -1522,6 +1522,7 @@ pub fn run() {
commands::check_provider_limits,
// Session usage sync
commands::sync_session_usage,
commands::rebuild_codex_usage,
commands::get_usage_data_sources,
// Stream health check
commands::stream_check_provider,