mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(desktop): wire persistedAgentMetrics notifier through archive sync
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 <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
co-authored by
Will Pfleger
parent
c537cbc1b8
commit
4f41319220
@@ -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();
|
||||
|
||||
@@ -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<unknown>;
|
||||
) => Promise<ArchiveBatchResult>;
|
||||
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);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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>,
|
||||
): 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<boolean> {
|
||||
export async function mergeSaveSubscriptionKinds(kind: number): Promise<void> {
|
||||
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<void> {
|
||||
export async function removeSaveSubscriptionKind(kind: number): Promise<void> {
|
||||
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<ArchiveBatchResult> {
|
||||
return invokeTauri<ArchiveBatchResult>("archive_events", {
|
||||
const raw = await invokeTauri<Partial<ArchiveBatchResult>>("archive_events", {
|
||||
candidates: candidates.map((c) => ({
|
||||
raw_event_json: c.rawEventJson,
|
||||
matched_scope: {
|
||||
@@ -218,6 +275,7 @@ export async function archiveEvents(
|
||||
},
|
||||
})),
|
||||
});
|
||||
return decodeArchiveBatchResult(raw);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user