From c537cbc1b8f9589a8b73b7c107f47ab90e9bbc14 Mon Sep 17 00:00:00 2001 From: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 Date: Thu, 16 Jul 2026 15:27:20 -0500 Subject: [PATCH] feat(desktop): add get_agent_usage_series NIP-AM archive backend Adds the Rust backend half of the NIP-AM local agent usage feature (Rev 3 frozen plan): a rebuildable agent_metric_index parsed from archived kind-44200 rows, a pure per-field accounting ladder (agent_usage.rs) computing token/cost deltas with adjacent-cumulative preference and direct-value fallback, and the get_agent_usage_series Tauri command wiring backfill, orphan repair, collection-enabled detection, and A13's hasArchivedEvidence into one series response. commit_archive now indexes kind-44200 rows in the same transaction as the canonical event insert, and upsert_archived_event/gc_orphaned_events return/enforce the new-row and cascade-delete invariants (A5/A6) that persistedAgentMetrics and the index's rebuildability depend on. get_agent_usage_series' SQLite core is split into a plain sync fn so it can be driven directly against an in-memory Connection in tests without a Tauri AppState. Co-authored-by: Will Pfleger Signed-off-by: Will Pfleger --- desktop/scripts/check-file-sizes.mjs | 8 +- desktop/src-tauri/src/archive/agent_usage.rs | 677 ++++++++++++++++++ .../src/archive/agent_usage_tests.rs | 592 +++++++++++++++ desktop/src-tauri/src/archive/metric_store.rs | 562 +++++++++++++++ .../src/archive/metric_store_tests.rs | 579 +++++++++++++++ desktop/src-tauri/src/archive/mod.rs | 100 +++ desktop/src-tauri/src/archive/mod_tests.rs | 224 ++++++ desktop/src-tauri/src/archive/pipeline.rs | 37 +- desktop/src-tauri/src/archive/store.rs | 118 ++- desktop/src-tauri/src/archive/store_tests.rs | 10 +- desktop/src-tauri/src/lib.rs | 1 + 11 files changed, 2880 insertions(+), 28 deletions(-) create mode 100644 desktop/src-tauri/src/archive/agent_usage.rs create mode 100644 desktop/src-tauri/src/archive/agent_usage_tests.rs create mode 100644 desktop/src-tauri/src/archive/metric_store.rs create mode 100644 desktop/src-tauri/src/archive/metric_store_tests.rs diff --git a/desktop/scripts/check-file-sizes.mjs b/desktop/scripts/check-file-sizes.mjs index 1963acc3b..b302f7c3f 100644 --- a/desktop/scripts/check-file-sizes.mjs +++ b/desktop/scripts/check-file-sizes.mjs @@ -75,7 +75,13 @@ const overrides = new Map([ // is test-only content; the override covers the test growth accumulated // across the local-archive + agent-metric-archive PR series. store_tests.rs // (~731 lines) is under 1000 so needs no override. - ["src-tauri/src/archive/mod_tests.rs", 1208], + // agent-usage-archive (Rev 3): persistedAgentMetrics decrypt-success + // assertion + re-ingest no-double-count test, plus three + // get_agent_usage_series integration tests (fresh ingest, agentPubkey + // filter + hasArchivedEvidence, pre-existing-row backfill) exercising the + // command's SQLite core end to end. Test-only content; ratcheted + // 1208 -> 1424. + ["src-tauri/src/archive/mod_tests.rs", 1424], // unified-agent-model 1A.1: profile reconcile split to agents_profile.rs, // ratcheting 1443 -> 1295. Queued to split further in the A2 fold. // global-agent-config: resolve_deploy_model_provider + visibility exports diff --git a/desktop/src-tauri/src/archive/agent_usage.rs b/desktop/src-tauri/src/archive/agent_usage.rs new file mode 100644 index 000000000..431c068a1 --- /dev/null +++ b/desktop/src-tauri/src/archive/agent_usage.rs @@ -0,0 +1,677 @@ +//! Pure NIP-AM usage accounting: request validation, wire types, and the +//! per-field cumulative/direct-fallback ladder. +//! +//! No Tauri or filesystem dependency — every function here takes already +//! loaded rows (or plain values) and returns plain values. The caller +//! (`mod.rs`'s `get_agent_usage_series` command, Phase 2) owns SQLite access +//! (`metric_store.rs`) and glues the two together: load window rows, compute +//! [`window_probe_keys`], load those exact-key rows, then call +//! [`compute_series`]. +//! +//! Accounting contract (Rev 3, frozen plan + amendments A1/A2/A4/A9/A11–A13): +//! see `docs/nips/NIP-AM.md:119-160` for the NIP itself. + +use std::collections::{HashMap, HashSet}; + +use serde::{Deserialize, Serialize}; + +use super::metric_store::AgentMetricIndexRow; + +// ── Request ────────────────────────────────────────────────────────────────── + +/// Request for [`compute_series`]'s caller. `bucket_boundaries` are exact +/// local-midnight Unix-second boundaries built by the frontend (inclusive +/// start / exclusive end per adjacent pair): 8 entries = 7 buckets, 31 +/// entries = 30 buckets. +#[derive(Debug, Clone, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct AgentUsageSeriesRequest { + pub bucket_boundaries: Vec, + pub agent_pubkey: Option, +} + +/// Widest interval NIP-AM query validation admits (A9): wide enough to admit +/// every real civil-day transition (ordinary DST, 30-minute-offset zones, +/// historical calendar skips) while still rejecting arbitrary bins. +const MAX_INTERVAL_SECS: i64 = 48 * 3600; + +/// Validate a request per the frozen contract + A9 (drops the 23–25h band +/// for a `> 0 && <= 48h` sanity band) and A13 pubkey normalization. +/// +/// Fails closed before any SQLite work. Returns the normalized (lowercased) +/// agent pubkey, if one was supplied. +pub(super) fn validate_request(req: &AgentUsageSeriesRequest) -> Result, String> { + let n = req.bucket_boundaries.len(); + if n != 8 && n != 31 { + return Err(format!( + "bucket_boundaries must have exactly 8 or 31 entries, got {n}" + )); + } + + for i in 0..n - 1 { + let (a, b) = (req.bucket_boundaries[i], req.bucket_boundaries[i + 1]); + if b <= a { + return Err(format!( + "bucket_boundaries must be strictly increasing (index {i}: {a} >= {b})" + )); + } + let interval = b - a; + if interval > MAX_INTERVAL_SECS { + return Err(format!( + "bucket_boundaries interval at index {i} is {interval}s, exceeds {MAX_INTERVAL_SECS}s" + )); + } + } + + // Finite Unix-second bounds: must be representable as an RFC 3339 + // instant (chrono's timestamp range), independent of local timezone. + for &t in &req.bucket_boundaries { + if chrono::DateTime::from_timestamp(t, 0).is_none() { + return Err(format!("bucket boundary {t} is out of representable range")); + } + } + + let normalized_pubkey = match &req.agent_pubkey { + None => None, + Some(pk) => { + if pk.len() != 64 || !pk.chars().all(|c| c.is_ascii_hexdigit()) { + return Err("agent_pubkey must be exactly 64 hex characters".to_string()); + } + Some(pk.to_lowercase()) + } + }; + + Ok(normalized_pubkey) +} + +// ── Wire types ─────────────────────────────────────────────────────────────── + +/// Per-field completeness (A2): `value: null` means no known increment in +/// scope; `incomplete: true` on a non-null value means the value is a +/// reported lower bound, not full coverage. +#[derive(Debug, Clone, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct UsageField { + pub value: Option, + pub incomplete: bool, +} + +#[derive(Debug, Clone, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct CostField { + pub value: Option, + pub incomplete: bool, +} + +#[derive(Debug, Clone, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct ReportedUsage { + pub input_tokens: UsageField, + pub output_tokens: UsageField, + pub total_tokens: UsageField, + pub estimated_cost_usd: CostField, +} + +#[derive(Debug, Clone, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct SeriesBucket { + pub start: i64, + pub end: i64, + pub usage: ReportedUsage, + pub report_count: i64, + pub has_unknown_usage: bool, +} + +#[derive(Debug, Clone, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct ModelUsage { + pub model: Option, + pub usage: ReportedUsage, + pub report_count: i64, + pub has_unknown_usage: bool, +} + +#[derive(Debug, Clone, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct AgentUsage { + pub agent_pubkey: String, + pub usage: ReportedUsage, + pub buckets: Vec, + pub models: Vec, + pub report_count: i64, + pub has_unknown_usage: bool, +} + +#[derive(Debug, Clone, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct Coverage { + pub first_archived_at: Option, + pub last_archived_at: Option, + pub first_reported_at: Option, + pub last_reported_at: Option, + pub report_count: i64, + pub invalid_report_count: i64, + pub has_unknown_usage: bool, +} + +#[derive(Debug, Clone, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct AgentUsageSeries { + pub collection_enabled: bool, + pub buckets: Vec, + pub agents: Vec, + pub coverage: Coverage, + /// A13: `null` when the request had no `agentPubkey` filter; otherwise + /// `true` iff at least one surviving `agent_metric_index` row (either + /// `parse_status`) exists for that author, with no bucket-boundary + /// restriction. Callers compute this (DB access) and pass it through. + pub has_archived_evidence: Option, +} + +// ── Per-event field ladder (A1, A4, A11, A12) ─────────────────────────────── + +#[derive(Debug, Clone, Copy)] +enum FieldValue { + Known(T), + Unknown, +} + +struct EventOutcome { + input: FieldValue, + output: FieldValue, + total: FieldValue, + cost: FieldValue, +} + +/// Per-field ladder for token counters (A1): adjacent nondecreasing cumulative +/// pair → diff; adjacent decreasing pair → unknown, terminal, no fallback; +/// no usable baseline for this field → `deltaReliable` direct value or +/// unknown. +fn ladder_token( + baseline_cumulative: Option, + current_cumulative: Option, + current_turn: Option, + delta_reliable: bool, +) -> FieldValue { + if let (Some(prev), Some(cur)) = (baseline_cumulative, current_cumulative) { + return if cur >= prev { + FieldValue::Known(cur - prev) + } else { + FieldValue::Unknown + }; + } + if delta_reliable { + if let Some(v) = current_turn { + return FieldValue::Known(v); + } + } + FieldValue::Unknown +} + +/// Same ladder for `costUsd` (f64). +fn ladder_cost( + baseline_cumulative: Option, + current_cumulative: Option, + current_turn: Option, + delta_reliable: bool, +) -> FieldValue { + if let (Some(prev), Some(cur)) = (baseline_cumulative, current_cumulative) { + return if cur >= prev { + FieldValue::Known(cur - prev) + } else { + FieldValue::Unknown + }; + } + if delta_reliable { + if let Some(v) = current_turn { + return FieldValue::Known(v); + } + } + FieldValue::Unknown +} + +/// Resolve the exact-`S-1` baseline row for `row`, or `None` if no usable +/// baseline exists (A11/A12): missing session/sequence key, duplicate row at +/// `row`'s own sequence (A4 — no cumulative delta for any row at a +/// duplicated sequence), sequence `0` (no predecessor, `checked_sub` +/// underflow), no predecessor row (gap), or duplicate rows at the +/// predecessor sequence (ambiguous baseline). +/// +/// `probe_by_key` must contain, for every key queried, ALL valid rows at +/// that exact `(agent, session, turnSeq)` — used for both this baseline +/// lookup and the A4/A11 duplicate-cardinality check, which is why absence +/// of a key from the map is treated identically to an empty group. +fn resolve_baseline<'a>( + row: &AgentMetricIndexRow, + probe_by_key: &HashMap<(String, String, u64), Vec<&'a AgentMetricIndexRow>>, +) -> Option<&'a AgentMetricIndexRow> { + let (agent, session, seq) = row.accounting_key()?; + + let own_group = probe_by_key.get(&(agent.clone(), session.clone(), seq))?; + if own_group.len() > 1 { + return None; // A4: duplicate at own sequence — no cumulative delta. + } + + let pred_seq = seq.checked_sub(1)?; // A12: seq == 0 has no baseline. + let pred_group = probe_by_key.get(&(agent, session, pred_seq))?; + if pred_group.len() != 1 { + return None; // Missing (gap) or ambiguous (duplicate) predecessor. + } + Some(pred_group[0]) +} + +fn compute_event_outcome( + row: &AgentMetricIndexRow, + probe_by_key: &HashMap<(String, String, u64), Vec<&AgentMetricIndexRow>>, +) -> EventOutcome { + let baseline = resolve_baseline(row, probe_by_key); + let delta_reliable = row.delta_reliable.unwrap_or(false); + + EventOutcome { + input: ladder_token( + baseline.and_then(|b| b.cumulative_input_tokens), + row.cumulative_input_tokens, + row.turn_input_tokens, + delta_reliable, + ), + output: ladder_token( + baseline.and_then(|b| b.cumulative_output_tokens), + row.cumulative_output_tokens, + row.turn_output_tokens, + delta_reliable, + ), + total: ladder_token( + baseline.and_then(|b| b.cumulative_total_tokens), + row.cumulative_total_tokens, + row.turn_total_tokens, + delta_reliable, + ), + cost: ladder_cost( + baseline.and_then(|b| b.cumulative_cost_usd), + row.cumulative_cost_usd, + row.turn_cost_usd, + delta_reliable, + ), + } +} + +/// The exact `(agent, session, turnSeq)` keys the caller must load via +/// `metric_store::load_rows_at_exact_keys` before calling [`compute_series`] +/// (A11): each in-window row's own key, plus its checked predecessor key +/// (`turnSeq - 1`) when one exists. Pure and DB-free so it is unit-testable +/// without SQLite. +pub(super) fn window_probe_keys( + window_rows: &[AgentMetricIndexRow], +) -> HashSet<(String, String, u64)> { + let mut keys = HashSet::new(); + for row in window_rows { + if let Some((agent, session, seq)) = row.accounting_key() { + if let Some(pred) = seq.checked_sub(1) { + keys.insert((agent.clone(), session.clone(), pred)); + } + keys.insert((agent, session, seq)); + } + } + keys +} + +// ── Accumulators ───────────────────────────────────────────────────────────── + +/// Sums known per-event increments with `checked_add`; an event with an +/// unknown value, or an overflow, marks the scope `incomplete` (A2) without +/// ever wrapping (overflow freezes the sum at its last valid value). +#[derive(Debug, Default, Clone)] +struct TokenAccumulator { + value: Option, + incomplete: bool, + overflowed: bool, +} + +impl TokenAccumulator { + fn add(&mut self, v: FieldValue) { + match v { + FieldValue::Unknown => self.incomplete = true, + FieldValue::Known(x) => { + if self.overflowed { + self.incomplete = true; + return; + } + self.value = Some(match self.value { + None => x, + Some(cur) => match cur.checked_add(x) { + Some(sum) => sum, + None => { + self.overflowed = true; + self.incomplete = true; + cur + } + }, + }); + } + } + } + + fn has_unknown(&self) -> bool { + self.incomplete + } + + fn finish(self) -> UsageField { + UsageField { + value: self.value.map(|v| v.to_string()), + incomplete: self.incomplete, + } + } +} + +/// Same contract as [`TokenAccumulator`] for `f64` costs: "checked finite +/// addition" means a sum that would become non-finite is rejected and the +/// scope freezes at its last valid value, flagged incomplete. +#[derive(Debug, Default, Clone)] +struct CostAccumulator { + value: Option, + incomplete: bool, + overflowed: bool, +} + +impl CostAccumulator { + fn add(&mut self, v: FieldValue) { + match v { + FieldValue::Unknown => self.incomplete = true, + FieldValue::Known(x) => { + if self.overflowed { + self.incomplete = true; + return; + } + let candidate = match self.value { + None => x, + Some(cur) => cur + x, + }; + if candidate.is_finite() { + self.value = Some(candidate); + } else { + self.overflowed = true; + self.incomplete = true; + } + } + } + } + + fn has_unknown(&self) -> bool { + self.incomplete + } + + fn finish(self) -> CostField { + CostField { + value: self.value, + incomplete: self.incomplete, + } + } +} + +#[derive(Debug, Default, Clone)] +struct UsageAccumulator { + input: TokenAccumulator, + output: TokenAccumulator, + total: TokenAccumulator, + cost: CostAccumulator, +} + +impl UsageAccumulator { + fn add(&mut self, outcome: &EventOutcome) { + self.input.add(outcome.input); + self.output.add(outcome.output); + self.total.add(outcome.total); + self.cost.add(outcome.cost); + } + + fn has_unknown(&self) -> bool { + self.input.has_unknown() + || self.output.has_unknown() + || self.total.has_unknown() + || self.cost.has_unknown() + } + + /// Raw total-token value, used only for the A2 ranking rule (sort by + /// known `totalTokens`, unknown-total scopes listed after). + fn total_tokens_value(&self) -> Option { + self.total.value + } + + fn finish(self) -> ReportedUsage { + ReportedUsage { + input_tokens: self.input.finish(), + output_tokens: self.output.finish(), + total_tokens: self.total.finish(), + estimated_cost_usd: self.cost.finish(), + } + } +} + +// ── Bucket assignment ──────────────────────────────────────────────────────── + +/// Which `[boundaries[i], boundaries[i+1])` bucket `t` falls in, or `None` +/// if outside every bucket (defensive; callers scope their row query to +/// `[boundaries[0], boundaries[last])` so this should never miss). +fn assign_bucket_index(boundaries: &[i64], t: i64) -> Option { + (0..boundaries.len().saturating_sub(1)).find(|&i| t >= boundaries[i] && t < boundaries[i + 1]) +} + +// ── compute_series ─────────────────────────────────────────────────────────── + +/// Per-agent accumulation scope, built while walking `window_rows` once. +struct AgentScope { + buckets: Vec, + bucket_counts: Vec, + total: UsageAccumulator, + report_count: i64, + models: HashMap, (UsageAccumulator, i64)>, +} + +/// Compute the full [`AgentUsageSeries`] from already-loaded rows. +/// +/// - `window_rows`: valid rows with `reported_at` in `[boundaries[0], +/// boundaries[last])` (optionally pre-filtered to one agent), from +/// `metric_store::load_window_valid_rows`. +/// - `probe_rows`: valid rows at the exact keys from [`window_probe_keys`], +/// from `metric_store::load_rows_at_exact_keys` — used for baseline +/// resolution and A4/A11 duplicate-cardinality checks, unrestricted by the +/// window. +/// - `invalid_report_count`: from `metric_store::count_invalid_rows_in_window`. +/// - `has_archived_evidence`: from `metric_store::has_archived_evidence`, +/// already resolved to `None` when the request has no `agentPubkey` filter +/// (A13) — this function does not decide that; it only carries the value. +pub(super) fn compute_series( + window_rows: &[AgentMetricIndexRow], + probe_rows: &[AgentMetricIndexRow], + invalid_report_count: i64, + boundaries: &[i64], + has_archived_evidence: Option, + collection_enabled: bool, +) -> AgentUsageSeries { + let bucket_count = boundaries.len().saturating_sub(1); + + let mut probe_by_key: HashMap<(String, String, u64), Vec<&AgentMetricIndexRow>> = + HashMap::new(); + for r in probe_rows { + if let Some(key) = r.accounting_key() { + probe_by_key.entry(key).or_default().push(r); + } + } + + let mut overall_buckets: Vec = (0..bucket_count) + .map(|_| UsageAccumulator::default()) + .collect(); + let mut overall_bucket_counts: Vec = vec![0; bucket_count]; + + let mut agents: HashMap = HashMap::new(); + + let mut first_reported_at: Option = None; + let mut last_reported_at: Option = None; + let mut first_archived_at: Option = None; + let mut last_archived_at: Option = None; + + for row in window_rows { + // Defensive: the loader already scopes to reported_at in-window and + // parse_status = 'valid'; a miss here means the caller passed rows + // it should not have, so skip rather than panic or miscount. + let Some(reported_at) = row.reported_at else { + continue; + }; + let Some(bucket_idx) = assign_bucket_index(boundaries, reported_at) else { + continue; + }; + + first_reported_at = Some(first_reported_at.map_or(reported_at, |v| v.min(reported_at))); + last_reported_at = Some(last_reported_at.map_or(reported_at, |v| v.max(reported_at))); + first_archived_at = + Some(first_archived_at.map_or(row.archived_at, |v| v.min(row.archived_at))); + last_archived_at = + Some(last_archived_at.map_or(row.archived_at, |v| v.max(row.archived_at))); + + let outcome = compute_event_outcome(row, &probe_by_key); + + overall_buckets[bucket_idx].add(&outcome); + overall_bucket_counts[bucket_idx] += 1; + + let scope = agents + .entry(row.agent_pubkey.clone()) + .or_insert_with(|| AgentScope { + buckets: (0..bucket_count) + .map(|_| UsageAccumulator::default()) + .collect(), + bucket_counts: vec![0; bucket_count], + total: UsageAccumulator::default(), + report_count: 0, + models: HashMap::new(), + }); + scope.buckets[bucket_idx].add(&outcome); + scope.bucket_counts[bucket_idx] += 1; + scope.total.add(&outcome); + scope.report_count += 1; + + let model_entry = scope + .models + .entry(row.model.clone()) + .or_insert_with(|| (UsageAccumulator::default(), 0i64)); + model_entry.0.add(&outcome); + model_entry.1 += 1; + } + + let overall_report_count: i64 = overall_bucket_counts.iter().sum(); + + let buckets: Vec = overall_buckets + .into_iter() + .zip(overall_bucket_counts) + .enumerate() + .map(|(i, (acc, count))| { + let has_unknown_usage = acc.has_unknown(); + SeriesBucket { + start: boundaries[i], + end: boundaries[i + 1], + usage: acc.finish(), + report_count: count, + has_unknown_usage, + } + }) + .collect(); + let any_overall_bucket_unknown = buckets.iter().any(|b| b.has_unknown_usage); + + // Build agent rows, then apply the A2 ranking rule: known totalTokens + // descending, unknown-total agents after, pubkey as the tiebreak/final + // key in both groups for determinism. + let mut agent_rows: Vec<(Option, String, AgentUsage)> = agents + .into_iter() + .map(|(agent_pubkey, scope)| { + let total_tokens_value = scope.total.total_tokens_value(); + let has_unknown_usage = scope.total.has_unknown(); + + let buckets: Vec = scope + .buckets + .into_iter() + .zip(scope.bucket_counts) + .enumerate() + .map(|(i, (acc, count))| { + let has_unknown_usage = acc.has_unknown(); + SeriesBucket { + start: boundaries[i], + end: boundaries[i + 1], + usage: acc.finish(), + report_count: count, + has_unknown_usage, + } + }) + .collect(); + + let mut model_rows: Vec<(Option, Option, ModelUsage)> = scope + .models + .into_iter() + .map(|(model, (acc, count))| { + let model_total = acc.total_tokens_value(); + let has_unknown_usage = acc.has_unknown(); + ( + model_total, + model.clone(), + ModelUsage { + model, + usage: acc.finish(), + report_count: count, + has_unknown_usage, + }, + ) + }) + .collect(); + model_rows.sort_by(|a, b| match (a.0, b.0) { + (Some(av), Some(bv)) => bv.cmp(&av).then_with(|| a.1.cmp(&b.1)), + (Some(_), None) => std::cmp::Ordering::Less, + (None, Some(_)) => std::cmp::Ordering::Greater, + (None, None) => a.1.cmp(&b.1), + }); + let models = model_rows.into_iter().map(|(_, _, m)| m).collect(); + + ( + total_tokens_value, + agent_pubkey.clone(), + AgentUsage { + agent_pubkey, + usage: scope.total.finish(), + buckets, + models, + report_count: scope.report_count, + has_unknown_usage, + }, + ) + }) + .collect(); + agent_rows.sort_by(|a, b| match (a.0, b.0) { + (Some(av), Some(bv)) => bv.cmp(&av).then_with(|| a.1.cmp(&b.1)), + (Some(_), None) => std::cmp::Ordering::Less, + (None, Some(_)) => std::cmp::Ordering::Greater, + (None, None) => a.1.cmp(&b.1), + }); + let any_agent_unknown = agent_rows.iter().any(|(_, _, a)| a.has_unknown_usage); + let agents: Vec = agent_rows.into_iter().map(|(_, _, a)| a).collect(); + + AgentUsageSeries { + collection_enabled, + buckets, + agents, + coverage: Coverage { + first_archived_at, + last_archived_at, + first_reported_at, + last_reported_at, + report_count: overall_report_count, + invalid_report_count, + has_unknown_usage: any_overall_bucket_unknown + || any_agent_unknown + || invalid_report_count > 0, + }, + has_archived_evidence, + } +} + +// ── Tests ──────────────────────────────────────────────────────────────────── + +#[cfg(test)] +#[path = "agent_usage_tests.rs"] +mod agent_usage_tests; diff --git a/desktop/src-tauri/src/archive/agent_usage_tests.rs b/desktop/src-tauri/src/archive/agent_usage_tests.rs new file mode 100644 index 000000000..f28c99a68 --- /dev/null +++ b/desktop/src-tauri/src/archive/agent_usage_tests.rs @@ -0,0 +1,592 @@ +//! Tests for the pure NIP-AM accounting ladder, accumulators, and request +//! validation. No SQLite — rows are constructed directly. + +use super::*; +use crate::archive::metric_store::ParseStatus; + +// ── Row builder ────────────────────────────────────────────────────────────── + +/// Build a fully-specified valid row for test scenarios. Defaults every +/// optional field to `None`/appropriate zero; callers override what a +/// scenario needs via struct-update syntax. +fn row(id: &str, agent: &str, session: &str, seq: u64, reported_at: i64) -> AgentMetricIndexRow { + AgentMetricIndexRow { + id: id.to_string(), + agent_pubkey: agent.to_string(), + event_created_at: reported_at, + archived_at: reported_at, + reported_at: Some(reported_at), + session_id: Some(session.to_string()), + turn_seq: Some(seq), + model: None, + delta_reliable: Some(true), + turn_input_tokens: None, + turn_output_tokens: None, + turn_total_tokens: None, + turn_cost_usd: None, + cumulative_input_tokens: None, + cumulative_output_tokens: None, + cumulative_total_tokens: None, + cumulative_cost_usd: None, + parse_status: ParseStatus::Valid, + } +} + +/// Standard 8-boundary window covering one day, seconds since epoch. +const DAY: i64 = 86_400; +fn boundaries_7() -> Vec { + (0..=7).map(|i| i * DAY).collect() +} + +fn probe_map( + rows: &[AgentMetricIndexRow], +) -> std::collections::HashMap<(String, String, u64), Vec<&AgentMetricIndexRow>> { + let mut m = std::collections::HashMap::new(); + for r in rows { + if let Some(key) = r.accounting_key() { + m.entry(key).or_insert_with(Vec::new).push(r); + } + } + m +} + +// ── Ladder: direct / cumulative / fallback ────────────────────────────────── + +#[test] +fn direct_turn_value_used_when_no_baseline() { + let r = AgentMetricIndexRow { + turn_input_tokens: Some(100), + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 1, 0) + }; + let __probe_rows = [r.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&r, &probes); + assert!(matches!(outcome.input, FieldValue::Known(100))); +} + +#[test] +fn adjacent_cumulative_preferred_over_direct() { + let prev = AgentMetricIndexRow { + cumulative_input_tokens: Some(1000), + ..row("e0", "agent1", "s1", 1, 0) + }; + let cur = AgentMetricIndexRow { + cumulative_input_tokens: Some(1300), + turn_input_tokens: Some(999), // deliberately wrong direct value + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 2, 10) + }; + let __probe_rows = [prev, cur.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&cur, &probes); + assert!(matches!(outcome.input, FieldValue::Known(300))); +} + +#[test] +fn direct_fallback_when_baseline_missing() { + // seq 5 has no row at seq 4 in the probe set → gap, direct-reliable only. + let cur = AgentMetricIndexRow { + cumulative_input_tokens: Some(500), + turn_input_tokens: Some(42), + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 5, 0) + }; + let __probe_rows = [cur.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&cur, &probes); + assert!(matches!(outcome.input, FieldValue::Known(42))); +} + +#[test] +fn unreliable_delta_with_no_baseline_is_unknown() { + let cur = AgentMetricIndexRow { + turn_input_tokens: Some(42), + delta_reliable: Some(false), + ..row("e1", "agent1", "s1", 5, 0) + }; + let __probe_rows = [cur.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&cur, &probes); + assert!(matches!(outcome.input, FieldValue::Unknown)); +} + +#[test] +fn sequence_gap_no_diff_but_direct_reliable_may_count() { + // predecessor exists at seq 1 but current is seq 3 (gap at seq 2) — + // predecessor lookup requires exact S-1, so seq 3's predecessor probe is + // for seq 2, which is absent. Direct fallback applies. + let baseline = AgentMetricIndexRow { + cumulative_input_tokens: Some(100), + ..row("e0", "agent1", "s1", 1, 0) + }; + let cur = AgentMetricIndexRow { + cumulative_input_tokens: Some(400), + turn_input_tokens: Some(50), + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 3, 10) + }; + let __probe_rows = [baseline, cur.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&cur, &probes); + assert!(matches!(outcome.input, FieldValue::Known(50))); +} + +// ── A1: counter decrease is terminal, no fallback ─────────────────────────── + +#[test] +fn adjacent_decrease_with_reliable_direct_present_is_unknown() { + // Required A1 test: decreasing cumulative pair + deltaReliable true + + // direct turn value present → field unknown, NOT the direct value. + let prev = AgentMetricIndexRow { + cumulative_input_tokens: Some(1000), + ..row("e0", "agent1", "s1", 1, 0) + }; + let cur = AgentMetricIndexRow { + cumulative_input_tokens: Some(600), // decreased + turn_input_tokens: Some(77), // present, but must NOT be used + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 2, 10) + }; + let __probe_rows = [prev, cur.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&cur, &probes); + assert!(matches!(outcome.input, FieldValue::Unknown)); +} + +#[test] +fn decrease_taints_only_the_affected_field() { + // Interpretation note in A1: a decrease on one field must not zero out + // sibling fields whose own adjacent pair is nondecreasing. + let prev = AgentMetricIndexRow { + cumulative_input_tokens: Some(1000), + cumulative_output_tokens: Some(200), + ..row("e0", "agent1", "s1", 1, 0) + }; + let cur = AgentMetricIndexRow { + cumulative_input_tokens: Some(600), // decreased + cumulative_output_tokens: Some(250), // increased, still valid + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 2, 10) + }; + let __probe_rows = [prev, cur.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&cur, &probes); + assert!(matches!(outcome.input, FieldValue::Unknown)); + assert!(matches!(outcome.output, FieldValue::Known(50))); +} + +// ── A4/A11: duplicate sequence quarantine ─────────────────────────────────── + +#[test] +fn duplicate_at_sequence_quarantines_successor() { + // Two rows at N, one at N+1: N+1 must not cumulative-diff against + // either N candidate. + let dup_a = AgentMetricIndexRow { + cumulative_input_tokens: Some(100), + ..row("e0a", "agent1", "s1", 5, 0) + }; + let dup_b = AgentMetricIndexRow { + cumulative_input_tokens: Some(150), + ..row("e0b", "agent1", "s1", 5, 1) + }; + let successor = AgentMetricIndexRow { + cumulative_input_tokens: Some(400), + turn_input_tokens: Some(30), + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 6, 10) + }; + let __probe_rows = [dup_a, dup_b, successor.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&successor, &probes); + // No usable baseline (ambiguous predecessor) → direct-reliable fallback. + assert!(matches!(outcome.input, FieldValue::Known(30))); +} + +#[test] +fn duplicate_row_itself_gets_no_cumulative_delta() { + let baseline = AgentMetricIndexRow { + cumulative_input_tokens: Some(100), + ..row("e_base", "agent1", "s1", 4, 0) + }; + let dup_a = AgentMetricIndexRow { + cumulative_input_tokens: Some(200), + turn_input_tokens: Some(99), + delta_reliable: Some(true), + ..row("e5a", "agent1", "s1", 5, 1) + }; + let dup_b = AgentMetricIndexRow { + cumulative_input_tokens: Some(250), + ..row("e5b", "agent1", "s1", 5, 2) + }; + let __probe_rows = [baseline, dup_a.clone(), dup_b]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&dup_a, &probes); + // Own sequence has >1 row → no cumulative delta; direct-reliable value + // for dup_a specifically still counts (A4: "only an independently + // reliable direct turn value may count"). + assert!(matches!(outcome.input, FieldValue::Known(99))); +} + +#[test] +fn duplicate_out_of_window_peer_still_quarantines_in_window_row() { + // A11: an out-of-window duplicate at the same sequence still poisons the + // in-window row's cumulative eligibility, because the probe set has no + // reported_at restriction. + let out_of_window_dup = AgentMetricIndexRow { + cumulative_input_tokens: Some(500), + ..row("e_old", "agent1", "s1", 5, -1000) + }; + let baseline = AgentMetricIndexRow { + cumulative_input_tokens: Some(100), + ..row("e_base", "agent1", "s1", 4, 0) + }; + let in_window = AgentMetricIndexRow { + cumulative_input_tokens: Some(200), + turn_input_tokens: Some(11), + delta_reliable: Some(true), + ..row("e5", "agent1", "s1", 5, 1) + }; + let __probe_rows = [out_of_window_dup, baseline, in_window.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&in_window, &probes); + assert!(matches!(outcome.input, FieldValue::Known(11))); +} + +// ── A12: checked_sub sequence arithmetic ──────────────────────────────────── + +#[test] +fn adjacent_pair_at_u64_max_computes_normally() { + let prev = AgentMetricIndexRow { + cumulative_input_tokens: Some(1000), + ..row("e_prev", "agent1", "s1", u64::MAX - 1, 0) + }; + let cur = AgentMetricIndexRow { + cumulative_input_tokens: Some(1500), + ..row("e_max", "agent1", "s1", u64::MAX, 10) + }; + let __probe_rows = [prev, cur.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&cur, &probes); + assert!(matches!(outcome.input, FieldValue::Known(500))); +} + +#[test] +fn duplicate_at_u64_max_needs_no_successor_probe() { + // u64::MAX has no successor sequence to quarantine — verify duplicate + // handling at MAX itself still works (own-sequence cardinality check). + let dup_a = AgentMetricIndexRow { + cumulative_input_tokens: Some(100), + turn_input_tokens: Some(5), + delta_reliable: Some(true), + ..row("e_max_a", "agent1", "s1", u64::MAX, 0) + }; + let dup_b = AgentMetricIndexRow { + cumulative_input_tokens: Some(150), + ..row("e_max_b", "agent1", "s1", u64::MAX, 1) + }; + let __probe_rows = [dup_a.clone(), dup_b]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&dup_a, &probes); + assert!(matches!(outcome.input, FieldValue::Known(5))); +} + +#[test] +fn seq_zero_has_no_baseline_underflow() { + let cur = AgentMetricIndexRow { + cumulative_input_tokens: Some(100), + turn_input_tokens: Some(100), + delta_reliable: Some(true), + ..row("e0", "agent1", "s1", 0, 0) + }; + let __probe_rows = [cur.clone()]; + let probes = probe_map(&__probe_rows); + // Must not panic (checked_sub) and must fall back to direct. + let outcome = compute_event_outcome(&cur, &probes); + assert!(matches!(outcome.input, FieldValue::Known(100))); +} + +// ── Null total / per-field independence ───────────────────────────────────── + +#[test] +fn null_total_tokens_stays_unknown_even_when_input_output_known() { + let r = AgentMetricIndexRow { + turn_input_tokens: Some(10), + turn_output_tokens: Some(20), + turn_total_tokens: None, + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 1, 0) + }; + let __probe_rows = [r.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&r, &probes); + assert!(matches!(outcome.input, FieldValue::Known(10))); + assert!(matches!(outcome.output, FieldValue::Known(20))); + assert!(matches!(outcome.total, FieldValue::Unknown)); +} + +// ── Cross-session / cross-agent isolation ─────────────────────────────────── + +#[test] +fn cumulative_diff_never_crosses_session_boundary() { + let other_session = AgentMetricIndexRow { + cumulative_input_tokens: Some(999_999), + ..row("e_other", "agent1", "s_other", 1, 0) + }; + let cur = AgentMetricIndexRow { + cumulative_input_tokens: Some(50), + turn_input_tokens: Some(50), + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 1, 10) + }; + let __probe_rows = [other_session, cur.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&cur, &probes); + // seq 1 has no predecessor (seq 0) in ANY session — direct fallback. + assert!(matches!(outcome.input, FieldValue::Known(50))); +} + +#[test] +fn cumulative_diff_never_crosses_agent_boundary() { + let other_agent = AgentMetricIndexRow { + cumulative_input_tokens: Some(999_999), + ..row("e_other", "agent2", "s1", 1, 0) + }; + let cur = AgentMetricIndexRow { + cumulative_input_tokens: Some(60), + turn_input_tokens: Some(60), + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 2, 10) + }; + let __probe_rows = [other_agent, cur.clone()]; + let probes = probe_map(&__probe_rows); + let outcome = compute_event_outcome(&cur, &probes); + // seq 2's predecessor (seq 1) does not exist for agent1 — direct. + assert!(matches!(outcome.input, FieldValue::Known(60))); +} + +// ── window_probe_keys ──────────────────────────────────────────────────────── + +#[test] +fn window_probe_keys_includes_own_and_predecessor() { + let r = row("e1", "agent1", "s1", 5, 0); + let keys = window_probe_keys(&[r]); + assert!(keys.contains(&("agent1".to_string(), "s1".to_string(), 5))); + assert!(keys.contains(&("agent1".to_string(), "s1".to_string(), 4))); + assert_eq!(keys.len(), 2); +} + +#[test] +fn window_probe_keys_skips_predecessor_at_seq_zero() { + let r = row("e1", "agent1", "s1", 0, 0); + let keys = window_probe_keys(&[r]); + assert_eq!(keys.len(), 1); + assert!(keys.contains(&("agent1".to_string(), "s1".to_string(), 0))); +} + +#[test] +fn window_probe_keys_skips_rows_without_session_or_seq() { + let r = AgentMetricIndexRow { + session_id: None, + turn_seq: None, + ..row("e1", "agent1", "s1", 5, 0) + }; + let keys = window_probe_keys(&[r]); + assert!(keys.is_empty()); +} + +// ── compute_series: bucketing, overflow, ranking, models ──────────────────── + +#[test] +fn compute_series_buckets_by_reported_at_not_created_at() { + let r = AgentMetricIndexRow { + turn_input_tokens: Some(10), + delta_reliable: Some(true), + event_created_at: 999_999_999, // deliberately wrong/misleading + ..row("e1", "agent1", "s1", 1, DAY + 100) // reported_at lands in bucket 1 + }; + let boundaries = boundaries_7(); + let series = compute_series( + std::slice::from_ref(&r), + std::slice::from_ref(&r), + 0, + &boundaries, + None, + true, + ); + assert_eq!(series.buckets[0].report_count, 0); + assert_eq!(series.buckets[1].report_count, 1); +} + +#[test] +fn compute_series_checked_add_overflow_marks_incomplete_without_wrapping() { + let r1 = AgentMetricIndexRow { + turn_input_tokens: Some(u64::MAX), + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 10, 0) + }; + let r2 = AgentMetricIndexRow { + turn_input_tokens: Some(5), + delta_reliable: Some(true), + ..row("e2", "agent1", "s2", 10, 1) // different session avoids adjacency + }; + let boundaries = boundaries_7(); + let rows = vec![r1, r2]; + let series = compute_series(&rows, &rows, 0, &boundaries, None, true); + let bucket = &series.buckets[0]; + assert!(bucket.has_unknown_usage); + // Value freezes at the last valid sum (u64::MAX) rather than wrapping. + assert_eq!(bucket.usage.input_tokens.value, Some(u64::MAX.to_string())); + assert!(bucket.usage.input_tokens.incomplete); +} + +#[test] +fn compute_series_cost_non_finite_marks_incomplete() { + let r1 = AgentMetricIndexRow { + turn_cost_usd: Some(f64::MAX), + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 10, 0) + }; + let r2 = AgentMetricIndexRow { + turn_cost_usd: Some(f64::MAX), + delta_reliable: Some(true), + ..row("e2", "agent1", "s2", 10, 1) + }; + let boundaries = boundaries_7(); + let rows = vec![r1, r2]; + let series = compute_series(&rows, &rows, 0, &boundaries, None, true); + assert!(series.buckets[0].usage.estimated_cost_usd.incomplete); +} + +#[test] +fn compute_series_ranks_known_total_before_unknown_total() { + let known = AgentMetricIndexRow { + turn_total_tokens: Some(500), + delta_reliable: Some(true), + ..row("e1", "agent_known", "s1", 1, 0) + }; + let unknown = AgentMetricIndexRow { + turn_total_tokens: None, + delta_reliable: Some(true), + ..row("e2", "agent_unknown", "s1", 1, 0) + }; + let boundaries = boundaries_7(); + let rows = vec![known, unknown]; + let series = compute_series(&rows, &rows, 0, &boundaries, None, true); + assert_eq!(series.agents[0].agent_pubkey, "agent_known"); + assert_eq!(series.agents[1].agent_pubkey, "agent_unknown"); +} + +#[test] +fn compute_series_model_breakdown_attributes_per_event_model() { + let r1 = AgentMetricIndexRow { + model: Some("model-a".to_string()), + turn_input_tokens: Some(10), + delta_reliable: Some(true), + ..row("e1", "agent1", "s1", 1, 0) + }; + let r2 = AgentMetricIndexRow { + model: Some("model-b".to_string()), + turn_input_tokens: Some(20), + delta_reliable: Some(true), + ..row("e2", "agent1", "s2", 1, 1) + }; + let boundaries = boundaries_7(); + let rows = vec![r1, r2]; + let series = compute_series(&rows, &rows, 0, &boundaries, None, true); + assert_eq!(series.agents.len(), 1); + assert_eq!(series.agents[0].models.len(), 2); +} + +#[test] +fn compute_series_invalid_report_count_passed_through_and_not_bucketed() { + let boundaries = boundaries_7(); + let series = compute_series(&[], &[], 3, &boundaries, None, true); + assert_eq!(series.coverage.invalid_report_count, 3); + assert_eq!(series.coverage.report_count, 0); + assert!(series.coverage.has_unknown_usage); +} + +#[test] +fn compute_series_zero_invalid_and_no_unknown_rows_is_not_unknown() { + let boundaries = boundaries_7(); + let series = compute_series(&[], &[], 0, &boundaries, None, true); + assert_eq!(series.coverage.invalid_report_count, 0); + assert!(!series.coverage.has_unknown_usage); +} + +#[test] +fn compute_series_collection_enabled_passthrough() { + let boundaries = boundaries_7(); + let series = compute_series(&[], &[], 0, &boundaries, None, false); + assert!(!series.collection_enabled); +} + +#[test] +fn compute_series_has_archived_evidence_passthrough() { + let boundaries = boundaries_7(); + let series = compute_series(&[], &[], 0, &boundaries, Some(true), true); + assert_eq!(series.has_archived_evidence, Some(true)); + let series_none = compute_series(&[], &[], 0, &boundaries, None, true); + assert_eq!(series_none.has_archived_evidence, None); +} + +// ── validate_request ───────────────────────────────────────────────────────── + +fn req(boundaries: Vec, agent_pubkey: Option) -> AgentUsageSeriesRequest { + AgentUsageSeriesRequest { + bucket_boundaries: boundaries, + agent_pubkey, + } +} + +#[test] +fn validate_request_accepts_8_and_31_boundaries() { + assert!(validate_request(&req(boundaries_7(), None)).is_ok()); + let b31: Vec = (0..=30).map(|i| i * DAY).collect(); + assert!(validate_request(&req(b31, None)).is_ok()); +} + +#[test] +fn validate_request_rejects_wrong_boundary_count() { + let bad: Vec = (0..=5).map(|i| i * DAY).collect(); + assert!(validate_request(&req(bad, None)).is_err()); +} + +#[test] +fn validate_request_rejects_non_increasing_boundaries() { + let mut b = boundaries_7(); + b[3] = b[2]; // zero-width interval + assert!(validate_request(&req(b, None)).is_err()); +} + +#[test] +fn validate_request_rejects_interval_over_48h() { + let mut b = boundaries_7(); + b[1] = b[0] + 49 * 3600; + assert!(validate_request(&req(b, None)).is_err()); +} + +#[test] +fn validate_request_accepts_30_minute_dst_interval() { + // Lord Howe Island: 30-minute DST offset. A day-boundary pair differing + // by 23.5h must be accepted under the 48h sanity band (A9). + let mut b = boundaries_7(); + b[1] = b[0] + 23 * 3600 + 1800; + assert!(validate_request(&req(b, None)).is_ok()); +} + +#[test] +fn validate_request_normalizes_pubkey_to_lowercase() { + let pk = "AB".repeat(32); + let result = validate_request(&req(boundaries_7(), Some(pk.clone()))); + assert_eq!(result.unwrap(), Some(pk.to_lowercase())); +} + +#[test] +fn validate_request_rejects_malformed_pubkey() { + let short = "ab".repeat(10); + assert!(validate_request(&req(boundaries_7(), Some(short))).is_err()); + let non_hex = "zz".repeat(32); + assert!(validate_request(&req(boundaries_7(), Some(non_hex))).is_err()); +} diff --git a/desktop/src-tauri/src/archive/metric_store.rs b/desktop/src-tauri/src/archive/metric_store.rs new file mode 100644 index 000000000..a2c739dc8 --- /dev/null +++ b/desktop/src-tauri/src/archive/metric_store.rs @@ -0,0 +1,562 @@ +//! Parsed index of kind 44200 (NIP-AM agent turn metric) archive rows. +//! +//! `agent_metric_index` is a derived, rebuildable cache of parsed NIP-AM +//! payload fields, keyed by `(identity_pubkey, relay_url, id)` exactly like +//! `archived_events`. It exists so the accounting algorithm in +//! `agent_usage.rs` can query parsed columns instead of re-parsing JSON on +//! every render. The canonical source of truth remains +//! `archived_events.raw_json`; every row here is reproducible from it alone +//! via [`AgentMetricIndexRow::from_payload`]. +//! +//! Kept in a sibling file (not `store.rs`) to keep that file under the +//! 1000-line gate, per the existing `pipeline.rs` precedent. + +use rusqlite::{params, Connection, OptionalExtension}; + +use buzz_core_pkg::agent_turn_metric::AgentTurnMetricPayload; + +// ── u64-safe sortable encoding ─────────────────────────────────────────────── + +/// Fixed-width digit count for the lexicographically order-preserving decimal +/// encoding of a `u64`. `u64::MAX` = 18446744073709551615 is 20 digits. +const U64_SORTABLE_WIDTH: usize = 20; + +/// Encode a `u64` as a fixed-width zero-padded decimal string so SQLite TEXT +/// ordering matches numeric ordering, and so the full `u64` range survives +/// SQLite's signed-`i64` INTEGER column type (rusqlite has no unsigned +/// binding). Used for both token counters and `turn_seq`. +pub(super) fn encode_u64_sortable(value: u64) -> String { + format!("{value:0U64_SORTABLE_WIDTH$}") +} + +/// Decode a value written by [`encode_u64_sortable`]. Returns `None` if the +/// string is not a well-formed same-width decimal `u64` — defensive only; +/// every value written by this module is always well-formed. +pub(super) fn decode_u64_sortable(text: &str) -> Option { + if text.len() != U64_SORTABLE_WIDTH { + return None; + } + text.parse::().ok() +} + +fn parse_rfc3339_secs(timestamp: &str) -> Option { + chrono::DateTime::parse_from_rfc3339(timestamp) + .ok() + .map(|dt| dt.timestamp()) +} + +// ── Row type ────────────────────────────────────────────────────────────── + +/// Parse status of a stored `agent_metric_index` row. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) enum ParseStatus { + Valid, + Invalid, +} + +impl ParseStatus { + fn as_str(self) -> &'static str { + match self { + ParseStatus::Valid => "valid", + ParseStatus::Invalid => "invalid", + } + } + + fn from_str(s: &str) -> Self { + match s { + "valid" => ParseStatus::Valid, + _ => ParseStatus::Invalid, + } + } +} + +/// One fully parsed `agent_metric_index` row. +#[derive(Debug, Clone, PartialEq)] +pub(super) struct AgentMetricIndexRow { + pub id: String, + pub agent_pubkey: String, + pub event_created_at: i64, + pub archived_at: i64, + pub reported_at: Option, + pub session_id: Option, + pub turn_seq: Option, + pub model: Option, + pub delta_reliable: Option, + pub turn_input_tokens: Option, + pub turn_output_tokens: Option, + pub turn_total_tokens: Option, + pub turn_cost_usd: Option, + pub cumulative_input_tokens: Option, + pub cumulative_output_tokens: Option, + pub cumulative_total_tokens: Option, + pub cumulative_cost_usd: Option, + pub parse_status: ParseStatus, +} + +impl AgentMetricIndexRow { + /// Parse a decrypted NIP-AM payload (the plaintext JSON already stored in + /// `archived_events.raw_json` for kind 44200 rows) into an index row. + /// + /// New ingest and backfill share this single parser so validity rules + /// cannot drift between the two paths (frozen plan requirement). + /// + /// "Invalid" per Rev 2 A5: the payload's JSON decodes but fails a + /// semantic check this layer owns: unparseable RFC3339 `timestamp`, or + /// `cumulative` present without both `sessionId` and `turnSeq` (NIP-AM + /// REQUIREs both whenever `cumulative` is present — a row missing either + /// cannot supply the complete `(agent, session, seq)` key needed to + /// compete as a cumulative snapshot). The upstream fail-closed ingest + /// path (`pipeline.rs::commit_archive`) has already decrypted, + /// deserialized, and numeric-validated (non-negative/finite `costUsd`) + /// before this row is ever produced — a raw-JSON parse failure here + /// would indicate on-disk corruption, not a normal producer error, but + /// is still handled fail-closed rather than panicking. + pub(super) fn from_payload( + raw_json: &str, + id: &str, + agent_pubkey: &str, + event_created_at: i64, + archived_at: i64, + ) -> Self { + let Ok(payload) = serde_json::from_str::(raw_json) else { + return Self::invalid(id, agent_pubkey, event_created_at, archived_at); + }; + + let reported_at = parse_rfc3339_secs(&payload.timestamp); + let cumulative_requires_session_seq = payload.cumulative.is_some(); + let has_session_and_seq = payload.session_id.is_some() && payload.turn_seq.is_some(); + + if reported_at.is_none() || (cumulative_requires_session_seq && !has_session_and_seq) { + return Self::invalid(id, agent_pubkey, event_created_at, archived_at); + } + + let turn = payload.turn.as_ref(); + let cumulative = payload.cumulative.as_ref(); + + Self { + id: id.to_string(), + agent_pubkey: agent_pubkey.to_string(), + event_created_at, + archived_at, + reported_at, + session_id: payload.session_id, + turn_seq: payload.turn_seq, + model: payload.model, + delta_reliable: Some(payload.delta_reliable), + turn_input_tokens: turn.and_then(|t| t.input_tokens), + turn_output_tokens: turn.and_then(|t| t.output_tokens), + turn_total_tokens: turn.and_then(|t| t.total_tokens), + turn_cost_usd: turn.and_then(|t| t.cost_usd), + cumulative_input_tokens: cumulative.and_then(|c| c.input_tokens), + cumulative_output_tokens: cumulative.and_then(|c| c.output_tokens), + cumulative_total_tokens: cumulative.and_then(|c| c.total_tokens), + cumulative_cost_usd: cumulative.and_then(|c| c.cost_usd), + parse_status: ParseStatus::Valid, + } + } + + fn invalid(id: &str, agent_pubkey: &str, event_created_at: i64, archived_at: i64) -> Self { + Self { + id: id.to_string(), + agent_pubkey: agent_pubkey.to_string(), + event_created_at, + archived_at, + reported_at: None, + session_id: None, + turn_seq: None, + model: None, + delta_reliable: None, + turn_input_tokens: None, + turn_output_tokens: None, + turn_total_tokens: None, + turn_cost_usd: None, + cumulative_input_tokens: None, + cumulative_output_tokens: None, + cumulative_total_tokens: None, + cumulative_cost_usd: None, + parse_status: ParseStatus::Invalid, + } + } + + /// The `(agent_pubkey, session_id, turn_seq)` cumulative-accounting key, + /// or `None` if this row cannot participate in cumulative delta + /// recomputation (missing session/sequence). + pub(super) fn accounting_key(&self) -> Option<(String, String, u64)> { + match (&self.session_id, self.turn_seq) { + (Some(sid), Some(seq)) => Some((self.agent_pubkey.clone(), sid.clone(), seq)), + _ => None, + } + } +} + +fn row_from_sql(row: &rusqlite::Row) -> rusqlite::Result { + let turn_seq_text: Option = row.get("turn_seq")?; + let turn_input_text: Option = row.get("turn_input_tokens")?; + let turn_output_text: Option = row.get("turn_output_tokens")?; + let turn_total_text: Option = row.get("turn_total_tokens")?; + let cum_input_text: Option = row.get("cumulative_input_tokens")?; + let cum_output_text: Option = row.get("cumulative_output_tokens")?; + let cum_total_text: Option = row.get("cumulative_total_tokens")?; + let delta_reliable_int: Option = row.get("delta_reliable")?; + let parse_status_str: String = row.get("parse_status")?; + + Ok(AgentMetricIndexRow { + id: row.get("id")?, + agent_pubkey: row.get("agent_pubkey")?, + event_created_at: row.get("event_created_at")?, + archived_at: row.get("archived_at")?, + reported_at: row.get("reported_at")?, + session_id: row.get("session_id")?, + turn_seq: turn_seq_text.as_deref().and_then(decode_u64_sortable), + model: row.get("model")?, + delta_reliable: delta_reliable_int.map(|v| v != 0), + turn_input_tokens: turn_input_text.as_deref().and_then(decode_u64_sortable), + turn_output_tokens: turn_output_text.as_deref().and_then(decode_u64_sortable), + turn_total_tokens: turn_total_text.as_deref().and_then(decode_u64_sortable), + turn_cost_usd: row.get("turn_cost_usd")?, + cumulative_input_tokens: cum_input_text.as_deref().and_then(decode_u64_sortable), + cumulative_output_tokens: cum_output_text.as_deref().and_then(decode_u64_sortable), + cumulative_total_tokens: cum_total_text.as_deref().and_then(decode_u64_sortable), + cumulative_cost_usd: row.get("cumulative_cost_usd")?, + parse_status: ParseStatus::from_str(&parse_status_str), + }) +} + +const ROW_COLUMNS: &str = "id, agent_pubkey, event_created_at, archived_at, reported_at, \ + session_id, turn_seq, model, delta_reliable, turn_input_tokens, turn_output_tokens, \ + turn_total_tokens, turn_cost_usd, cumulative_input_tokens, cumulative_output_tokens, \ + cumulative_total_tokens, cumulative_cost_usd, parse_status"; + +// ── Write path ─────────────────────────────────────────────────────────────── + +/// Insert one metric index row inside the caller's transaction. +/// +/// Called from `pipeline::commit_archive` ONLY when the corresponding +/// `archived_events` row was newly inserted (never for a duplicate), and +/// from the backfill driver for pre-existing unindexed rows. `INSERT OR +/// IGNORE` on the shared PK makes a second call for the same +/// `(identity, relay, id)` a safe no-op (defensive; callers already guard +/// against re-indexing). +pub(super) fn insert_metric_index_row( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, + row: &AgentMetricIndexRow, +) -> Result { + let turn_seq = row.turn_seq.map(encode_u64_sortable); + let turn_input = row.turn_input_tokens.map(encode_u64_sortable); + let turn_output = row.turn_output_tokens.map(encode_u64_sortable); + let turn_total = row.turn_total_tokens.map(encode_u64_sortable); + let cum_input = row.cumulative_input_tokens.map(encode_u64_sortable); + let cum_output = row.cumulative_output_tokens.map(encode_u64_sortable); + let cum_total = row.cumulative_total_tokens.map(encode_u64_sortable); + let delta_reliable = row.delta_reliable.map(|b| b as i64); + + let affected = conn + .execute( + "INSERT INTO agent_metric_index + (identity_pubkey, relay_url, id, agent_pubkey, event_created_at, + archived_at, reported_at, session_id, turn_seq, model, + delta_reliable, turn_input_tokens, turn_output_tokens, + turn_total_tokens, turn_cost_usd, cumulative_input_tokens, + cumulative_output_tokens, cumulative_total_tokens, + cumulative_cost_usd, parse_status) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, + ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19, ?20) + ON CONFLICT (identity_pubkey, relay_url, id) DO NOTHING", + params![ + identity_pubkey, + relay_url, + row.id, + row.agent_pubkey, + row.event_created_at, + row.archived_at, + row.reported_at, + row.session_id, + turn_seq, + row.model, + delta_reliable, + turn_input, + turn_output, + turn_total, + row.turn_cost_usd, + cum_input, + cum_output, + cum_total, + row.cumulative_cost_usd, + row.parse_status.as_str(), + ], + ) + .map_err(|e| format!("failed to insert agent_metric_index row: {e}"))?; + Ok(affected > 0) +} + +// ── Backfill ───────────────────────────────────────────────────────────────── + +/// Backfill existing `archived_events` kind-44200 rows that have no matching +/// `agent_metric_index` row yet, for the given identity + relay. +/// +/// Runs in bounded chunks (~500 rows per transaction) so a large existing +/// archive never holds one unbounded write lock; each chunk is independently +/// atomic and the whole backfill is idempotent (anti-join + index PK is the +/// source of truth) and restartable — interruption between chunks loses +/// nothing, and a later run simply resumes against the still-missing rows. +/// +/// Returns the total number of newly indexed rows. +pub(super) fn backfill_agent_metric_index( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, +) -> Result { + const CHUNK_SIZE: i64 = 500; + let mut total = 0usize; + + loop { + let mut stmt = conn + .prepare( + "SELECT ae.id, ae.pubkey, ae.created_at, ae.archived_at, ae.raw_json + FROM archived_events ae + WHERE ae.identity_pubkey = ?1 + AND ae.relay_url = ?2 + AND ae.kind = 44200 + AND ae.id NOT IN ( + SELECT id FROM agent_metric_index + WHERE identity_pubkey = ?1 + AND relay_url = ?2 + ) + ORDER BY ae.created_at ASC, ae.id ASC + LIMIT ?3", + ) + .map_err(|e| format!("prepare backfill_agent_metric_index select: {e}"))?; + + let chunk: Vec<(String, String, i64, i64, String)> = stmt + .query_map(params![identity_pubkey, relay_url, CHUNK_SIZE], |row| { + Ok(( + row.get::<_, String>(0)?, + row.get::<_, String>(1)?, + row.get::<_, i64>(2)?, + row.get::<_, i64>(3)?, + row.get::<_, String>(4)?, + )) + }) + .map_err(|e| format!("query backfill_agent_metric_index select: {e}"))? + .collect::, _>>() + .map_err(|e| format!("read backfill_agent_metric_index row: {e}"))?; + drop(stmt); + + if chunk.is_empty() { + break; + } + let chunk_len = chunk.len(); + + let tx = conn + .unchecked_transaction() + .map_err(|e| format!("failed to begin backfill chunk transaction: {e}"))?; + for (id, pubkey, created_at, archived_at, raw_json) in &chunk { + let parsed = + AgentMetricIndexRow::from_payload(raw_json, id, pubkey, *created_at, *archived_at); + insert_metric_index_row(&tx, identity_pubkey, relay_url, &parsed)?; + } + tx.commit() + .map_err(|e| format!("failed to commit backfill chunk: {e}"))?; + + total += chunk_len; + if (chunk_len as i64) < CHUNK_SIZE { + break; + } + } + + Ok(total) +} + +// ── GC / orphan repair ─────────────────────────────────────────────────────── + +/// Delete `agent_metric_index` rows whose canonical `archived_events` row no +/// longer exists. Called from `store::gc_orphaned_events` inside the SAME +/// SQLite transaction as the canonical delete (A6) so the index can never +/// observe a canonical row as gone while its own row survives. +pub(super) fn delete_orphaned_metric_index_rows( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, +) -> Result { + let affected = conn + .execute( + "DELETE FROM agent_metric_index + WHERE identity_pubkey = ?1 + AND relay_url = ?2 + AND id NOT IN ( + SELECT id FROM archived_events + WHERE identity_pubkey = ?1 + AND relay_url = ?2 + )", + params![identity_pubkey, relay_url], + ) + .map_err(|e| format!("failed to gc orphaned agent_metric_index rows: {e}"))?; + Ok(affected) +} + +/// Read-time orphan repair: same anti-join delete as +/// [`delete_orphaned_metric_index_rows`], run defensively before every read +/// so a planted/legacy orphan is self-healed even if a future deletion path +/// forgets to call the GC cascade. Defense in depth, not an atomicity +/// substitute for A6. +pub(super) fn repair_orphaned_metric_index_rows( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, +) -> Result { + delete_orphaned_metric_index_rows(conn, identity_pubkey, relay_url) +} + +// ── Read path ──────────────────────────────────────────────────────────────── + +/// Load all VALID rows whose `reported_at` falls in `[start, end)`, optionally +/// filtered to one agent author. This is the exact set of rows that may ever +/// be counted into a bucket (A11 step 1). +pub(super) fn load_window_valid_rows( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, + start: i64, + end: i64, + agent_pubkey: Option<&str>, +) -> Result, String> { + let sql = format!( + "SELECT {ROW_COLUMNS} FROM agent_metric_index + WHERE identity_pubkey = ?1 AND relay_url = ?2 AND parse_status = 'valid' + AND reported_at >= ?3 AND reported_at < ?4 + AND (?5 IS NULL OR agent_pubkey = ?5) + ORDER BY reported_at ASC, id ASC" + ); + let mut stmt = stmt_prepare(conn, &sql)?; + let rows = stmt + .query_map( + params![identity_pubkey, relay_url, start, end, agent_pubkey], + row_from_sql, + ) + .map_err(|e| format!("query load_window_valid_rows: {e}"))?; + rows.collect::, _>>() + .map_err(|e| format!("read load_window_valid_rows row: {e}")) +} + +/// Count INVALID rows whose `event_created_at` falls in `[start, end)` +/// (invalid rows have no trustworthy `reported_at`, so coarse signed +/// `created_at` is the only available time signal), optionally filtered to +/// one agent author. +pub(super) fn count_invalid_rows_in_window( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, + start: i64, + end: i64, + agent_pubkey: Option<&str>, +) -> Result { + conn.query_row( + "SELECT COUNT(*) FROM agent_metric_index + WHERE identity_pubkey = ?1 AND relay_url = ?2 AND parse_status = 'invalid' + AND event_created_at >= ?3 AND event_created_at < ?4 + AND (?5 IS NULL OR agent_pubkey = ?5)", + params![identity_pubkey, relay_url, start, end, agent_pubkey], + |row| row.get(0), + ) + .map_err(|e| format!("count_invalid_rows_in_window: {e}")) +} + +/// For the given set of exact `(agent_pubkey, session_id, turn_seq)` keys, +/// load ALL valid rows matching those exact keys with NO `reported_at` +/// restriction (A11). Used both for duplicate-cardinality checks at a +/// sequence and for exact-predecessor baseline lookups — a single probe +/// covers both, since the predecessor key is included in the request set +/// alongside each window row's own key. +/// +/// Keys are grouped by `(agent_pubkey, session_id)` and queried with a +/// `turn_seq IN (...)` clause per group (typically few groups per window), +/// served by `idx_agent_metric_session`. +pub(super) fn load_rows_at_exact_keys( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, + keys: &std::collections::HashSet<(String, String, u64)>, +) -> Result, String> { + use std::collections::HashMap; + + // Group by (agent, session) so each group becomes one IN-list query. + let mut groups: HashMap<(String, String), Vec> = HashMap::new(); + for (agent, session, seq) in keys { + groups + .entry((agent.clone(), session.clone())) + .or_default() + .push(*seq); + } + + let mut out = Vec::new(); + for ((agent, session), seqs) in groups { + let encoded: Vec = seqs.into_iter().map(encode_u64_sortable).collect(); + let sql = format!( + "SELECT {ROW_COLUMNS} FROM agent_metric_index + WHERE identity_pubkey = ?1 AND relay_url = ?2 AND parse_status = 'valid' + AND agent_pubkey = ?3 AND session_id = ?4 + AND turn_seq IN ({})", + encoded + .iter() + .enumerate() + .map(|(i, _)| format!("?{}", i + 5)) + .collect::>() + .join(",") + ); + let mut stmt = stmt_prepare(conn, &sql)?; + let mut bound: Vec> = vec![ + Box::new(identity_pubkey.to_owned()), + Box::new(relay_url.to_owned()), + Box::new(agent.clone()), + Box::new(session.clone()), + ]; + for e in &encoded { + bound.push(Box::new(e.clone())); + } + let refs: Vec<&dyn rusqlite::ToSql> = bound.iter().map(|b| b.as_ref()).collect(); + let rows = stmt + .query_map(refs.as_slice(), row_from_sql) + .map_err(|e| format!("query load_rows_at_exact_keys: {e}"))?; + for r in rows { + out.push(r.map_err(|e| format!("read load_rows_at_exact_keys row: {e}"))?); + } + } + + Ok(out) +} + +/// `hasArchivedEvidence` (A13): does at least one surviving `agent_metric_index` +/// row (either `parse_status`) exist for this author under the active +/// identity+relay, with NO bucket-boundary restriction? Computed after +/// backfill + orphan repair by the caller. +pub(super) fn has_archived_evidence( + conn: &Connection, + identity_pubkey: &str, + relay_url: &str, + agent_pubkey: &str, +) -> Result { + let exists: Option = conn + .query_row( + "SELECT 1 FROM agent_metric_index + WHERE identity_pubkey = ?1 AND relay_url = ?2 AND agent_pubkey = ?3 + LIMIT 1", + params![identity_pubkey, relay_url, agent_pubkey], + |row| row.get(0), + ) + .optional() + .map_err(|e| format!("has_archived_evidence: {e}"))?; + Ok(exists.is_some()) +} + +fn stmt_prepare<'a>(conn: &'a Connection, sql: &str) -> Result, String> { + conn.prepare(sql) + .map_err(|e| format!("prepare failed: {e} — sql: {sql}")) +} + +// ── Tests ──────────────────────────────────────────────────────────────────── + +#[cfg(test)] +#[path = "metric_store_tests.rs"] +mod metric_store_tests; diff --git a/desktop/src-tauri/src/archive/metric_store_tests.rs b/desktop/src-tauri/src/archive/metric_store_tests.rs new file mode 100644 index 000000000..48fa091fd --- /dev/null +++ b/desktop/src-tauri/src/archive/metric_store_tests.rs @@ -0,0 +1,579 @@ +//! Unit tests for `archive/metric_store.rs`. +//! +//! Kept in a sibling file so `metric_store.rs` stays under the file-size +//! gate; `#[path]`-included from there. + +use super::*; +use crate::archive::store::{self, SCHEMA}; + +fn in_memory() -> Connection { + let conn = Connection::open_in_memory().unwrap(); + conn.pragma_update(None, "journal_mode", "WAL").unwrap(); + conn.pragma_update(None, "busy_timeout", 5000).unwrap(); + conn.execute_batch(SCHEMA).unwrap(); + conn +} + +fn valid_payload_json(session_id: &str, seq: u64, timestamp: &str) -> String { + format!( + r#"{{"harness":"goose","model":"claude","channelId":null,"sessionId":"{session_id}","turnId":null,"turnSeq":{seq},"timestamp":"{timestamp}","turn":{{"inputTokens":10,"outputTokens":20,"totalTokens":30,"costUsd":0.01}},"cumulative":{{"inputTokens":100,"outputTokens":200,"totalTokens":300,"costUsd":0.1}},"deltaReliable":true,"stopReason":"end_turn"}}"# + ) +} + +#[allow(clippy::too_many_arguments)] +fn insert_archived_event( + conn: &Connection, + identity: &str, + relay: &str, + id: &str, + kind: i64, + pubkey: &str, + created_at: i64, + raw_json: &str, + archived_at: i64, +) { + store::upsert_archived_event( + conn, + identity, + relay, + id, + kind, + pubkey, + created_at, + raw_json, + archived_at, + ) + .unwrap(); +} + +// ── u64 sortable encoding ──────────────────────────────────────────────────── + +#[test] +fn u64_sortable_round_trips_zero_and_max() { + for v in [0u64, 1, 12345, u64::MAX - 1, u64::MAX] { + let encoded = encode_u64_sortable(v); + assert_eq!(encoded.len(), U64_SORTABLE_WIDTH); + assert_eq!(decode_u64_sortable(&encoded), Some(v)); + } +} + +#[test] +fn u64_sortable_encoding_preserves_numeric_order_across_i64_max() { + let below = i64::MAX as u64; + let above = (i64::MAX as u64) + 1; + let e_below = encode_u64_sortable(below); + let e_above = encode_u64_sortable(above); + assert!( + e_below < e_above, + "lexicographic order must match numeric order" + ); +} + +#[test] +fn decode_rejects_wrong_width() { + assert_eq!(decode_u64_sortable("123"), None); + assert_eq!(decode_u64_sortable(""), None); +} + +// ── from_payload parsing ───────────────────────────────────────────────────── + +#[test] +fn from_payload_parses_valid_row() { + let json = valid_payload_json("s1", 7, "2026-07-01T20:11:03.213Z"); + let row = AgentMetricIndexRow::from_payload(&json, "eid1", "agent1", 100, 200); + assert_eq!(row.parse_status, ParseStatus::Valid); + assert_eq!(row.session_id, Some("s1".to_string())); + assert_eq!(row.turn_seq, Some(7)); + assert_eq!(row.turn_input_tokens, Some(10)); + assert_eq!(row.cumulative_input_tokens, Some(100)); + assert_eq!(row.model, Some("claude".to_string())); +} + +#[test] +fn from_payload_marks_unparseable_json_invalid() { + let row = AgentMetricIndexRow::from_payload("not json", "eid1", "agent1", 100, 200); + assert_eq!(row.parse_status, ParseStatus::Invalid); + assert_eq!(row.turn_input_tokens, None); +} + +#[test] +fn from_payload_marks_unparseable_timestamp_invalid() { + let json = r#"{"harness":"goose","timestamp":"not-a-timestamp"}"#; + let row = AgentMetricIndexRow::from_payload(json, "eid1", "agent1", 100, 200); + assert_eq!(row.parse_status, ParseStatus::Invalid); +} + +#[test] +fn from_payload_marks_cumulative_without_session_seq_invalid() { + // cumulative present but sessionId/turnSeq missing — semantic-invalid per A5. + let json = r#"{"harness":"goose","timestamp":"2026-07-01T20:11:03Z","cumulative":{"inputTokens":1,"outputTokens":null,"totalTokens":null,"costUsd":null}}"#; + let row = AgentMetricIndexRow::from_payload(json, "eid1", "agent1", 100, 200); + assert_eq!(row.parse_status, ParseStatus::Invalid); +} + +#[test] +fn from_payload_accepts_missing_cumulative_without_session_seq() { + // No cumulative object at all — session/seq are not required. + let json = r#"{"harness":"goose","timestamp":"2026-07-01T20:11:03Z","turn":{"inputTokens":5,"outputTokens":null,"totalTokens":null,"costUsd":null}}"#; + let row = AgentMetricIndexRow::from_payload(json, "eid1", "agent1", 100, 200); + assert_eq!(row.parse_status, ParseStatus::Valid); + assert_eq!(row.turn_input_tokens, Some(5)); +} + +// ── insert / idempotence ────────────────────────────────────────────────────── + +#[test] +fn insert_metric_index_row_is_idempotent_on_pk() { + let conn = in_memory(); + let row = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"), + "eid1", + "agent1", + 100, + 200, + ); + let first = insert_metric_index_row(&conn, "id", "relay", &row).unwrap(); + let second = insert_metric_index_row(&conn, "id", "relay", &row).unwrap(); + assert!(first); + assert!(!second, "second insert of the same PK must be a no-op"); + + let count: i64 = conn + .query_row("SELECT COUNT(*) FROM agent_metric_index", [], |r| r.get(0)) + .unwrap(); + assert_eq!(count, 1); +} + +#[test] +fn insert_metric_index_row_round_trips_u64_max() { + let conn = in_memory(); + let row = AgentMetricIndexRow { + turn_seq: Some(u64::MAX), + turn_input_tokens: Some(u64::MAX), + cumulative_input_tokens: Some(u64::MAX), + ..AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"), + "eid1", + "agent1", + 100, + 200, + ) + }; + insert_metric_index_row(&conn, "id", "relay", &row).unwrap(); + + let loaded = load_window_valid_rows(&conn, "id", "relay", 0, i64::MAX, None).unwrap(); + assert_eq!(loaded.len(), 1); + assert_eq!(loaded[0].turn_seq, Some(u64::MAX)); + assert_eq!(loaded[0].turn_input_tokens, Some(u64::MAX)); + assert_eq!(loaded[0].cumulative_input_tokens, Some(u64::MAX)); +} + +#[test] +fn insert_invalid_row_preserves_null_parsed_columns() { + let conn = in_memory(); + let row = AgentMetricIndexRow::from_payload("bad json", "eid1", "agent1", 100, 200); + insert_metric_index_row(&conn, "id", "relay", &row).unwrap(); + + let count = count_invalid_rows_in_window(&conn, "id", "relay", 0, i64::MAX, None).unwrap(); + assert_eq!(count, 1); +} + +// ── Identity/relay isolation ───────────────────────────────────────────────── + +#[test] +fn load_window_valid_rows_scoped_to_identity_and_relay() { + let conn = in_memory(); + let row_a = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"), + "eidA", + "agent1", + 100, + 200, + ); + let row_b = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"), + "eidB", + "agent1", + 100, + 200, + ); + insert_metric_index_row(&conn, "identityA", "relay1", &row_a).unwrap(); + insert_metric_index_row(&conn, "identityB", "relay1", &row_b).unwrap(); + + let loaded_a = load_window_valid_rows(&conn, "identityA", "relay1", 0, i64::MAX, None).unwrap(); + assert_eq!(loaded_a.len(), 1); + assert_eq!(loaded_a[0].id, "eidA"); + + let loaded_b = load_window_valid_rows(&conn, "identityB", "relay1", 0, i64::MAX, None).unwrap(); + assert_eq!(loaded_b.len(), 1); + assert_eq!(loaded_b[0].id, "eidB"); +} + +#[test] +fn load_window_valid_rows_filters_by_agent_pubkey() { + let conn = in_memory(); + let row_a = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"), + "eidA", + "agentA", + 100, + 200, + ); + let row_b = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"), + "eidB", + "agentB", + 100, + 200, + ); + insert_metric_index_row(&conn, "id", "relay", &row_a).unwrap(); + insert_metric_index_row(&conn, "id", "relay", &row_b).unwrap(); + + let loaded = load_window_valid_rows(&conn, "id", "relay", 0, i64::MAX, Some("agentA")).unwrap(); + assert_eq!(loaded.len(), 1); + assert_eq!(loaded[0].agent_pubkey, "agentA"); +} + +#[test] +fn load_window_valid_rows_excludes_out_of_window_reported_at() { + let conn = in_memory(); + let in_window = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "2026-01-02T00:00:00Z"), + "eid_in", + "agent1", + 0, + 0, + ); + let out_of_window = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 2, "2020-01-01T00:00:00Z"), + "eid_out", + "agent1", + 0, + 0, + ); + insert_metric_index_row(&conn, "id", "relay", &in_window).unwrap(); + insert_metric_index_row(&conn, "id", "relay", &out_of_window).unwrap(); + + let start = chrono::DateTime::parse_from_rfc3339("2026-01-01T00:00:00Z") + .unwrap() + .timestamp(); + let end = chrono::DateTime::parse_from_rfc3339("2026-01-03T00:00:00Z") + .unwrap() + .timestamp(); + let loaded = load_window_valid_rows(&conn, "id", "relay", start, end, None).unwrap(); + assert_eq!(loaded.len(), 1); + assert_eq!(loaded[0].id, "eid_in"); +} + +// ── load_rows_at_exact_keys ─────────────────────────────────────────────────── + +#[test] +fn load_rows_at_exact_keys_matches_multiple_groups() { + let conn = in_memory(); + let r1 = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 5, "2026-07-01T00:00:00Z"), + "e1", + "agent1", + 0, + 0, + ); + let r2 = AgentMetricIndexRow::from_payload( + &valid_payload_json("s2", 9, "2026-07-01T00:00:00Z"), + "e2", + "agent1", + 0, + 0, + ); + let r3_not_requested = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 99, "2026-07-01T00:00:00Z"), + "e3", + "agent1", + 0, + 0, + ); + insert_metric_index_row(&conn, "id", "relay", &r1).unwrap(); + insert_metric_index_row(&conn, "id", "relay", &r2).unwrap(); + insert_metric_index_row(&conn, "id", "relay", &r3_not_requested).unwrap(); + + let mut keys = std::collections::HashSet::new(); + keys.insert(("agent1".to_string(), "s1".to_string(), 5u64)); + keys.insert(("agent1".to_string(), "s2".to_string(), 9u64)); + + let loaded = load_rows_at_exact_keys(&conn, "id", "relay", &keys).unwrap(); + let mut ids: Vec<&str> = loaded.iter().map(|r| r.id.as_str()).collect(); + ids.sort(); + assert_eq!(ids, vec!["e1", "e2"]); +} + +#[test] +fn load_rows_at_exact_keys_returns_all_duplicates_at_one_key() { + let conn = in_memory(); + let dup_a = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 5, "2026-07-01T00:00:00Z"), + "eA", + "agent1", + 0, + 0, + ); + let dup_b = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 5, "2026-07-01T00:00:01Z"), + "eB", + "agent1", + 0, + 1, + ); + insert_metric_index_row(&conn, "id", "relay", &dup_a).unwrap(); + insert_metric_index_row(&conn, "id", "relay", &dup_b).unwrap(); + + let mut keys = std::collections::HashSet::new(); + keys.insert(("agent1".to_string(), "s1".to_string(), 5u64)); + let loaded = load_rows_at_exact_keys(&conn, "id", "relay", &keys).unwrap(); + assert_eq!(loaded.len(), 2); +} + +// ── has_archived_evidence ───────────────────────────────────────────────────── + +#[test] +fn has_archived_evidence_true_for_either_parse_status() { + let conn = in_memory(); + let invalid_row = AgentMetricIndexRow::from_payload("bad", "eid1", "agentX", 0, 0); + insert_metric_index_row(&conn, "id", "relay", &invalid_row).unwrap(); + + assert!(has_archived_evidence(&conn, "id", "relay", "agentX").unwrap()); + assert!(!has_archived_evidence(&conn, "id", "relay", "agentY").unwrap()); +} + +#[test] +fn has_archived_evidence_ignores_bucket_boundaries() { + let conn = in_memory(); + // A very old row (outside any realistic window) still counts as evidence. + let old_row = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "1999-01-01T00:00:00Z"), + "eid1", + "agentX", + 0, + 0, + ); + insert_metric_index_row(&conn, "id", "relay", &old_row).unwrap(); + assert!(has_archived_evidence(&conn, "id", "relay", "agentX").unwrap()); +} + +// ── Backfill ────────────────────────────────────────────────────────────────── + +#[test] +fn backfill_indexes_existing_unindexed_kind_44200_rows() { + let conn = in_memory(); + let json = valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"); + insert_archived_event( + &conn, "id", "relay", "eid1", 44200, "agent1", 100, &json, 200, + ); + + let indexed = backfill_agent_metric_index(&conn, "id", "relay").unwrap(); + assert_eq!(indexed, 1); + + let loaded = load_window_valid_rows(&conn, "id", "relay", 0, i64::MAX, None).unwrap(); + assert_eq!(loaded.len(), 1); + assert_eq!(loaded[0].id, "eid1"); +} + +#[test] +fn backfill_is_idempotent_anti_join() { + let conn = in_memory(); + let json = valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"); + insert_archived_event( + &conn, "id", "relay", "eid1", 44200, "agent1", 100, &json, 200, + ); + + backfill_agent_metric_index(&conn, "id", "relay").unwrap(); + let second_run = backfill_agent_metric_index(&conn, "id", "relay").unwrap(); + assert_eq!(second_run, 0, "second backfill run must index nothing new"); + + let count: i64 = conn + .query_row("SELECT COUNT(*) FROM agent_metric_index", [], |r| r.get(0)) + .unwrap(); + assert_eq!(count, 1); +} + +#[test] +fn backfill_processes_chunks_larger_than_500_rows() { + let conn = in_memory(); + // 501 rows exercises the CHUNK_SIZE=500 boundary. + for i in 0..501 { + let json = valid_payload_json(&format!("s{i}"), 1, "2026-07-01T00:00:00Z"); + insert_archived_event( + &conn, + "id", + "relay", + &format!("eid{i}"), + 44200, + "agent1", + 100 + i as i64, + &json, + 200, + ); + } + let indexed = backfill_agent_metric_index(&conn, "id", "relay").unwrap(); + assert_eq!(indexed, 501); + + let count: i64 = conn + .query_row("SELECT COUNT(*) FROM agent_metric_index", [], |r| r.get(0)) + .unwrap(); + assert_eq!(count, 501); +} + +#[test] +fn backfill_ignores_non_44200_rows() { + let conn = in_memory(); + insert_archived_event(&conn, "id", "relay", "eid1", 1, "author1", 100, "{}", 200); + let indexed = backfill_agent_metric_index(&conn, "id", "relay").unwrap(); + assert_eq!(indexed, 0); +} + +// ── GC / orphan repair ──────────────────────────────────────────────────────── + +#[test] +fn delete_orphaned_metric_index_rows_removes_rows_with_no_canonical_event() { + let conn = in_memory(); + // Planted orphan: index row with no matching archived_events row. + let orphan = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"), + "orphan_id", + "agent1", + 0, + 0, + ); + insert_metric_index_row(&conn, "id", "relay", &orphan).unwrap(); + + let deleted = delete_orphaned_metric_index_rows(&conn, "id", "relay").unwrap(); + assert_eq!(deleted, 1); + + let count: i64 = conn + .query_row("SELECT COUNT(*) FROM agent_metric_index", [], |r| r.get(0)) + .unwrap(); + assert_eq!(count, 0); +} + +#[test] +fn delete_orphaned_metric_index_rows_preserves_rows_with_canonical_event() { + let conn = in_memory(); + let json = valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"); + insert_archived_event( + &conn, "id", "relay", "eid1", 44200, "agent1", 100, &json, 200, + ); + let row = AgentMetricIndexRow::from_payload(&json, "eid1", "agent1", 100, 200); + insert_metric_index_row(&conn, "id", "relay", &row).unwrap(); + + let deleted = delete_orphaned_metric_index_rows(&conn, "id", "relay").unwrap(); + assert_eq!(deleted, 0); +} + +#[test] +fn repair_orphaned_metric_index_rows_self_heals_planted_orphan_before_read() { + let conn = in_memory(); + let orphan = AgentMetricIndexRow::from_payload( + &valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"), + "orphan_id", + "agent1", + 0, + 0, + ); + insert_metric_index_row(&conn, "id", "relay", &orphan).unwrap(); + + repair_orphaned_metric_index_rows(&conn, "id", "relay").unwrap(); + let loaded = load_window_valid_rows(&conn, "id", "relay", 0, i64::MAX, None).unwrap(); + assert!( + loaded.is_empty(), + "planted orphan must never be reported after repair" + ); +} + +#[test] +fn gc_orphaned_events_cascades_to_metric_index_atomically() { + let conn = in_memory(); + let json = valid_payload_json("s1", 1, "2026-07-01T00:00:00Z"); + insert_archived_event( + &conn, "id", "relay", "eid1", 44200, "agent1", 100, &json, 200, + ); + let row = AgentMetricIndexRow::from_payload(&json, "eid1", "agent1", 100, 200); + insert_metric_index_row(&conn, "id", "relay", &row).unwrap(); + + // Remove the last scope row so the event becomes orphaned, then GC. + // (No scope row was ever added in this test, so the event is already + // orphaned by construction — gc_orphaned_events should delete both the + // canonical row and its index row in one transaction.) + store::gc_orphaned_events(&conn, "id", "relay").unwrap(); + + let event_count: i64 = conn + .query_row("SELECT COUNT(*) FROM archived_events", [], |r| r.get(0)) + .unwrap(); + let index_count: i64 = conn + .query_row("SELECT COUNT(*) FROM agent_metric_index", [], |r| r.get(0)) + .unwrap(); + assert_eq!(event_count, 0); + assert_eq!( + index_count, 0, + "index row must not outlive its canonical event" + ); +} + +// ── EXPLAIN QUERY PLAN index assertions (A7) ───────────────────────────────── + +fn query_plan(conn: &Connection, sql: &str, params: &[&dyn rusqlite::ToSql]) -> String { + let explain_sql = format!("EXPLAIN QUERY PLAN {sql}"); + let mut stmt = conn.prepare(&explain_sql).unwrap(); + let mut rows = stmt.query(params).unwrap(); + let mut plan = String::new(); + while let Some(row) = rows.next().unwrap() { + let detail: String = row.get(3).unwrap(); + plan.push_str(&detail); + plan.push('\n'); + } + plan +} + +#[test] +fn backfill_anti_join_uses_partial_index() { + let conn = in_memory(); + let plan = query_plan( + &conn, + "SELECT ae.id FROM archived_events ae + WHERE ae.identity_pubkey = ?1 AND ae.relay_url = ?2 AND ae.kind = 44200 + AND ae.id NOT IN (SELECT id FROM agent_metric_index WHERE identity_pubkey = ?1 AND relay_url = ?2)", + &[&"id", &"relay"], + ); + assert!( + plan.contains("idx_archived_events_agent_metric"), + "backfill anti-join must use the partial index, plan was:\n{plan}" + ); +} + +#[test] +fn window_scan_uses_reported_index() { + let conn = in_memory(); + let plan = query_plan( + &conn, + "SELECT * FROM agent_metric_index + WHERE identity_pubkey = ?1 AND relay_url = ?2 AND parse_status = 'valid' + AND reported_at >= ?3 AND reported_at < ?4", + &[&"id", &"relay", &0i64, &100i64], + ); + assert!( + plan.contains("idx_agent_metric_reported"), + "window scan must use the reported-time index, plan was:\n{plan}" + ); +} + +#[test] +fn predecessor_lookup_uses_session_index() { + let conn = in_memory(); + let plan = query_plan( + &conn, + "SELECT * FROM agent_metric_index + WHERE identity_pubkey = ?1 AND relay_url = ?2 AND parse_status = 'valid' + AND agent_pubkey = ?3 AND session_id = ?4 AND turn_seq IN (?5)", + &[&"id", &"relay", &"agent1", &"s1", &"00000000000000000005"], + ); + assert!( + plan.contains("idx_agent_metric_session"), + "predecessor lookup must use the session index, plan was:\n{plan}" + ); +} diff --git a/desktop/src-tauri/src/archive/mod.rs b/desktop/src-tauri/src/archive/mod.rs index a65b126a4..0f56c432b 100644 --- a/desktop/src-tauri/src/archive/mod.rs +++ b/desktop/src-tauri/src/archive/mod.rs @@ -17,6 +17,8 @@ //! validation (sig/id + kind + p-tag + agent tag + frame=telemetry + author //! == agent) is applied fail-closed. +mod agent_usage; +mod metric_store; mod pipeline; pub mod store; @@ -116,9 +118,17 @@ pub struct MatchedScope { /// Result of a batch archive call. #[derive(Debug, Serialize)] +#[serde(rename_all = "camelCase")] pub struct ArchiveBatchResult { /// Events successfully written to the store (event + scope rows). pub persisted: u32, + /// Newly-indexed `agent_metric_index` rows (valid or invalid) written in + /// this call — the count the frontend uses to decide whether an + /// agent-usage query needs to be invalidated. Distinct from `persisted`: + /// a re-ingested duplicate can be `persisted` (event/scope rows upserted, + /// no-op) without incrementing this counter, since the index row for + /// that id was already inserted by whichever earlier batch first saw it. + pub persisted_agent_metrics: u32, /// Events dropped due to access denial or invalid payload (not an error). pub dropped: u32, } @@ -666,6 +676,96 @@ pub async fn read_archived_events( .await } +// ── get_agent_usage_series ─────────────────────────────────────────────────── + +/// Compute the locally archived NIP-AM usage series for one identity/relay. +/// +/// Synchronous SQLite core of [`get_agent_usage_series`], split out so tests +/// can drive it directly against an in-memory `Connection` without a Tauri +/// `AppState`. Backfills any unindexed kind-44200 rows, repairs orphaned +/// index rows (defense in depth alongside A6's GC-time cascade), validates +/// the request, loads the window + exact-key probe rows, and hands them to +/// the pure `agent_usage::compute_series`. +fn agent_usage_series( + conn: &Connection, + identity_pk: &str, + relay_url: &str, + request: &agent_usage::AgentUsageSeriesRequest, +) -> Result { + // Fail-closed request validation happens before any SQLite work. + let agent_pubkey = agent_usage::validate_request(request)?; + + metric_store::backfill_agent_metric_index(conn, identity_pk, relay_url)?; + metric_store::repair_orphaned_metric_index_rows(conn, identity_pk, relay_url)?; + + let collection_enabled = { + let kinds_json = + store::get_subscription_kinds(conn, identity_pk, relay_url, "owner_p", identity_pk)? + .unwrap_or_else(|| "[]".to_string()); + let kinds: Vec = serde_json::from_str(&kinds_json).unwrap_or_default(); + kinds.contains(&(KIND_AGENT_TURN_METRIC as u64)) + }; + + // A13: only meaningful when the request is scoped to one author. + let has_archived_evidence = match &agent_pubkey { + None => None, + Some(pk) => Some(metric_store::has_archived_evidence( + conn, + identity_pk, + relay_url, + pk, + )?), + }; + + let start = request.bucket_boundaries[0]; + let end = *request + .bucket_boundaries + .last() + .expect("validate_request already rejected fewer than 8 boundaries"); + + let window_rows = metric_store::load_window_valid_rows( + conn, + identity_pk, + relay_url, + start, + end, + agent_pubkey.as_deref(), + )?; + let invalid_report_count = metric_store::count_invalid_rows_in_window( + conn, + identity_pk, + relay_url, + start, + end, + agent_pubkey.as_deref(), + )?; + let probe_keys = agent_usage::window_probe_keys(&window_rows); + let probe_rows = + metric_store::load_rows_at_exact_keys(conn, identity_pk, relay_url, &probe_keys)?; + + Ok(agent_usage::compute_series( + &window_rows, + &probe_rows, + invalid_report_count, + &request.bucket_boundaries, + has_archived_evidence, + collection_enabled, + )) +} + +/// Compute the locally archived NIP-AM usage series for the active identity +/// + relay (Rev 3 frozen contract). See [`agent_usage_series`] for the logic. +#[tauri::command] +pub async fn get_agent_usage_series( + state: State<'_, AppState>, + request: agent_usage::AgentUsageSeriesRequest, +) -> Result { + let identity_pk = identity_pubkey(&state)?; + let relay_url = relay_ws_url_with_override(&state); + run_archive_db_task(move |conn| agent_usage_series(conn, &identity_pk, &relay_url, &request)) + .await +} + // ── Tests ──────────────────────────────────────────────────────────────────── #[cfg(test)] diff --git a/desktop/src-tauri/src/archive/mod_tests.rs b/desktop/src-tauri/src/archive/mod_tests.rs index 288d2ab34..0addb8b6c 100644 --- a/desktop/src-tauri/src/archive/mod_tests.rs +++ b/desktop/src-tauri/src/archive/mod_tests.rs @@ -740,6 +740,10 @@ fn test_turn_metric_decrypt_success_stores_plaintext() { assert_eq!(result.persisted, 1, "event must be persisted"); assert_eq!(result.dropped, 0, "no drops on successful decrypt"); + assert_eq!( + result.persisted_agent_metrics, 1, + "one newly-indexed agent_metric_index row on first ingest" + ); // The stored raw_json must be plaintext JSON, not NIP-44 ciphertext. let raw_json: String = conn @@ -798,6 +802,226 @@ fn test_turn_metric_decrypt_fail_drops_fail_closed() { ); } +/// Re-ingesting a batch containing an already-archived kind-44200 event must +/// no-op the metric index insert: `persisted` still counts the (idempotent) +/// event/scope upsert, but `persisted_agent_metrics` must be 0 for the +/// duplicate — the row was already indexed by the first ingest (A5: this is +/// exactly the signal the frontend uses to skip a redundant query +/// invalidation). +#[test] +fn test_reingest_of_same_metric_event_does_not_double_count_persisted_agent_metrics() { + let conn = in_memory(); + let owner_keys = Keys::generate(); + let agent_keys = Keys::generate(); + let owner_pk = owner_keys.public_key().to_hex(); + let relay_url = "wss://relay.example"; + add_sub(&conn, &owner_pk, relay_url, "owner_p", &owner_pk, "[44200]"); + + let ev = make_turn_metric_event(&owner_keys, &agent_keys); + let cand1 = candidate(&ev, ScopeType::OwnerP, &owner_pk); + + let first = run_batch_sync_with_keys( + vec![cand1], + &owner_pk, + relay_url, + &conn, + vec![ev.clone()], + &owner_keys, + ); + assert_eq!(first.persisted, 1); + assert_eq!(first.persisted_agent_metrics, 1); + + // Same event re-ingested in a second batch (e.g. relay redelivery). + let cand2 = candidate(&ev, ScopeType::OwnerP, &owner_pk); + let second = run_batch_sync_with_keys( + vec![cand2], + &owner_pk, + relay_url, + &conn, + vec![ev.clone()], + &owner_keys, + ); + assert_eq!( + second.persisted, 1, + "re-ingest of a duplicate is still an accepted (idempotent) write" + ); + assert_eq!( + second.persisted_agent_metrics, 0, + "re-ingest must NOT double-count the metric index row" + ); + + let index_count: i64 = conn + .query_row("SELECT COUNT(*) FROM agent_metric_index", [], |r| r.get(0)) + .unwrap(); + assert_eq!(index_count, 1, "exactly one index row must exist total"); +} + +// ── get_agent_usage_series integration ────────────────────────────────────── +// +// Exercises `agent_usage_series` (the sync core `get_agent_usage_series` +// delegates to) end to end: ingest a real encrypted turn-metric event +// through the full archive pipeline, then read it back through the command +// core, proving backfill/indexing, collection-enabled detection, and the +// pure accounting ladder are wired together correctly — not just each in +// isolation. + +/// A freshly ingested single turn-metric event surfaces in the series with +/// its direct (delta-reliable, no-baseline) token counts, and +/// `collectionEnabled` reflects the active owner_p/44200 subscription. +#[test] +fn test_agent_usage_series_surfaces_freshly_ingested_event() { + let conn = in_memory(); + let owner_keys = Keys::generate(); + let agent_keys = Keys::generate(); + let owner_pk = owner_keys.public_key().to_hex(); + let agent_pk = agent_keys.public_key().to_hex(); + let relay_url = "wss://relay.example"; + add_sub(&conn, &owner_pk, relay_url, "owner_p", &owner_pk, "[44200]"); + + let ev = make_turn_metric_event(&owner_keys, &agent_keys); + let cand = candidate(&ev, ScopeType::OwnerP, &owner_pk); + let batch = run_batch_sync_with_keys( + vec![cand], + &owner_pk, + relay_url, + &conn, + vec![ev.clone()], + &owner_keys, + ); + assert_eq!(batch.persisted_agent_metrics, 1, "event must be indexed"); + + // `make_turn_metric_event`'s payload timestamp is 2026-07-01T00:00:00Z. + const EVENT_DAY_START: i64 = 1_782_864_000; + let boundaries: Vec = (0..=7).map(|i| EVENT_DAY_START + i * 86_400).collect(); + let request = agent_usage::AgentUsageSeriesRequest { + bucket_boundaries: boundaries, + agent_pubkey: None, + }; + + let series = agent_usage_series(&conn, &owner_pk, relay_url, &request).unwrap(); + + assert!( + series.collection_enabled, + "owner_p subscription includes kind 44200" + ); + assert_eq!(series.coverage.report_count, 1); + assert_eq!(series.coverage.invalid_report_count, 0); + assert_eq!(series.agents.len(), 1, "exactly one agent reported usage"); + let agent = &series.agents[0]; + assert_eq!(agent.agent_pubkey, agent_pk); + // No baseline row exists, so the ladder falls back to the direct + // (delta-reliable) turn values from the payload: 100/50/150. + assert_eq!(agent.usage.input_tokens.value.as_deref(), Some("100")); + assert_eq!(agent.usage.output_tokens.value.as_deref(), Some("50")); + assert_eq!(agent.usage.total_tokens.value.as_deref(), Some("150")); + assert!(!agent.usage.input_tokens.incomplete); + assert_eq!( + series.has_archived_evidence, None, + "no agentPubkey filter was supplied" + ); +} + +/// Filtering by `agentPubkey` scopes both the returned series and +/// `hasArchivedEvidence` (A13) to that one author; an unrelated agent's +/// events must not leak into either. +#[test] +fn test_agent_usage_series_filters_by_agent_pubkey_and_sets_has_archived_evidence() { + let conn = in_memory(); + let owner_keys = Keys::generate(); + let target_agent = Keys::generate(); + let other_agent = Keys::generate(); + let owner_pk = owner_keys.public_key().to_hex(); + let target_pk = target_agent.public_key().to_hex(); + let relay_url = "wss://relay.example"; + add_sub(&conn, &owner_pk, relay_url, "owner_p", &owner_pk, "[44200]"); + + let target_ev = make_turn_metric_event(&owner_keys, &target_agent); + let other_ev = make_turn_metric_event(&owner_keys, &other_agent); + let cands = vec![ + candidate(&target_ev, ScopeType::OwnerP, &owner_pk), + candidate(&other_ev, ScopeType::OwnerP, &owner_pk), + ]; + let batch = run_batch_sync_with_keys( + cands, + &owner_pk, + relay_url, + &conn, + vec![target_ev.clone(), other_ev.clone()], + &owner_keys, + ); + assert_eq!(batch.persisted_agent_metrics, 2); + + const EVENT_DAY_START: i64 = 1_782_864_000; + let boundaries: Vec = (0..=7).map(|i| EVENT_DAY_START + i * 86_400).collect(); + let request = agent_usage::AgentUsageSeriesRequest { + bucket_boundaries: boundaries, + agent_pubkey: Some(target_pk.clone()), + }; + + let series = agent_usage_series(&conn, &owner_pk, relay_url, &request).unwrap(); + + assert_eq!( + series.agents.len(), + 1, + "only the filtered agent's usage must be returned" + ); + assert_eq!(series.agents[0].agent_pubkey, target_pk); + assert_eq!( + series.has_archived_evidence, + Some(true), + "A13: evidence exists for the filtered author" + ); +} + +/// An unindexed pre-existing kind-44200 row (simulating an event archived +/// by a prior build before `agent_metric_index` existed) is picked up by +/// the command's backfill step before the window is read. +#[test] +fn test_agent_usage_series_backfills_unindexed_row_before_reading() { + let conn = in_memory(); + let owner_keys = Keys::generate(); + let agent_keys = Keys::generate(); + let owner_pk = owner_keys.public_key().to_hex(); + let relay_url = "wss://relay.example"; + + // Insert directly into `archived_events`, bypassing `commit_archive`, so + // no `agent_metric_index` row is created — the exact state a fresh + // backfill must repair. + let ev = make_turn_metric_event(&owner_keys, &agent_keys); + let plaintext = r#"{"harness":"test-harness","model":"test-model","sessionId":"sess-1","turnId":"turn-1","turnSeq":1,"timestamp":"2026-07-01T00:00:00Z","turn":{"inputTokens":100,"outputTokens":50,"totalTokens":150,"costUsd":0.001},"deltaReliable":true}"#; + store::upsert_archived_event( + &conn, + &owner_pk, + relay_url, + &ev.id.to_hex(), + 44200, + &agent_keys.public_key().to_hex(), + ev.created_at.as_secs() as i64, + plaintext, + 0, + ) + .unwrap(); + + let index_count_before: i64 = conn + .query_row("SELECT COUNT(*) FROM agent_metric_index", [], |r| r.get(0)) + .unwrap(); + assert_eq!(index_count_before, 0, "no index row before backfill"); + + const EVENT_DAY_START: i64 = 1_782_864_000; + let boundaries: Vec = (0..=7).map(|i| EVENT_DAY_START + i * 86_400).collect(); + let request = agent_usage::AgentUsageSeriesRequest { + bucket_boundaries: boundaries, + agent_pubkey: None, + }; + + let series = agent_usage_series(&conn, &owner_pk, relay_url, &request).unwrap(); + + assert_eq!( + series.coverage.report_count, 1, + "backfill must index the pre-existing row before the window read" + ); +} + // ── Real-relay integration tests ────────────────────────────────────────── // // Gated on `#[cfg(not(target_os = "windows"))]` because `build_app_state()` diff --git a/desktop/src-tauri/src/archive/pipeline.rs b/desktop/src-tauri/src/archive/pipeline.rs index 2bd149dce..98ff64dff 100644 --- a/desktop/src-tauri/src/archive/pipeline.rs +++ b/desktop/src-tauri/src/archive/pipeline.rs @@ -268,6 +268,7 @@ pub(super) fn commit_archive( conn: &Connection, ) -> Result { let mut persisted: u32 = 0; + let mut persisted_agent_metrics: u32 = 0; let mut dropped: u32 = pre_dropped; // Collect writes; count drops first, then execute inside a single @@ -399,6 +400,36 @@ pub(super) fn commit_archive( &w.scope_value, now, )?; + + // Index kind-44200 rows in the SAME transaction as the canonical + // insert (Rev 2 F5): the plaintext payload was already decrypted + // above into `w.raw_json`, so this is parse-only, no re-decrypt. + // `insert_metric_index_row`'s own `ON CONFLICT DO NOTHING` makes + // a duplicate call for an already-indexed id (e.g. re-ingest of + // a row seen in an earlier batch) a safe no-op — its `bool` + // return tells us whether this call actually inserted a new + // index row, which is exactly what `persisted_agent_metrics` + // counts (A5: newly-indexed rows, valid or invalid, not raw + // write attempts). + if w.kind == super::KIND_AGENT_TURN_METRIC as i64 { + let index_row = super::metric_store::AgentMetricIndexRow::from_payload( + &w.raw_json, + &w.eid, + &w.pubkey, + w.created_at, + now, + ); + let index_inserted = super::metric_store::insert_metric_index_row( + &tx, + identity_pk, + relay_url, + &index_row, + )?; + if index_inserted { + persisted_agent_metrics += 1; + } + } + persisted += 1; } @@ -450,5 +481,9 @@ pub(super) fn commit_archive( .map_err(|e| format!("failed to commit archive transaction: {e}"))?; } - Ok(ArchiveBatchResult { persisted, dropped }) + Ok(ArchiveBatchResult { + persisted, + persisted_agent_metrics, + dropped, + }) } diff --git a/desktop/src-tauri/src/archive/store.rs b/desktop/src-tauri/src/archive/store.rs index ae0ef92e4..d3e469c6b 100644 --- a/desktop/src-tauri/src/archive/store.rs +++ b/desktop/src-tauri/src/archive/store.rs @@ -75,6 +75,59 @@ CREATE TABLE IF NOT EXISTS archive_migrations ( name TEXT PRIMARY KEY, applied_at INTEGER NOT NULL ); + +-- Parsed index of kind 44200 (NIP-AM agent turn metric) archive rows. +-- +-- Rebuildable from `archived_events.raw_json` — never the source of truth. +-- Every archived kind-44200 row gets exactly one row here, keyed by +-- (identity, relay, id), with `parse_status` 'valid' or 'invalid'. Token +-- counters are stored as fixed-width 20-digit zero-padded decimal TEXT +-- (order-preserving lexicographically) because SQLite INTEGER is signed +-- i64 and NIP-AM counters are full-range u64. `turn_seq` uses the same +-- encoding so it survives sequence values above i64::MAX. +CREATE TABLE IF NOT EXISTS agent_metric_index ( + identity_pubkey TEXT NOT NULL, + relay_url TEXT NOT NULL, + id TEXT NOT NULL, + agent_pubkey TEXT NOT NULL, + event_created_at INTEGER NOT NULL, + archived_at INTEGER NOT NULL, + reported_at INTEGER, + session_id TEXT, + turn_seq TEXT, + model TEXT, + delta_reliable INTEGER, + turn_input_tokens TEXT, + turn_output_tokens TEXT, + turn_total_tokens TEXT, + turn_cost_usd REAL, + cumulative_input_tokens TEXT, + cumulative_output_tokens TEXT, + cumulative_total_tokens TEXT, + cumulative_cost_usd REAL, + parse_status TEXT NOT NULL CHECK (parse_status IN ('valid','invalid')), + PRIMARY KEY (identity_pubkey, relay_url, id) +); + +-- Backfill/anti-join source: scoped partial index so the per-read backfill +-- scan over archived_events only touches kind-44200 rows. +CREATE INDEX IF NOT EXISTS idx_archived_events_agent_metric + ON archived_events (identity_pubkey, relay_url, id) + WHERE kind = 44200; + +-- Predecessor/baseline lookup: exact-sequence probes (A11) and duplicate +-- cardinality checks key on this prefix. +CREATE INDEX IF NOT EXISTS idx_agent_metric_session + ON agent_metric_index (identity_pubkey, relay_url, agent_pubkey, session_id, turn_seq, id); + +-- Window scan by reported time. +CREATE INDEX IF NOT EXISTS idx_agent_metric_reported + ON agent_metric_index (identity_pubkey, relay_url, reported_at); + +-- Coarse-time scan for invalid-row coverage (invalid rows lack a trustworthy +-- reported_at, so their window membership is judged by event_created_at). +CREATE INDEX IF NOT EXISTS idx_agent_metric_created + ON agent_metric_index (identity_pubkey, relay_url, event_created_at, parse_status); "; // ── Open / init ───────────────────────────────────────────────────────────── @@ -442,6 +495,8 @@ pub fn get_subscription_kinds( /// Upsert an event row (idempotent on the PK). /// /// Does nothing if the event is already archived (same identity/relay/id). +/// Returns `true` iff this call inserted a new row (`false` if the row +/// already existed and the `ON CONFLICT DO NOTHING` no-op'd). // Args mirror the archived_events columns; a params struct would just rename them. #[allow(clippy::too_many_arguments)] pub fn upsert_archived_event( @@ -454,25 +509,26 @@ pub fn upsert_archived_event( created_at: i64, raw_json: &str, archived_at: i64, -) -> Result<(), String> { - conn.execute( - "INSERT INTO archived_events - (identity_pubkey, relay_url, id, kind, pubkey, created_at, raw_json, archived_at) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) - ON CONFLICT (identity_pubkey, relay_url, id) DO NOTHING", - params![ - identity_pubkey, - relay_url, - event_id, - kind, - pubkey, - created_at, - raw_json, - archived_at - ], - ) - .map_err(|e| format!("failed to upsert archived event: {e}"))?; - Ok(()) +) -> Result { + let affected = conn + .execute( + "INSERT INTO archived_events + (identity_pubkey, relay_url, id, kind, pubkey, created_at, raw_json, archived_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) + ON CONFLICT (identity_pubkey, relay_url, id) DO NOTHING", + params![ + identity_pubkey, + relay_url, + event_id, + kind, + pubkey, + created_at, + raw_json, + archived_at + ], + ) + .map_err(|e| format!("failed to upsert archived event: {e}"))?; + Ok(affected > 0) } /// Upsert a scope membership row for an event. @@ -771,17 +827,28 @@ pub fn read_archived_observer_events_for_channel( .map_err(|e| format!("read read_archived_observer_events_for_channel row: {e}")) } -/// GC: delete orphaned event rows whose last scope row was just removed. +/// GC: delete orphaned event rows whose last scope row was just removed, and +/// atomically cascade-delete any `agent_metric_index` rows whose canonical +/// `archived_events` row no longer exists. /// -/// Called after any batch deletion of scope rows. Uses a LEFT JOIN so only -/// events with zero remaining scope rows are deleted. +/// Both deletes run inside ONE SQLite transaction (not two autocommit +/// statements) so the derived index can never observe a canonical row as +/// gone while the index row it produced still exists — an index row must +/// never outlive the event it was parsed from (A6). +/// +/// Called after any batch deletion of scope rows. Uses a LEFT JOIN-equivalent +/// anti-join so only events with zero remaining scope rows are deleted. #[allow(dead_code)] // Used by P4 purge commands; not yet wired to a Tauri command. pub fn gc_orphaned_events( conn: &Connection, identity_pubkey: &str, relay_url: &str, ) -> Result { - let affected = conn + let tx = conn + .unchecked_transaction() + .map_err(|e| format!("failed to begin gc_orphaned_events transaction: {e}"))?; + + let affected = tx .execute( "DELETE FROM archived_events WHERE identity_pubkey = ?1 @@ -794,6 +861,11 @@ pub fn gc_orphaned_events( params![identity_pubkey, relay_url], ) .map_err(|e| format!("failed to gc orphaned events: {e}"))?; + + super::metric_store::delete_orphaned_metric_index_rows(&tx, identity_pubkey, relay_url)?; + + tx.commit() + .map_err(|e| format!("failed to commit gc_orphaned_events transaction: {e}"))?; Ok(affected) } diff --git a/desktop/src-tauri/src/archive/store_tests.rs b/desktop/src-tauri/src/archive/store_tests.rs index 0a07811e7..bbd15391e 100644 --- a/desktop/src-tauri/src/archive/store_tests.rs +++ b/desktop/src-tauri/src/archive/store_tests.rs @@ -285,9 +285,13 @@ fn test_remove_owner_p_kind_noop_when_kind_absent() { #[test] fn test_upsert_archived_event_is_idempotent() { let conn = in_memory(); - upsert_archived_event(&conn, "pk", "wss://r", "id1", 1, "author", 100, "{}", 200).unwrap(); - // Second call must not error or duplicate. - upsert_archived_event(&conn, "pk", "wss://r", "id1", 1, "author", 100, "{}", 201).unwrap(); + let first = + upsert_archived_event(&conn, "pk", "wss://r", "id1", 1, "author", 100, "{}", 200).unwrap(); + assert!(first, "first insert of a new id must report newly-inserted"); + // Second call must not error or duplicate, and must report false (no new row). + let second = + upsert_archived_event(&conn, "pk", "wss://r", "id1", 1, "author", 100, "{}", 201).unwrap(); + assert!(!second, "duplicate insert must report NOT newly-inserted"); let count: i64 = conn .query_row("SELECT COUNT(*) FROM archived_events", [], |r| r.get(0)) .unwrap(); diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index feddd4dad..979b721e6 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -876,6 +876,7 @@ pub fn run() { archive::read_archived_observer_events_for_channel, archive::index_observer_channel_id, archive::read_unindexed_observer_rows, + archive::get_agent_usage_series, is_auto_update_supported, set_window_vibrancy, ])