mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
refactor(desktop): drop agent-usage file-size overrides and dedupe date formatter
The branch had ratcheted the desktop file-size gate four times. One of those entries re-declared `src-tauri/src/lib.rs` at 1001 while a pre-existing entry at line 58 sets 1013 — the override map is a JS `Map`, so the later key won and the branch silently lowered a ceiling it never needed to touch (lib.rs is 920). Splitting at existing seams removes the other two: the M1 schema migration moves to `archive/store_migrations.rs` behind one `pub(super)` entry point, and the kind-44200 archive tests move to `archive/mod_agent_metric_tests.rs`, `#[path]`-included from `mod_tests.rs` so they keep reaching the shared fixtures through `use super::*`. `mod_tests.rs` returns to its pre-branch 1208 ceiling rather than dropping it, since that debt predates this work. `formatCoverageDate` was duplicated verbatim in both usage components; it now lives beside the other formatters in `lib/agentUsage.ts`. Co-authored-by: Will Pfleger <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
co-authored by
Will Pfleger
parent
194ceb64f5
commit
e1274be394
@@ -75,18 +75,7 @@ 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.
|
||||
// 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],
|
||||
// agent-usage-harness: M1 migration (atomic transaction wrapper, PRAGMA
|
||||
// guard, canonical-store worklist) + rebuild_metric_index_in_tx helper
|
||||
// + apply_schema_migrations + store_migration_tests.rs #[path] include.
|
||||
// Load-bearing correctness fix (Thufir IMPORTANT). M1 tests split out.
|
||||
["src-tauri/src/archive/store.rs", 1044],
|
||||
["src-tauri/src/archive/mod_tests.rs", 1208],
|
||||
// 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
|
||||
@@ -532,13 +521,6 @@ const overrides = new Map([
|
||||
// runtimeSupportsLlmProviderSelection guard on discovery provider (codex fix);
|
||||
// hideProviderIds computation for Databricks v1 gate. Queued to split.
|
||||
["src/features/agents/ui/AgentDefinitionDialog.tsx", 1035],
|
||||
// agent-usage-archive (Rev 3): Phase 2's single invoke_handler registration
|
||||
// line for get_agent_usage_series (826d79221) pre-existed under the prior
|
||||
// 1000-line ceiling; rebasing this branch onto a later main (which grew
|
||||
// lib.rs by 2 lines upstream, unrelated to this feature) tipped the file
|
||||
// to 1001. Not generic debt growth from this PR — one line, pre-existing.
|
||||
// Queued to split with the rest of this list.
|
||||
["src-tauri/src/lib.rs", 1001],
|
||||
]);
|
||||
|
||||
await runFileSizeCheck({
|
||||
|
||||
@@ -21,6 +21,7 @@ mod agent_usage;
|
||||
mod metric_store;
|
||||
mod pipeline;
|
||||
pub mod store;
|
||||
mod store_migrations;
|
||||
|
||||
use pipeline::{commit_archive, plan_archive, query_buckets};
|
||||
|
||||
|
||||
@@ -0,0 +1,410 @@
|
||||
//! Kind-44200 (NIP-AM agent turn metric) archive and `get_agent_usage_series`
|
||||
//! integration tests for `archive/mod.rs`.
|
||||
//!
|
||||
//! Kept in a sibling file so `mod_tests.rs` stays under the 1000-line gate;
|
||||
//! `#[path]`-included from there so the shared fixtures (`in_memory`,
|
||||
//! `add_sub`, `candidate`, `make_observer_frame`, `run_batch_sync_with_keys`)
|
||||
//! stay private to `mod_tests`.
|
||||
|
||||
use super::*;
|
||||
|
||||
// ── Kind-44200 agent-turn-metric archive tests ───────────────────────────
|
||||
|
||||
fn make_turn_metric_event(owner_keys: &Keys, agent_keys: &Keys) -> Event {
|
||||
use buzz_core_pkg::agent_turn_metric::{
|
||||
encrypt_agent_turn_metric, AgentTurnMetricPayload, TokenCounts,
|
||||
};
|
||||
let owner_pk = owner_keys.public_key().to_hex();
|
||||
let payload = AgentTurnMetricPayload {
|
||||
harness: "test-harness".to_string(),
|
||||
model: Some("test-model".to_string()),
|
||||
channel_id: None,
|
||||
session_id: Some("sess-1".to_string()),
|
||||
turn_id: Some("turn-1".to_string()),
|
||||
turn_seq: Some(1),
|
||||
timestamp: "2026-07-01T00:00:00Z".to_string(),
|
||||
turn: Some(TokenCounts {
|
||||
input_tokens: Some(100),
|
||||
output_tokens: Some(50),
|
||||
total_tokens: Some(150),
|
||||
cost_usd: Some(0.001),
|
||||
cache_read_tokens: None,
|
||||
cache_write_tokens: None,
|
||||
}),
|
||||
cumulative: None,
|
||||
delta_reliable: true,
|
||||
stop_reason: None,
|
||||
};
|
||||
let ciphertext =
|
||||
encrypt_agent_turn_metric(agent_keys, &owner_keys.public_key(), &payload).unwrap();
|
||||
let tags = vec![
|
||||
Tag::parse(["p", &owner_pk]).unwrap(),
|
||||
Tag::parse(["agent", &agent_keys.public_key().to_hex()]).unwrap(),
|
||||
];
|
||||
EventBuilder::new(Kind::Custom(44200), &ciphertext)
|
||||
.tags(tags)
|
||||
.sign_with_keys(agent_keys)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
/// A kind-44200 event with `owner_p` scope must route to the persistent
|
||||
/// (relay-query) path, NOT the ephemeral path.
|
||||
#[test]
|
||||
fn test_owner_p_44200_routes_to_persistent_path() {
|
||||
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";
|
||||
// Subscription for kind 44200 under owner_p.
|
||||
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 plan = plan_archive(vec![cand], &owner_pk, relay_url, &conn).unwrap();
|
||||
|
||||
// Must be in persistent buckets, NOT ephemeral list.
|
||||
assert_eq!(plan.buckets.len(), 1, "kind-44200 must land in a bucket");
|
||||
assert_eq!(
|
||||
plan.ephemeral.len(),
|
||||
0,
|
||||
"kind-44200 must NOT be on the ephemeral path"
|
||||
);
|
||||
assert_eq!(
|
||||
plan.buckets[0].scope_type_str, "owner_p",
|
||||
"bucket scope_type must be owner_p"
|
||||
);
|
||||
}
|
||||
|
||||
/// A kind-24200 event with `owner_p` scope must still route to ephemeral.
|
||||
#[test]
|
||||
fn test_owner_p_24200_still_routes_to_ephemeral() {
|
||||
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, "[24200]");
|
||||
|
||||
let ev = make_observer_frame(&owner_keys, &agent_keys, OBSERVER_FRAME_TELEMETRY);
|
||||
let cand = candidate(&ev, ScopeType::OwnerP, &owner_pk);
|
||||
|
||||
let plan = plan_archive(vec![cand], &owner_pk, relay_url, &conn).unwrap();
|
||||
|
||||
assert_eq!(
|
||||
plan.buckets.len(),
|
||||
0,
|
||||
"kind-24200 must NOT land in a bucket"
|
||||
);
|
||||
assert_eq!(
|
||||
plan.ephemeral.len(),
|
||||
1,
|
||||
"kind-24200 must be on the ephemeral path"
|
||||
);
|
||||
}
|
||||
|
||||
/// Decrypt success: plaintext payload JSON is stored, not raw ciphertext.
|
||||
#[test]
|
||||
fn test_turn_metric_decrypt_success_stores_plaintext() {
|
||||
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 cand = candidate(&ev, ScopeType::OwnerP, &owner_pk);
|
||||
let result = run_batch_sync_with_keys(
|
||||
vec![cand],
|
||||
&owner_pk,
|
||||
relay_url,
|
||||
&conn,
|
||||
vec![ev.clone()],
|
||||
&owner_keys,
|
||||
);
|
||||
|
||||
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
|
||||
.query_row("SELECT raw_json FROM archived_events", [], |r| r.get(0))
|
||||
.unwrap();
|
||||
// Plaintext JSON should be a valid object with "harness" key.
|
||||
let parsed: serde_json::Value =
|
||||
serde_json::from_str(&raw_json).expect("stored raw_json must be valid JSON");
|
||||
assert_eq!(
|
||||
parsed["harness"], "test-harness",
|
||||
"stored plaintext must decode to AgentTurnMetricPayload"
|
||||
);
|
||||
// Sanity: must NOT be the original NIP-44 ciphertext (which is not JSON).
|
||||
assert_ne!(
|
||||
raw_json, ev.content,
|
||||
"stored content must differ from original ciphertext"
|
||||
);
|
||||
}
|
||||
|
||||
/// Decrypt fail: event is dropped, nothing written to the store (fail-closed).
|
||||
#[test]
|
||||
fn test_turn_metric_decrypt_fail_drops_fail_closed() {
|
||||
let conn = in_memory();
|
||||
let owner_keys = Keys::generate();
|
||||
let wrong_keys = Keys::generate(); // wrong owner key — decrypt will fail
|
||||
let agent_keys = Keys::generate();
|
||||
let owner_pk = owner_keys.public_key().to_hex();
|
||||
let relay_url = "wss://relay.example";
|
||||
// Register subscription under owner_pk so the event passes plan-phase,
|
||||
// but use `wrong_keys` in commit so decrypt fails.
|
||||
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 result = run_batch_sync_with_keys(
|
||||
vec![cand],
|
||||
&owner_pk,
|
||||
relay_url,
|
||||
&conn,
|
||||
vec![ev.clone()],
|
||||
&wrong_keys, // wrong key → decrypt fails
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
result.persisted, 0,
|
||||
"decrypt failure must not persist the event"
|
||||
);
|
||||
assert_eq!(result.dropped, 1, "decrypt failure must count as dropped");
|
||||
|
||||
let event_count: i64 = conn
|
||||
.query_row("SELECT COUNT(*) FROM archived_events", [], |r| r.get(0))
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
event_count, 0,
|
||||
"no rows must be written to archived_events on decrypt failure"
|
||||
);
|
||||
}
|
||||
|
||||
/// 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"
|
||||
);
|
||||
}
|
||||
@@ -621,406 +621,11 @@ fn test_commit_archive_rolls_back_when_scope_write_would_fail() {
|
||||
);
|
||||
}
|
||||
|
||||
// ── Kind-44200 agent-turn-metric archive tests ───────────────────────────
|
||||
|
||||
fn make_turn_metric_event(owner_keys: &Keys, agent_keys: &Keys) -> Event {
|
||||
use buzz_core_pkg::agent_turn_metric::{
|
||||
encrypt_agent_turn_metric, AgentTurnMetricPayload, TokenCounts,
|
||||
};
|
||||
let owner_pk = owner_keys.public_key().to_hex();
|
||||
let payload = AgentTurnMetricPayload {
|
||||
harness: "test-harness".to_string(),
|
||||
model: Some("test-model".to_string()),
|
||||
channel_id: None,
|
||||
session_id: Some("sess-1".to_string()),
|
||||
turn_id: Some("turn-1".to_string()),
|
||||
turn_seq: Some(1),
|
||||
timestamp: "2026-07-01T00:00:00Z".to_string(),
|
||||
turn: Some(TokenCounts {
|
||||
input_tokens: Some(100),
|
||||
output_tokens: Some(50),
|
||||
total_tokens: Some(150),
|
||||
cost_usd: Some(0.001),
|
||||
cache_read_tokens: None,
|
||||
cache_write_tokens: None,
|
||||
}),
|
||||
cumulative: None,
|
||||
delta_reliable: true,
|
||||
stop_reason: None,
|
||||
};
|
||||
let ciphertext =
|
||||
encrypt_agent_turn_metric(agent_keys, &owner_keys.public_key(), &payload).unwrap();
|
||||
let tags = vec![
|
||||
Tag::parse(["p", &owner_pk]).unwrap(),
|
||||
Tag::parse(["agent", &agent_keys.public_key().to_hex()]).unwrap(),
|
||||
];
|
||||
EventBuilder::new(Kind::Custom(44200), &ciphertext)
|
||||
.tags(tags)
|
||||
.sign_with_keys(agent_keys)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
/// A kind-44200 event with `owner_p` scope must route to the persistent
|
||||
/// (relay-query) path, NOT the ephemeral path.
|
||||
#[test]
|
||||
fn test_owner_p_44200_routes_to_persistent_path() {
|
||||
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";
|
||||
// Subscription for kind 44200 under owner_p.
|
||||
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 plan = plan_archive(vec![cand], &owner_pk, relay_url, &conn).unwrap();
|
||||
|
||||
// Must be in persistent buckets, NOT ephemeral list.
|
||||
assert_eq!(plan.buckets.len(), 1, "kind-44200 must land in a bucket");
|
||||
assert_eq!(
|
||||
plan.ephemeral.len(),
|
||||
0,
|
||||
"kind-44200 must NOT be on the ephemeral path"
|
||||
);
|
||||
assert_eq!(
|
||||
plan.buckets[0].scope_type_str, "owner_p",
|
||||
"bucket scope_type must be owner_p"
|
||||
);
|
||||
}
|
||||
|
||||
/// A kind-24200 event with `owner_p` scope must still route to ephemeral.
|
||||
#[test]
|
||||
fn test_owner_p_24200_still_routes_to_ephemeral() {
|
||||
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, "[24200]");
|
||||
|
||||
let ev = make_observer_frame(&owner_keys, &agent_keys, OBSERVER_FRAME_TELEMETRY);
|
||||
let cand = candidate(&ev, ScopeType::OwnerP, &owner_pk);
|
||||
|
||||
let plan = plan_archive(vec![cand], &owner_pk, relay_url, &conn).unwrap();
|
||||
|
||||
assert_eq!(
|
||||
plan.buckets.len(),
|
||||
0,
|
||||
"kind-24200 must NOT land in a bucket"
|
||||
);
|
||||
assert_eq!(
|
||||
plan.ephemeral.len(),
|
||||
1,
|
||||
"kind-24200 must be on the ephemeral path"
|
||||
);
|
||||
}
|
||||
|
||||
/// Decrypt success: plaintext payload JSON is stored, not raw ciphertext.
|
||||
#[test]
|
||||
fn test_turn_metric_decrypt_success_stores_plaintext() {
|
||||
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 cand = candidate(&ev, ScopeType::OwnerP, &owner_pk);
|
||||
let result = run_batch_sync_with_keys(
|
||||
vec![cand],
|
||||
&owner_pk,
|
||||
relay_url,
|
||||
&conn,
|
||||
vec![ev.clone()],
|
||||
&owner_keys,
|
||||
);
|
||||
|
||||
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
|
||||
.query_row("SELECT raw_json FROM archived_events", [], |r| r.get(0))
|
||||
.unwrap();
|
||||
// Plaintext JSON should be a valid object with "harness" key.
|
||||
let parsed: serde_json::Value =
|
||||
serde_json::from_str(&raw_json).expect("stored raw_json must be valid JSON");
|
||||
assert_eq!(
|
||||
parsed["harness"], "test-harness",
|
||||
"stored plaintext must decode to AgentTurnMetricPayload"
|
||||
);
|
||||
// Sanity: must NOT be the original NIP-44 ciphertext (which is not JSON).
|
||||
assert_ne!(
|
||||
raw_json, ev.content,
|
||||
"stored content must differ from original ciphertext"
|
||||
);
|
||||
}
|
||||
|
||||
/// Decrypt fail: event is dropped, nothing written to the store (fail-closed).
|
||||
#[test]
|
||||
fn test_turn_metric_decrypt_fail_drops_fail_closed() {
|
||||
let conn = in_memory();
|
||||
let owner_keys = Keys::generate();
|
||||
let wrong_keys = Keys::generate(); // wrong owner key — decrypt will fail
|
||||
let agent_keys = Keys::generate();
|
||||
let owner_pk = owner_keys.public_key().to_hex();
|
||||
let relay_url = "wss://relay.example";
|
||||
// Register subscription under owner_pk so the event passes plan-phase,
|
||||
// but use `wrong_keys` in commit so decrypt fails.
|
||||
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 result = run_batch_sync_with_keys(
|
||||
vec![cand],
|
||||
&owner_pk,
|
||||
relay_url,
|
||||
&conn,
|
||||
vec![ev.clone()],
|
||||
&wrong_keys, // wrong key → decrypt fails
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
result.persisted, 0,
|
||||
"decrypt failure must not persist the event"
|
||||
);
|
||||
assert_eq!(result.dropped, 1, "decrypt failure must count as dropped");
|
||||
|
||||
let event_count: i64 = conn
|
||||
.query_row("SELECT COUNT(*) FROM archived_events", [], |r| r.get(0))
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
event_count, 0,
|
||||
"no rows must be written to archived_events on decrypt failure"
|
||||
);
|
||||
}
|
||||
|
||||
/// 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"
|
||||
);
|
||||
}
|
||||
// Kind-44200 agent-turn-metric coverage lives in a sibling file to keep this
|
||||
// one under the 1000-line gate; nested here (not in `mod.rs`) so it inherits
|
||||
// the shared fixtures above through `use super::*`.
|
||||
#[path = "mod_agent_metric_tests.rs"]
|
||||
mod agent_metric;
|
||||
|
||||
// ── Real-relay integration tests ──────────────────────────────────────────
|
||||
//
|
||||
|
||||
@@ -13,6 +13,8 @@ use std::path::Path;
|
||||
use rusqlite::{params, Connection, OptionalExtension};
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use super::store_migrations::apply_schema_migrations;
|
||||
|
||||
// ── Schema ─────────────────────────────────────────────────────────────────
|
||||
|
||||
pub(super) const SCHEMA: &str = "
|
||||
@@ -159,167 +161,6 @@ pub fn open_archive_db(path: &Path) -> Result<Connection, String> {
|
||||
Ok(conn)
|
||||
}
|
||||
|
||||
/// One-shot, idempotent schema migrations recorded in `archive_migrations`.
|
||||
///
|
||||
/// Each migration is guarded by a `SELECT 1 FROM archive_migrations WHERE
|
||||
/// name = '...'` check so it is safe to call on every open — a migration
|
||||
/// that already ran is a no-op.
|
||||
fn apply_schema_migrations(conn: &Connection) -> Result<(), String> {
|
||||
migrate_add_harness_to_metric_index(conn)
|
||||
}
|
||||
|
||||
/// M1: add `harness TEXT` column to `agent_metric_index` and rebuild index
|
||||
/// rows so pre-existing rows gain harness populated from `archived_events`.
|
||||
///
|
||||
/// The entire migration — column addition, full DELETE of stale index rows,
|
||||
/// fresh re-insertion from `archived_events.raw_json`, and the
|
||||
/// `archive_migrations` marker — is wrapped in a single `unchecked_transaction`
|
||||
/// (BEGIN DEFERRED). A crash anywhere before COMMIT leaves the DB in its
|
||||
/// pre-migration state and the marker absent, so the next open re-runs the
|
||||
/// migration from scratch.
|
||||
///
|
||||
/// The worklist is built from `archived_events` (the canonical store) rather
|
||||
/// than from `agent_metric_index` so that a partially-rebuilt or entirely
|
||||
/// empty index never causes us to forget which (identity, relay) scopes exist.
|
||||
fn migrate_add_harness_to_metric_index(conn: &Connection) -> Result<(), String> {
|
||||
let already_run: bool = conn
|
||||
.query_row(
|
||||
"SELECT COUNT(*) FROM archive_migrations WHERE name = 'add_harness_to_metric_index'",
|
||||
[],
|
||||
|r| r.get::<_, i64>(0),
|
||||
)
|
||||
.map_err(|e| format!("migration M1: guard check: {e}"))?
|
||||
> 0;
|
||||
|
||||
if already_run {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// All steps run inside one transaction so a crash before COMMIT leaves the
|
||||
// DB fully pre-migration and the marker absent — the next open re-runs.
|
||||
let tx = conn
|
||||
.unchecked_transaction()
|
||||
.map_err(|e| format!("migration M1: begin transaction: {e}"))?;
|
||||
|
||||
// 1. Add the `harness` column only when it is genuinely absent, so we
|
||||
// propagate unexpected DDL errors instead of swallowing them.
|
||||
let harness_exists: bool = {
|
||||
let mut stmt = tx
|
||||
.prepare("PRAGMA table_info(agent_metric_index)")
|
||||
.map_err(|e| format!("migration M1: PRAGMA table_info prepare: {e}"))?;
|
||||
let names: Vec<String> = stmt
|
||||
.query_map([], |row| row.get::<_, String>(1))
|
||||
.map_err(|e| format!("migration M1: PRAGMA table_info query: {e}"))?
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.map_err(|e| format!("migration M1: PRAGMA table_info read: {e}"))?;
|
||||
names.iter().any(|n| n == "harness")
|
||||
};
|
||||
if !harness_exists {
|
||||
tx.execute_batch("ALTER TABLE agent_metric_index ADD COLUMN harness TEXT")
|
||||
.map_err(|e| format!("migration M1: ALTER TABLE failed: {e}"))?;
|
||||
}
|
||||
|
||||
// 2. Collect all (identity, relay) scopes from `archived_events` — the
|
||||
// canonical store — so the worklist is correct even if the index was
|
||||
// partially rebuilt or empty before this migration runs.
|
||||
let scopes: Vec<(String, String)> = {
|
||||
let mut stmt = tx
|
||||
.prepare(
|
||||
"SELECT DISTINCT identity_pubkey, relay_url
|
||||
FROM archived_events
|
||||
WHERE kind = 44200",
|
||||
)
|
||||
.map_err(|e| format!("migration M1: prepare scope query: {e}"))?;
|
||||
let result = stmt
|
||||
.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
|
||||
.map_err(|e| format!("migration M1: query scopes: {e}"))?
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.map_err(|e| format!("migration M1: read scopes: {e}"))?;
|
||||
result
|
||||
};
|
||||
|
||||
if !scopes.is_empty() {
|
||||
// 3. Delete all stale index rows; re-insert them with harness populated
|
||||
// from the shared `from_payload` parser below. Both the DELETE and
|
||||
// all inserts run inside the same outer transaction — no intermediate
|
||||
// commit, so readers never observe an empty index.
|
||||
tx.execute_batch("DELETE FROM agent_metric_index")
|
||||
.map_err(|e| format!("migration M1: delete index rows: {e}"))?;
|
||||
|
||||
for (identity, relay) in &scopes {
|
||||
rebuild_metric_index_in_tx(&tx, identity, relay)?;
|
||||
}
|
||||
}
|
||||
|
||||
// 4. Record migration as applied — written last so the marker is only
|
||||
// present in a fully committed transaction.
|
||||
tx.execute(
|
||||
"INSERT OR IGNORE INTO archive_migrations (name, applied_at) \
|
||||
VALUES ('add_harness_to_metric_index', ?1)",
|
||||
params![std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs() as i64],
|
||||
)
|
||||
.map_err(|e| format!("migration M1: record marker: {e}"))?;
|
||||
|
||||
tx.commit()
|
||||
.map_err(|e| format!("migration M1: commit: {e}"))?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Re-insert all kind-44200 rows for `(identity, relay)` from
|
||||
/// `archived_events.raw_json` through the shared `from_payload` parser.
|
||||
///
|
||||
/// Unlike `metric_store::backfill_agent_metric_index`, this function runs
|
||||
/// entirely inside the caller's transaction — no nested `BEGIN`/`COMMIT`.
|
||||
/// It is used exclusively by M1 where the outer transaction provides atomicity.
|
||||
fn rebuild_metric_index_in_tx(
|
||||
conn: &Connection,
|
||||
identity_pubkey: &str,
|
||||
relay_url: &str,
|
||||
) -> Result<(), String> {
|
||||
let mut stmt = conn
|
||||
.prepare(
|
||||
"SELECT id, pubkey, created_at, archived_at, raw_json
|
||||
FROM archived_events
|
||||
WHERE identity_pubkey = ?1
|
||||
AND relay_url = ?2
|
||||
AND kind = 44200",
|
||||
)
|
||||
.map_err(|e| format!("migration M1: prepare rebuild select: {e}"))?;
|
||||
|
||||
let rows: Vec<(String, String, i64, i64, String)> = stmt
|
||||
.query_map(params![identity_pubkey, relay_url], |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!("migration M1: query rebuild rows: {e}"))?
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.map_err(|e| format!("migration M1: read rebuild row: {e}"))?;
|
||||
drop(stmt);
|
||||
|
||||
for (id, pubkey, created_at, archived_at, raw_json) in &rows {
|
||||
let parsed = super::metric_store::AgentMetricIndexRow::from_payload(
|
||||
raw_json,
|
||||
id,
|
||||
pubkey,
|
||||
*created_at,
|
||||
*archived_at,
|
||||
);
|
||||
super::metric_store::insert_metric_index_row(conn, identity_pubkey, relay_url, &parsed)
|
||||
.map_err(|e| format!("migration M1: insert rebuilt row: {e}"))?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn set_wal_mode(conn: &Connection) -> Result<(), String> {
|
||||
let deadline = Instant::now() + Duration::from_secs(5);
|
||||
loop {
|
||||
|
||||
@@ -0,0 +1,171 @@
|
||||
//! One-shot schema migrations for the archive database.
|
||||
//!
|
||||
//! Applied by `store::open_archive_db` on every open and guarded by markers in
|
||||
//! `archive_migrations`, so a migration that already ran is a no-op.
|
||||
//!
|
||||
//! Kept in a sibling file (not `store.rs`) to keep that file under the
|
||||
//! 1000-line gate, per the existing `metric_store.rs` / `pipeline.rs`
|
||||
//! precedent.
|
||||
|
||||
use rusqlite::{params, Connection};
|
||||
|
||||
/// One-shot, idempotent schema migrations recorded in `archive_migrations`.
|
||||
///
|
||||
/// Each migration is guarded by a `SELECT 1 FROM archive_migrations WHERE
|
||||
/// name = '...'` check so it is safe to call on every open — a migration
|
||||
/// that already ran is a no-op.
|
||||
pub(super) fn apply_schema_migrations(conn: &Connection) -> Result<(), String> {
|
||||
migrate_add_harness_to_metric_index(conn)
|
||||
}
|
||||
|
||||
/// M1: add `harness TEXT` column to `agent_metric_index` and rebuild index
|
||||
/// rows so pre-existing rows gain harness populated from `archived_events`.
|
||||
///
|
||||
/// The entire migration — column addition, full DELETE of stale index rows,
|
||||
/// fresh re-insertion from `archived_events.raw_json`, and the
|
||||
/// `archive_migrations` marker — is wrapped in a single `unchecked_transaction`
|
||||
/// (BEGIN DEFERRED). A crash anywhere before COMMIT leaves the DB in its
|
||||
/// pre-migration state and the marker absent, so the next open re-runs the
|
||||
/// migration from scratch.
|
||||
///
|
||||
/// The worklist is built from `archived_events` (the canonical store) rather
|
||||
/// than from `agent_metric_index` so that a partially-rebuilt or entirely
|
||||
/// empty index never causes us to forget which (identity, relay) scopes exist.
|
||||
fn migrate_add_harness_to_metric_index(conn: &Connection) -> Result<(), String> {
|
||||
let already_run: bool = conn
|
||||
.query_row(
|
||||
"SELECT COUNT(*) FROM archive_migrations WHERE name = 'add_harness_to_metric_index'",
|
||||
[],
|
||||
|r| r.get::<_, i64>(0),
|
||||
)
|
||||
.map_err(|e| format!("migration M1: guard check: {e}"))?
|
||||
> 0;
|
||||
|
||||
if already_run {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// All steps run inside one transaction so a crash before COMMIT leaves the
|
||||
// DB fully pre-migration and the marker absent — the next open re-runs.
|
||||
let tx = conn
|
||||
.unchecked_transaction()
|
||||
.map_err(|e| format!("migration M1: begin transaction: {e}"))?;
|
||||
|
||||
// 1. Add the `harness` column only when it is genuinely absent, so we
|
||||
// propagate unexpected DDL errors instead of swallowing them.
|
||||
let harness_exists: bool = {
|
||||
let mut stmt = tx
|
||||
.prepare("PRAGMA table_info(agent_metric_index)")
|
||||
.map_err(|e| format!("migration M1: PRAGMA table_info prepare: {e}"))?;
|
||||
let names: Vec<String> = stmt
|
||||
.query_map([], |row| row.get::<_, String>(1))
|
||||
.map_err(|e| format!("migration M1: PRAGMA table_info query: {e}"))?
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.map_err(|e| format!("migration M1: PRAGMA table_info read: {e}"))?;
|
||||
names.iter().any(|n| n == "harness")
|
||||
};
|
||||
if !harness_exists {
|
||||
tx.execute_batch("ALTER TABLE agent_metric_index ADD COLUMN harness TEXT")
|
||||
.map_err(|e| format!("migration M1: ALTER TABLE failed: {e}"))?;
|
||||
}
|
||||
|
||||
// 2. Collect all (identity, relay) scopes from `archived_events` — the
|
||||
// canonical store — so the worklist is correct even if the index was
|
||||
// partially rebuilt or empty before this migration runs.
|
||||
let scopes: Vec<(String, String)> = {
|
||||
let mut stmt = tx
|
||||
.prepare(
|
||||
"SELECT DISTINCT identity_pubkey, relay_url
|
||||
FROM archived_events
|
||||
WHERE kind = 44200",
|
||||
)
|
||||
.map_err(|e| format!("migration M1: prepare scope query: {e}"))?;
|
||||
let result = stmt
|
||||
.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
|
||||
.map_err(|e| format!("migration M1: query scopes: {e}"))?
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.map_err(|e| format!("migration M1: read scopes: {e}"))?;
|
||||
result
|
||||
};
|
||||
|
||||
if !scopes.is_empty() {
|
||||
// 3. Delete all stale index rows; re-insert them with harness populated
|
||||
// from the shared `from_payload` parser below. Both the DELETE and
|
||||
// all inserts run inside the same outer transaction — no intermediate
|
||||
// commit, so readers never observe an empty index.
|
||||
tx.execute_batch("DELETE FROM agent_metric_index")
|
||||
.map_err(|e| format!("migration M1: delete index rows: {e}"))?;
|
||||
|
||||
for (identity, relay) in &scopes {
|
||||
rebuild_metric_index_in_tx(&tx, identity, relay)?;
|
||||
}
|
||||
}
|
||||
|
||||
// 4. Record migration as applied — written last so the marker is only
|
||||
// present in a fully committed transaction.
|
||||
tx.execute(
|
||||
"INSERT OR IGNORE INTO archive_migrations (name, applied_at) \
|
||||
VALUES ('add_harness_to_metric_index', ?1)",
|
||||
params![std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs() as i64],
|
||||
)
|
||||
.map_err(|e| format!("migration M1: record marker: {e}"))?;
|
||||
|
||||
tx.commit()
|
||||
.map_err(|e| format!("migration M1: commit: {e}"))?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Re-insert all kind-44200 rows for `(identity, relay)` from
|
||||
/// `archived_events.raw_json` through the shared `from_payload` parser.
|
||||
///
|
||||
/// Unlike `metric_store::backfill_agent_metric_index`, this function runs
|
||||
/// entirely inside the caller's transaction — no nested `BEGIN`/`COMMIT`.
|
||||
/// It is used exclusively by M1 where the outer transaction provides atomicity.
|
||||
fn rebuild_metric_index_in_tx(
|
||||
conn: &Connection,
|
||||
identity_pubkey: &str,
|
||||
relay_url: &str,
|
||||
) -> Result<(), String> {
|
||||
let mut stmt = conn
|
||||
.prepare(
|
||||
"SELECT id, pubkey, created_at, archived_at, raw_json
|
||||
FROM archived_events
|
||||
WHERE identity_pubkey = ?1
|
||||
AND relay_url = ?2
|
||||
AND kind = 44200",
|
||||
)
|
||||
.map_err(|e| format!("migration M1: prepare rebuild select: {e}"))?;
|
||||
|
||||
let rows: Vec<(String, String, i64, i64, String)> = stmt
|
||||
.query_map(params![identity_pubkey, relay_url], |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!("migration M1: query rebuild rows: {e}"))?
|
||||
.collect::<Result<Vec<_>, _>>()
|
||||
.map_err(|e| format!("migration M1: read rebuild row: {e}"))?;
|
||||
drop(stmt);
|
||||
|
||||
for (id, pubkey, created_at, archived_at, raw_json) in &rows {
|
||||
let parsed = super::metric_store::AgentMetricIndexRow::from_payload(
|
||||
raw_json,
|
||||
id,
|
||||
pubkey,
|
||||
*created_at,
|
||||
*archived_at,
|
||||
);
|
||||
super::metric_store::insert_metric_index_row(conn, identity_pubkey, relay_url, &parsed)
|
||||
.map_err(|e| format!("migration M1: insert rebuilt row: {e}"))?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -5,6 +5,7 @@ import {
|
||||
bigintRatio,
|
||||
buildLocalDayBoundaries,
|
||||
deriveUsageIngressTrailing,
|
||||
formatCoverageDate,
|
||||
formatEstimatedCostUsd,
|
||||
formatTokenCountCompact,
|
||||
formatTokenCountExact,
|
||||
@@ -300,6 +301,21 @@ test("formatEstimatedCostUsd renders two-decimal USD currency", () => {
|
||||
assert.equal(formatEstimatedCostUsd(0), "$0.00");
|
||||
});
|
||||
|
||||
// ── formatCoverageDate ───────────────────────────────────────────────────────
|
||||
|
||||
test("formatCoverageDate renders unknown for a missing timestamp", () => {
|
||||
assert.equal(formatCoverageDate(null), "unknown");
|
||||
});
|
||||
|
||||
test("formatCoverageDate renders a timestamp as its local month and day without a year", () => {
|
||||
const unixSeconds = 1_737_849_600; // 2025-01-26T00:00:00Z
|
||||
const localDate = new Date(unixSeconds * 1000);
|
||||
const formatted = formatCoverageDate(unixSeconds);
|
||||
|
||||
assert.match(formatted, new RegExp(`\\b${localDate.getDate()}\\b`));
|
||||
assert.doesNotMatch(formatted, new RegExp(`${localDate.getFullYear()}`));
|
||||
});
|
||||
|
||||
// ── bigintRatio ──────────────────────────────────────────────────────────────
|
||||
|
||||
test("bigintRatio computes a bounded ratio without losing bigint precision on large magnitudes", () => {
|
||||
|
||||
@@ -165,6 +165,15 @@ export function formatEstimatedCostUsd(value: number): string {
|
||||
}).format(value);
|
||||
}
|
||||
|
||||
/** Short coverage-date display, e.g. `1737849600` -> "Jan 25". `null` renders "unknown". */
|
||||
export function formatCoverageDate(unixSeconds: number | null): string {
|
||||
if (unixSeconds === null) return "unknown";
|
||||
return new Date(unixSeconds * 1000).toLocaleDateString(undefined, {
|
||||
month: "short",
|
||||
day: "numeric",
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Bigint-safe ratio in `[0, 1]` for a relative bar, e.g. `part` tokens against
|
||||
* `whole` tokens. Never converts the full magnitude through `Number(...)`;
|
||||
|
||||
@@ -14,6 +14,7 @@ import { Skeleton } from "@/shared/ui/skeleton";
|
||||
import { Tabs, TabsList, TabsTrigger } from "@/shared/ui/tabs";
|
||||
import { useAgentUsageSeries } from "../hooks";
|
||||
import {
|
||||
formatCoverageDate,
|
||||
formatEstimatedCostUsd,
|
||||
formatTokenCountCompact,
|
||||
formatTokenCountExact,
|
||||
@@ -331,14 +332,6 @@ function UsageStat({
|
||||
);
|
||||
}
|
||||
|
||||
function formatCoverageDate(unixSeconds: number | null): string {
|
||||
if (unixSeconds === null) return "unknown";
|
||||
return new Date(unixSeconds * 1000).toLocaleDateString(undefined, {
|
||||
month: "short",
|
||||
day: "numeric",
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Human-readable coverage range for the focused view's footer, from the
|
||||
* exact first/last reported timestamps the backend already computes
|
||||
|
||||
@@ -21,6 +21,7 @@ import { Tabs, TabsList, TabsTrigger } from "@/shared/ui/tabs";
|
||||
import { useAgentUsageSeries } from "../hooks";
|
||||
import {
|
||||
bigintRatio,
|
||||
formatCoverageDate,
|
||||
formatTokenCountCompact,
|
||||
isPartialField,
|
||||
isUnknownField,
|
||||
@@ -245,14 +246,6 @@ function AgentUsageCard({
|
||||
);
|
||||
}
|
||||
|
||||
function formatCoverageDate(unixSeconds: number | null): string {
|
||||
if (unixSeconds === null) return "unknown";
|
||||
return new Date(unixSeconds * 1000).toLocaleDateString(undefined, {
|
||||
month: "short",
|
||||
day: "numeric",
|
||||
});
|
||||
}
|
||||
|
||||
function AgentUsageRow({
|
||||
agent,
|
||||
days,
|
||||
|
||||
Reference in New Issue
Block a user