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 <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7
2026-07-20 23:20:29 -07:00
co-authored by Will Pfleger
parent 8a65d81fab
commit c537cbc1b8
11 changed files with 2880 additions and 28 deletions
+7 -1
View File
@@ -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
@@ -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<i64>,
pub agent_pubkey: Option<String>,
}
/// 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<Option<String>, 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<String>,
pub incomplete: bool,
}
#[derive(Debug, Clone, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct CostField {
pub value: Option<f64>,
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<String>,
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<SeriesBucket>,
pub models: Vec<ModelUsage>,
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<i64>,
pub last_archived_at: Option<i64>,
pub first_reported_at: Option<i64>,
pub last_reported_at: Option<i64>,
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<SeriesBucket>,
pub agents: Vec<AgentUsage>,
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<bool>,
}
// ── Per-event field ladder (A1, A4, A11, A12) ───────────────────────────────
#[derive(Debug, Clone, Copy)]
enum FieldValue<T> {
Known(T),
Unknown,
}
struct EventOutcome {
input: FieldValue<u64>,
output: FieldValue<u64>,
total: FieldValue<u64>,
cost: FieldValue<f64>,
}
/// 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<u64>,
current_cumulative: Option<u64>,
current_turn: Option<u64>,
delta_reliable: bool,
) -> FieldValue<u64> {
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<f64>,
current_cumulative: Option<f64>,
current_turn: Option<f64>,
delta_reliable: bool,
) -> FieldValue<f64> {
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<u64>,
incomplete: bool,
overflowed: bool,
}
impl TokenAccumulator {
fn add(&mut self, v: FieldValue<u64>) {
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<f64>,
incomplete: bool,
overflowed: bool,
}
impl CostAccumulator {
fn add(&mut self, v: FieldValue<f64>) {
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<u64> {
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<usize> {
(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<UsageAccumulator>,
bucket_counts: Vec<i64>,
total: UsageAccumulator,
report_count: i64,
models: HashMap<Option<String>, (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<bool>,
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<UsageAccumulator> = (0..bucket_count)
.map(|_| UsageAccumulator::default())
.collect();
let mut overall_bucket_counts: Vec<i64> = vec![0; bucket_count];
let mut agents: HashMap<String, AgentScope> = HashMap::new();
let mut first_reported_at: Option<i64> = None;
let mut last_reported_at: Option<i64> = None;
let mut first_archived_at: Option<i64> = None;
let mut last_archived_at: Option<i64> = 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<SeriesBucket> = 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<u64>, 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<SeriesBucket> = 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<u64>, Option<String>, 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<AgentUsage> = 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;
@@ -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<i64> {
(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<i64>, agent_pubkey: Option<String>) -> 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<i64> = (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<i64> = (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());
}
@@ -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<u64> {
if text.len() != U64_SORTABLE_WIDTH {
return None;
}
text.parse::<u64>().ok()
}
fn parse_rfc3339_secs(timestamp: &str) -> Option<i64> {
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<i64>,
pub session_id: Option<String>,
pub turn_seq: Option<u64>,
pub model: Option<String>,
pub delta_reliable: Option<bool>,
pub turn_input_tokens: Option<u64>,
pub turn_output_tokens: Option<u64>,
pub turn_total_tokens: Option<u64>,
pub turn_cost_usd: Option<f64>,
pub cumulative_input_tokens: Option<u64>,
pub cumulative_output_tokens: Option<u64>,
pub cumulative_total_tokens: Option<u64>,
pub cumulative_cost_usd: Option<f64>,
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::<AgentTurnMetricPayload>(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<AgentMetricIndexRow> {
let turn_seq_text: Option<String> = row.get("turn_seq")?;
let turn_input_text: Option<String> = row.get("turn_input_tokens")?;
let turn_output_text: Option<String> = row.get("turn_output_tokens")?;
let turn_total_text: Option<String> = row.get("turn_total_tokens")?;
let cum_input_text: Option<String> = row.get("cumulative_input_tokens")?;
let cum_output_text: Option<String> = row.get("cumulative_output_tokens")?;
let cum_total_text: Option<String> = row.get("cumulative_total_tokens")?;
let delta_reliable_int: Option<i64> = 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<bool, String> {
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<usize, String> {
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::<Result<Vec<_>, _>>()
.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<usize, String> {
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<usize, String> {
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<Vec<AgentMetricIndexRow>, 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::<Result<Vec<_>, _>>()
.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<i64, String> {
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<Vec<AgentMetricIndexRow>, String> {
use std::collections::HashMap;
// Group by (agent, session) so each group becomes one IN-list query.
let mut groups: HashMap<(String, String), Vec<u64>> = 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<String> = 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::<Vec<_>>()
.join(",")
);
let mut stmt = stmt_prepare(conn, &sql)?;
let mut bound: Vec<Box<dyn rusqlite::ToSql>> = 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<bool, String> {
let exists: Option<i64> = 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<rusqlite::Statement<'a>, 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;
@@ -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}"
);
}
+100
View File
@@ -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<agent_usage::AgentUsageSeries, String> {
// 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<u64> = 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<agent_usage::AgentUsageSeries, String> {
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)]
+224
View File
@@ -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<i64> = (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<i64> = (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<i64> = (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()`
+36 -1
View File
@@ -268,6 +268,7 @@ pub(super) fn commit_archive(
conn: &Connection,
) -> Result<ArchiveBatchResult, String> {
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,
})
}
+95 -23
View File
@@ -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<bool, String> {
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<usize, String> {
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)
}
+7 -3
View File
@@ -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();
+1
View File
@@ -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,
])