mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
Add opt-in archive performance diagnostics
Log 10-second archive traffic and event-loop stall windows only when the diagnostic build flag is enabled. Distinguish observer frames from turn metrics so current archive preferences can be evaluated objectively. Co-authored-by: Carl <c7ebe626f000404285d3686e1dc74cc07cc60a9754a150041ba132e14bd3e2ec@buzz.block.builderlab.xyz> Signed-off-by: Wes <wesbillman@users.noreply.github.com>
This commit is contained in:
@@ -1023,6 +1023,103 @@ test("manager_notifies_agent_metrics_changed_on_destroy_flush", async () => {
|
||||
off();
|
||||
});
|
||||
|
||||
test("manager_reports_and_resets_archive_perf_diagnostics", async () => {
|
||||
const relay = makeFakeRelayClient();
|
||||
const archive = makeFakeArchive();
|
||||
archive.setSubs([
|
||||
{
|
||||
scopeType: "owner_p",
|
||||
scopeValue: "agent-pk",
|
||||
kinds: [24200, 44200],
|
||||
identityPubkey: "pk",
|
||||
relayUrl: "wss://r",
|
||||
createdAt: 0,
|
||||
},
|
||||
]);
|
||||
const logs = [];
|
||||
const nowValues = [1_000, 11_000, 21_000];
|
||||
const mgr = makeManager(relay, archive, {
|
||||
enableArchivePerfDiagnostics: true,
|
||||
flushBatchSize: 2,
|
||||
perfNow: () => nowValues.shift(),
|
||||
perfLog: (message) => logs.push(message),
|
||||
});
|
||||
await mgr.start();
|
||||
|
||||
const filter = JSON.parse([...relay.subs.keys()][0]);
|
||||
const observerEvent = {
|
||||
id: "observer-1",
|
||||
kind: 24200,
|
||||
pubkey: "agent-pk",
|
||||
created_at: 1,
|
||||
content: "frame",
|
||||
tags: [],
|
||||
};
|
||||
const metricEvent = {
|
||||
id: "metric-1",
|
||||
kind: 44200,
|
||||
pubkey: "agent-pk",
|
||||
created_at: 2,
|
||||
content: "metric",
|
||||
tags: [],
|
||||
};
|
||||
relay.push(filter, observerEvent);
|
||||
relay.push(filter, metricEvent);
|
||||
|
||||
// biome-ignore lint/complexity/useLiteralKeys: deterministic private diagnostic seam
|
||||
mgr["reportPerfWindow"]();
|
||||
const report = logs.at(-1);
|
||||
assert.match(report, /windowMs=10000/);
|
||||
assert.match(report, /observerEvents=1/);
|
||||
assert.match(
|
||||
report,
|
||||
new RegExp(`observerBytes=${JSON.stringify(observerEvent).length}`),
|
||||
);
|
||||
assert.match(report, /metricEvents=1/);
|
||||
assert.match(
|
||||
report,
|
||||
new RegExp(`metricBytes=${JSON.stringify(metricEvent).length}`),
|
||||
);
|
||||
assert.match(report, /archiveBatches=1 archiveEvents=2/);
|
||||
|
||||
// biome-ignore lint/complexity/useLiteralKeys: deterministic private diagnostic seam
|
||||
mgr["reportPerfWindow"]();
|
||||
assert.match(
|
||||
logs.at(-1),
|
||||
/observerEvents=0 observerBytes=0 metricEvents=0 metricBytes=0 archiveBatches=0 archiveEvents=0 archiveBytes=0/,
|
||||
);
|
||||
mgr.destroy();
|
||||
});
|
||||
|
||||
test("manager_clears_archive_perf_diagnostic_timers_on_destroy", async () => {
|
||||
const relay = makeFakeRelayClient();
|
||||
const archive = makeFakeArchive();
|
||||
const mgr = makeManager(relay, archive, {
|
||||
enableArchivePerfDiagnostics: true,
|
||||
perfNow: () => 0,
|
||||
perfLog: () => {},
|
||||
});
|
||||
await mgr.start();
|
||||
|
||||
assert.notEqual(
|
||||
// biome-ignore lint/complexity/useLiteralKeys: intentional lifecycle assertion
|
||||
mgr["perfReportTimer"],
|
||||
null,
|
||||
);
|
||||
assert.notEqual(
|
||||
// biome-ignore lint/complexity/useLiteralKeys: intentional lifecycle assertion
|
||||
mgr["eventLoopProbeTimer"],
|
||||
null,
|
||||
);
|
||||
|
||||
mgr.destroy();
|
||||
|
||||
// biome-ignore lint/complexity/useLiteralKeys: intentional lifecycle assertion
|
||||
assert.equal(mgr["perfReportTimer"], null);
|
||||
// biome-ignore lint/complexity/useLiteralKeys: intentional lifecycle assertion
|
||||
assert.equal(mgr["eventLoopProbeTimer"], null);
|
||||
});
|
||||
|
||||
test("manager_flushes_buffer_on_destroy", async () => {
|
||||
const relay = makeFakeRelayClient();
|
||||
const archive = makeFakeArchive();
|
||||
|
||||
@@ -1,7 +1,10 @@
|
||||
import { relayClient as defaultRelayClient } from "@/shared/api/relayClient";
|
||||
import type { RelaySubscriptionFilter } from "@/shared/api/relayClientShared";
|
||||
import type { RelayEvent } from "@/shared/api/types";
|
||||
import { KIND_AGENT_TURN_METRIC } from "@/shared/constants/kinds";
|
||||
import {
|
||||
KIND_AGENT_OBSERVER_FRAME,
|
||||
KIND_AGENT_TURN_METRIC,
|
||||
} from "@/shared/constants/kinds";
|
||||
import {
|
||||
archiveEvents as defaultArchiveEvents,
|
||||
listSaveSubscriptions as defaultListSaveSubscriptions,
|
||||
@@ -18,6 +21,11 @@ const FLUSH_BATCH_SIZE = 25;
|
||||
const FLUSH_IDLE_MS = 2_000;
|
||||
const DISABLE_AGENT_METRIC_ARCHIVE =
|
||||
import.meta.env?.VITE_BUZZ_DISABLE_AGENT_METRIC_ARCHIVE === "1";
|
||||
const ENABLE_ARCHIVE_PERF_DIAGNOSTICS =
|
||||
import.meta.env?.VITE_BUZZ_ARCHIVE_PERF_DIAGNOSTICS === "1";
|
||||
const PERF_REPORT_INTERVAL_MS = 10_000;
|
||||
const EVENT_LOOP_PROBE_INTERVAL_MS = 250;
|
||||
const EVENT_LOOP_STALL_THRESHOLD_MS = 50;
|
||||
|
||||
// ── Types ─────────────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -38,6 +46,9 @@ export interface ArchiveSyncDeps {
|
||||
) => Promise<ArchiveBatchResult>;
|
||||
onSubscriptionChange: (listener: () => void) => () => void;
|
||||
disableAgentMetricArchive?: boolean;
|
||||
enableArchivePerfDiagnostics?: boolean;
|
||||
perfNow?: () => number;
|
||||
perfLog?: (message: string) => void;
|
||||
flushBatchSize?: number;
|
||||
flushIdleMs?: number;
|
||||
}
|
||||
@@ -87,10 +98,18 @@ export class ArchiveSyncManager {
|
||||
private readonly deps: Required<
|
||||
Omit<
|
||||
ArchiveSyncDeps,
|
||||
"disableAgentMetricArchive" | "flushBatchSize" | "flushIdleMs"
|
||||
| "disableAgentMetricArchive"
|
||||
| "enableArchivePerfDiagnostics"
|
||||
| "perfNow"
|
||||
| "perfLog"
|
||||
| "flushBatchSize"
|
||||
| "flushIdleMs"
|
||||
>
|
||||
>;
|
||||
private readonly disableAgentMetricArchive: boolean;
|
||||
private readonly enableArchivePerfDiagnostics: boolean;
|
||||
private readonly perfNow: () => number;
|
||||
private readonly perfLog: (message: string) => void;
|
||||
private readonly flushBatchSize: number;
|
||||
private readonly flushIdleMs: number;
|
||||
|
||||
@@ -108,6 +127,22 @@ export class ArchiveSyncManager {
|
||||
matchedScope: { scopeType: ScopeType; scopeValue: string };
|
||||
}> = [];
|
||||
private flushTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
private perfReportTimer: ReturnType<typeof setInterval> | null = null;
|
||||
private eventLoopProbeTimer: ReturnType<typeof setInterval> | null = null;
|
||||
private perfWindowStartedAt = 0;
|
||||
private nextEventLoopProbeAt = 0;
|
||||
private perfCounters = {
|
||||
observerEvents: 0,
|
||||
observerBytes: 0,
|
||||
metricEvents: 0,
|
||||
metricBytes: 0,
|
||||
archiveBatches: 0,
|
||||
archiveEvents: 0,
|
||||
archiveBytes: 0,
|
||||
stalls: 0,
|
||||
stallMs: 0,
|
||||
maxStallMs: 0,
|
||||
};
|
||||
private destroyed = false;
|
||||
private offSubscriptionChange: (() => void) | null = null;
|
||||
|
||||
@@ -122,11 +157,16 @@ export class ArchiveSyncManager {
|
||||
};
|
||||
this.disableAgentMetricArchive =
|
||||
deps?.disableAgentMetricArchive ?? DISABLE_AGENT_METRIC_ARCHIVE;
|
||||
this.enableArchivePerfDiagnostics =
|
||||
deps?.enableArchivePerfDiagnostics ?? ENABLE_ARCHIVE_PERF_DIAGNOSTICS;
|
||||
this.perfNow = deps?.perfNow ?? (() => performance.now());
|
||||
this.perfLog = deps?.perfLog ?? ((message) => console.info(message));
|
||||
this.flushBatchSize = deps?.flushBatchSize ?? FLUSH_BATCH_SIZE;
|
||||
this.flushIdleMs = deps?.flushIdleMs ?? FLUSH_IDLE_MS;
|
||||
}
|
||||
|
||||
async start(): Promise<void> {
|
||||
this.startPerfDiagnostics();
|
||||
// Register the change listener before the initial load so that any
|
||||
// subscription change arriving while the first pass is running sets
|
||||
// reloadPending and gets picked up by the coalescing loop.
|
||||
@@ -138,6 +178,14 @@ export class ArchiveSyncManager {
|
||||
|
||||
destroy(): void {
|
||||
this.destroyed = true;
|
||||
if (this.perfReportTimer !== null) {
|
||||
clearInterval(this.perfReportTimer);
|
||||
this.perfReportTimer = null;
|
||||
}
|
||||
if (this.eventLoopProbeTimer !== null) {
|
||||
clearInterval(this.eventLoopProbeTimer);
|
||||
this.eventLoopProbeTimer = null;
|
||||
}
|
||||
if (this.flushTimer !== null) {
|
||||
clearTimeout(this.flushTimer);
|
||||
this.flushTimer = null;
|
||||
@@ -279,14 +327,77 @@ export class ArchiveSyncManager {
|
||||
}
|
||||
}
|
||||
|
||||
private startPerfDiagnostics(): void {
|
||||
if (!this.enableArchivePerfDiagnostics || this.perfReportTimer !== null) {
|
||||
return;
|
||||
}
|
||||
|
||||
this.perfWindowStartedAt = this.perfNow();
|
||||
this.nextEventLoopProbeAt =
|
||||
this.perfWindowStartedAt + EVENT_LOOP_PROBE_INTERVAL_MS;
|
||||
this.perfLog(
|
||||
`[archive-perf] started metricArchiveDisabled=${this.disableAgentMetricArchive}`,
|
||||
);
|
||||
|
||||
this.eventLoopProbeTimer = setInterval(() => {
|
||||
const now = this.perfNow();
|
||||
const delayMs = Math.max(0, now - this.nextEventLoopProbeAt);
|
||||
this.nextEventLoopProbeAt = now + EVENT_LOOP_PROBE_INTERVAL_MS;
|
||||
if (delayMs >= EVENT_LOOP_STALL_THRESHOLD_MS) {
|
||||
this.perfCounters.stalls++;
|
||||
this.perfCounters.stallMs += delayMs;
|
||||
this.perfCounters.maxStallMs = Math.max(
|
||||
this.perfCounters.maxStallMs,
|
||||
delayMs,
|
||||
);
|
||||
}
|
||||
}, EVENT_LOOP_PROBE_INTERVAL_MS);
|
||||
|
||||
this.perfReportTimer = setInterval(() => {
|
||||
this.reportPerfWindow();
|
||||
}, PERF_REPORT_INTERVAL_MS);
|
||||
}
|
||||
|
||||
private reportPerfWindow(): void {
|
||||
const now = this.perfNow();
|
||||
const durationMs = Math.round(now - this.perfWindowStartedAt);
|
||||
const counters = this.perfCounters;
|
||||
this.perfLog(
|
||||
`[archive-perf] windowMs=${durationMs} metricArchiveDisabled=${this.disableAgentMetricArchive} observerEvents=${counters.observerEvents} observerBytes=${counters.observerBytes} metricEvents=${counters.metricEvents} metricBytes=${counters.metricBytes} archiveBatches=${counters.archiveBatches} archiveEvents=${counters.archiveEvents} archiveBytes=${counters.archiveBytes} stalls=${counters.stalls} stallMs=${Math.round(counters.stallMs)} maxStallMs=${Math.round(counters.maxStallMs)}`,
|
||||
);
|
||||
this.perfWindowStartedAt = now;
|
||||
this.perfCounters = {
|
||||
observerEvents: 0,
|
||||
observerBytes: 0,
|
||||
metricEvents: 0,
|
||||
metricBytes: 0,
|
||||
archiveBatches: 0,
|
||||
archiveEvents: 0,
|
||||
archiveBytes: 0,
|
||||
stalls: 0,
|
||||
stallMs: 0,
|
||||
maxStallMs: 0,
|
||||
};
|
||||
}
|
||||
|
||||
private enqueue(
|
||||
event: RelayEvent,
|
||||
scopeType: ScopeType,
|
||||
scopeValue: string,
|
||||
): void {
|
||||
if (this.destroyed) return;
|
||||
const rawEventJson = JSON.stringify(event);
|
||||
if (this.enableArchivePerfDiagnostics) {
|
||||
if (event.kind === KIND_AGENT_OBSERVER_FRAME) {
|
||||
this.perfCounters.observerEvents++;
|
||||
this.perfCounters.observerBytes += rawEventJson.length;
|
||||
} else if (event.kind === KIND_AGENT_TURN_METRIC) {
|
||||
this.perfCounters.metricEvents++;
|
||||
this.perfCounters.metricBytes += rawEventJson.length;
|
||||
}
|
||||
}
|
||||
this.buffer.push({
|
||||
rawEventJson: JSON.stringify(event),
|
||||
rawEventJson,
|
||||
matchedScope: { scopeType, scopeValue },
|
||||
});
|
||||
if (this.buffer.length >= this.flushBatchSize) {
|
||||
@@ -328,6 +439,14 @@ export class ArchiveSyncManager {
|
||||
}>,
|
||||
errLabel: string,
|
||||
): void {
|
||||
if (this.enableArchivePerfDiagnostics) {
|
||||
this.perfCounters.archiveBatches++;
|
||||
this.perfCounters.archiveEvents += batch.length;
|
||||
this.perfCounters.archiveBytes += batch.reduce(
|
||||
(total, candidate) => total + candidate.rawEventJson.length,
|
||||
0,
|
||||
);
|
||||
}
|
||||
void this.deps
|
||||
.archiveEvents(batch)
|
||||
.then((result) => {
|
||||
|
||||
Reference in New Issue
Block a user