diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index b9ae43f..f6f75b5 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -39,7 +39,6 @@ permissions: # If there's a prerelease-style suffix to the version, then the release(s) # will be marked as a prerelease. on: - pull_request: push: tags: - '**[0-9]+.[0-9]+.[0-9]+*' diff --git a/Cargo.lock b/Cargo.lock index 3899396..ffdd553 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2421,6 +2421,7 @@ dependencies = [ "anyhow", "base64", "chrono", + "regex", "rusqlite", "serde", "serde_json", diff --git a/README.md b/README.md index 78ed69b..809b920 100644 --- a/README.md +++ b/README.md @@ -160,7 +160,7 @@ dashboard. Raw usage events and complete archived conversation records stay local and are never included in hosted sync. StatsAI does not upload full prompts, full responses, or raw provider logs. When hosted task sync is explicitly enabled -with `statsai sync --include-tasks`, the current `sync_batch.v3` payload may +with `statsai sync --include-tasks`, the current `sync_batch.v5` payload may include bounded conversation-derived task titles, summary previews, and todo excerpts. You can inspect the exact sync contract in [`docs/sync-contract.md`](docs/sync-contract.md) and verify the resolved sync diff --git a/crates/statsai-adapters/Cargo.toml b/crates/statsai-adapters/Cargo.toml index 885b813..7ed3bac 100644 --- a/crates/statsai-adapters/Cargo.toml +++ b/crates/statsai-adapters/Cargo.toml @@ -20,6 +20,7 @@ anyhow.workspace = true base64.workspace = true chrono.workspace = true rusqlite.workspace = true +regex.workspace = true serde.workspace = true serde_json.workspace = true url.workspace = true diff --git a/crates/statsai-adapters/src/lib.rs b/crates/statsai-adapters/src/lib.rs index 0a071ab..fdb48c5 100644 --- a/crates/statsai-adapters/src/lib.rs +++ b/crates/statsai-adapters/src/lib.rs @@ -2,22 +2,29 @@ use anyhow::{Context, Result}; use chrono::{DateTime, NaiveDate, TimeZone, Utc}; +use regex::Regex; use rusqlite::{Connection, OpenFlags}; use serde::Deserialize; use serde_json::Value; use statsai_core::{ - branch_family, canonical_display, display_path, expand_home_path, extract_issue_keys, - hash_text, home_dir, normalize_git_remote, normalize_task_title, path_hash, project_bucket_key, - semantic_event_id, summarize_task_text, summary_id, task_preview_from_prompt, task_span_id, + account_identity_observation_id, account_plan_observation_id, branch_family, canonical_display, + conversation_account_binding_id, display_path, expand_home_path, extract_issue_keys, hash_text, + home_dir, normalize_email, normalize_git_remote, normalize_plan_name, normalize_task_title, + path_hash, project_bucket_key, provider_account_id_from_identity, semantic_event_id, + summarize_task_text, summary_id, task_preview_from_prompt, task_span_id, task_title_from_prompt, task_title_is_generic, task_title_is_weak_signal, - task_title_signal_score, title_topic_tokens, BillingPeriod, Confidence, CostAccumulator, - CostInfo, EventId, EventSource, IdentitySource, LatencySource, LocationOrigin, MetricStats, - ModelInfo, ParseEvidence, PrivacyInfo, PrivacyMode, ProjectInfo, QuotaCreditsV1, - QuotaObservationRecordV1, QuotaObservationV1, QuotaStatusV1, QuotaUsageLinkKind, - QuotaWindowObservationV1, ReasoningLevel, RuntimeInfo, SessionInfo, SourceKind, SourceLocation, - SubscriptionStatus, SummaryMetadata, SummaryMetrics, TaskSpan, UsageCounts, UsageEvent, - UsageSummary, QUOTA_OBSERVATION_SCHEMA_VERSION, QUOTA_WINDOW_OBSERVATION_SCHEMA_VERSION, - TASK_SPAN_SCHEMA_VERSION, USAGE_EVENT_SCHEMA_VERSION, USAGE_SUMMARY_SCHEMA_VERSION, + task_title_signal_score, title_topic_tokens, AccountEvidenceCheckpointV1, AccountEvidenceKind, + AccountIdentityObservationV1, AccountPlanObservationV1, Confidence, + ConversationAccountBindingV1, CostAccumulator, CostInfo, EventId, EventSource, IdentitySource, + LatencySource, LocationOrigin, MetricStats, ModelInfo, ParseEvidence, PrivacyInfo, PrivacyMode, + ProjectInfo, ProviderAccountId, QuotaCreditsV1, QuotaObservationRecordV1, QuotaObservationV1, + QuotaStatusV1, QuotaUsageLinkKind, QuotaWindowObservationV1, ReasoningLevel, RuntimeInfo, + SessionInfo, SourceKind, SourceLocation, SummaryMetadata, SummaryMetrics, TaskSpan, + UsageCounts, UsageEvent, UsageSummary, ACCOUNT_EVIDENCE_CHECKPOINT_SCHEMA_VERSION, + ACCOUNT_IDENTITY_OBSERVATION_SCHEMA_VERSION, ACCOUNT_PLAN_OBSERVATION_SCHEMA_VERSION, + CONVERSATION_ACCOUNT_BINDING_SCHEMA_VERSION, QUOTA_OBSERVATION_SCHEMA_VERSION, + QUOTA_WINDOW_OBSERVATION_SCHEMA_VERSION, TASK_SPAN_SCHEMA_VERSION, USAGE_EVENT_SCHEMA_VERSION, + USAGE_SUMMARY_SCHEMA_VERSION, }; use statsai_pricing::{ estimate_cost_at, normalize_model_name, pricing_changes_between, unknown_cost, @@ -41,7 +48,10 @@ const PROVIDER_RECORD_EVENT_KEY_VERSION: &str = "provider_record_usage_event.v1" const SCAN_CACHE_SIGNATURE_VERSION: &str = "scan-cache.v1"; // Invalidate unchanged-file scan cache entries whenever provider parsing semantics change, // so historical sessions get rescanned for runtime, pricing, and project context updates. -const CODEX_SCAN_CACHE_PARSER_REVISION: &str = "quota-history.v27"; +// session-identity.v28: usage events adopt the session_meta id (the telemetry +// `conversation.id`) as their session identity; cached files must reparse or +// conversation-to-account bindings can never reach previously scanned events. +const CODEX_SCAN_CACHE_PARSER_REVISION: &str = "session-identity.v28"; const CLAUDE_SCAN_CACHE_PARSER_REVISION: &str = "task-spans.v23"; const OPENCODE_SCAN_CACHE_PARSER_REVISION: &str = "task-spans.v15"; const GROK_BUILD_SCAN_CACHE_PARSER_REVISION: &str = "task-spans.v20"; @@ -60,8 +70,86 @@ const CLAUDE_SETTINGS_AUTH_OVERRIDE_KEYS: &[&str] = &[ "CLAUDE_CODE_USE_ANTHROPIC_AWS", ]; const CODEX_TASK_PREVIEW_RAW_BYTES: usize = 24 * 1024; +const CODEX_ACCOUNT_EVIDENCE_PARSER_VERSION: &str = "codex-account-evidence.v2"; pub(crate) const MAX_JSONL_RECORD_BYTES: usize = 16 * 1024 * 1024; +#[derive(Debug, Clone, PartialEq, Eq)] +struct CodexTelemetryCursor { + maximum_row_id: i64, + checkpoint_row_fingerprint: Option, + database_size: i64, + database_modified_nanos: i64, + wal_size: i64, + wal_modified_nanos: i64, +} + +fn codex_telemetry_file_state( + path: &Path, + maximum_row_id: i64, + checkpoint_row_fingerprint: Option, +) -> CodexTelemetryCursor { + let metadata = std::fs::metadata(path).ok(); + let wal_metadata = std::fs::metadata(path.with_extension("sqlite-wal")).ok(); + let modified_nanos = |metadata: Option<&std::fs::Metadata>| { + metadata + .and_then(|value| value.modified().ok()) + .and_then(|value| value.duration_since(UNIX_EPOCH).ok()) + .map_or(0, |value| { + i64::try_from(value.as_nanos()).unwrap_or(i64::MAX) + }) + }; + CodexTelemetryCursor { + maximum_row_id, + checkpoint_row_fingerprint, + database_size: metadata + .as_ref() + .map_or(0, |value| i64::try_from(value.len()).unwrap_or(i64::MAX)), + database_modified_nanos: modified_nanos(metadata.as_ref()), + wal_size: wal_metadata + .as_ref() + .map_or(0, |value| i64::try_from(value.len()).unwrap_or(i64::MAX)), + wal_modified_nanos: modified_nanos(wal_metadata.as_ref()), + } +} + +fn codex_telemetry_cursor(checkpoint: &AccountEvidenceCheckpointV1) -> CodexTelemetryCursor { + CodexTelemetryCursor { + maximum_row_id: checkpoint.maximum_row_id, + checkpoint_row_fingerprint: checkpoint.checkpoint_row_fingerprint.clone(), + database_size: checkpoint.database_size, + database_modified_nanos: checkpoint.database_modified_nanos, + wal_size: checkpoint.wal_size, + wal_modified_nanos: checkpoint.wal_modified_nanos, + } +} + +fn codex_telemetry_checkpoint_row_fingerprint( + connection: &Connection, + row_id: i64, +) -> Option { + if row_id == 0 { + return Some(hash_text("codex-telemetry-checkpoint-row.v1:empty")); + } + let (seconds, nanos, body) = connection + .query_row( + "SELECT ts, ts_nanos, feedback_log_body FROM logs WHERE id = ?1", + [row_id], + |row| { + Ok(( + row.get::<_, i64>(0)?, + row.get::<_, i64>(1)?, + row.get::<_, Option>(2)?, + )) + }, + ) + .ok()?; + Some(hash_text(&format!( + "codex-telemetry-checkpoint-row.v1:{row_id}:{seconds}:{nanos}:{}", + body.as_deref() + .map_or_else(|| "none".to_string(), hash_text,) + ))) +} + #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum BoundedLineRead { Eof, @@ -203,6 +291,111 @@ pub struct AdapterScan { pub verified_source_state: Option, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ObservedProviderAccount { + pub provider_user_id: Option, + pub email: Option, + pub plan_name: Option, + pub observed_at: DateTime, +} + +#[derive(Debug, Clone, Default)] +pub struct AccountEvidenceScan { + pub accounts: Vec, + pub identity_observations: Vec, + pub plan_observations: Vec, + pub conversation_bindings: Vec, + pub checkpoints: Vec, +} + +/// Rewrites collected account references and their account-dependent deterministic IDs. +/// +/// Callers use this with already-known aliases before Store deduplication, and again after an +/// identity upsert when a newly discovered identity changes the canonical account reference. +pub fn remap_account_evidence_account_ids( + evidence: &mut AccountEvidenceScan, + canonical_ids: &HashMap, +) { + for observation in &mut evidence.identity_observations { + if let Some(canonical_id) = observation + .provider_account_id + .as_ref() + .and_then(|account_id| canonical_ids.get(account_id)) + { + observation.provider_account_id = Some(canonical_id.clone()); + } + } + for observation in &mut evidence.plan_observations { + if let Some(canonical_id) = observation + .provider_account_id + .as_ref() + .and_then(|account_id| canonical_ids.get(account_id)) + { + observation.provider_account_id = Some(canonical_id.clone()); + observation.observation_id = account_plan_observation_id( + &observation.source_id, + Some(canonical_id), + &observation.raw_plan_name, + observation.observed_at, + observation.evidence_kind, + ); + } + } + for binding in &mut evidence.conversation_bindings { + if let Some(canonical_id) = canonical_ids.get(&binding.provider_account_id) { + binding.provider_account_id = canonical_id.clone(); + binding.binding_id = conversation_account_binding_id( + &binding.source_id, + &binding.conversation_id_hash, + binding.turn_id_hash.as_deref(), + canonical_id, + ); + } + } +} + +/// Drops discovered accounts whose evidence was already filtered out by the Store. +/// +/// Call this after `Store::retain_unseen_account_evidence` so unrelated usage-file changes do not +/// refresh an unchanged account and mark it pending for cloud sync again. +pub fn retain_accounts_referenced_by_account_evidence( + provider: &str, + canonical_ids: &HashMap, + evidence: &mut AccountEvidenceScan, +) { + let referenced_account_ids = evidence + .identity_observations + .iter() + .filter_map(|observation| observation.provider_account_id.clone()) + .chain( + evidence + .plan_observations + .iter() + .filter_map(|observation| observation.provider_account_id.clone()), + ) + .chain( + evidence + .conversation_bindings + .iter() + .map(|binding| binding.provider_account_id.clone()), + ) + .collect::>(); + evidence.accounts.retain(|observed| { + provider_account_id_from_identity( + provider, + observed.provider_user_id.as_deref(), + observed.email.as_deref(), + ) + .map(|account_id| { + canonical_ids + .get(&account_id) + .cloned() + .unwrap_or(account_id) + }) + .is_some_and(|account_id| referenced_account_ids.contains(&account_id)) + }); +} + pub trait ProviderAdapter: Send + Sync { fn id(&self) -> &'static str; fn version(&self) -> &'static str; @@ -237,6 +430,13 @@ pub trait ProviderAdapter: Send + Sync { ) -> bool { false } + fn collect_account_evidence( + &self, + _source: &SourceLocation, + _checkpoints: &[AccountEvidenceCheckpointV1], + ) -> Result { + Ok(AccountEvidenceScan::default()) + } fn scan(&self, source: &SourceLocation, options: &ScanOptions) -> Result; fn collect_archive( @@ -399,6 +599,17 @@ impl ProviderAdapter for CodexAdapter { .unwrap_or(VerifiedSourceObservation::Unavailable)) } + fn collect_account_evidence( + &self, + source: &SourceLocation, + checkpoints: &[AccountEvidenceCheckpointV1], + ) -> Result { + let Some(root) = source_root_path(source) else { + return Ok(AccountEvidenceScan::default()); + }; + collect_codex_account_evidence(source, &codex_source_root(&root), checkpoints) + } + fn scan(&self, source: &SourceLocation, options: &ScanOptions) -> Result { scan_codex_source(self, source, options) } @@ -2594,7 +2805,12 @@ fn parse_codex_file( let mut current_project: Option = None; let mut current_title: Option = None; let mut current_thread_id: Option = None; - let session_raw = codex_session_id(usage_root, path); + // Seeded from the file path, replaced by the session's own id as soon as + // the `session_meta` line declares one. The embedded id is the same UUID + // Codex reports as `conversation.id` in telemetry, and conversation-to- + // account bindings hash that UUID — a path-derived identity can never meet + // them, which left every binding unable to attribute a single event. + let mut session_raw = codex_session_id(usage_root, path); let mut records = Vec::new(); let mut quota_observation_indices = HashMap::new(); let mut project_cache = ProjectContextCache::new(); @@ -2835,11 +3051,15 @@ fn parse_codex_file( }; if is_codex_session_meta(&value) { + let declared_session_id = value + .pointer("/payload/id") + .and_then(Value::as_str) + .map(ToOwned::to_owned); + if let Some(declared_session_id) = declared_session_id.clone() { + session_raw = declared_session_id; + } if collect_tasks { - current_thread_id = value - .pointer("/payload/id") - .and_then(Value::as_str) - .map(ToOwned::to_owned); + current_thread_id = declared_session_id; let session_id = current_thread_id .clone() .or_else(|| Some(session_raw.clone())); @@ -5534,6 +5754,14 @@ fn open_sqlite_readonly(path: &Path) -> Result { .with_context(|| format!("open sqlite {}", path.display())) } +fn sqlite_table_exists(connection: &Connection, table: &str) -> Result { + Ok(connection.query_row( + "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?1)", + [table], + |row| row.get(0), + )?) +} + fn sqlite_column_exists(connection: &Connection, table: &str, column: &str) -> Result { let mut statement = connection.prepare(&format!("PRAGMA table_info({table})"))?; let mut rows = statement.query([])?; @@ -7627,9 +7855,20 @@ fn claude_profile_timestamp(value: &Value) -> Option> { } } -fn codex_auth_snapshot(root: &Path) -> Option { - let auth_path = root.join("auth.json"); - let value = std::fs::read_to_string(&auth_path).ok()?; +#[derive(Debug)] +struct CodexAuthClaims { + provider_user_id: Option, + email: Option, + plan_type: Option, + auth_mode: Option, + authenticated_at: Option>, + subscription_checked_at: Option>, + active_from: Option>, + active_until: Option>, +} + +fn codex_auth_claims(auth_path: &Path) -> Option { + let value = std::fs::read_to_string(auth_path).ok()?; let value: Value = serde_json::from_str(&value).ok()?; let payload = string_at_any( &value, @@ -7671,206 +7910,960 @@ fn codex_auth_snapshot(root: &Path) -> Option { } let plan_type = auth.and_then(|auth| string_at_any(auth, &["chatgpt_plan_type"])); - let plan_name = plan_type.as_deref().map(display_codex_plan_name); + let auth_mode = string_at_any(&value, &["auth_mode", "authMode"]); let authenticated_at = payload .as_ref() .and_then(|payload| timestamp_at_any(payload, &["auth_time", "iat"])) - .or_else(|| file_modified_at(&auth_path)); + .or_else(|| file_modified_at(auth_path)); // An auth-file mtime or a fresh ID token proves a refreshed local session, not // that the embedded subscription claims were refreshed at the same time. let subscription_checked_at = auth.and_then(|auth| timestamp_at_any(auth, &["chatgpt_subscription_last_checked"])); - let verified_at = subscription_checked_at.or(authenticated_at); - let paid_at = + let active_from = auth.and_then(|auth| timestamp_at_any(auth, &["chatgpt_subscription_active_start"])); - let current_period_ends_at = + let active_until = auth.and_then(|auth| timestamp_at_any(auth, &["chatgpt_subscription_active_until"])); - let subscription = plan_type.as_deref().and_then(|plan_type| { - codex_verified_subscription( - plan_type, - paid_at, - current_period_ends_at, - subscription_checked_at, - ) - }); - - Some(VerifiedSourceState { + Some(CodexAuthClaims { provider_user_id, email, - account_label: None, - plan_name, + plan_type, + auth_mode, authenticated_at, - verified_at, - subscription, + subscription_checked_at, + active_from, + active_until, }) } -fn file_modified_at(path: &Path) -> Option> { - let modified = std::fs::metadata(path).ok()?.modified().ok()?; - Some(DateTime::::from(modified)) -} - -fn string_at_any(value: &Value, keys: &[&str]) -> Option { - keys.iter() - .filter_map(|key| { - if key.starts_with('/') { - value.pointer(key) - } else { - value.get(*key) - } - }) - .find_map(Value::as_str) - .map(str::trim) - .filter(|text| !text.is_empty()) - .map(ToOwned::to_owned) +fn collect_codex_account_evidence( + source: &SourceLocation, + root: &Path, + checkpoints: &[AccountEvidenceCheckpointV1], +) -> Result { + let mut scan = AccountEvidenceScan::default(); + collect_codex_auth_evidence(source, root, &mut scan); + collect_codex_telemetry_evidence(source, root, checkpoints, &mut scan)?; + collect_codex_reset_history_evidence(source, root, &mut scan); + collect_codex_login_evidence(source, root, &mut scan)?; + dedupe_codex_account_evidence(&mut scan); + Ok(scan) } -fn timestamp_at_any(value: &Value, keys: &[&str]) -> Option> { - keys.iter() - .filter_map(|key| { - if key.starts_with('/') { - value.pointer(key) - } else { - value.get(*key) - } +fn collect_codex_auth_evidence( + source: &SourceLocation, + root: &Path, + scan: &mut AccountEvidenceScan, +) { + push_codex_auth_snapshot(source, &root.join("auth.json"), true, scan); + // Multi-account setups keep the other logins next to the live one as + // `auth-