From 4f413192208f62cd001d5871013d49820c1ca9a3 Mon Sep 17 00:00:00 2001 From: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 Date: Thu, 16 Jul 2026 15:46:28 -0500 Subject: [PATCH] feat(desktop): wire persistedAgentMetrics notifier through archive sync MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Completes Phase 2 of the NIP-AM local agent usage feature. tauriArchive decodes the backend's persistedAgentMetrics count defensively (missing field -> 0, per A5/M2's backward-compatible wire contract) and exposes an onAgentMetricsChanged notifier, fired only when a subscription mutation touches kind 44200 specifically. ArchiveSyncManager routes both its idle/size-triggered and destroy-time flushes through one sendBatch helper that notifies usage listeners only when the backend reports persistedAgentMetrics > 0 — rejections, duplicate-only batches, and non-metric successes never notify, keeping the refresh signal backend-authoritative. Co-authored-by: Will Pfleger Signed-off-by: Will Pfleger --- .../local-archive/archiveSyncManager.test.mjs | 185 +++++++++++++++++- .../local-archive/archiveSyncManager.ts | 38 +++- desktop/src/shared/api/tauriArchive.ts | 60 +++++- 3 files changed, 274 insertions(+), 9 deletions(-) diff --git a/desktop/src/features/local-archive/archiveSyncManager.test.mjs b/desktop/src/features/local-archive/archiveSyncManager.test.mjs index 0bcefa8e6..a7feaac8e 100644 --- a/desktop/src/features/local-archive/archiveSyncManager.test.mjs +++ b/desktop/src/features/local-archive/archiveSyncManager.test.mjs @@ -1,5 +1,6 @@ import assert from "node:assert/strict"; import test from "node:test"; +import { onAgentMetricsChanged } from "@/shared/api/tauriArchive"; import { ArchiveSyncManager } from "./archiveSyncManager.ts"; // ── Fakes ──────────────────────────────────────────────────────────────────── @@ -41,6 +42,7 @@ function makeFakeArchive() { let subs = []; const archiveCalls = []; const listeners = new Set(); + let nextPersistedAgentMetrics = 0; return { async listSaveSubscriptions() { @@ -81,7 +83,13 @@ function makeFakeArchive() { }, async archiveEvents(candidates) { archiveCalls.push(candidates); - return { persisted: candidates.length, dropped: 0 }; + const persistedAgentMetrics = nextPersistedAgentMetrics; + nextPersistedAgentMetrics = 0; + return { + persisted: candidates.length, + persistedAgentMetrics, + dropped: 0, + }; }, onSubscriptionChange(listener) { listeners.add(listener); @@ -91,6 +99,10 @@ function makeFakeArchive() { setSubs(s) { subs = s; }, + /** Next archiveEvents() call reports this many newly persisted agent metrics. */ + setNextPersistedAgentMetrics(n) { + nextPersistedAgentMetrics = n; + }, }; } @@ -803,6 +815,177 @@ test("manager_handles_out_of_order_list_resolution", async () => { ); }); +test("manager_notifies_agent_metrics_changed_when_persisted_agent_metrics_positive", async () => { + const relay = makeFakeRelayClient(); + const archive = makeFakeArchive(); + archive.setSubs([ + { + scopeType: "owner_p", + scopeValue: "agent-pk", + kinds: [44200], + identityPubkey: "pk", + relayUrl: "wss://r", + createdAt: 0, + }, + ]); + const mgr = makeManager(relay, archive, { flushBatchSize: 1 }); + await mgr.start(); + await tick(); + + let fired = 0; + const off = onAgentMetricsChanged(() => { + fired++; + }); + + archive.setNextPersistedAgentMetrics(1); + const filter = JSON.parse([...relay.subs.keys()][0]); + relay.push(filter, { + id: "metric-1", + kind: 44200, + pubkey: "agent-pk", + created_at: 1, + content: "encrypted", + tags: [], + }); + await tick(); + + assert.equal(fired, 1, "notifier fires once when persistedAgentMetrics > 0"); + off(); + mgr.destroy(); +}); + +test("manager_does_not_notify_agent_metrics_changed_when_persisted_agent_metrics_zero", async () => { + const relay = makeFakeRelayClient(); + const archive = makeFakeArchive(); + archive.setSubs([ + { + scopeType: "channel_h", + scopeValue: "chan-1", + kinds: [9], + identityPubkey: "pk", + relayUrl: "wss://r", + createdAt: 0, + }, + ]); + const mgr = makeManager(relay, archive, { flushBatchSize: 1 }); + await mgr.start(); + await tick(); + + let fired = 0; + const off = onAgentMetricsChanged(() => { + fired++; + }); + + // Non-metric event; archiveEvents defaults to persistedAgentMetrics: 0. + const filter = JSON.parse([...relay.subs.keys()][0]); + relay.push(filter, { + id: "ev1", + kind: 9, + pubkey: "pk", + created_at: 1, + content: "hi", + tags: [], + }); + await tick(); + + assert.equal(fired, 0, "notifier must not fire for persistedAgentMetrics: 0"); + off(); + mgr.destroy(); +}); + +test("manager_does_not_notify_agent_metrics_changed_on_archive_events_rejection", async () => { + const relay = makeFakeRelayClient(); + const archive = makeFakeArchive(); + archive.setSubs([ + { + scopeType: "owner_p", + scopeValue: "agent-pk", + kinds: [44200], + identityPubkey: "pk", + relayUrl: "wss://r", + createdAt: 0, + }, + ]); + const mgr = makeManager(relay, archive, { + flushBatchSize: 1, + archiveEvents: async () => { + throw new Error("simulated archive_events failure"); + }, + }); + await mgr.start(); + await tick(); + + let fired = 0; + const off = onAgentMetricsChanged(() => { + fired++; + }); + + const filter = JSON.parse([...relay.subs.keys()][0]); + relay.push(filter, { + id: "metric-1", + kind: 44200, + pubkey: "agent-pk", + created_at: 1, + content: "encrypted", + tags: [], + }); + await tick(); + + assert.equal(fired, 0, "a rejected archiveEvents call must never notify"); + off(); + mgr.destroy(); +}); + +test("manager_notifies_agent_metrics_changed_on_destroy_flush", async () => { + const relay = makeFakeRelayClient(); + const archive = makeFakeArchive(); + archive.setSubs([ + { + scopeType: "owner_p", + scopeValue: "agent-pk", + kinds: [44200], + identityPubkey: "pk", + relayUrl: "wss://r", + createdAt: 0, + }, + ]); + // Large batch size / idle so the event stays buffered until destroy() flushes it. + const mgr = makeManager(relay, archive, { + flushBatchSize: 100, + flushIdleMs: 10000, + }); + await mgr.start(); + await tick(); + + let fired = 0; + const off = onAgentMetricsChanged(() => { + fired++; + }); + + archive.setNextPersistedAgentMetrics(1); + const filter = JSON.parse([...relay.subs.keys()][0]); + relay.push(filter, { + id: "metric-1", + kind: 44200, + pubkey: "agent-pk", + created_at: 1, + content: "encrypted", + tags: [], + }); + assert.equal(archive.archiveCalls.length, 0, "buffered, not yet flushed"); + + mgr.destroy(); + await tick(); + + assert.equal(archive.archiveCalls.length, 1, "destroy() flushes the buffer"); + assert.equal( + fired, + 1, + "destroy flush notifies when persistedAgentMetrics > 0", + ); + off(); +}); + test("manager_flushes_buffer_on_destroy", async () => { const relay = makeFakeRelayClient(); const archive = makeFakeArchive(); diff --git a/desktop/src/features/local-archive/archiveSyncManager.ts b/desktop/src/features/local-archive/archiveSyncManager.ts index c97364169..3672989a7 100644 --- a/desktop/src/features/local-archive/archiveSyncManager.ts +++ b/desktop/src/features/local-archive/archiveSyncManager.ts @@ -4,7 +4,9 @@ import type { RelayEvent } from "@/shared/api/types"; import { archiveEvents as defaultArchiveEvents, listSaveSubscriptions as defaultListSaveSubscriptions, + notifyAgentMetricsChanged, onSubscriptionChange as defaultOnSubscriptionChange, + type ArchiveBatchResult, type SaveSubscription, type ScopeType, } from "@/shared/api/tauriArchive"; @@ -30,7 +32,7 @@ export interface ArchiveSyncDeps { rawEventJson: string; matchedScope: { scopeType: ScopeType; scopeValue: string }; }>, - ) => Promise; + ) => Promise; onSubscriptionChange: (listener: () => void) => () => void; flushBatchSize?: number; flushIdleMs?: number; @@ -135,9 +137,7 @@ export class ArchiveSyncManager { // Flush any buffered events before tearing down. if (this.buffer.length > 0) { const toFlush = this.buffer.splice(0); - void this.deps.archiveEvents(toFlush).catch((err: unknown) => { - console.warn("[archiveSyncManager] flush on destroy failed:", err); - }); + this.sendBatch(toFlush, "flush on destroy failed"); } for (const [, unsub] of this.active) { void unsub(); @@ -292,9 +292,33 @@ export class ArchiveSyncManager { } if (this.buffer.length === 0) return; const batch = this.buffer.splice(0); - void this.deps.archiveEvents(batch).catch((err: unknown) => { - console.warn("[archiveSyncManager] archive_events failed:", err); - }); + this.sendBatch(batch, "archive_events failed"); + } + + /** + * Fire-and-forget `archiveEvents(batch)`, shared by the idle/size-triggered + * flush and the destroy-time flush. Notifies `onAgentMetricsChanged` + * subscribers only when the backend confirms `persistedAgentMetrics > 0` — + * the backend is authoritative, so a rejected call, a duplicate-only batch, + * or a batch with no kind-44200 events never notifies. + */ + private sendBatch( + batch: Array<{ + rawEventJson: string; + matchedScope: { scopeType: ScopeType; scopeValue: string }; + }>, + errLabel: string, + ): void { + void this.deps + .archiveEvents(batch) + .then((result) => { + if (result.persistedAgentMetrics > 0) { + notifyAgentMetricsChanged(); + } + }) + .catch((err: unknown) => { + console.warn(`[archiveSyncManager] ${errLabel}:`, err); + }); } } diff --git a/desktop/src/shared/api/tauriArchive.ts b/desktop/src/shared/api/tauriArchive.ts index 0e4be6c9a..dc0b704d6 100644 --- a/desktop/src/shared/api/tauriArchive.ts +++ b/desktop/src/shared/api/tauriArchive.ts @@ -1,3 +1,5 @@ +import { KIND_AGENT_TURN_METRIC } from "@/shared/constants/kinds"; + import { invokeTauri } from "./tauri"; // ── Wire-shape types (raw Tauri responses) ─────────────────────────────────── @@ -32,9 +34,33 @@ export type SaveSubscription = { export type ArchiveBatchResult = { persisted: number; + /** + * Newly-indexed `agent_metric_index` rows (valid or invalid) written by + * this call. A re-ingested duplicate event does not increment this even + * when `persisted` counts it, because the index row for that id was + * already written by whichever earlier batch first saw it. Missing on + * the wire (older/mocked responses) decodes as `0` — see + * `decodeArchiveBatchResult`. + */ + persistedAgentMetrics: number; dropped: number; }; +/** + * Rust sends camelCase (`#[serde(rename_all = "camelCase")]` on + * `ArchiveBatchResult`), but decode defensively rather than trust every + * caller (including mocks/tests) to supply every field. + */ +function decodeArchiveBatchResult( + raw: Partial, +): ArchiveBatchResult { + return { + persisted: raw.persisted ?? 0, + persistedAgentMetrics: raw.persistedAgentMetrics ?? 0, + dropped: raw.dropped ?? 0, + }; +} + // ── Subscription-change notifier ───────────────────────────────────────────── /** @@ -55,6 +81,28 @@ function notifySubscriptionChange(): void { } } +// ── Agent-metrics-change notifier ──────────────────────────────────────────── + +/** + * Module-level notifier for newly persisted agent turn metrics (kind 44200). + * `useAgentUsageSeries` subscribes to this to invalidate its query without + * polling. Fired only when the backend confirms `persistedAgentMetrics > 0` + * for a successful `archiveEvents` call, or when a kind-44200 subscription + * mutation succeeds (`collectionEnabled` is part of the usage query result). + */ +const agentMetricsChangeListeners = new Set<() => void>(); + +export function onAgentMetricsChanged(listener: () => void): () => void { + agentMetricsChangeListeners.add(listener); + return () => agentMetricsChangeListeners.delete(listener); +} + +export function notifyAgentMetricsChanged(): void { + for (const listener of agentMetricsChangeListeners) { + listener(); + } +} + // ── Decoder ────────────────────────────────────────────────────────────────── function decodeRawSubscription(raw: RawSaveSubscription): SaveSubscription { @@ -124,6 +172,12 @@ export async function agentMetricArchiveDefaultEnabled(): Promise { export async function mergeSaveSubscriptionKinds(kind: number): Promise { await invokeTauri("merge_save_subscription_kinds", { kind }); notifySubscriptionChange(); + // `collectionEnabled` is part of the usage query result — toggling kind + // 44200 on must invalidate mounted usage queries. Other kinds don't affect + // usage state. + if (kind === KIND_AGENT_TURN_METRIC) { + notifyAgentMetricsChanged(); + } } /** @@ -142,6 +196,9 @@ export async function mergeSaveSubscriptionKinds(kind: number): Promise { export async function removeSaveSubscriptionKind(kind: number): Promise { await invokeTauri("remove_save_subscription_kind", { kind }); notifySubscriptionChange(); + if (kind === KIND_AGENT_TURN_METRIC) { + notifyAgentMetricsChanged(); + } } /** @@ -209,7 +266,7 @@ export async function archiveEvents( matchedScope: { scopeType: ScopeType; scopeValue: string }; }>, ): Promise { - return invokeTauri("archive_events", { + const raw = await invokeTauri>("archive_events", { candidates: candidates.map((c) => ({ raw_event_json: c.rawEventJson, matched_scope: { @@ -218,6 +275,7 @@ export async function archiveEvents( }, })), }); + return decodeArchiveBatchResult(raw); } /**