From b0a098eb8c78fbaa3bd481684154871e3686b073 Mon Sep 17 00:00:00 2001 From: Amir Ghorbani Date: Tue, 29 Sep 2026 12:15:40 -0400 Subject: [PATCH 1/6] docs(observability): make the metrics table the single source of truth MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Spec 10 §7 and the telemetry guide now list every instrument with its type, unit, attributes and when it is recorded, including the three notification counters and the three process gauges, and name the guard test that keeps the tables, the catalogue and the recorded metrics in step (spec 09 §3.2). --- docs/guide/telemetry.md | 47 ++++++++++++--------- specs/09-testing.md | 2 +- specs/10-error-handling-and-telemetry.md | 54 +++++++++++++----------- 3 files changed, 56 insertions(+), 47 deletions(-) diff --git a/docs/guide/telemetry.md b/docs/guide/telemetry.md index c4cb4b6..c0fcb11 100644 --- a/docs/guide/telemetry.md +++ b/docs/guide/telemetry.md @@ -40,30 +40,35 @@ If the MCP client sends a W3C `traceparent` in the tool call's `_meta`, the tool ### Metrics -| Instrument | Type | Attributes | -|---|---|---| -| `browserhive.tool_calls` | counter | `tool`, `ok`, `error_code`, `harness` | -| `browserhive.tool_call.duration` | histogram (ms) | `tool` | -| `browserhive.sessions.active` | up-down counter | `harness` | -| `browserhive.session.launch.duration` | histogram (ms) | `channel`, `stealth` | -| `browserhive.session.lifetime` | histogram (ms) | `closed_reason` | -| `browserhive.ws.connections` | up-down counter | | -| `browserhive.ws.buffered_bytes` | gauge | `connection_id` | -| `browserhive.ws.frames_dropped` | counter | `channel` | -| `browserhive.db.write_queue.depth` | gauge | | -| `browserhive.db.dropped_writes` | counter | `table` | -| `browserhive.db.size_bytes` | gauge | | -| `browserhive.browser.rss_bytes` | gauge | `session_id` | -| `browserhive.attention.open` | up-down counter | `kind` | -| `browserhive.attention.wait` | histogram (ms) | `status` | -| `browserhive.vault.fills` | counter | `result` | -| `browserhive.blocklist.hits` | counter | `source` | -| `browserhive.retention.pruned_rows` | counter | `table` | -| `browserhive.process.*` | gauges | rss, heap, event-loop lag | +| Instrument | Type | Unit | Attributes | What it measures | +|---|---|---|---|---| +| `browserhive.tool_calls` | counter | `{call}` | `tool`, `ok`, `error_code` (failed calls only), `harness` | Tool calls. | +| `browserhive.tool_call.duration` | histogram | `ms` | `tool` | How long each tool call took. | +| `browserhive.sessions.active` | up-down counter | `{session}` | `harness` | Live sessions. | +| `browserhive.session.launch.duration` | histogram | `ms` | `channel`, `stealth` | From the create request until the browser is up, per successful launch. | +| `browserhive.session.lifetime` | histogram | `ms` | `closed_reason` | How long each closed session lived. | +| `browserhive.ws.connections` | up-down counter | `{connection}` | | Open dashboard WebSocket connections. | +| `browserhive.ws.buffered_bytes` | gauge | `By` | `connection_id` | Bytes waiting to be sent, for the 50 most backed-up connections. | +| `browserhive.ws.frames_dropped` | counter | `{frame}` | `channel` (`screencast`, `logs`, `feed`) | Frames not delivered because a connection was congested. | +| `browserhive.db.write_queue.depth` | gauge | `{write}` | | Writes waiting to be saved. | +| `browserhive.db.dropped_writes` | counter | `{write}` | `table` | Writes lost because the queue was full or the write failed. | +| `browserhive.db.size_bytes` | gauge | `By` | | Database size. | +| `browserhive.browser.rss_bytes` | gauge | `By` | `session_id` | Memory of each session's browser, all its processes together, sampled every 10 seconds. Linux and macOS only. | +| `browserhive.attention.open` | up-down counter | `{request}` | `kind` (`attention`, `vault_confirm`) | Operator requests waiting for an answer. | +| `browserhive.attention.wait` | histogram | `ms` | `status` | How long each attention request waited until it was answered, timed out or cancelled. | +| `browserhive.vault.fills` | counter | `{fill}` | `result` | Vault fills. | +| `browserhive.blocklist.hits` | counter | `{hit}` | `source` | Blocked requests. | +| `browserhive.retention.pruned_rows` | counter | `{row}` | `table` | Rows removed by retention. | +| `browserhive.notifications.deliveries` | counter | `{delivery}` | `channel_kind`, `status` | Notification deliveries by outcome: `sent`, `retrying`, `dead`, `suppressed`, `superseded`. | +| `browserhive.notifications.reports` | counter | `{report}` | `kind`, `outcome` | Digest and anomaly-alert decisions: `sent`, `late`, `empty`, `skipped`, `manual`, `resolved`, `in_app`. | +| `browserhive.notifications.actions` | counter | `{press}` | `channel_kind`, `outcome` | Presses of act buttons in Telegram, Discord and ntfy, by outcome (`unknown` for buttons BrowserHive did not create). | +| `browserhive.process.rss_bytes` | gauge | `By` | | Memory of the BrowserHive process. | +| `browserhive.process.heap_bytes` | gauge | `By` | | JavaScript heap in use. | +| `browserhive.process.event_loop_lag` | gauge | `ms` | | Event-loop delay (99th percentile since the previous export). | `harness` on metrics is one of the known harness names, `unknown` or `other` (every name BrowserHive doesn't know is folded into `other`), so it adds a bounded number of series. The model, workspace and extra labels are never metric attributes. -Metrics are exported every 30 seconds. +Metrics are exported every 30 seconds. Gauges and the WebSocket and write-queue totals are read at export time; with telemetry off nothing is measured. ### Logs diff --git a/specs/09-testing.md b/specs/09-testing.md index 9299d5a..afafb71 100644 --- a/specs/09-testing.md +++ b/specs/09-testing.md @@ -92,7 +92,7 @@ Domain and application (no I/O, fakes only): - Retention (`app/observability/retention.test.ts`): every table has a retention class asserted over the table list; terminal `notification_deliveries` and settled `notification_channel_messages` older than 30 d are pruned while pending jobs and pending TTLs are kept; `notification_channels` is never pruned; `purge` inventories a database from an older schema without failing on tables it does not have; age prune; byte cap converges without `VACUUM` in a loop; artifact deletion via the outbox; `pending_attention`/operator-request rows are pruned once terminal; failures are counted and surfaced, never thrown out of the timer. - Error registry (`kernel/errors/*.test.ts`): every code projects to MCP, problem+json, and WS frames; `docs/errors.md` generation is deterministic; `INTERNAL_ERROR` never carries host paths in the public message. - Logger (`infra/logging/*.test.ts`): JSON and pretty renderers with fixed clock and width; field priority order; redaction of key names and registered secrets; `child()` bindings; per-module level spec; ring buffer size and eviction; stdio guard routes `console.*` to stderr. -- Request context / telemetry (`infra/telemetry/*.test.ts`): spans nest under `AsyncLocalStorage`; log records carry `trace_id`/`span_id`; the no-op provider adds no fields when `otel=false`. +- Request context / telemetry (`infra/telemetry/*.test.ts`): spans nest under `AsyncLocalStorage`; log records carry `trace_id`/`span_id`; the no-op provider adds no fields when `otel=false`. Metrics guard (`browserhive/test/composition/metrics-guard.test.ts`): the metric tables of spec 10 §7 and `docs/guide/telemetry.md` and the instrument catalogue must list the same instruments with the same type, unit and attributes, and every instrument must reach a real OTLP/HTTP receiver when its source is driven (bus events, the realtime hub's totals, the write queue's dropped writes, the browser-memory sampler over a real process tree, the event-loop monitor, and the notification outbox, act buttons and report scheduler as the composition root builds them), with only documented attribute keys; a metric that is documented or defined but never recorded fails it. - WS hub (`interface/ws/hub.test.ts`, fake socket): envelope shape; `subscribe` with a cursor replays buffered events in order; buffer bounds by count, bytes, and age; overflow yields `resync_required`; screencast is latest-wins per connection and drops frames when the fake socket reports backpressure while feed events are never dropped; auth invalidation closes 4401; five protocol violations close 4400; stale reap 1001. - HTTP routes (`interface/http/routes/*.test.ts`): the real Hono app with an in-process `:memory:` database and in-memory browser driver, no port; every route has at least one success test and one validation-failure test (400 problem+json with field errors); an auth matrix test iterates all routes × {no cookie → 401, must-change-password → 403 except the allowed three, valid → 2xx/4xx}; response bodies parsed with the contracts schema; response goldens for list routes with a seeded fixture dataset and fixed clock/ids (deep equality). - Auth (`app/auth/*.test.ts`): provider chain order; present-but-invalid never falls through; password sessions idle/absolute expiry; token hashing and constant-time lookup; grants single-use and revoked with the parent; `Authorizer.can` table; login rate limit and lockout; password change revokes other sessions. diff --git a/specs/10-error-handling-and-telemetry.md b/specs/10-error-handling-and-telemetry.md index 7310dee..703f137 100644 --- a/specs/10-error-handling-and-telemetry.md +++ b/specs/10-error-handling-and-telemetry.md @@ -339,31 +339,35 @@ All spans use `@opentelemetry/api`'s tracer `browserhive`; attributes use the `b ## 7. Metrics catalogue -| Instrument | Type | Attributes | -|---|---|---| -| `browserhive.tool_calls` | counter | `tool`, `ok`, `error_code`, `harness` (folded, below) | -| `browserhive.tool_call.duration` | histogram (ms) | `tool` | -| `browserhive.sessions.active` | up-down counter | `harness` (folded, below) | -| `browserhive.session.launch.duration` | histogram (ms) | `channel`, `stealth` | -| `browserhive.session.lifetime` | histogram (ms) | `closed_reason` | -| `browserhive.ws.connections` | up-down counter | — | -| `browserhive.ws.buffered_bytes` | gauge (observable) | `connection_id` capped to 50 series | -| `browserhive.ws.frames_dropped` | counter | `channel` | -| `browserhive.db.write_queue.depth` | gauge | — | -| `browserhive.db.dropped_writes` | counter | `table` | -| `browserhive.db.size_bytes` | gauge | — | -| `browserhive.browser.rss_bytes` | gauge (observable; sampled every 10 s from the process tree) | `session_id` | -| `browserhive.attention.open` | up-down counter | `kind` | -| `browserhive.attention.wait` | histogram (ms) | `status` | -| `browserhive.vault.fills` | counter | `result` | -| `browserhive.blocklist.hits` | counter | `source` | -| `browserhive.retention.pruned_rows` | counter | `table` | -| `browserhive.notifications.deliveries` | counter | `channel_kind`, `status` (`sent`, `retrying`, `dead`, `suppressed`, `superseded`) — one increment per finished or rescheduled outbox job (03 §9.4) | -| `browserhive.notifications.reports` | counter | `kind` (`digest.daily`, `digest.weekly`, `report.anomaly`), `outcome` (`sent` for a produced report, `late`, `empty`, `skipped` per skipped window, `manual`, `resolved` for an anomaly alert that cleared, `in_app` for a new in-app copy, D-45) — one increment per report decision of the scheduler (03 §9.7) | -| `browserhive.notifications.actions` | counter | `channel_kind`, `outcome` (`done`, `failed`, `not_allowed`, `used`, `expired`, `stale`, `wrong_channel`, `disabled`, `unknown`) — one increment per act-button press (03 §9.6); `unknown` counts presses of tokens BrowserHive never minted, which are not audited | -| `browserhive.process.*` | gauges: rss, heap, event-loop lag (sampled) | — | - -The same registry backs `/api/v1/system` figures; with `--otel` off, the in-process meter provider is the SDK's no-op. +This table is the single source of truth: `infra/telemetry/metrics.ts` defines exactly these instruments, the guard test (09 §3.2) fails when a row here, a row of `docs/guide/telemetry.md` and the catalogue disagree on name, type, unit or attributes, or when an instrument is never recorded. Attribute values in parentheses are the closed set that can occur. Types are the OTLP data a collector receives: *counter* is a monotonic sum, *up-down counter* a non-monotonic sum, *gauge* a gauge; "observable" instruments are read by a callback at each export (every 30 s) instead of being written on each event. + +| Instrument | Type | Unit | Attributes | Recorded | +|---|---|---|---|---| +| `browserhive.tool_calls` | counter | `{call}` | `tool`, `ok`, `error_code` (failed calls only), `harness` (folded, below) | once per terminal tool outcome (`tool.called`) | +| `browserhive.tool_call.duration` | histogram | `ms` | `tool` | once per terminal tool outcome (`tool.called`) | +| `browserhive.sessions.active` | up-down counter | `{session}` | `harness` (folded, below) | +1 on `session.opened`, −1 on `session.closed` | +| `browserhive.session.launch.duration` | histogram | `ms` | `channel`, `stealth` | once per successful launch: create request to browser attached (the session's `launch_ms`, first seen on `session.updated`) | +| `browserhive.session.lifetime` | histogram | `ms` | `closed_reason` | once per `session.closed`, opened to closed | +| `browserhive.ws.connections` | up-down counter (observable) | `{connection}` | — | open dashboard WebSocket connections of the realtime hub (http transport only) | +| `browserhive.ws.buffered_bytes` | gauge (observable) | `By` | `connection_id` (the 50 connections with the most buffered bytes) | bytes the socket has not flushed yet, per connection (http transport only) | +| `browserhive.ws.frames_dropped` | counter (observable) | `{frame}` | `channel` (`screencast`, `logs`, `feed`) | the hub's running totals: a screencast frame replaced by a newer one or refused by a congested socket; a log record skipped for a congested `logs` subscriber; a feed frame the socket refused (the connection closes with 1013 and replays from its cursor) | +| `browserhive.db.write_queue.depth` | gauge (observable) | `{write}` | — | recorder writes waiting in the queue | +| `browserhive.db.dropped_writes` | counter (observable) | `{write}` | `table` (the recorder table: `sessions`, `tool_calls`, `pages`, `screenshots`, `blocked_requests`, `vault_access`, `events`, `logs`) | the write queue's running totals: writes refused (queue full or closed) or failed inside their drain | +| `browserhive.db.size_bytes` | gauge (observable) | `By` | — | database file size, read in the background and reported at the next export | +| `browserhive.browser.rss_bytes` | gauge (observable) | `By` | `session_id` | resident memory of each live session's browser process tree (the browser process and every descendant: renderers, GPU, utilities), sampled every 10 s; Linux and macOS only, no data points on Windows | +| `browserhive.attention.open` | up-down counter | `{request}` | `kind` (`attention`, `vault_confirm`) | +1 when an operator request opens, −1 when it settles | +| `browserhive.attention.wait` | histogram | `ms` | `status` (`resolved`, `rejected`, `timeout`, `cancelled`) | once per settled attention request, its `waited_ms` | +| `browserhive.vault.fills` | counter | `{fill}` | `result` | once per `vault.access` | +| `browserhive.blocklist.hits` | counter | `{hit}` | `source` | once per `blocklist.hit` | +| `browserhive.retention.pruned_rows` | counter | `{row}` | `table` | once per retention pass and table that lost rows | +| `browserhive.notifications.deliveries` | counter | `{delivery}` | `channel_kind`, `status` (`sent`, `retrying`, `dead`, `suppressed`, `superseded`) | one increment per finished or rescheduled outbox job (03 §9.4); `channel_kind` is `unknown` for a job of a channel removed meanwhile | +| `browserhive.notifications.reports` | counter | `{report}` | `kind` (`digest.daily`, `digest.weekly`, `report.anomaly`), `outcome` (`sent` for a produced report, `late`, `empty`, `skipped` per skipped window, `manual` for an on-demand digest sent from the dashboard, `resolved` for an anomaly alert that cleared, `in_app` for a new in-app copy, D-45) | one increment per report decision of the scheduler (03 §9.7); a silent revision of an open anomaly alert is not a decision and is not counted | +| `browserhive.notifications.actions` | counter | `{press}` | `channel_kind`, `outcome` (`done`, `failed`, `not_allowed`, `used`, `expired`, `stale`, `wrong_channel`, `disabled`, `unknown`) | one increment per act-button press (03 §9.6); `unknown` counts presses of tokens BrowserHive never minted, which are not audited | +| `browserhive.process.rss_bytes` | gauge (observable) | `By` | — | resident memory of the BrowserHive process | +| `browserhive.process.heap_bytes` | gauge (observable) | `By` | — | JavaScript heap in use | +| `browserhive.process.event_loop_lag` | gauge (observable) | `ms` | — | 99th percentile event-loop delay since the previous export | + +Instruments are only created and fed when `--otel` is on with the `metrics` signal: the consumers, the 10 s browser-memory sampler and the event-loop monitor are wired by the composition root in that case only, so with `--otel` off nothing is subscribed, sampled or timed. The sources themselves only keep plain running totals whether or not telemetry is on (the hub's dropped frames per channel, the write queue's dropped writes per table, next to the per-connection counters `/api/v1/system/realtime` already shows); the observable instruments read those at export time. `/api/v1/system` reads the same sources directly, never the meter. **Cardinality rule for client identity (D-30).** `harness` is a metric attribute only on `browserhive.tool_calls` and `browserhive.sessions.active`, and only as `metricHarness(slug)`: the known slug table of `contracts/harness` (including `unknown` and `other`), with every other value folded into `other`, so it adds at most that many series per existing combination. It is never an attribute of `browserhive.tool_call.duration` (a histogram already split by 43 tools). `model`, `workspace`, `harness_source`, client names, `User-Agent` and meta-bag keys or values are never metric attributes; spans carry `browserhive.harness`, `browserhive.harness_source` and `browserhive.model` because a span is one event. The `tool called` log line carries `harness` as a field. From 9ff350b54a5d040cdc525748a75efc1b3dcf58de Mon Sep 17 00:00:00 2001 From: Amir Ghorbani Date: Tue, 29 Sep 2026 12:32:14 -0400 Subject: [PATCH 2/6] fix(observability): record every documented metric The adapter only fed eleven of the catalogue's instruments. It now records the rest from their sources: - session launch duration (channel, stealth) from the first session.updated carrying launch_ms; lifetime by closed_reason - attention.open by kind (attention and vault confirmations) and the attention wait by status, from the broker's events - retention pruned rows by table, carried internally on retention.completed - WebSocket connections, buffered bytes (top 50) and frames dropped by channel, read from the realtime hub's running totals at export - dropped writes by table, read from the write queue's totals - each live session's browser process-tree RSS, sampled every 10 s from the browser pid (one browser-level DevTools read, cached) and /proc or ps - the p99 event-loop delay since the previous export Every instrument now carries the unit of its catalogue row. The report scheduler counts on-demand digests as manual and no longer counts a silent revision of an in-app anomaly alert. With telemetry off nothing is subscribed, sampled or timed. --- .changeset/otel-documented-metrics.md | 10 + .../composition/adapters/browser-memory.ts | 77 +++++ .../src/composition/adapters/metrics.test.ts | 245 +++++++++++++- .../src/composition/adapters/metrics.ts | 203 ++++++++++-- .../browserhive/src/composition/context.ts | 6 +- .../src/composition/phases/build-domain.ts | 6 +- .../src/composition/phases/listeners-http.ts | 2 + .../src/composition/phases/wire-observers.ts | 48 ++- packages/core/src/app/events/catalog.ts | 5 +- .../app/maintenance/retention-scheduler.ts | 1 + .../src/app/notifications/report-scheduler.ts | 7 +- .../core/src/infra/browsers/session-handle.ts | 27 ++ packages/core/src/infra/host/event-loop.ts | 39 +++ packages/core/src/infra/host/process-tree.ts | 123 +++++++ .../core/src/infra/persistence/write-queue.ts | 17 +- packages/core/src/infra/telemetry/metrics.ts | 312 +++++++++++------- packages/core/src/interface/ws/hub.ts | 30 +- packages/core/src/interface/ws/index.ts | 2 +- packages/core/src/ports/browser-driver.ts | 5 + .../core/src/ports/persistence/write-queue.ts | 5 + packages/core/src/public/runtime.ts | 13 +- packages/core/test/helpers/in-memory-repos.ts | 15 +- 22 files changed, 1008 insertions(+), 190 deletions(-) create mode 100644 .changeset/otel-documented-metrics.md create mode 100644 packages/browserhive/src/composition/adapters/browser-memory.ts create mode 100644 packages/core/src/infra/host/event-loop.ts create mode 100644 packages/core/src/infra/host/process-tree.ts diff --git a/.changeset/otel-documented-metrics.md b/.changeset/otel-documented-metrics.md new file mode 100644 index 0000000..42cd468 --- /dev/null +++ b/.changeset/otel-documented-metrics.md @@ -0,0 +1,10 @@ +--- +"browserhive": patch +--- + +OpenTelemetry now exports the notification, attention, session-launch, WebSocket, write-queue and browser-memory metrics the docs describe. + +- **Metrics that were documented but never sent** now reach your collector: `browserhive.attention.wait`, `browserhive.session.launch.duration`, `browserhive.ws.connections`, `browserhive.ws.buffered_bytes`, `browserhive.ws.frames_dropped`, `browserhive.db.dropped_writes`, `browserhive.browser.rss_bytes` (each session's browser with all its processes, every 10 seconds; Linux and macOS) and `browserhive.process.event_loop_lag`. +- **Attributes the docs promised** are now set: `closed_reason` on `browserhive.session.lifetime`, `kind` on `browserhive.attention.open` (vault confirmations count too), `table` on `browserhive.retention.pruned_rows`. +- **The notification metrics** `browserhive.notifications.deliveries`, `.actions` and `.reports` are now in the [telemetry guide](https://browserhive.ai/docs/guide/telemetry). `.reports` now counts on-demand digests as `manual`, as documented, and no longer counts a silent revision of an open in-app anomaly alert. +- **Every metric has a unit** (`ms`, `By`, or a count such as `{call}`), and the guide's table lists each one with its type, unit and attributes. With telemetry off nothing is measured, as before. diff --git a/packages/browserhive/src/composition/adapters/browser-memory.ts b/packages/browserhive/src/composition/adapters/browser-memory.ts new file mode 100644 index 0000000..f5d04d0 --- /dev/null +++ b/packages/browserhive/src/composition/adapters/browser-memory.ts @@ -0,0 +1,77 @@ +/** @module composition/adapters/browser-memory — samples the resident memory of each live session's browser process tree every 10 s for `browserhive.browser.rss_bytes` (spec 10 §7). Only started when telemetry is on. */ + +import type { ProcessTreeReader } from '@browserhive/core/runtime'; +import { BROWSER_RSS_SAMPLE_INTERVAL_MS } from '@browserhive/core/runtime'; +import type { EveryFn } from './timers.ts'; + +/** A live session as the sampler sees it. */ +export interface SampledSession { + readonly id: string; + /** The browser's main pid (cached by the handle); absent before launch and in fakes. */ + readonly browserPid?: () => Promise; +} + +/** Dependencies of {@link startBrowserMemorySampler}. */ +export interface BrowserMemorySamplerDeps { + /** The sessions to sample now. */ + readonly sessions: () => Iterable; + /** `null` where the platform has no reader (Windows): the sampler then never runs. */ + readonly reader: ProcessTreeReader | null; + /** Repeating timer (`createTimers().every`). */ + readonly repeat: EveryFn; + readonly intervalMs?: number; + readonly onError?: (error: unknown) => void; +} + +/** A running sampler. */ +export interface BrowserMemorySampler { + /** RSS bytes per session id from the latest sample. */ + latest(): ReadonlyMap; + /** Takes one sample now (the timer calls this). */ + sample(): Promise; + stop(): void; +} + +/** Starts sampling; the first sample is taken at once so the first export has data. */ +export function startBrowserMemorySampler(deps: BrowserMemorySamplerDeps): BrowserMemorySampler { + let latest: ReadonlyMap = new Map(); + let running: Promise | undefined; + const reader = deps.reader; + const take = async (): Promise => { + if (reader === null) return; + const roots = new Map(); + for (const session of deps.sessions()) { + const pid = await session.browserPid?.().catch(() => null); + if (pid !== null && pid !== undefined) roots.set(pid, session.id); + } + const sums = await reader.rssOfTrees([...roots.keys()]); + const next = new Map(); + for (const [pid, bytes] of sums) { + const id = roots.get(pid); + if (id !== undefined) next.set(id, bytes); + } + latest = next; + }; + const sample = (): Promise => { + // Coalesce: a slow table read never stacks samples. + running ??= take() + .catch((err: unknown) => deps.onError?.(err)) + .finally(() => { + running = undefined; + }); + return running; + }; + const cancel = + reader === null + ? () => undefined + : deps.repeat(() => void sample(), deps.intervalMs ?? BROWSER_RSS_SAMPLE_INTERVAL_MS); + void sample(); + return { + latest: () => latest, + sample, + stop() { + cancel(); + latest = new Map(); + }, + }; +} diff --git a/packages/browserhive/src/composition/adapters/metrics.test.ts b/packages/browserhive/src/composition/adapters/metrics.test.ts index 78deb70..5e0efb1 100644 --- a/packages/browserhive/src/composition/adapters/metrics.test.ts +++ b/packages/browserhive/src/composition/adapters/metrics.test.ts @@ -1,8 +1,8 @@ -/** @module composition/adapters/metrics.test — the `harness` metric attribute is folded to the known slug table and never carries model, workspace or meta (spec 10 §7, D-30). */ +/** @module composition/adapters/metrics.test — what each bus event and observable source records (spec 10 §7): the `harness` attribute folded to the known slug table and never model, workspace or meta (D-30); launch duration once per session; lifetime by closed reason; operator requests by kind and the attention wait by status; pruned rows by table; the write queue, hub, browser-memory and event-loop callbacks; unwiring. */ import { describe, expect, it } from 'bun:test'; import type { DomainEvents, EventBus, Instruments } from '@browserhive/core/runtime'; -import { wireMetrics } from './metrics.ts'; +import { type MetricSources, wireMetrics } from './metrics.ts'; type Handler = (event: { name: string; at: number; payload: unknown }) => void; @@ -13,7 +13,12 @@ function fakeBus(): { bus: EventBus; emit(name: string, payload: u subscribeAll: () => () => undefined, subscribe: (name: string, handler: Handler) => { handlers.set(name, [...(handlers.get(name) ?? []), handler]); - return () => undefined; + return () => { + handlers.set( + name, + (handlers.get(name) ?? []).filter((h) => h !== handler), + ); + }; }, } as unknown as EventBus; return { @@ -24,31 +29,84 @@ function fakeBus(): { bus: EventBus; emit(name: string, payload: u }; } +type Observed = { value: number; attrs: Record }; +type Callback = (result: { observe(value: number, attrs?: Record): void }) => void; + function fakeInstruments() { const adds: { instrument: string; value: number; attrs: Record }[] = []; - const counter = (instrument: string) => ({ + const callbacks = new Map(); + const cache = new Map(); + const instrument = (name: string) => ({ add: (value: number, attrs: Record = {}) => - adds.push({ instrument, value, attrs }), + adds.push({ instrument: name, value, attrs }), record: (value: number, attrs: Record = {}) => - adds.push({ instrument, value, attrs }), - addCallback: () => undefined, - removeCallback: () => undefined, + adds.push({ instrument: name, value, attrs }), + addCallback: (cb: Callback) => callbacks.set(name, [...(callbacks.get(name) ?? []), cb]), + removeCallback: (cb: Callback) => + callbacks.set( + name, + (callbacks.get(name) ?? []).filter((c) => c !== cb), + ), }); const instruments = new Proxy( {}, - { get: (_target, name) => counter(String(name)) }, + { + get: (_target, name) => { + const key = String(name); + if (!cache.has(key)) cache.set(key, instrument(key)); + return cache.get(key); + }, + }, ) as unknown as Instruments; - return { instruments, adds }; + const collect = (name: string): Observed[] => { + const out: Observed[] = []; + for (const cb of callbacks.get(name) ?? []) { + cb({ observe: (value, attrs = {}) => out.push({ value, attrs }) }); + } + return out; + }; + const of = (name: string) => + adds.filter((a) => a.instrument === name).map((a) => [a.value, a.attrs]); + return { instruments, adds, callbacks, collect, of }; +} + +function sources(overrides: Partial = {}): MetricSources { + return { + queue: { depth: 0, droppedWritesByTable: new Map() }, + analytics: { databaseSize: async () => 0 }, + eventLoop: () => null, + ...overrides, + }; +} + +function summary(id: string, extra: Record = {}) { + return { + session_id: id, + created_at: 1_000, + harness: 'claude-code', + channel: 'chromium', + stealth: false, + ...extra, + }; +} + +function request(id: string, kind: string, extra: Record = {}) { + return { + request_id: id, + kind, + status: 'pending', + created_at: 1_000, + resolved_at: null, + waited_ms: null, + ...extra, + }; } describe('wireMetrics: harness attribute', () => { it('folds unknown slugs into other and keeps known ones', () => { const { bus, emit } = fakeBus(); const { instruments, adds } = fakeInstruments(); - const stop = wireMetrics(instruments, bus, { - queue: { depth: 0 }, - analytics: { databaseSize: async () => 0 }, - }); + const stop = wireMetrics(instruments, bus, sources()); const call = (harness: string) => emit('tool.called', { observation: { tool: 'navigate', ok: true, errorCode: null, durationMs: 5, harness }, @@ -56,7 +114,7 @@ describe('wireMetrics: harness attribute', () => { call('claude-code'); call('nightly-scraper'); emit('session.opened', { session: { session_id: 's-1', created_at: 1, harness: 'my-bot' } }); - emit('session.closed', { session_id: 's-1', closed_at: 2 }); + emit('session.closed', { session_id: 's-1', closed_at: 2, reason: 'user' }); stop(); const calls = adds.filter((a) => a.instrument === 'toolCalls').map((a) => a.attrs); expect(calls).toEqual([ @@ -74,4 +132,161 @@ describe('wireMetrics: harness attribute', () => { [-1, { harness: 'other' }], ]); }); + + it('adds error_code only on failed calls', () => { + const { bus, emit } = fakeBus(); + const { instruments, of } = fakeInstruments(); + wireMetrics(instruments, bus, sources()); + emit('tool.called', { + observation: { + tool: 'click', + ok: false, + errorCode: 'ELEMENT_NOT_FOUND', + durationMs: 9, + harness: 'unknown', + }, + }); + expect(of('toolCalls')).toEqual([ + [1, { tool: 'click', ok: false, harness: 'unknown', error_code: 'ELEMENT_NOT_FOUND' }], + ]); + }); +}); + +describe('wireMetrics: sessions', () => { + it('records the launch duration once per session, by channel and stealth', () => { + const { bus, emit } = fakeBus(); + const { instruments, of } = fakeInstruments(); + wireMetrics(instruments, bus, sources()); + emit('session.opened', { session: summary('s-1', { channel: 'chrome', stealth: true }) }); + // Reserved, then launched, then later patches that repeat `launchMs`. + emit('session.updated', { session: summary('s-1'), patch: { launchMs: null } }); + emit('session.updated', { + session: summary('s-1', { channel: 'chrome', stealth: true }), + patch: { launchMs: 850 }, + }); + emit('session.updated', { session: summary('s-1'), patch: { launchMs: 850 } }); + // A session this process never saw open (no baseline) is not a launch. + emit('session.updated', { session: summary('s-9'), patch: { launchMs: 5 } }); + expect(of('sessionLaunchDuration')).toEqual([[850, { channel: 'chrome', stealth: true }]]); + }); + + it('records the lifetime by closed reason', () => { + const { bus, emit } = fakeBus(); + const { instruments, of } = fakeInstruments(); + wireMetrics(instruments, bus, sources()); + emit('session.opened', { session: summary('s-1') }); + emit('session.closed', { session_id: 's-1', closed_at: 61_000, reason: 'lease_expired' }); + expect(of('sessionLifetime')).toEqual([[60_000, { closed_reason: 'lease_expired' }]]); + }); +}); + +describe('wireMetrics: operator requests', () => { + it('counts open requests by kind and records the attention wait by status', () => { + const { bus, emit } = fakeBus(); + const { instruments, of } = fakeInstruments(); + wireMetrics(instruments, bus, sources()); + emit('attention.created', { request: request('a-1', 'attention') }); + emit('vault.confirm.created', { request: request('v-1', 'vault_confirm') }); + emit('attention.resolved', { + request: request('a-1', 'attention', { + status: 'timeout', + resolved_at: 4_000, + waited_ms: 3_000, + }), + }); + emit('vault.confirm.resolved', { request: request('v-1', 'vault_confirm') }); + // Settled by the startup reconcile, opened by a previous run: no negative count. + emit('attention.resolved', { + request: request('a-0', 'attention', { status: 'rejected', resolved_at: 1_500 }), + }); + expect(of('attentionOpen')).toEqual([ + [1, { kind: 'attention' }], + [1, { kind: 'vault_confirm' }], + [-1, { kind: 'attention' }], + [-1, { kind: 'vault_confirm' }], + ]); + expect(of('attentionWait')).toEqual([ + [3_000, { status: 'timeout' }], + [500, { status: 'rejected' }], + ]); + }); +}); + +describe('wireMetrics: retention', () => { + it('adds pruned rows per table, skipping tables that lost none', () => { + const { bus, emit } = fakeBus(); + const { instruments, of } = fakeInstruments(); + wireMetrics(instruments, bus, sources()); + emit('retention.completed', { + pruned_rows: 7, + prunedByTable: { tool_calls: 5, pages: 2, logs: 0 }, + }); + expect(of('retentionPrunedRows')).toEqual([ + [5, { table: 'tool_calls' }], + [2, { table: 'pages' }], + ]); + }); +}); + +describe('wireMetrics: observable instruments', () => { + it('reads the queue, the hub, the browser sampler, the event loop and this process', async () => { + const { bus } = fakeBus(); + const { instruments, collect } = fakeInstruments(); + let hubOpen = false; + const sockets = Array.from({ length: 60 }, (_, i) => ({ + connectionId: `c-${i}`, + bufferedBytes: i, + })); + wireMetrics( + instruments, + bus, + sources({ + queue: { depth: 3, droppedWritesByTable: new Map([['tool_calls', 2]]) }, + analytics: { databaseSize: async () => 4096 }, + realtime: () => + hubOpen + ? { sockets: () => sockets, droppedFrames: () => ({ screencast: 4, logs: 1, feed: 0 }) } + : undefined, + browserMemory: () => new Map([['s-1', 123_456]]), + eventLoop: () => ({ takeP99Ms: () => 12.5, stop: () => undefined }), + }), + ); + await Promise.resolve(); + expect(collect('dbWriteQueueDepth')).toEqual([{ value: 3, attrs: {} }]); + expect(collect('dbDroppedWrites')).toEqual([{ value: 2, attrs: { table: 'tool_calls' } }]); + expect(collect('dbSizeBytes')).toEqual([{ value: 4096, attrs: {} }]); + expect(collect('browserRssBytes')).toEqual([{ value: 123_456, attrs: { session_id: 's-1' } }]); + expect(collect('processEventLoopLag')).toEqual([{ value: 12.5, attrs: {} }]); + expect(collect('processRssBytes')[0]?.value).toBeGreaterThan(0); + expect(collect('processHeapBytes')[0]?.value).toBeGreaterThan(0); + // Before the listeners open (and under stdio) the hub instruments have no data points. + expect(collect('wsConnections')).toEqual([]); + hubOpen = true; + expect(collect('wsConnections')).toEqual([{ value: 60, attrs: {} }]); + const buffered = collect('wsBufferedBytes'); + expect(buffered).toHaveLength(50); + expect(buffered[0]).toEqual({ value: 59, attrs: { connection_id: 'c-59' } }); + expect(collect('wsFramesDropped')).toEqual([ + { value: 4, attrs: { channel: 'screencast' } }, + { value: 1, attrs: { channel: 'logs' } }, + { value: 0, attrs: { channel: 'feed' } }, + ]); + }); + + it('unwiring removes every callback and subscription and stops the event-loop monitor', () => { + const { bus, emit } = fakeBus(); + const { instruments, callbacks, adds } = fakeInstruments(); + let stopped = 0; + const unwire = wireMetrics( + instruments, + bus, + sources({ eventLoop: () => ({ takeP99Ms: () => 1, stop: () => void stopped++ }) }), + ); + expect([...callbacks.values()].flat().length).toBe(10); + unwire(); + expect([...callbacks.values()].flat().length).toBe(0); + expect(stopped).toBe(1); + emit('attention.created', { request: request('a-1', 'attention') }); + expect(adds).toEqual([]); + }); }); diff --git a/packages/browserhive/src/composition/adapters/metrics.ts b/packages/browserhive/src/composition/adapters/metrics.ts index 785efcd..32450bd 100644 --- a/packages/browserhive/src/composition/adapters/metrics.ts +++ b/packages/browserhive/src/composition/adapters/metrics.ts @@ -1,22 +1,78 @@ -/** @module composition/adapters/metrics — OTel metrics consumers: bus events → counters/histograms, observable gauges over queue depth, DB size and process memory (spec 10 §7). Only wired when telemetry is on. */ +/** @module composition/adapters/metrics — OTel metrics consumers (spec 10 §7): bus events → counters/histograms; observable instruments over the write queue, the database size, the realtime hub, the browser-memory sampler and this process. Only wired when telemetry is on. */ import { metricHarness } from '@browserhive/contracts/harness'; import type { AnalyticsQueries, DomainEvents, EventBus, + EventLoopLagMonitor, Instruments, WriteQueue, } from '@browserhive/core/runtime'; -import { readProcessMemory } from '@browserhive/core/runtime'; +import { + readProcessMemory, + startEventLoopLagMonitor, + WS_BUFFERED_BYTES_MAX_SERIES, +} from '@browserhive/core/runtime'; +import type { ActionCounter, DeliveryCounter, ReportCounter } from '@browserhive/core/server'; + +/** What the realtime hub exposes to the metrics (http transport only). */ +export interface RealtimeMetricsSource { + /** Open connections and the bytes each socket has not flushed yet. */ + sockets(): readonly { readonly connectionId: string; readonly bufferedBytes: number }[]; + /** Frames dropped since start, by channel. */ + droppedFrames(): Readonly>; +} + +/** The metrics view of the realtime hub, as `listeners-http` exposes it on `ctx.listeners.realtime`. */ +export function realtimeMetrics(hub: { + connections(): readonly { readonly connection_id: string; readonly buffered_bytes: number }[]; + droppedFrames(): Readonly>; +}): RealtimeMetricsSource { + return { + sockets: () => + hub + .connections() + .map((c) => ({ connectionId: c.connection_id, bufferedBytes: c.buffered_bytes })), + droppedFrames: () => hub.droppedFrames(), + }; +} -/** Sources of the observable gauges. */ +/** Sources of the observable instruments. */ export interface MetricSources { - readonly queue: Pick; + readonly queue: Pick; readonly analytics: Pick; + /** The hub once the listeners are open; `undefined` before that and under stdio. */ + readonly realtime?: () => RealtimeMetricsSource | undefined; + /** RSS bytes per session id from the browser-memory sampler's latest sample. */ + readonly browserMemory?: () => ReadonlyMap; + /** Starts the event-loop monitor (default the `perf_hooks` one); `null` disables the gauge. */ + readonly eventLoop?: () => EventLoopLagMonitor | null; +} + +interface Observer { + observe(value: number, attributes?: Record): void; } -/** Subscribes the instrument consumers; returns the unsubscribe function. */ +interface ObservableInstrument { + addCallback(callback: (result: Observer) => void): void; + removeCallback(callback: (result: Observer) => void): void; +} + +/** The notification services' counters, exactly as the composition root passes them (spec 10 §7). */ +export function notificationCounters(instruments: Instruments): { + readonly deliveryCounter: DeliveryCounter; + readonly actionCounter: ActionCounter; + readonly reportCounter: ReportCounter; +} { + return { + deliveryCounter: instruments.notificationDeliveries, + actionCounter: instruments.notificationActions, + reportCounter: instruments.notificationReports, + }; +} + +/** Subscribes the instrument consumers and registers the callbacks; returns the unwire function. */ export function wireMetrics( instruments: Instruments, bus: EventBus, @@ -24,8 +80,22 @@ export function wireMetrics( ): () => void { const offs: (() => void)[] = []; const openedAt = new Map(); + const launched = new Set(); + // Operator requests seen opening: a settlement of one opened before this process started (the + // startup reconcile) never drives the up-down counter below zero. + const openRequests = new Map(); // `harness` is folded to the known slug table (spec 10 §7 cardinality rule); never model/workspace/meta. const harnessOf = new Map(); + const requestOpened = (request: { readonly request_id: string; readonly kind: string }) => { + openRequests.set(request.request_id, request.kind); + instruments.attentionOpen.add(1, { kind: request.kind }); + }; + const requestSettled = (request: { readonly request_id: string }) => { + const kind = openRequests.get(request.request_id); + if (kind === undefined) return; + openRequests.delete(request.request_id); + instruments.attentionOpen.add(-1, { kind }); + }; offs.push( bus.subscribe('tool.called', ({ payload }) => { const o = payload.observation; @@ -44,54 +114,131 @@ export function wireMetrics( harnessOf.set(payload.session.session_id, harness); openedAt.set(payload.session.session_id, payload.session.created_at); }), + bus.subscribe('session.updated', ({ payload }) => { + // `launch_ms` is set once, when the browser is attached; later patches repeat it. + const launchMs = payload.patch.launchMs; + const id = payload.session.session_id; + if (typeof launchMs !== 'number' || launched.has(id) || !openedAt.has(id)) return; + launched.add(id); + instruments.sessionLaunchDuration.record(launchMs, { + channel: payload.session.channel, + stealth: payload.session.stealth, + }); + }), bus.subscribe('session.closed', ({ payload }) => { instruments.sessionsActive.add(-1, { harness: harnessOf.get(payload.session_id) ?? metricHarness(null), }); harnessOf.delete(payload.session_id); + launched.delete(payload.session_id); const started = openedAt.get(payload.session_id); openedAt.delete(payload.session_id); if (started !== undefined) { - instruments.sessionLifetime.record(Math.max(0, payload.closed_at - started)); + instruments.sessionLifetime.record(Math.max(0, payload.closed_at - started), { + closed_reason: payload.reason, + }); + } + }), + bus.subscribe('attention.created', ({ payload }) => requestOpened(payload.request)), + bus.subscribe('vault.confirm.created', ({ payload }) => requestOpened(payload.request)), + bus.subscribe('vault.confirm.resolved', ({ payload }) => requestSettled(payload.request)), + bus.subscribe('attention.resolved', ({ payload }) => { + const r = payload.request; + requestSettled(r); + const waited = r.waited_ms ?? (r.resolved_at === null ? null : r.resolved_at - r.created_at); + if (waited !== null) { + instruments.attentionWait.record(Math.max(0, waited), { status: r.status }); } }), - bus.subscribe('attention.created', () => instruments.attentionOpen.add(1)), - bus.subscribe('attention.resolved', () => instruments.attentionOpen.add(-1)), bus.subscribe('vault.access', ({ payload }) => instruments.vaultFills.add(1, { result: payload.row.result }), ), bus.subscribe('blocklist.hit', ({ payload }) => instruments.blocklistHits.add(1, { source: payload.row.source }), ), - bus.subscribe('retention.completed', ({ payload }) => - instruments.retentionPrunedRows.add(payload.pruned_rows), - ), + bus.subscribe('retention.completed', ({ payload }) => { + for (const [table, rows] of Object.entries(payload.prunedByTable ?? {})) { + if (rows > 0) instruments.retentionPrunedRows.add(rows, { table }); + } + }), ); - let dbBytes = 0; - const depth = (result: { observe(value: number): void }) => result.observe(sources.queue.depth); - const size = (result: { observe(value: number): void }) => { - result.observe(dbBytes); + // The size is read in the background and reported at the next export (never blocks a collect). + let dbBytes: number | undefined; + const refreshSize = () => void sources.analytics.databaseSize().then( (bytes) => { dbBytes = bytes; }, () => undefined, ); - }; - const rss = (result: { observe(value: number): void }) => - result.observe(readProcessMemory().rssBytes); - const heap = (result: { observe(value: number): void }) => - result.observe(readProcessMemory().heapUsedBytes); - instruments.dbWriteQueueDepth.addCallback(depth); - instruments.dbSizeBytes.addCallback(size); - instruments.processRssBytes.addCallback(rss); - instruments.processHeapBytes.addCallback(heap); + refreshSize(); + const monitor = (sources.eventLoop ?? startEventLoopLagMonitor)(); + + const hub = () => sources.realtime?.(); + const callbacks: readonly [ObservableInstrument, (result: Observer) => void][] = [ + [instruments.dbWriteQueueDepth, (r) => r.observe(sources.queue.depth)], + [ + instruments.dbDroppedWrites, + (r) => { + for (const [table, count] of sources.queue.droppedWritesByTable) { + r.observe(count, { table }); + } + }, + ], + [ + instruments.dbSizeBytes, + (r) => { + if (dbBytes !== undefined) r.observe(dbBytes); + refreshSize(); + }, + ], + [instruments.processRssBytes, (r) => r.observe(readProcessMemory().rssBytes)], + [instruments.processHeapBytes, (r) => r.observe(readProcessMemory().heapUsedBytes)], + [ + instruments.processEventLoopLag, + (r) => { + const p99 = monitor?.takeP99Ms() ?? null; + if (p99 !== null) r.observe(p99); + }, + ], + [ + instruments.wsConnections, + (r) => { + const h = hub(); + if (h !== undefined) r.observe(h.sockets().length); + }, + ], + [ + instruments.wsBufferedBytes, + (r) => { + const top = [...(hub()?.sockets() ?? [])] + .sort((a, b) => b.bufferedBytes - a.bufferedBytes) + .slice(0, WS_BUFFERED_BYTES_MAX_SERIES); + for (const s of top) r.observe(s.bufferedBytes, { connection_id: s.connectionId }); + }, + ], + [ + instruments.wsFramesDropped, + (r) => { + for (const [channel, count] of Object.entries(hub()?.droppedFrames() ?? {})) { + r.observe(count, { channel }); + } + }, + ], + [ + instruments.browserRssBytes, + (r) => { + for (const [sessionId, bytes] of sources.browserMemory?.() ?? []) { + r.observe(bytes, { session_id: sessionId }); + } + }, + ], + ]; + for (const [instrument, callback] of callbacks) instrument.addCallback(callback); offs.push(() => { - instruments.dbWriteQueueDepth.removeCallback(depth); - instruments.dbSizeBytes.removeCallback(size); - instruments.processRssBytes.removeCallback(rss); - instruments.processHeapBytes.removeCallback(heap); + for (const [instrument, callback] of callbacks) instrument.removeCallback(callback); + monitor?.stop(); }); return () => { for (const off of offs.splice(0)) off(); diff --git a/packages/browserhive/src/composition/context.ts b/packages/browserhive/src/composition/context.ts index cde0b64..45e0f2f 100644 --- a/packages/browserhive/src/composition/context.ts +++ b/packages/browserhive/src/composition/context.ts @@ -53,6 +53,7 @@ import type { VaultService, } from '@browserhive/core/server'; import type { DegradationRelay } from './adapters/degradation-relay.ts'; +import type { RealtimeMetricsSource } from './adapters/metrics.ts'; import type { DataDirLayout } from './data-dir.ts'; import type { PhaseTracker } from './health.ts'; import type { DataDirLock } from './lock-file.ts'; @@ -157,7 +158,10 @@ export interface ListenersPart { /** Open MCP Streamable HTTP sessions (0 under stdio). */ mcpConnections(): number; /** WS hub counters; absent under stdio. */ - readonly realtime?: { connections(): number; activeScreencasts(): number }; + readonly realtime?: RealtimeMetricsSource & { + connections(): number; + activeScreencasts(): number; + }; } /** Boot-time record shared by every phase. */ diff --git a/packages/browserhive/src/composition/phases/build-domain.ts b/packages/browserhive/src/composition/phases/build-domain.ts index 69f08dc..276cfbc 100644 --- a/packages/browserhive/src/composition/phases/build-domain.ts +++ b/packages/browserhive/src/composition/phases/build-domain.ts @@ -24,6 +24,7 @@ import { createPlaywrightPageActions, InProcessEventBus, } from '@browserhive/core/server'; +import { notificationCounters } from '../adapters/metrics.ts'; import { asyncTick, createTimers } from '../adapters/timers.ts'; import { createAuthStack } from '../auth-stack.ts'; import { type BootContext, part, type SeedNotice } from '../context.ts'; @@ -177,6 +178,8 @@ async function buildDomain( }, }); const ops = buildOps({ + // `browserhive.notifications.*` (spec 10 §7): no-ops unless telemetry is on. + ...notificationCounters(telemetry.instruments), config, repos, queue: storage.queue, @@ -193,7 +196,6 @@ async function buildDomain( env: ctx.input.env, registerSecret: (literal) => secrets.add(literal), dashboardUrl: () => ctx.listeners?.url ?? `http://${config.host}:${config.port}`, - deliveryCounter: telemetry.instruments.notificationDeliveries, channelFactories: channelFactories({ images, telegramUpdates, @@ -205,14 +207,12 @@ async function buildDomain( telegram: createTelegramSetup({ updates: telegramUpdates }), discord: createDiscordSetup({ gateway: discordGateway }), actionExecutors, - actionCounter: telemetry.instruments.notificationActions, probe: createUrlProbe(), instanceId, capacity: () => { const status = sessions.serverStatus(); return { live: status.count, max: status.limit ?? 0 }; }, - reportCounter: telemetry.instruments.notificationReports, snapshots: createNotificationSnapshots({ sessions, screenshots: repos.screenshots, diff --git a/packages/browserhive/src/composition/phases/listeners-http.ts b/packages/browserhive/src/composition/phases/listeners-http.ts index 934c204..5bdf9be 100644 --- a/packages/browserhive/src/composition/phases/listeners-http.ts +++ b/packages/browserhive/src/composition/phases/listeners-http.ts @@ -34,6 +34,7 @@ import { notificationsPort, preferencesPort, } from '../adapters/http-ports.ts'; +import { realtimeMetrics } from '../adapters/metrics.ts'; import { createTimers } from '../adapters/timers.ts'; import { type BootContext, part } from '../context.ts'; import { resolveDashboardDir } from '../dashboard-dir.ts'; @@ -297,6 +298,7 @@ export async function openHttpListener( realtime: { connections: () => realtime.hub.connections().length, activeScreencasts: () => realtime.hub.activeScreencasts(), + ...realtimeMetrics(realtime.hub), }, }; log.info('listening', { url }); diff --git a/packages/browserhive/src/composition/phases/wire-observers.ts b/packages/browserhive/src/composition/phases/wire-observers.ts index 05014aa..f945334 100644 --- a/packages/browserhive/src/composition/phases/wire-observers.ts +++ b/packages/browserhive/src/composition/phases/wire-observers.ts @@ -2,9 +2,11 @@ import { createRequire } from 'node:module'; import { listBackups } from '@browserhive/core/persistence'; -import { serializeError } from '@browserhive/core/runtime'; +import { createProcessTreeReader, type Logger, serializeError } from '@browserhive/core/runtime'; import { pinnedPlaywrightVersion, SystemStatusService } from '@browserhive/core/server'; +import { startBrowserMemorySampler } from '../adapters/browser-memory.ts'; import { wireMetrics } from '../adapters/metrics.ts'; +import { createTimers } from '../adapters/timers.ts'; import { type BootContext, part } from '../context.ts'; import type { PhaseHandle } from '../unwind.ts'; @@ -22,6 +24,43 @@ export function patchrightVersion(): string | null { return null; } +/** + * The OTel metrics consumers and the browser-memory sampler (spec 10 §7); only called with + * telemetry on, so with `--otel` off nothing is subscribed, sampled or timed. + */ +function wireMetricsFor(ctx: BootContext, logger: Logger): () => void { + const { telemetry } = part(ctx.observability, 'observability'); + const storage = part(ctx.storage, 'storage'); + const domain = part(ctx.domain, 'domain'); + const report = (err: unknown) => + logger.warn('metrics sampler failed', { err: serializeError(err) }); + const sampler = startBrowserMemorySampler({ + sessions: () => + domain.sessions.listAll().map((session) => { + const handle = session.handle; + return { + id: session.id, + ...(handle?.browserPid !== undefined && { + browserPid: () => handle.browserPid?.() ?? Promise.resolve(null), + }), + }; + }), + reader: createProcessTreeReader(), + repeat: createTimers(report).every, + onError: report, + }); + const unwire = wireMetrics(telemetry.instruments, domain.bus, { + queue: storage.queue, + analytics: storage.analytics, + realtime: () => ctx.listeners?.realtime, + browserMemory: () => sampler.latest(), + }); + return () => { + unwire(); + sampler.stop(); + }; +} + /** Phase `wire-observers`. */ export async function wireObserversPhase(ctx: BootContext): Promise { ctx.health.enter('wire-observers'); @@ -116,12 +155,7 @@ export async function wireObserversPhase(ctx: BootContext): Promise .refreshLastBackup() .catch((err: unknown) => logger.warn('backup scan failed', { err: serializeError(err) })); domain.sweeper.start(); - const unwireMetrics = telemetry.enabled - ? wireMetrics(telemetry.instruments, domain.bus, { - queue: storage.queue, - analytics: storage.analytics, - }) - : () => undefined; + const unwireMetrics = telemetry.enabled ? wireMetricsFor(ctx, logger) : () => undefined; ctx.observers = { status }; return { diff --git a/packages/core/src/app/events/catalog.ts b/packages/core/src/app/events/catalog.ts index a074676..74868c4 100644 --- a/packages/core/src/app/events/catalog.ts +++ b/packages/core/src/app/events/catalog.ts @@ -148,7 +148,10 @@ export type DomainEvents = { readonly status: 'starting' | 'ready' | 'degraded' | 'stopping'; readonly at: number; }; - readonly 'retention.completed': z.infer; + readonly 'retention.completed': z.infer & { + /** Rows deleted per table in this pass (internal: the `retention.pruned_rows{table}` metric). */ + readonly prunedByTable?: Readonly>; + }; // notifications readonly 'notification.created': z.infer; readonly 'notification.updated': z.infer; diff --git a/packages/core/src/app/maintenance/retention-scheduler.ts b/packages/core/src/app/maintenance/retention-scheduler.ts index 5ccff95..35d7f99 100644 --- a/packages/core/src/app/maintenance/retention-scheduler.ts +++ b/packages/core/src/app/maintenance/retention-scheduler.ts @@ -136,6 +136,7 @@ export class RetentionScheduler { type: 'retention.completed', at: run.at, pruned_rows: run.prunedRows, + prunedByTable: run.result?.prunedRows ?? {}, result: run.outcome, ...(run.outcome !== 'ok' && { severity: run.outcome === 'failed' ? 'error' : 'warn' }), }); diff --git a/packages/core/src/app/notifications/report-scheduler.ts b/packages/core/src/app/notifications/report-scheduler.ts index d32a950..0408093 100644 --- a/packages/core/src/app/notifications/report-scheduler.ts +++ b/packages/core/src/app/notifications/report-scheduler.ts @@ -662,6 +662,8 @@ export class ReportScheduler { */ async storeManualCopy(built: BuiltReport): Promise { const now = this.deps.clock.now(); + // Only the send path stores a copy, so this is where an on-demand digest is counted. + this.count(built.message.kind, 'manual'); const thread = manualPeriodThread(built.ctx.zone, built.window); const inApp = await this.inAppDigest(thread, built.rule, built.facts, now, built.ctx); let found: { id: string | null; inserted: NotificationRecord | null } = { @@ -860,7 +862,10 @@ export class ReportScheduler { if (step.revised !== null) revised.push(step.revised); if (step.outcome !== 'none') { decisions++; - this.count('report.anomaly', step.outcome === 'sent' ? 'in_app' : step.outcome); + // A silent revision is not a decision (spec 10 §7), as for a channel's alert. + if (step.outcome !== 'revised') { + this.count('report.anomaly', step.outcome === 'sent' ? 'in_app' : step.outcome); + } } } if (rows.length === 0 && revised.length === 0 && sameWatches(stored, next)) return 0; diff --git a/packages/core/src/infra/browsers/session-handle.ts b/packages/core/src/infra/browsers/session-handle.ts index c341c14..01232a2 100644 --- a/packages/core/src/infra/browsers/session-handle.ts +++ b/packages/core/src/infra/browsers/session-handle.ts @@ -52,6 +52,7 @@ export class PlaywrightSessionHandle implements SessionHandle { private readonly crashListeners = new Set<(reason: string) => void>(); private closing = false; private closed: Promise | undefined; + private pid: Promise | undefined; constructor(input: SessionHandleInput) { this.sessionId = input.sessionId; @@ -151,6 +152,32 @@ export class PlaywrightSessionHandle implements SessionHandle { return warnings; } + /** + * The browser's main process id, from a browser-level DevTools session opened once + * (`SystemInfo.getProcessInfo`); no page target is touched. Cached, including a `null`. + */ + browserPid(): Promise { + this.pid ??= this.readBrowserPid(); + return this.pid; + } + + private async readBrowserPid(): Promise { + const browser = this.browser ?? this.context.browser(); + if (browser === null || this.closing) return null; + try { + const cdp = await browser.newBrowserCDPSession(); + try { + const info = await cdp.send('SystemInfo.getProcessInfo'); + return info.processInfo.find((p) => p.type === 'browser')?.id ?? null; + } finally { + await cdp.detach().catch(() => undefined); + } + } catch (err) { + this.logger.debug('browser pid unavailable', { err: serializeError(err) }); + return null; + } + } + /** Last resort on overrun: the process, not the protocol, is what holds the resources. */ private killProcess(): void { const browser = this.browser ?? this.context.browser(); diff --git a/packages/core/src/infra/host/event-loop.ts b/packages/core/src/infra/host/event-loop.ts new file mode 100644 index 0000000..ffb74ac --- /dev/null +++ b/packages/core/src/infra/host/event-loop.ts @@ -0,0 +1,39 @@ +/** @module infra/host/event-loop — event-loop delay of this process over a window (spec 10 §7 `browserhive.process.event_loop_lag`). */ + +import { monitorEventLoopDelay } from 'node:perf_hooks'; + +/** A running event-loop delay monitor. */ +export interface EventLoopLagMonitor { + /** 99th percentile delay in ms since the previous call (or the start), then starts a new window; `null` without samples. */ + takeP99Ms(): number | null; + /** Stops sampling. Idempotent. */ + stop(): void; +} + +/** + * Starts sampling the event-loop delay every `resolutionMs`. Returns `null` where the runtime has + * no `monitorEventLoopDelay`; nothing runs until this is called. + */ +export function startEventLoopLagMonitor(resolutionMs = 20): EventLoopLagMonitor | null { + let histogram: ReturnType; + try { + histogram = monitorEventLoopDelay({ resolution: resolutionMs }); + histogram.enable(); + } catch { + return null; + } + let stopped = false; + return { + takeP99Ms() { + if (histogram.count === 0) return null; + const p99 = histogram.percentile(99) / 1e6; + histogram.reset(); + return Number.isFinite(p99) ? p99 : null; + }, + stop() { + if (stopped) return; + stopped = true; + histogram.disable(); + }, + }; +} diff --git a/packages/core/src/infra/host/process-tree.ts b/packages/core/src/infra/host/process-tree.ts new file mode 100644 index 0000000..8816ed6 --- /dev/null +++ b/packages/core/src/infra/host/process-tree.ts @@ -0,0 +1,123 @@ +/** @module infra/host/process-tree — resident memory of process trees (a browser and every descendant), read from `/proc` on Linux and `ps` on macOS (spec 10 §7 `browserhive.browser.rss_bytes`). */ + +import { execFile } from 'node:child_process'; +import { readdir, readFile } from 'node:fs/promises'; + +/** One process of the host's table. */ +export interface ProcessEntry { + readonly pid: number; + readonly ppid: number; + readonly rssBytes: number; +} + +/** Reads the RSS of whole process trees. */ +export interface ProcessTreeReader { + /** + * Sums the RSS of each root and all its descendants. Roots that no longer exist are absent from + * the result. Never rejects: an unreadable table yields an empty map. + */ + rssOfTrees(roots: readonly number[]): Promise>; +} + +/** Sums each root's tree over a process table (pure). */ +export function sumProcessTrees( + table: readonly ProcessEntry[], + roots: readonly number[], +): Map { + const children = new Map(); + const byPid = new Map(); + for (const entry of table) { + byPid.set(entry.pid, entry); + const siblings = children.get(entry.ppid); + if (siblings === undefined) children.set(entry.ppid, [entry]); + else siblings.push(entry); + } + const sums = new Map(); + for (const root of roots) { + const top = byPid.get(root); + if (top === undefined) continue; + let total = 0; + const seen = new Set(); + const stack: ProcessEntry[] = [top]; + for (let next = stack.pop(); next !== undefined; next = stack.pop()) { + if (seen.has(next.pid)) continue; + seen.add(next.pid); + total += next.rssBytes; + for (const child of children.get(next.pid) ?? []) stack.push(child); + } + sums.set(root, total); + } + return sums; +} + +/** Parses `/proc//status` into the parent pid and `VmRSS` (kernel threads have none: 0). */ +export function parseProcStatus(pid: number, text: string): ProcessEntry | null { + const ppid = /^PPid:\s+(\d+)/m.exec(text)?.[1]; + if (ppid === undefined) return null; + const rssKb = /^VmRSS:\s+(\d+)\s+kB/m.exec(text)?.[1]; + return { pid, ppid: Number(ppid), rssBytes: rssKb === undefined ? 0 : Number(rssKb) * 1024 }; +} + +/** Parses `ps -A -o pid=,ppid=,rss=` output (RSS in KiB). */ +export function parsePsTable(text: string): ProcessEntry[] { + const entries: ProcessEntry[] = []; + for (const line of text.split('\n')) { + const [pid, ppid, rss] = line.trim().split(/\s+/); + if (pid === undefined || ppid === undefined || rss === undefined) continue; + const entry = { pid: Number(pid), ppid: Number(ppid), rssBytes: Number(rss) * 1024 }; + if ( + Number.isFinite(entry.pid) && + Number.isFinite(entry.ppid) && + Number.isFinite(entry.rssBytes) + ) + entries.push(entry); + } + return entries; +} + +async function linuxTable(): Promise { + const names = await readdir('/proc'); + const reads = names + .filter((name) => /^\d+$/.test(name)) + .map(async (name) => { + try { + return parseProcStatus(Number(name), await readFile(`/proc/${name}/status`, 'utf8')); + } catch { + // The process exited between the listing and the read. + return null; + } + }); + return (await Promise.all(reads)).filter((e): e is ProcessEntry => e !== null); +} + +function psTable(): Promise { + return new Promise((resolve) => { + execFile( + 'ps', + ['-A', '-o', 'pid=,ppid=,rss='], + { timeout: 5_000, maxBuffer: 8 * 1024 * 1024 }, + (err, stdout) => resolve(err === null ? parsePsTable(stdout) : []), + ); + }); +} + +/** + * The reader for `platform`: `/proc` on Linux, `ps` on macOS, `null` elsewhere (Windows has no + * cheap equivalent; the metric then has no data points, as spec 10 §7 says). + */ +export function createProcessTreeReader( + platform: NodeJS.Platform = process.platform, +): ProcessTreeReader | null { + const table = platform === 'linux' ? linuxTable : platform === 'darwin' ? psTable : null; + if (table === null) return null; + return { + async rssOfTrees(roots) { + if (roots.length === 0) return new Map(); + try { + return sumProcessTrees(await table(), roots); + } catch { + return new Map(); + } + }, + }; +} diff --git a/packages/core/src/infra/persistence/write-queue.ts b/packages/core/src/infra/persistence/write-queue.ts index 36d605c..bfb2792 100644 --- a/packages/core/src/infra/persistence/write-queue.ts +++ b/packages/core/src/infra/persistence/write-queue.ts @@ -35,6 +35,7 @@ export class SqliteWriteQueue implements WriteQueue { #scheduled = false; #closed = false; #dropped = 0; + readonly #droppedByTable = new Map(); constructor(options: WriteQueueOptions) { this.#uow = options.uow; @@ -51,9 +52,13 @@ export class SqliteWriteQueue implements WriteQueue { return this.#dropped; } + get droppedWritesByTable(): ReadonlyMap { + return this.#droppedByTable; + } + enqueue(operation: string, job: WriteJob): boolean { if (this.#closed || this.#queue.length >= this.#maxDepth) { - this.#dropped += 1; + this.#drop(operation, 1); this.#logger.warn('write dropped', { operation, reason: this.#closed ? 'closed' : 'full' }); return false; } @@ -93,7 +98,7 @@ export class SqliteWriteQueue implements WriteQueue { try { await item.job(repos); } catch (error) { - this.#dropped += 1; + this.#drop(item.operation, 1); this.#logger.error('write failed', { operation: item.operation, error: String(error), @@ -103,13 +108,19 @@ export class SqliteWriteQueue implements WriteQueue { }); } catch (error) { // The transaction itself failed (lock timeout, I/O): every job of the batch is lost. - this.#dropped += batch.length; + for (const item of batch) this.#drop(item.operation, 1); this.#logger.error('drain failed', { jobs: batch.length, error: String(error) }); } }, ); } + #drop(operation: string, count: number): void { + this.#dropped += count; + const table = operation.split('.', 1)[0] || 'unknown'; + this.#droppedByTable.set(table, (this.#droppedByTable.get(table) ?? 0) + count); + } + async close(): Promise { if (this.#closed) return; await this.drain(); diff --git a/packages/core/src/infra/telemetry/metrics.ts b/packages/core/src/infra/telemetry/metrics.ts index 5cdf498..ccf11bc 100644 --- a/packages/core/src/infra/telemetry/metrics.ts +++ b/packages/core/src/infra/telemetry/metrics.ts @@ -1,6 +1,14 @@ -/** @module infra/telemetry/metrics — the instrument catalogue, created lazily from a Meter (spec 10 §7). */ +/** @module infra/telemetry/metrics — the instrument catalogue (spec 10 §7, the single source of truth with the docs table), created lazily from a Meter. */ -import type { Counter, Histogram, Meter, ObservableGauge, UpDownCounter } from '@opentelemetry/api'; +import type { + Counter, + Histogram, + Meter, + ObservableCounter, + ObservableGauge, + ObservableUpDownCounter, + UpDownCounter, +} from '@opentelemetry/api'; /** Instrument names from the catalogue. */ export const METRIC = { @@ -29,9 +37,138 @@ export const METRIC = { PROCESS_EVENT_LOOP_LAG: 'browserhive.process.event_loop_lag', } as const; +/** A catalogue name. */ +export type MetricName = (typeof METRIC)[keyof typeof METRIC]; + /** Cap on `connection_id` series for the buffered-bytes gauge (spec 10 §7). */ export const WS_BUFFERED_BYTES_MAX_SERIES = 50; +/** How often the browser-memory sampler reads the process trees (spec 10 §7). */ +export const BROWSER_RSS_SAMPLE_INTERVAL_MS = 10_000; + +/** The OTLP data a collector receives: a monotonic sum, a non-monotonic sum, a histogram, a gauge. */ +export type MetricKind = 'counter' | 'up-down counter' | 'histogram' | 'gauge'; + +/** One row of the spec 10 §7 table. */ +export interface MetricDefinition { + readonly name: MetricName; + readonly kind: MetricKind; + /** Read by a callback at each export rather than written per event. */ + readonly observable: boolean; + /** UCUM unit (`ms`, `By`) or a `{annotation}` count. */ + readonly unit: string; + /** Every attribute key a data point may carry. */ + readonly attributes: readonly string[]; + readonly description: string; +} + +const def = ( + name: MetricName, + kind: MetricKind, + unit: string, + attributes: readonly string[], + description: string, + observable = false, +): MetricDefinition => ({ name, kind, observable, unit, attributes, description }); + +/** The catalogue, in the order of the spec 10 §7 table. */ +export const METRIC_DEFINITIONS: readonly MetricDefinition[] = [ + def( + METRIC.TOOL_CALLS, + 'counter', + '{call}', + ['tool', 'ok', 'error_code', 'harness'], + 'Tool calls by tool, ok, error_code and harness', + ), + def(METRIC.TOOL_CALL_DURATION, 'histogram', 'ms', ['tool'], 'Tool call duration'), + def(METRIC.SESSIONS_ACTIVE, 'up-down counter', '{session}', ['harness'], 'Live sessions'), + def( + METRIC.SESSION_LAUNCH_DURATION, + 'histogram', + 'ms', + ['channel', 'stealth'], + 'Create request to browser attached, per successful launch', + ), + def(METRIC.SESSION_LIFETIME, 'histogram', 'ms', ['closed_reason'], 'Session lifetime'), + def( + METRIC.WS_CONNECTIONS, + 'up-down counter', + '{connection}', + [], + 'Open WebSocket connections', + true, + ), + def( + METRIC.WS_BUFFERED_BYTES, + 'gauge', + 'By', + ['connection_id'], + 'Per-connection buffered bytes', + true, + ), + def( + METRIC.WS_FRAMES_DROPPED, + 'counter', + '{frame}', + ['channel'], + 'Frames dropped by channel', + true, + ), + def(METRIC.DB_WRITE_QUEUE_DEPTH, 'gauge', '{write}', [], 'Pending writes', true), + def(METRIC.DB_DROPPED_WRITES, 'counter', '{write}', ['table'], 'Writes dropped by table', true), + def(METRIC.DB_SIZE_BYTES, 'gauge', 'By', [], 'Database size', true), + def( + METRIC.BROWSER_RSS_BYTES, + 'gauge', + 'By', + ['session_id'], + 'Browser process-tree RSS by session', + true, + ), + def( + METRIC.ATTENTION_OPEN, + 'up-down counter', + '{request}', + ['kind'], + 'Open operator requests by kind', + ), + def(METRIC.ATTENTION_WAIT, 'histogram', 'ms', ['status'], 'Attention wait by status'), + def(METRIC.VAULT_FILLS, 'counter', '{fill}', ['result'], 'Vault fills by result'), + def(METRIC.BLOCKLIST_HITS, 'counter', '{hit}', ['source'], 'Blocked URLs by source'), + def(METRIC.RETENTION_PRUNED_ROWS, 'counter', '{row}', ['table'], 'Rows pruned by table'), + def( + METRIC.NOTIFICATION_DELIVERIES, + 'counter', + '{delivery}', + ['channel_kind', 'status'], + 'Notification deliveries by channel_kind and status', + ), + def( + METRIC.NOTIFICATION_REPORTS, + 'counter', + '{report}', + ['kind', 'outcome'], + 'Scheduled report decisions by kind and outcome', + ), + def( + METRIC.NOTIFICATION_ACTIONS, + 'counter', + '{press}', + ['channel_kind', 'outcome'], + 'Act-button presses by channel_kind and outcome', + ), + def(METRIC.PROCESS_RSS_BYTES, 'gauge', 'By', [], 'Process RSS', true), + def(METRIC.PROCESS_HEAP_BYTES, 'gauge', 'By', [], 'Process heap used', true), + def( + METRIC.PROCESS_EVENT_LOOP_LAG, + 'gauge', + 'ms', + [], + 'Event-loop delay, p99 since the previous export', + true, + ), +]; + /** Every instrument in the catalogue. Each is created on first access. */ export interface Instruments { readonly toolCalls: Counter; @@ -39,11 +176,11 @@ export interface Instruments { readonly sessionsActive: UpDownCounter; readonly sessionLaunchDuration: Histogram; readonly sessionLifetime: Histogram; - readonly wsConnections: UpDownCounter; + readonly wsConnections: ObservableUpDownCounter; readonly wsBufferedBytes: ObservableGauge; - readonly wsFramesDropped: Counter; + readonly wsFramesDropped: ObservableCounter; readonly dbWriteQueueDepth: ObservableGauge; - readonly dbDroppedWrites: Counter; + readonly dbDroppedWrites: ObservableCounter; readonly dbSizeBytes: ObservableGauge; readonly browserRssBytes: ObservableGauge; readonly attentionOpen: UpDownCounter; @@ -62,178 +199,105 @@ export interface Instruments { readonly processEventLoopLag: ObservableGauge; } +interface InstrumentOptions { + readonly unit: string; + readonly description: string; +} + /** - * Builds the lazily-created {@link Instruments} over `meter`. With telemetry off, `meter` is the - * API's no-op meter and every instrument is a no-op; the same object backs `/api/v1/system` - * figures through the counters' owners. + * Builds the lazily-created {@link Instruments} over `meter`, each with the name, unit and + * description of its {@link METRIC_DEFINITIONS} row. With telemetry off, `meter` is the API's + * no-op meter and every instrument is a no-op. */ export function createInstruments(meter: Meter): Instruments { const cache = new Map(); - const lazy = (key: string, make: () => T): T => { - const hit = cache.get(key); + const byName = new Map(METRIC_DEFINITIONS.map((d) => [d.name, d])); + const lazy = (name: MetricName, make: (options: InstrumentOptions) => T): T => { + const hit = cache.get(name); if (hit !== undefined) return hit as T; - const made = make(); - cache.set(key, made); + const d = byName.get(name); + if (d === undefined) throw new Error(`metric ${name} is not in the catalogue`); + const made = make({ unit: d.unit, description: d.description }); + cache.set(name, made); return made; }; - const ms = { unit: 'ms' } as const; - const bytes = { unit: 'By' } as const; + const counter = (name: MetricName) => lazy(name, (o) => meter.createCounter(name, o)); + const upDown = (name: MetricName) => lazy(name, (o) => meter.createUpDownCounter(name, o)); + const histogram = (name: MetricName) => lazy(name, (o) => meter.createHistogram(name, o)); + const gauge = (name: MetricName) => lazy(name, (o) => meter.createObservableGauge(name, o)); + const observableCounter = (name: MetricName) => + lazy(name, (o) => meter.createObservableCounter(name, o)); + const observableUpDown = (name: MetricName) => + lazy(name, (o) => meter.createObservableUpDownCounter(name, o)); return { get toolCalls() { - return lazy(METRIC.TOOL_CALLS, () => - meter.createCounter(METRIC.TOOL_CALLS, { - description: 'Tool calls by tool, ok, error_code', - }), - ); + return counter(METRIC.TOOL_CALLS); }, get toolCallDuration() { - return lazy(METRIC.TOOL_CALL_DURATION, () => - meter.createHistogram(METRIC.TOOL_CALL_DURATION, { - ...ms, - description: 'Tool call duration', - }), - ); + return histogram(METRIC.TOOL_CALL_DURATION); }, get sessionsActive() { - return lazy(METRIC.SESSIONS_ACTIVE, () => - meter.createUpDownCounter(METRIC.SESSIONS_ACTIVE, { - description: 'Live sessions by state', - }), - ); + return upDown(METRIC.SESSIONS_ACTIVE); }, get sessionLaunchDuration() { - return lazy(METRIC.SESSION_LAUNCH_DURATION, () => - meter.createHistogram(METRIC.SESSION_LAUNCH_DURATION, { - ...ms, - description: 'Session launch duration', - }), - ); + return histogram(METRIC.SESSION_LAUNCH_DURATION); }, get sessionLifetime() { - return lazy(METRIC.SESSION_LIFETIME, () => - meter.createHistogram(METRIC.SESSION_LIFETIME, { ...ms, description: 'Session lifetime' }), - ); + return histogram(METRIC.SESSION_LIFETIME); }, get wsConnections() { - return lazy(METRIC.WS_CONNECTIONS, () => - meter.createUpDownCounter(METRIC.WS_CONNECTIONS, { - description: 'Open WebSocket connections', - }), - ); + return observableUpDown(METRIC.WS_CONNECTIONS); }, get wsBufferedBytes() { - return lazy(METRIC.WS_BUFFERED_BYTES, () => - meter.createObservableGauge(METRIC.WS_BUFFERED_BYTES, { - ...bytes, - description: 'Per-connection buffered bytes', - }), - ); + return gauge(METRIC.WS_BUFFERED_BYTES); }, get wsFramesDropped() { - return lazy(METRIC.WS_FRAMES_DROPPED, () => - meter.createCounter(METRIC.WS_FRAMES_DROPPED, { description: 'Frames dropped by channel' }), - ); + return observableCounter(METRIC.WS_FRAMES_DROPPED); }, get dbWriteQueueDepth() { - return lazy(METRIC.DB_WRITE_QUEUE_DEPTH, () => - meter.createObservableGauge(METRIC.DB_WRITE_QUEUE_DEPTH, { description: 'Pending writes' }), - ); + return gauge(METRIC.DB_WRITE_QUEUE_DEPTH); }, get dbDroppedWrites() { - return lazy(METRIC.DB_DROPPED_WRITES, () => - meter.createCounter(METRIC.DB_DROPPED_WRITES, { description: 'Writes dropped by table' }), - ); + return observableCounter(METRIC.DB_DROPPED_WRITES); }, get dbSizeBytes() { - return lazy(METRIC.DB_SIZE_BYTES, () => - meter.createObservableGauge(METRIC.DB_SIZE_BYTES, { - ...bytes, - description: 'Database size', - }), - ); + return gauge(METRIC.DB_SIZE_BYTES); }, get browserRssBytes() { - return lazy(METRIC.BROWSER_RSS_BYTES, () => - meter.createObservableGauge(METRIC.BROWSER_RSS_BYTES, { - ...bytes, - description: 'Browser RSS by session', - }), - ); + return gauge(METRIC.BROWSER_RSS_BYTES); }, get attentionOpen() { - return lazy(METRIC.ATTENTION_OPEN, () => - meter.createUpDownCounter(METRIC.ATTENTION_OPEN, { - description: 'Open operator requests by kind', - }), - ); + return upDown(METRIC.ATTENTION_OPEN); }, get attentionWait() { - return lazy(METRIC.ATTENTION_WAIT, () => - meter.createHistogram(METRIC.ATTENTION_WAIT, { - ...ms, - description: 'Attention wait by status', - }), - ); + return histogram(METRIC.ATTENTION_WAIT); }, get vaultFills() { - return lazy(METRIC.VAULT_FILLS, () => - meter.createCounter(METRIC.VAULT_FILLS, { description: 'Vault fills by result' }), - ); + return counter(METRIC.VAULT_FILLS); }, get blocklistHits() { - return lazy(METRIC.BLOCKLIST_HITS, () => - meter.createCounter(METRIC.BLOCKLIST_HITS, { description: 'Blocked URLs by source' }), - ); + return counter(METRIC.BLOCKLIST_HITS); }, get retentionPrunedRows() { - return lazy(METRIC.RETENTION_PRUNED_ROWS, () => - meter.createCounter(METRIC.RETENTION_PRUNED_ROWS, { description: 'Rows pruned by table' }), - ); + return counter(METRIC.RETENTION_PRUNED_ROWS); }, get notificationDeliveries() { - return lazy(METRIC.NOTIFICATION_DELIVERIES, () => - meter.createCounter(METRIC.NOTIFICATION_DELIVERIES, { - description: 'Notification deliveries by channel_kind and status', - }), - ); + return counter(METRIC.NOTIFICATION_DELIVERIES); }, get notificationActions() { - return lazy(METRIC.NOTIFICATION_ACTIONS, () => - meter.createCounter(METRIC.NOTIFICATION_ACTIONS, { - description: 'Act-button presses by channel_kind and outcome', - }), - ); + return counter(METRIC.NOTIFICATION_ACTIONS); }, get notificationReports() { - return lazy(METRIC.NOTIFICATION_REPORTS, () => - meter.createCounter(METRIC.NOTIFICATION_REPORTS, { - description: 'Scheduled report decisions by kind and outcome', - }), - ); + return counter(METRIC.NOTIFICATION_REPORTS); }, get processRssBytes() { - return lazy(METRIC.PROCESS_RSS_BYTES, () => - meter.createObservableGauge(METRIC.PROCESS_RSS_BYTES, { - ...bytes, - description: 'Process RSS', - }), - ); + return gauge(METRIC.PROCESS_RSS_BYTES); }, get processHeapBytes() { - return lazy(METRIC.PROCESS_HEAP_BYTES, () => - meter.createObservableGauge(METRIC.PROCESS_HEAP_BYTES, { - ...bytes, - description: 'Process heap used', - }), - ); + return gauge(METRIC.PROCESS_HEAP_BYTES); }, get processEventLoopLag() { - return lazy(METRIC.PROCESS_EVENT_LOOP_LAG, () => - meter.createObservableGauge(METRIC.PROCESS_EVENT_LOOP_LAG, { - ...ms, - description: 'Event-loop lag', - }), - ); + return gauge(METRIC.PROCESS_EVENT_LOOP_LAG); }, }; } diff --git a/packages/core/src/interface/ws/hub.ts b/packages/core/src/interface/ws/hub.ts index 2c8d9f2..cf8a8bd 100644 --- a/packages/core/src/interface/ws/hub.ts +++ b/packages/core/src/interface/ws/hub.ts @@ -57,10 +57,19 @@ export interface RealtimeHubDeps { readonly degradations?: DegradationReporter; } +/** Where a frame was dropped (spec 10 §7 `browserhive.ws.frames_dropped{channel}`). */ +export type DroppedFrameChannel = 'screencast' | 'logs' | 'feed'; + /** The realtime hub. One instance per process. */ export class RealtimeHub implements HubCommandHost { readonly liveView: LiveViewPort; private readonly conns = new Set(); + /** Frames dropped since start, by channel (plain counters; read by the metrics at export). */ + private readonly dropped: Record = { + screencast: 0, + logs: 0, + feed: 0, + }; private readonly feed: FeedBuffer; private readonly limits: HubLimits; private readonly log: Logger; @@ -160,7 +169,10 @@ export class RealtimeHub implements HubCommandHost { const text = eventFrame(this.feed.head, this.now(), 'logs', { type: 'log.record', record }); for (const conn of this.conns) { if (!conn.topics.has('logs') || !matchesLogFilter(conn, entry)) continue; - if (conn.socket.bufferedAmount() > this.limits.screencastDropBytes) continue; + if (conn.socket.bufferedAmount() > this.limits.screencastDropBytes) { + this.dropped.logs += 1; + continue; + } conn.send(text); } } @@ -210,6 +222,11 @@ export class RealtimeHub implements HubCommandHost { })); } + /** Frames dropped since start, by channel: running totals across every connection. */ + droppedFrames(): Readonly> { + return { ...this.dropped }; + } + /** Sessions currently screencasting. */ activeScreencasts(): number { return this.liveView.activeCount; @@ -298,7 +315,10 @@ export class RealtimeHub implements HubCommandHost { bytes.set(header, 0); bytes.set(jpeg, header.byteLength); if (conn.congested || conn.socket.bufferedAmount() > this.limits.screencastDropBytes) { - if (sub.pending !== undefined) conn.droppedFrames += 1; + if (sub.pending !== undefined) { + conn.droppedFrames += 1; + this.dropped.screencast += 1; + } sub.pending = bytes; return; } @@ -309,7 +329,10 @@ export class RealtimeHub implements HubCommandHost { private writeFrame(conn: WsConnection, bytes: Uint8Array): void { const status = conn.send(bytes); - if (status === 0) conn.droppedFrames += 1; + if (status === 0) { + conn.droppedFrames += 1; + this.dropped.screencast += 1; + } if (status <= 0) conn.congested = true; } @@ -318,6 +341,7 @@ export class RealtimeHub implements HubCommandHost { if (status === 0) { // Bun dropped the frame: the feed must never lose events, so the client reconnects and // replays from its cursor. + this.dropped.feed += 1; this.overloaded(conn); return; } diff --git a/packages/core/src/interface/ws/index.ts b/packages/core/src/interface/ws/index.ts index e039c8f..d0ddd88 100644 --- a/packages/core/src/interface/ws/index.ts +++ b/packages/core/src/interface/ws/index.ts @@ -18,7 +18,7 @@ export { } from './cdp-bridge.ts'; export { WsConnection } from './connection.ts'; export { topicsForEvent, wireFeed } from './feed.ts'; -export { RealtimeHub, type RealtimeHubDeps } from './hub.ts'; +export { type DroppedFrameChannel, RealtimeHub, type RealtimeHubDeps } from './hub.ts'; export { DEFAULT_HUB_LIMITS, type HubLimits } from './hub-support.ts'; export { INPUT_AUDIT_WINDOW_MS, diff --git a/packages/core/src/ports/browser-driver.ts b/packages/core/src/ports/browser-driver.ts index edeeb5d..6ce4766 100644 --- a/packages/core/src/ports/browser-driver.ts +++ b/packages/core/src/ports/browser-driver.ts @@ -208,6 +208,11 @@ export interface SessionHandle { onCrash(listener: (reason: string) => void): () => void; /** Closes everything within `deadlineMs`; never throws (findings become warnings). */ close(deadlineMs: number, signal?: AbortSignal): Promise; + /** + * The OS pid of the browser's main process (the root of its process tree), read once and cached; + * `null` when it cannot be read. Only the browser-memory metric asks (spec 10 §7). Absent in fakes. + */ + browserPid?(): Promise; } /** Strategy for spawning a new isolated browser session (spec 01 §9 engine seam). */ diff --git a/packages/core/src/ports/persistence/write-queue.ts b/packages/core/src/ports/persistence/write-queue.ts index e1c818d..f0a7cb7 100644 --- a/packages/core/src/ports/persistence/write-queue.ts +++ b/packages/core/src/ports/persistence/write-queue.ts @@ -21,6 +21,11 @@ export interface WriteQueue { readonly depth: number; /** Writes refused or failed since start (`dropped_writes_total`). */ readonly droppedWrites: number; + /** + * {@link droppedWrites} split by table: the part of each dropped write's `operation` before the + * first `.` (`tool_calls.insert` → `tool_calls`). Only tables that lost a write appear. + */ + readonly droppedWritesByTable: ReadonlyMap; /** Drains and refuses further writes. Idempotent. */ close(): Promise; } diff --git a/packages/core/src/public/runtime.ts b/packages/core/src/public/runtime.ts index 024e570..2faf8cf 100644 --- a/packages/core/src/public/runtime.ts +++ b/packages/core/src/public/runtime.ts @@ -27,7 +27,9 @@ export { createCredentialsFile } from '../infra/auth/credentials-file.ts'; export { createWebCryptoRandom } from '../infra/auth/web-crypto-random.ts'; export { createSystemClock } from '../infra/clock/system-clock.ts'; export { createNodeFileSystem } from '../infra/fs/node-file-system.ts'; +export { type EventLoopLagMonitor, startEventLoopLagMonitor } from '../infra/host/event-loop.ts'; export { readHostMemory, readProcessMemory } from '../infra/host/memory.ts'; +export { createProcessTreeReader, type ProcessTreeReader } from '../infra/host/process-tree.ts'; export { createNanoidIdGenerator } from '../infra/ids/nanoid-id-generator.ts'; export { resolveColor } from '../infra/logging/color.ts'; export { redirectConsoleToLogger } from '../infra/logging/console-redirect.ts'; @@ -41,7 +43,16 @@ export { createLogger, type RootLogger } from '../infra/logging/logger.ts'; export { createRingBuffer, type LogRingBuffer } from '../infra/logging/ring-buffer.ts'; export type { LogSink } from '../infra/logging/sinks.ts'; export { createBunProcessRunner } from '../infra/process/bun-process-runner.ts'; -export type { Instruments } from '../infra/telemetry/metrics.ts'; +export { + BROWSER_RSS_SAMPLE_INTERVAL_MS, + type Instruments, + METRIC, + METRIC_DEFINITIONS, + type MetricDefinition, + type MetricKind, + type MetricName, + WS_BUFFERED_BYTES_MAX_SERIES, +} from '../infra/telemetry/metrics.ts'; export { createOtelLogSink } from '../infra/telemetry/otel-log-sink.ts'; export { createTelemetry, type Telemetry } from '../infra/telemetry/telemetry.ts'; export { AppError, isAppError } from '../kernel/errors/app-error.ts'; diff --git a/packages/core/test/helpers/in-memory-repos.ts b/packages/core/test/helpers/in-memory-repos.ts index 15251a9..76d4b23 100644 --- a/packages/core/test/helpers/in-memory-repos.ts +++ b/packages/core/test/helpers/in-memory-repos.ts @@ -368,6 +368,7 @@ export class InMemoryWriteQueue implements WriteQueue { private readonly jobs: Array<{ operation: string; job: WriteJob }> = []; private closed = false; private dropped = 0; + private readonly droppedByTable = new Map(); private draining: Promise | undefined; constructor( @@ -377,7 +378,7 @@ export class InMemoryWriteQueue implements WriteQueue { enqueue(operation: string, job: WriteJob): boolean { if (this.closed) { - this.dropped++; + this.drop(operation); return false; } this.operations.push(operation); @@ -402,11 +403,21 @@ export class InMemoryWriteQueue implements WriteQueue { return this.dropped; } + get droppedWritesByTable(): ReadonlyMap { + return this.droppedByTable; + } + async close(): Promise { await this.drain(); this.closed = true; } + private drop(operation: string): void { + this.dropped++; + const table = operation.split('.', 1)[0] || 'unknown'; + this.droppedByTable.set(table, (this.droppedByTable.get(table) ?? 0) + 1); + } + private async runAll(): Promise { while (this.jobs.length > 0) { const next = this.jobs.shift(); @@ -414,7 +425,7 @@ export class InMemoryWriteQueue implements WriteQueue { try { await next.job(this.repos); } catch (error) { - this.dropped++; + this.drop(next.operation); this.failures.push({ operation: next.operation, error }); } } From 55705aaacb7823400e2b2cfd4d2b0332aa68b81e Mon Sep 17 00:00:00 2001 From: Amir Ghorbani Date: Tue, 29 Sep 2026 12:32:19 -0400 Subject: [PATCH 3/6] test(observability): guard that every documented metric is recorded MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A table-driven guard parses the metric tables of spec 10 §7 and the telemetry guide and requires them to match the instrument catalogue row for row, then drives every source the composition root wires (bus consumers, the SQLite write queue, the realtime hub, the browser-memory sampler over a real process tree, the event-loop monitor, and the notification outbox, act buttons and report scheduler built by buildOps) and requires each instrument to reach a real OTLP/HTTP receiver with its documented type, unit and attribute keys. Unit tests cover the process-tree reader, the event-loop monitor, the sampler, the hub's and the write queue's totals, the catalogue units and the report outcomes; an integration test reads each real browser's pid and process-tree memory. --- .../adapters/browser-memory.test.ts | 118 +++ .../test/composition/metrics-guard.test.ts | 706 ++++++++++++++++++ .../maintenance/retention-scheduler.test.ts | 8 +- .../report-scheduler.in-app.test.ts | 12 + .../notifications/report-scheduler.test.ts | 4 + .../core/src/infra/host/event-loop.test.ts | 29 + .../core/src/infra/host/process-tree.test.ts | 74 ++ .../src/infra/persistence/write-queue.test.ts | 8 +- .../core/src/infra/telemetry/metrics.test.ts | 75 +- packages/core/src/interface/ws/hub.test.ts | 9 + .../core/test/integration/isolation.test.ts | 21 +- 11 files changed, 1039 insertions(+), 25 deletions(-) create mode 100644 packages/browserhive/src/composition/adapters/browser-memory.test.ts create mode 100644 packages/browserhive/test/composition/metrics-guard.test.ts create mode 100644 packages/core/src/infra/host/event-loop.test.ts create mode 100644 packages/core/src/infra/host/process-tree.test.ts diff --git a/packages/browserhive/src/composition/adapters/browser-memory.test.ts b/packages/browserhive/src/composition/adapters/browser-memory.test.ts new file mode 100644 index 0000000..ed00d45 --- /dev/null +++ b/packages/browserhive/src/composition/adapters/browser-memory.test.ts @@ -0,0 +1,118 @@ +/** @module composition/adapters/browser-memory.test — the 10 s browser-memory sampler behind `browserhive.browser.rss_bytes` (spec 10 §7): one tree per live session keyed by session id, sessions without a pid skipped, overlapping samples coalesced, no timer without a reader, stop. */ + +import { describe, expect, it } from 'bun:test'; +import type { ProcessTreeReader } from '@browserhive/core/runtime'; +import { startBrowserMemorySampler } from './browser-memory.ts'; + +function manualTimer() { + const timers: { fn: () => void; ms: number; cancelled: boolean }[] = []; + return { + timers, + repeat: (fn: () => void, ms: number) => { + const t = { fn, ms, cancelled: false }; + timers.push(t); + return () => { + t.cancelled = true; + }; + }, + }; +} + +describe('startBrowserMemorySampler', () => { + it('samples at once and every 10 s, one process tree per session', async () => { + const reads: number[][] = []; + const reader: ProcessTreeReader = { + rssOfTrees: async (roots) => { + reads.push([...roots]); + return new Map(roots.filter((pid) => pid !== 30).map((pid) => [pid, pid * 1000])); + }, + }; + let sessions = [ + { id: 's-a', browserPid: async () => 10 }, + { id: 's-b', browserPid: async () => 20 }, + { id: 's-launching' }, + { id: 's-unreadable', browserPid: async () => null }, + { id: 's-gone', browserPid: async () => 30 }, + ]; + const timer = manualTimer(); + const sampler = startBrowserMemorySampler({ + sessions: () => sessions, + reader, + repeat: timer.repeat, + }); + await sampler.sample(); + expect(timer.timers.map((t) => t.ms)).toEqual([10_000]); + expect(reads[0]).toEqual([10, 20, 30]); + expect([...sampler.latest()]).toEqual([ + ['s-a', 10_000], + ['s-b', 20_000], + ]); + // A closed session disappears at the next sample. + sessions = sessions.slice(0, 1); + timer.timers[0]?.fn(); + await sampler.sample(); + expect([...sampler.latest()]).toEqual([['s-a', 10_000]]); + sampler.stop(); + expect(timer.timers[0]?.cancelled).toBe(true); + expect(sampler.latest().size).toBe(0); + }); + + it('coalesces a sample requested while one is running', async () => { + let release: () => void = () => undefined; + let calls = 0; + const reader: ProcessTreeReader = { + rssOfTrees: () => { + calls++; + return new Promise((resolve) => { + release = () => resolve(new Map([[10, 1]])); + }); + }, + }; + const sampler = startBrowserMemorySampler({ + sessions: () => [{ id: 's-a', browserPid: async () => 10 }], + reader, + repeat: manualTimer().repeat, + }); + const second = sampler.sample(); + await new Promise((resolve) => setTimeout(resolve, 5)); + release(); + await second; + expect(calls).toBe(1); + sampler.stop(); + }); + + it('never starts a timer without a reader (Windows)', async () => { + const timer = manualTimer(); + const sampler = startBrowserMemorySampler({ + sessions: () => [{ id: 's-a', browserPid: async () => 10 }], + reader: null, + repeat: timer.repeat, + }); + await sampler.sample(); + expect(timer.timers).toEqual([]); + expect(sampler.latest().size).toBe(0); + sampler.stop(); + }); + + it('reports a failing read and keeps the previous sample', async () => { + const errors: unknown[] = []; + let fail = false; + const sampler = startBrowserMemorySampler({ + sessions: () => [{ id: 's-a', browserPid: async () => 10 }], + reader: { + rssOfTrees: async () => { + if (fail) throw new Error('boom'); + return new Map([[10, 5]]); + }, + }, + repeat: manualTimer().repeat, + onError: (err) => errors.push(err), + }); + await sampler.sample(); + fail = true; + await sampler.sample(); + expect(errors).toHaveLength(1); + expect([...sampler.latest()]).toEqual([['s-a', 5]]); + sampler.stop(); + }); +}); diff --git a/packages/browserhive/test/composition/metrics-guard.test.ts b/packages/browserhive/test/composition/metrics-guard.test.ts new file mode 100644 index 0000000..0242090 --- /dev/null +++ b/packages/browserhive/test/composition/metrics-guard.test.ts @@ -0,0 +1,706 @@ +/** + * @module test/composition/metrics-guard.test — documented == defined == recorded (spec 10 §7, spec + * 09 §3.2). The metric tables of spec 10 §7 and `docs/guide/telemetry.md` must match the instrument + * catalogue row for row (name, type, unit, attributes); then every source the composition root + * wires is driven for real (the bus consumers, the write queue on SQLite, the realtime hub, the + * browser-memory sampler over a real process tree, the event-loop monitor, and the notification + * outbox, act buttons and report scheduler built by `buildOps` with the counters `build-domain` + * passes) and every instrument must reach a real OTLP/HTTP receiver with the documented type, unit + * and attribute keys. A metric documented or defined but never recorded fails here. + */ + +import { afterAll, beforeAll, describe, expect, it } from 'bun:test'; +import { mkdtempSync, readFileSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join, resolve } from 'node:path'; +import { CHANNEL_RENDERERS, channelFactories } from '@browserhive/core/notifications'; +import { + openDatabase, + SqliteAnalyticsQueries, + SqliteMaintenanceService, + SqliteUnitOfWork, + SqliteWriteQueue, +} from '@browserhive/core/persistence'; +import { + createNanoidIdGenerator, + createProcessTreeReader, + createRedactor, + createSystemClock, + createTelemetry, + DegradationService, + type DomainEvents, + type Logger, + METRIC_DEFINITIONS, + type MetricDefinition, + type MetricName, + RetentionScheduler, + retentionPolicyFromConfig, + type Telemetry, +} from '@browserhive/core/runtime'; +import { createRealtimeHub, InProcessEventBus, type Realtime } from '@browserhive/core/server'; +import { startBrowserMemorySampler } from '../../src/composition/adapters/browser-memory.ts'; +import { + notificationCounters, + realtimeMetrics, + wireMetrics, +} from '../../src/composition/adapters/metrics.ts'; +import { buildOps, type OpsParts } from '../../src/composition/phases/domain-ops.ts'; +import { resolvedFor } from './support.ts'; + +const ROOT = resolve(import.meta.dir, '../../../..'); + +// ------------------------------------------------------------------------------------------------ +// The documented tables +// ------------------------------------------------------------------------------------------------ + +interface DocumentedMetric { + readonly name: string; + readonly kind: string; + readonly observable: boolean; + readonly unit: string; + readonly attributes: readonly string[]; +} + +const withoutParentheses = (cell: string) => cell.replace(/\([^)]*\)/g, ''); +const ticked = (cell: string) => [...cell.matchAll(/`([^`]+)`/g)].map((m) => m[1] ?? ''); + +/** The rows of the metric table between `start` and `end` in a Markdown file. */ +function metricTable(file: string, start: string, end: string): DocumentedMetric[] { + const text = readFileSync(join(ROOT, file), 'utf8'); + const from = text.indexOf(start); + const to = text.indexOf(end, from + start.length); + if (from < 0 || to < 0) throw new Error(`${file}: cannot find the metric table`); + return text + .slice(from, to) + .split('\n') + .filter((line) => line.startsWith('| `browserhive.')) + .map((line) => { + const [name = '', type = '', unit = '', attributes = ''] = line + .split('|') + .slice(1, -1) + .map((c) => c.trim()); + return { + name: ticked(name)[0] ?? '', + kind: withoutParentheses(type).trim(), + observable: type.includes('(observable)'), + unit: ticked(unit)[0] ?? '', + attributes: ticked(withoutParentheses(attributes)), + }; + }); +} + +const SPEC = metricTable( + 'specs/10-error-handling-and-telemetry.md', + '## 7. Metrics catalogue', + '## 8. OpenTelemetry wiring', +); +const DOCS = metricTable('docs/guide/telemetry.md', '### Metrics', '### Logs'); +const CATALOGUE = METRIC_DEFINITIONS.map((d: MetricDefinition) => ({ + name: d.name, + kind: d.kind, + observable: d.observable, + unit: d.unit, + attributes: d.attributes, +})); + +describe('documented == defined', () => { + it('spec 10 §7 lists exactly the catalogue, row for row', () => { + expect(SPEC).toEqual(CATALOGUE); + }); + + it('the telemetry guide lists exactly the catalogue (types without the observable note)', () => { + expect(DOCS).toEqual(CATALOGUE.map((d) => ({ ...d, observable: false }))); + }); +}); + +// ------------------------------------------------------------------------------------------------ +// A real OTLP/HTTP receiver +// ------------------------------------------------------------------------------------------------ + +interface OtlpAttribute { + readonly key: string; + readonly value: Record; +} +interface OtlpPoint { + readonly attributes?: readonly OtlpAttribute[]; + readonly asInt?: number | string; + readonly asDouble?: number; + readonly count?: number | string; + readonly sum?: number; +} +interface OtlpMetric { + readonly name: string; + readonly unit?: string; + readonly sum?: { readonly dataPoints: readonly OtlpPoint[]; readonly isMonotonic?: boolean }; + readonly gauge?: { readonly dataPoints: readonly OtlpPoint[] }; + readonly histogram?: { readonly dataPoints: readonly OtlpPoint[] }; +} +interface OtlpBody { + readonly resourceMetrics?: readonly { + readonly scopeMetrics?: readonly { readonly metrics?: readonly OtlpMetric[] }[]; + }[]; +} + +/** What a collector would store for one instrument. */ +interface Received { + readonly kind: string; + readonly unit: string; + readonly points: readonly { readonly value: number; readonly attrs: Record }[]; +} + +function receivedMetrics(bodies: readonly OtlpBody[]): Map { + const out = new Map(); + for (const body of bodies) { + for (const rm of body.resourceMetrics ?? []) { + for (const sm of rm.scopeMetrics ?? []) { + for (const m of sm.metrics ?? []) { + const kind = + m.histogram !== undefined + ? 'histogram' + : m.gauge !== undefined + ? 'gauge' + : m.sum?.isMonotonic === true + ? 'counter' + : 'up-down counter'; + const raw = m.histogram?.dataPoints ?? m.gauge?.dataPoints ?? m.sum?.dataPoints ?? []; + const points = raw.map((p) => ({ + value: Number(p.asDouble ?? p.asInt ?? p.sum ?? 0), + attrs: Object.fromEntries( + (p.attributes ?? []).map((a) => [a.key, Object.values(a.value)[0]]), + ), + })); + // Cumulative temporality: the latest export holds every series. + out.set(m.name, { kind, unit: m.unit ?? '', points }); + } + } + } + } + return out; +} + +// ------------------------------------------------------------------------------------------------ +// The world: every source wired as the composition root wires it +// ------------------------------------------------------------------------------------------------ + +const quiet: Logger = { + child: () => quiet, + isLevelEnabled: () => false, + error: () => undefined, + warn: () => undefined, + info: () => undefined, + debug: () => undefined, + trace: () => undefined, +} as unknown as Logger; + +interface World { + readonly bus: InProcessEventBus; + readonly queue: SqliteWriteQueue; + readonly realtime: Realtime; + readonly ops: OpsParts; + readonly retention: RetentionScheduler; + readonly sampler: ReturnType; + readonly channelId: string; + readonly hooks: string[]; +} + +const SESSION = 'guard-00000001'; + +function requestRow(id: string, kind: 'attention' | 'vault_confirm', extra = {}) { + return { + request_id: id, + kind, + session_id: SESSION, + session_slug: 'guard', + owner: 'admin', + reason: 'solve the captcha', + mode: kind === 'attention' ? ('takeover' as const) : null, + options: null, + status: 'pending' as const, + message: null, + resolved_by: null, + resolution_reason: null, + created_at: Date.now() - 2_000, + resolved_at: null, + deadline_at: null, + waited_ms: null, + page_url: 'https://example.test/login', + tool: 'request_attention', + event_id: null, + entry_name: kind === 'vault_confirm' ? 'shop' : null, + ...extra, + }; +} + +function sessionSummary(extra = {}) { + return { + session_id: SESSION, + slug: 'guard', + created_at: Date.now() - 60_000, + harness: 'nightly-scraper', + channel: 'chromium', + stealth: true, + ...extra, + }; +} + +/** One source of the composition root, what it drives, and the metrics it must produce. */ +interface Source { + readonly name: string; + readonly drives: readonly MetricName[]; + run(w: World): Promise; +} + +const publish = (w: World, name: N, payload: unknown) => + w.bus.publish(name, payload as DomainEvents[N]); + +/** Table-driven: every catalogue metric must appear in exactly one `drives` list. */ +const SOURCES: readonly Source[] = [ + { + name: 'MCP dispatcher (tool.called)', + drives: ['browserhive.tool_calls', 'browserhive.tool_call.duration'], + async run(w) { + const observation = (ok: boolean) => ({ + tool: ok ? 'navigate' : 'click', + ok, + errorCode: ok ? null : 'ELEMENT_NOT_FOUND', + durationMs: ok ? 120 : 40, + harness: 'claude-code', + }); + publish(w, 'tool.called', { type: 'tool.called', observation: observation(true) }); + publish(w, 'tool.called', { type: 'tool.called', observation: observation(false) }); + }, + }, + { + name: 'session service (session.opened / updated / closed)', + drives: [ + 'browserhive.sessions.active', + 'browserhive.session.launch.duration', + 'browserhive.session.lifetime', + ], + async run(w) { + publish(w, 'session.opened', { type: 'session.opened', session: sessionSummary() }); + publish(w, 'session.updated', { + type: 'session.updated', + session: sessionSummary(), + patch: { launchMs: 900 }, + }); + publish(w, 'session.closed', { + type: 'session.closed', + session_id: SESSION, + closed_at: Date.now(), + reason: 'user', + }); + }, + }, + { + name: 'operator-request broker (attention.*, vault.confirm.*)', + drives: ['browserhive.attention.open', 'browserhive.attention.wait'], + async run(w) { + // The notification producers listen from here on, as `wire-observers` starts them. + w.ops.notifications.start(); + publish(w, 'attention.created', { + type: 'attention.created', + request: requestRow('a-000000000001', 'attention'), + }); + publish(w, 'vault.confirm.created', { + type: 'vault.confirm.created', + request: requestRow('v-000000000001', 'vault_confirm'), + }); + publish(w, 'attention.resolved', { + type: 'attention.resolved', + request: requestRow('a-000000000001', 'attention', { + status: 'resolved', + resolved_at: Date.now(), + waited_ms: 2_000, + resolved_by: 'admin', + }), + }); + }, + }, + { + name: 'vault broker and blocklist (vault.access, blocklist.hit)', + drives: ['browserhive.vault.fills', 'browserhive.blocklist.hits'], + async run(w) { + publish(w, 'vault.access', { type: 'vault.access', row: { result: 'filled' } }); + publish(w, 'blocklist.hit', { type: 'blocklist.hit', row: { source: 'config' } }); + }, + }, + { + name: 'retention scheduler (retention.completed)', + drives: ['browserhive.retention.pruned_rows'], + async run(w) { + await w.retention.tick(); + }, + }, + { + name: 'write queue (SQLite)', + drives: [ + 'browserhive.db.write_queue.depth', + 'browserhive.db.dropped_writes', + 'browserhive.db.size_bytes', + ], + async run(w) { + w.queue.enqueue('tool_calls.insert', async () => { + throw new Error('constraint'); + }); + await w.queue.drain(); + }, + }, + { + name: 'realtime hub', + drives: [ + 'browserhive.ws.connections', + 'browserhive.ws.buffered_bytes', + 'browserhive.ws.frames_dropped', + ], + async run(w) { + const socket = { + buffered: 0, + send: () => 1, + bufferedAmount: () => socket.buffered, + close: () => undefined, + }; + const principal = { + subject: 'admin', + kind: 'operator', + display: 'admin', + auth: { method: 'password-session', sessionId: 'as-1' }, + scopes: [], + tenantId: null, + mustChangePassword: false, + } as unknown as Parameters[1]; + const conn = w.realtime.hub.open(socket, principal); + await w.realtime.hub.message(conn, JSON.stringify({ type: 'logs.tail' })); + socket.buffered = 3 * 1024 * 1024; + w.realtime.hub.publishLog({ + seq: 1, + record: { ts: Date.now(), level: 'info', msg: 'x', module: 'http' }, + }); + }, + }, + { + name: 'browser-memory sampler (process tree)', + drives: ['browserhive.browser.rss_bytes'], + async run(w) { + await w.sampler.sample(); + }, + }, + { + name: 'this process', + drives: [ + 'browserhive.process.rss_bytes', + 'browserhive.process.heap_bytes', + 'browserhive.process.event_loop_lag', + ], + async run() { + await new Promise((r) => setTimeout(r, 60)); + }, + }, + { + name: 'notification outbox (the attention request above, to a webhook channel)', + drives: ['browserhive.notifications.deliveries'], + async run(w) { + await w.ops.notifications.idle(); + await w.ops.notificationOutbox.tick(); + }, + }, + { + name: 'act buttons (a press of a token never minted)', + drives: ['browserhive.notifications.actions'], + async run(w) { + await w.ops.actions.press({ + token: 'a'.repeat(22), + origin: null, + actor: { platform: 'ntfy', id: null, name: null }, + }); + }, + }, + { + name: 'report scheduler (an on-demand digest)', + drives: ['browserhive.notifications.reports'], + async run(w) { + const record = await w.ops.channels.channels().find((c) => c.record.channelId === w.channelId) + ?.record; + if (record === undefined) throw new Error('no channel'); + const built = await w.ops.reports.manualDigest(record); + await w.ops.reports.storeManualCopy(built); + }, + }, +]; + +describe('documented == recorded', () => { + const dir = mkdtempSync(join(tmpdir(), 'bh-metrics-guard-')); + const bodies: OtlpBody[] = []; + const hooks: string[] = []; + let server: ReturnType | undefined; + let telemetry: Telemetry | undefined; + let received = new Map(); + const stops: (() => unknown)[] = []; + + beforeAll(async () => { + server = Bun.serve({ + hostname: '127.0.0.1', + port: 0, + async fetch(request) { + const path = new URL(request.url).pathname; + if (path === '/v1/metrics') bodies.push((await request.json()) as OtlpBody); + if (path === '/hook') hooks.push(await request.text()); + return Response.json({}); + }, + }); + const base = `http://127.0.0.1:${server.port}`; + telemetry = await createTelemetry({ + enabled: true, + registerGlobals: false, + endpoint: base, + protocol: 'http/json', + signals: { traces: false, metrics: true, logs: false }, + metricsIntervalMs: 3_600_000, + }); + const clock = createSystemClock(); + const ids = createNanoidIdGenerator({ clock }); + const handle = await openDatabase({ + path: join(dir, 'browserhive.db'), + dataDir: dir, + appVersion: '0.0.0-test', + clock, + logger: quiet, + }); + stops.push(() => handle.close()); + const uow = new SqliteUnitOfWork(handle.db); + const queue = new SqliteWriteQueue({ uow, logger: quiet }); + const analytics = new SqliteAnalyticsQueries(handle.db, uow.repos, queue); + const maintenance = new SqliteMaintenanceService({ + handle, + clock, + logger: quiet, + appVersion: '0.0.0-test', + queue, + }); + const bus = new InProcessEventBus({ clock, logger: quiet }); + const degradations = new DegradationService({ + repo: uow.repos.systemEvents, + bus, + clock, + ids, + logger: quiet, + }); + const config = resolvedFor({ port: 0, dataDir: dir }).config; + const env = { BHTEST_WEBHOOK_URL: `${base}/hook` }; + const ops = buildOps({ + // Exactly what `build-domain` passes. + ...notificationCounters(telemetry.instruments), + config, + repos: uow.repos, + queue, + handle, + analytics, + maintenance, + bus, + clock, + ids, + logger: quiet, + redactor: createRedactor(), + degradations, + uow, + env, + registerSecret: () => undefined, + dashboardUrl: () => base, + channelFactories: channelFactories({ images: { read: async () => null } }), + renderers: CHANNEL_RENDERERS, + probe: async () => ({ ok: true }) as never, + instanceId: 'guard', + capacity: () => ({ live: 0, max: 4 }), + }); + const channelId = 'nc-000000guard'; + await uow.repos.notificationChannels.upsert({ + channelId, + name: 'guard-hook', + kind: 'webhook', + mode: null, + source: 'db', + status: 'active', + target: {}, + secretRefs: { url: 'BHTEST_WEBHOOK_URL' }, + rules: {}, + failureCount: 0, + lastError: null, + lastOkAt: null, + lastFailureAt: null, + createdAt: clock.now(), + updatedAt: clock.now(), + }); + await ops.channels.load([]); + stops.push(() => ops.notifications.stop()); + // The session the operator requests and their notifications belong to. + const now = clock.now(); + await uow.repos.sessions.insert({ + sessionId: SESSION, + slug: 'guard', + owner: 'admin', + tenantId: null, + connectionId: null, + engine: 'chromium', + channel: 'chromium', + headless: true, + incognito: false, + persistenceMode: 'memory', + disableEvaluate: false, + vaultEnabled: false, + stealth: true, + fingerprint: false, + humanize: false, + identity: null, + proxyLabel: null, + state: 'live', + createdAt: now - 60_000, + launchedAt: now - 59_000, + lastActivityAt: now, + leaseExpiresAt: now + 600_000, + leasePausedAt: null, + closedAt: null, + closedReason: null, + archivedAt: null, + lastUrl: null, + launchMs: 900, + config: { channel: 'chromium', headless: true }, + harness: null, + sandboxed: null, + browserVersion: null, + }); + const retention = new RetentionScheduler({ + maintenance: { + retentionSweep: async () => ({ + startedAt: clock.now(), + durationMs: 1, + prunedRows: { tool_calls: 4, pages: 0 }, + artifactsEnqueued: 0, + bytesBefore: 0, + bytesAfter: 0, + failures: [], + }), + }, + policy: retentionPolicyFromConfig(config), + clock, + logger: quiet, + degradations, + bus, + }); + const realtime = createRealtimeHub({ + clock, + ids, + logger: quiet, + auth: { touchSession: async () => true }, + attention: { isInputPermitted: () => false }, + serverVersion: '0.0.0-test', + epoch: 'guard', + bus, + sessions: { peek: () => undefined }, + pageOf: () => { + throw new Error('no pages in the guard'); + }, + bridges: (() => { + throw new Error('no bridges in the guard'); + }) as never, + schedule: () => () => undefined, + every: () => () => undefined, + }); + // A "browser" whose process tree is this test process: the real reader, the real sampler. + const sampler = startBrowserMemorySampler({ + sessions: () => [{ id: SESSION, browserPid: async () => process.pid }], + reader: createProcessTreeReader(), + repeat: () => () => undefined, + }); + stops.push(() => sampler.stop()); + // Exactly what `wire-observers` wires. + stops.push( + wireMetrics(telemetry.instruments, bus, { + queue, + analytics, + realtime: () => realtimeMetrics(realtime.hub), + browserMemory: () => sampler.latest(), + }), + ); + const world: World = { bus, queue, realtime, ops, retention, sampler, channelId, hooks }; + for (const source of SOURCES) await source.run(world); + // The size gauge reports the value read at the previous collection. + await telemetry.forceFlush(); + await telemetry.forceFlush(); + received = receivedMetrics(bodies); + }); + + afterAll(async () => { + for (const stop of stops.reverse()) await stop(); + await telemetry?.shutdown(); + await server?.stop(true); + rmSync(dir, { recursive: true, force: true }); + }); + + it('every catalogue metric is driven by exactly one source of the table', () => { + const driven = SOURCES.flatMap((s) => s.drives); + expect(new Set(driven).size).toBe(driven.length); + expect([...driven].sort()).toEqual(METRIC_DEFINITIONS.map((d) => d.name).sort()); + }); + + it('the webhook channel was really called', () => { + expect(hooks).toHaveLength(2); + }); + + for (const d of METRIC_DEFINITIONS) { + it(`${d.name} reaches the collector as a ${d.kind} in ${d.unit} with its documented attributes`, () => { + const got = received.get(d.name); + expect(got, `${d.name} was never exported`).toBeDefined(); + if (got === undefined) return; + expect({ kind: got.kind, unit: got.unit }).toEqual({ kind: d.kind, unit: d.unit }); + expect(got.points.length).toBeGreaterThan(0); + const keys = new Set(got.points.flatMap((p) => Object.keys(p.attrs))); + expect([...keys].sort()).toEqual([...d.attributes].sort()); + }); + } + + it('carries the values of the driven sources', () => { + const points = (name: MetricName) => received.get(name)?.points ?? []; + expect(points('browserhive.tool_calls')).toContainEqual({ + value: 1, + attrs: { tool: 'click', ok: false, harness: 'claude-code', error_code: 'ELEMENT_NOT_FOUND' }, + }); + expect(points('browserhive.sessions.active')).toContainEqual({ + value: 0, + attrs: { harness: 'other' }, + }); + expect(points('browserhive.session.launch.duration')[0]?.attrs).toEqual({ + channel: 'chromium', + stealth: true, + }); + expect(points('browserhive.attention.open')).toContainEqual({ + value: 1, + attrs: { kind: 'vault_confirm' }, + }); + expect(points('browserhive.retention.pruned_rows')).toEqual([ + { value: 4, attrs: { table: 'tool_calls' } }, + ]); + expect(points('browserhive.db.dropped_writes')).toEqual([ + { value: 1, attrs: { table: 'tool_calls' } }, + ]); + expect(points('browserhive.ws.connections')).toContainEqual({ value: 1, attrs: {} }); + expect(points('browserhive.ws.frames_dropped')).toContainEqual({ + value: 1, + attrs: { channel: 'logs' }, + }); + expect(points('browserhive.browser.rss_bytes')[0]?.value).toBeGreaterThan(1024 * 1024); + expect(points('browserhive.db.size_bytes')[0]?.value).toBeGreaterThan(0); + // The attention request and the vault confirm, each one message to the webhook. + expect(points('browserhive.notifications.deliveries')).toContainEqual({ + value: 2, + attrs: { channel_kind: 'webhook', status: 'sent' }, + }); + expect(points('browserhive.notifications.actions')).toContainEqual({ + value: 1, + attrs: { channel_kind: 'ntfy', outcome: 'unknown' }, + }); + expect(points('browserhive.notifications.reports')).toContainEqual({ + value: 1, + attrs: { kind: 'digest.daily', outcome: 'manual' }, + }); + }); +}); diff --git a/packages/core/src/app/maintenance/retention-scheduler.test.ts b/packages/core/src/app/maintenance/retention-scheduler.test.ts index 5e88925..9f26cba 100644 --- a/packages/core/src/app/maintenance/retention-scheduler.test.ts +++ b/packages/core/src/app/maintenance/retention-scheduler.test.ts @@ -103,7 +103,13 @@ describe('RetentionScheduler', () => { expect(recovered).toEqual([]); const completed = bus.published.find((p) => p.name === 'retention.completed') ?.payload as DomainEvents['retention.completed']; - expect(completed).toMatchObject({ result: 'partial', severity: 'warn', pruned_rows: 4 }); + expect(completed).toMatchObject({ + result: 'partial', + severity: 'warn', + pruned_rows: 4, + // Internal, for `retention.pruned_rows{table}` (spec 10 §7); the WS feed schema strips it. + prunedByTable: { tool_calls: 3, pages: 1 }, + }); }); it('status matches the contracts RetentionStatus and coalesces overlapping ticks', async () => { diff --git a/packages/core/src/app/notifications/report-scheduler.in-app.test.ts b/packages/core/src/app/notifications/report-scheduler.in-app.test.ts index 635bcc9..db354a1 100644 --- a/packages/core/src/app/notifications/report-scheduler.in-app.test.ts +++ b/packages/core/src/app/notifications/report-scheduler.in-app.test.ts @@ -104,6 +104,7 @@ async function setup(opts: Setup = {}) { }; const intervals = new ManualIntervals(); const announced: { op: string; notification: Notification }[] = []; + const counted: { kind: string; outcome: string; n: number }[] = []; const make = () => new ReportScheduler({ registry, @@ -122,6 +123,7 @@ async function setup(opts: Setup = {}) { inbox: (op, notification) => void announced.push({ op, notification }), redactor: createRedactor(), scheduler: intervals, + counter: { add: (n, a) => void counted.push({ ...a, n }) }, }); const scheduler = make(); const rows = () => [...repos.notifications.rows.values()]; @@ -138,6 +140,7 @@ async function setup(opts: Setup = {}) { make, intervals, announced, + counted, inApp, copies, service, @@ -400,6 +403,11 @@ describe('anomaly watches (D-45)', () => { ['created', 'open'], ['updated', 'resolved'], ]); + // `browserhive.notifications.reports`: the new in-app alert, then its resolution (spec 10 §7). + expect(t.counted).toEqual([ + { kind: 'report.anomaly', outcome: 'in_app', n: 1 }, + { kind: 'report.anomaly', outcome: 'resolved', n: 1 }, + ]); }); it('closes the open alert of a watch no longer wanted', async () => { @@ -436,6 +444,10 @@ describe('on-demand digests (D-45)', () => { const built = await t.scheduler.manualDigest(record); const id = await t.scheduler.storeManualCopy(built); expect(t.inApp().map((r) => r.notificationId)).toEqual([id ?? 'missing']); + expect(t.counted).toEqual([ + { kind: 'digest.daily', outcome: 'manual', n: 1 }, + { kind: 'digest.daily', outcome: 'in_app', n: 1 }, + ]); expect(decodeMessage(t.inApp()[0]?.messageJson ?? null)?.report?.manual).toBe(true); // Asking twice for the same instant finds the same copy. expect(await t.scheduler.storeManualCopy(built)).toBe(id); diff --git a/packages/core/src/app/notifications/report-scheduler.test.ts b/packages/core/src/app/notifications/report-scheduler.test.ts index 4140802..972bed3 100644 --- a/packages/core/src/app/notifications/report-scheduler.test.ts +++ b/packages/core/src/app/notifications/report-scheduler.test.ts @@ -323,6 +323,10 @@ describe('ReportScheduler: digests (D-43)', () => { expect(built.window).toEqual({ since: START - DAY, until: START }); expect(built.message.report?.manual).toBe(true); expect(await t.repos.notificationCursors.get(digestCursorKey(CHANNEL))).toBe(before); + // A preview is not a decision; sending (which stores the copy) counts `manual` (spec 10 §7). + expect(t.counted.filter((c) => c.outcome === 'manual')).toEqual([]); + await t.scheduler.storeManualCopy(built); + expect(t.counted).toContainEqual({ kind: 'digest.daily', outcome: 'manual', n: 1 }); }); it('shows the next run in the channel zone', async () => { diff --git a/packages/core/src/infra/host/event-loop.test.ts b/packages/core/src/infra/host/event-loop.test.ts new file mode 100644 index 0000000..64d1381 --- /dev/null +++ b/packages/core/src/infra/host/event-loop.test.ts @@ -0,0 +1,29 @@ +/** @module infra/host/event-loop.test — the event-loop delay monitor behind `browserhive.process.event_loop_lag` (spec 10 §7): a blocked loop shows up in the window's p99, and each read starts a new window. */ + +import { describe, expect, it } from 'bun:test'; +import { startEventLoopLagMonitor } from './event-loop.ts'; + +const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); + +describe('startEventLoopLagMonitor', () => { + it('reports the p99 delay of the window in ms, then resets', async () => { + const monitor = startEventLoopLagMonitor(5); + expect(monitor).not.toBeNull(); + await sleep(30); + // Block the loop for ~60 ms. + const until = performance.now() + 60; + while (performance.now() < until) { + // busy wait + } + await sleep(30); + const p99 = monitor?.takeP99Ms() ?? null; + expect(p99).toBeGreaterThan(20); + expect(p99).toBeLessThan(10_000); + await sleep(30); + const next = monitor?.takeP99Ms() ?? null; + // The next window did not block, so its p99 is far below the blocked one. + expect(next === null || next < (p99 ?? 0)).toBe(true); + monitor?.stop(); + monitor?.stop(); + }); +}); diff --git a/packages/core/src/infra/host/process-tree.test.ts b/packages/core/src/infra/host/process-tree.test.ts new file mode 100644 index 0000000..27c14db --- /dev/null +++ b/packages/core/src/infra/host/process-tree.test.ts @@ -0,0 +1,74 @@ +/** @module infra/host/process-tree.test — process-tree RSS for `browserhive.browser.rss_bytes` (spec 10 §7): tree sums, `/proc` and `ps` parsing, and the real reader over this process. */ + +import { describe, expect, it } from 'bun:test'; +import { + createProcessTreeReader, + parseProcStatus, + parsePsTable, + sumProcessTrees, +} from './process-tree.ts'; + +describe('sumProcessTrees', () => { + const table = [ + { pid: 1, ppid: 0, rssBytes: 1_000 }, + { pid: 10, ppid: 1, rssBytes: 100 }, // browser A + { pid: 11, ppid: 10, rssBytes: 20 }, // zygote + { pid: 12, ppid: 11, rssBytes: 3 }, // renderer under the zygote + { pid: 13, ppid: 10, rssBytes: 4 }, // GPU + { pid: 20, ppid: 1, rssBytes: 200 }, // browser B + { pid: 21, ppid: 20, rssBytes: 5 }, + ]; + + it('sums each root with every descendant, never a sibling tree or the parent', () => { + expect([...sumProcessTrees(table, [10, 20])]).toEqual([ + [10, 127], + [20, 205], + ]); + }); + + it('leaves out roots that are gone and survives a ppid cycle', () => { + expect([...sumProcessTrees(table, [99])]).toEqual([]); + const cycle = [ + { pid: 5, ppid: 6, rssBytes: 1 }, + { pid: 6, ppid: 5, rssBytes: 2 }, + ]; + expect(sumProcessTrees(cycle, [5]).get(5)).toBe(3); + }); +}); + +describe('parsers', () => { + it('reads PPid and VmRSS from /proc//status (kernel threads have no VmRSS)', () => { + const status = 'Name:\tchrome\nState:\tS (sleeping)\nPPid:\t4242\nVmRSS:\t 90524 kB\n'; + expect(parseProcStatus(7, status)).toEqual({ pid: 7, ppid: 4242, rssBytes: 90_524 * 1024 }); + expect(parseProcStatus(2, 'Name:\tkthreadd\nPPid:\t0\n')).toEqual({ + pid: 2, + ppid: 0, + rssBytes: 0, + }); + expect(parseProcStatus(3, 'garbage')).toBeNull(); + }); + + it('reads `ps -A -o pid=,ppid=,rss=` rows in KiB', () => { + expect(parsePsTable(' 1 0 1200\n 501 1 64\n\n bad row\n')).toEqual([ + { pid: 1, ppid: 0, rssBytes: 1200 * 1024 }, + { pid: 501, ppid: 1, rssBytes: 64 * 1024 }, + ]); + }); +}); + +describe('createProcessTreeReader', () => { + it('has no reader on Windows', () => { + expect(createProcessTreeReader('win32')).toBeNull(); + }); + + it.if(process.platform === 'linux' || process.platform === 'darwin')( + 'reads the tree of this process', + async () => { + const reader = createProcessTreeReader(); + const sums = await reader?.rssOfTrees([process.pid, 2 ** 30]); + expect(sums?.get(process.pid)).toBeGreaterThan(1024 * 1024); + expect(sums?.has(2 ** 30)).toBe(false); + expect((await reader?.rssOfTrees([]))?.size).toBe(0); + }, + ); +}); diff --git a/packages/core/src/infra/persistence/write-queue.test.ts b/packages/core/src/infra/persistence/write-queue.test.ts index 55add9a..cbfe24c 100644 --- a/packages/core/src/infra/persistence/write-queue.test.ts +++ b/packages/core/src/infra/persistence/write-queue.test.ts @@ -54,6 +54,7 @@ describe('SqliteWriteQueue', () => { }); await queue.drain(); expect(queue.droppedWrites).toBe(1); + expect([...queue.droppedWritesByTable]).toEqual([['bad', 1]]); expect(await t.repos.sessions.get('shop-a1b2c3d4')).not.toBeNull(); expect(t.logger.records.some((r) => r.msg === 'write failed')).toBe(true); await t.close(); @@ -66,8 +67,13 @@ describe('SqliteWriteQueue', () => { expect(queue.enqueue('two', async () => undefined)).toBe(false); expect(queue.droppedWrites).toBe(1); await queue.close(); - expect(queue.enqueue('late', async () => undefined)).toBe(false); + expect(queue.enqueue('tool_calls.insert', async () => undefined)).toBe(false); expect(queue.droppedWrites).toBe(2); + // By table: the operation up to the first dot (spec 10 §7 `browserhive.db.dropped_writes`). + expect([...queue.droppedWritesByTable]).toEqual([ + ['two', 1], + ['tool_calls', 1], + ]); await queue.close(); await t.close(); }); diff --git a/packages/core/src/infra/telemetry/metrics.test.ts b/packages/core/src/infra/telemetry/metrics.test.ts index c8917fc..3d3c25a 100644 --- a/packages/core/src/infra/telemetry/metrics.test.ts +++ b/packages/core/src/infra/telemetry/metrics.test.ts @@ -1,34 +1,48 @@ -/** @module infra/telemetry/metrics.test — instruments are created lazily once, with catalogue names. */ +/** @module infra/telemetry/metrics.test — instruments are created lazily once, with the catalogue's name, unit, description and instrument kind (spec 10 §7). */ import { describe, expect, it } from 'bun:test'; import { metrics } from '@opentelemetry/api'; -import { createInstruments, METRIC } from './metrics.ts'; +import { createInstruments, type Instruments, METRIC, METRIC_DEFINITIONS } from './metrics.ts'; + +type Created = { method: string; name: string; options: { unit?: string; description?: string } }; + +function spyMeter() { + const created: Created[] = []; + const meter = metrics.getMeter('test'); + const spy = new Proxy(meter, { + get(target, prop, receiver) { + const value = Reflect.get(target, prop, receiver); + if (typeof value === 'function' && String(prop).startsWith('create')) { + return (name: string, options: Created['options'] = {}) => { + created.push({ method: String(prop), name, options }); + return Reflect.apply(value, target, [name, options]); + }; + } + return value; + }, + }); + return { created, instruments: createInstruments(spy) }; +} + +/** The meter method each catalogue kind maps to. */ +function methodOf(kind: string, observable: boolean): string { + if (kind === 'histogram') return 'createHistogram'; + if (kind === 'gauge') return 'createObservableGauge'; + if (kind === 'counter') return observable ? 'createObservableCounter' : 'createCounter'; + return observable ? 'createObservableUpDownCounter' : 'createUpDownCounter'; +} describe('createInstruments', () => { it('creates each instrument on first access, once', () => { - const names: string[] = []; - const meter = metrics.getMeter('test'); - const spy = new Proxy(meter, { - get(target, prop, receiver) { - const value = Reflect.get(target, prop, receiver); - if (typeof value === 'function' && String(prop).startsWith('create')) { - return (name: string, ...rest: unknown[]) => { - names.push(name); - return Reflect.apply(value, target, [name, ...rest]); - }; - } - return value; - }, - }); - const instruments = createInstruments(spy); - expect(names).toEqual([]); + const { created, instruments } = spyMeter(); + expect(created).toEqual([]); instruments.toolCalls.add(1, { tool: 'navigate' }); instruments.toolCalls.add(1, { tool: 'click' }); instruments.toolCallDuration.record(5, { tool: 'navigate' }); - instruments.sessionsActive.add(1, { state: 'live' }); + instruments.sessionsActive.add(1, { harness: 'unknown' }); instruments.wsBufferedBytes.addCallback(() => undefined); instruments.processEventLoopLag.addCallback(() => undefined); - expect(names).toEqual([ + expect(created.map((c) => c.name)).toEqual([ METRIC.TOOL_CALLS, METRIC.TOOL_CALL_DURATION, METRIC.SESSIONS_ACTIVE, @@ -37,7 +51,24 @@ describe('createInstruments', () => { ]); }); - it('every catalogue name starts with browserhive.', () => { - for (const name of Object.values(METRIC)) expect(name.startsWith('browserhive.')).toBe(true); + it('gives every instrument the kind, unit and description of its catalogue row', () => { + const { created, instruments } = spyMeter(); + for (const key of Object.keys(instruments) as (keyof Instruments)[]) void instruments[key]; + const byName = new Map(created.map((c) => [c.name, c])); + expect(created).toHaveLength(METRIC_DEFINITIONS.length); + for (const d of METRIC_DEFINITIONS) { + expect({ name: d.name, ...byName.get(d.name) }).toEqual({ + name: d.name, + method: methodOf(d.kind, d.observable), + options: { unit: d.unit, description: d.description }, + }); + } + }); + + it('the catalogue lists every name once, each starting with browserhive.', () => { + const names = METRIC_DEFINITIONS.map((d) => d.name); + expect(new Set(names).size).toBe(names.length); + expect([...names].sort()).toEqual([...Object.values(METRIC)].sort()); + for (const name of names) expect(name.startsWith('browserhive.')).toBe(true); }); }); diff --git a/packages/core/src/interface/ws/hub.test.ts b/packages/core/src/interface/ws/hub.test.ts index 433e200..78323c4 100644 --- a/packages/core/src/interface/ws/hub.test.ts +++ b/packages/core/src/interface/ws/hub.test.ts @@ -183,6 +183,7 @@ describe('feed', () => { socket.sendResults.push(0); hub.publish('system', event(2)); expect(socket.closed?.code).toBe(WS_CLOSE.OVERLOADED); + expect(hub.droppedFrames()).toEqual({ screencast: 0, logs: 0, feed: 1 }); }); it('closes 1013 when bufferedAmount stays above the overload bound past the grace', async () => { @@ -236,6 +237,9 @@ describe('screencast', () => { expect(conn.droppedFrames).toBe(1); hub.drain(conn); expect(socket.binaries.map((b) => b[16])).toEqual([1, 3]); + // The hub's running total outlives the connection (spec 10 §7 `ws.frames_dropped`). + hub.close(conn); + expect(hub.droppedFrames().screencast).toBe(1); }); it('drops frames while bufferedAmount exceeds the screencast bound', async () => { @@ -487,5 +491,10 @@ describe('connection lifecycle', () => { m.kind === 'event' && m.payload.type === 'log.record' ? m.payload.record.msg : null, ), ).toEqual(['b']); + expect(hub.droppedFrames().logs).toBe(0); + // A congested subscriber skips the record and the hub counts it. + socket.buffered = 2 * 1024 * 1024; + hub.publishLog({ seq: 4, record: { ts: 1, level: 'error', msg: 'd', module: 'auth' } }); + expect(hub.droppedFrames()).toEqual({ screencast: 0, logs: 1, feed: 0 }); }); }); diff --git a/packages/core/test/integration/isolation.test.ts b/packages/core/test/integration/isolation.test.ts index c3eec8a..08b1c48 100644 --- a/packages/core/test/integration/isolation.test.ts +++ b/packages/core/test/integration/isolation.test.ts @@ -1,6 +1,7 @@ -/** @module test/integration/isolation — two sessions never share cookies, localStorage, sessionStorage, IndexedDB or service workers (the project's core guarantee). */ +/** @module test/integration/isolation — two sessions never share cookies, localStorage, sessionStorage, IndexedDB or service workers (the project's core guarantee), and each runs its own browser process tree (whose memory `browserhive.browser.rss_bytes` reports). */ import { describe, expect, it } from 'bun:test'; +import { createProcessTreeReader } from '../../src/infra/host/process-tree.ts'; import { useDriverFixtures } from './driver-fixture.ts'; /** What `/storage` renders into `#storage` after `window.__setStorage`. */ @@ -153,4 +154,22 @@ describe('session isolation', () => { expect(await b.page.locator('#title').textContent()).toBe('fixture'); expect(crashes).toEqual([]); }); + + it('each session runs its own browser process tree, whose memory can be read', async () => { + const a = await state.launch({ slug: 'treea' }); + const b = await state.launch({ slug: 'treeb', persistenceMode: 'persistent' }); + const pidA = await a.handle.browserPid?.(); + const pidB = await b.handle.browserPid?.(); + expect(pidA).toBeGreaterThan(0); + expect(pidB).toBeGreaterThan(0); + expect(pidA).not.toBe(pidB); + // Cached: the second read opens no new DevTools session. + expect(await a.handle.browserPid?.()).toBe(pidA); + const reader = createProcessTreeReader(); + if (reader === null || pidA == null || pidB == null) return; // Windows: no reader (spec 10 §7). + const rss = await reader.rssOfTrees([pidA, pidB]); + // A whole Chromium tree (browser, zygotes, renderer, GPU, network) is well above 10 MB. + expect(rss.get(pidA)).toBeGreaterThan(10 * 1024 * 1024); + expect(rss.get(pidB)).toBeGreaterThan(10 * 1024 * 1024); + }); }); From b2ab504a40d4fd5062004109371db1dae07b407a Mon Sep 17 00:00:00 2001 From: Amir Ghorbani Date: Tue, 29 Sep 2026 12:46:16 -0400 Subject: [PATCH 4/6] fix(observability): drop closed sessions from gauges, report every recorder table MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The SDK repeats the last value of an observable gauge's series that is no longer observed, so a closed session's browser memory or a closed WebSocket's buffered bytes was exported forever. Observable gauges are now collected with delta temporality, which OTLP gauges do not carry, so each export holds only what exists; counters stay cumulative. browserhive.db.dropped_writes now reports every recorder table from the start at 0, so a rate over the series works before the first drop. Spec 10 §7 and the guide say both, and that browser memory is the sum of the processes' RSS. The guard test checks that a closed session's series disappears and that the spec's recorder tables match. --- .changeset/otel-documented-metrics.md | 1 + docs/guide/telemetry.md | 4 +- .../src/composition/adapters/metrics.test.ts | 6 ++- .../src/composition/adapters/metrics.ts | 7 ++- .../test/composition/metrics-guard.test.ts | 48 +++++++++++++++---- packages/core/src/infra/telemetry/metrics.ts | 15 ++++++ .../core/src/infra/telemetry/telemetry.ts | 30 ++++++++++-- packages/core/src/public/runtime.ts | 1 + specs/10-error-handling-and-telemetry.md | 6 +-- 9 files changed, 97 insertions(+), 21 deletions(-) diff --git a/.changeset/otel-documented-metrics.md b/.changeset/otel-documented-metrics.md index 42cd468..5f2e1e5 100644 --- a/.changeset/otel-documented-metrics.md +++ b/.changeset/otel-documented-metrics.md @@ -7,4 +7,5 @@ OpenTelemetry now exports the notification, attention, session-launch, WebSocket - **Metrics that were documented but never sent** now reach your collector: `browserhive.attention.wait`, `browserhive.session.launch.duration`, `browserhive.ws.connections`, `browserhive.ws.buffered_bytes`, `browserhive.ws.frames_dropped`, `browserhive.db.dropped_writes`, `browserhive.browser.rss_bytes` (each session's browser with all its processes, every 10 seconds; Linux and macOS) and `browserhive.process.event_loop_lag`. - **Attributes the docs promised** are now set: `closed_reason` on `browserhive.session.lifetime`, `kind` on `browserhive.attention.open` (vault confirmations count too), `table` on `browserhive.retention.pruned_rows`. - **The notification metrics** `browserhive.notifications.deliveries`, `.actions` and `.reports` are now in the [telemetry guide](https://browserhive.ai/docs/guide/telemetry). `.reports` now counts on-demand digests as `manual`, as documented, and no longer counts a silent revision of an open in-app anomaly alert. +- **Gauges only report what exists**: a closed session's browser memory or a closed WebSocket's buffered bytes disappears from the next export instead of repeating its last value. `browserhive.db.dropped_writes` reports every recorder table from the start, at 0. - **Every metric has a unit** (`ms`, `By`, or a count such as `{call}`), and the guide's table lists each one with its type, unit and attributes. With telemetry off nothing is measured, as before. diff --git a/docs/guide/telemetry.md b/docs/guide/telemetry.md index c0fcb11..7d5af53 100644 --- a/docs/guide/telemetry.md +++ b/docs/guide/telemetry.md @@ -53,7 +53,7 @@ If the MCP client sends a W3C `traceparent` in the tool call's `_meta`, the tool | `browserhive.db.write_queue.depth` | gauge | `{write}` | | Writes waiting to be saved. | | `browserhive.db.dropped_writes` | counter | `{write}` | `table` | Writes lost because the queue was full or the write failed. | | `browserhive.db.size_bytes` | gauge | `By` | | Database size. | -| `browserhive.browser.rss_bytes` | gauge | `By` | `session_id` | Memory of each session's browser, all its processes together, sampled every 10 seconds. Linux and macOS only. | +| `browserhive.browser.rss_bytes` | gauge | `By` | `session_id` | Memory of each session's browser, all its processes added up (memory they share counts once per process), sampled every 10 seconds. Linux and macOS only. | | `browserhive.attention.open` | up-down counter | `{request}` | `kind` (`attention`, `vault_confirm`) | Operator requests waiting for an answer. | | `browserhive.attention.wait` | histogram | `ms` | `status` | How long each attention request waited until it was answered, timed out or cancelled. | | `browserhive.vault.fills` | counter | `{fill}` | `result` | Vault fills. | @@ -68,7 +68,7 @@ If the MCP client sends a W3C `traceparent` in the tool call's `_meta`, the tool `harness` on metrics is one of the known harness names, `unknown` or `other` (every name BrowserHive doesn't know is folded into `other`), so it adds a bounded number of series. The model, workspace and extra labels are never metric attributes. -Metrics are exported every 30 seconds. Gauges and the WebSocket and write-queue totals are read at export time; with telemetry off nothing is measured. +Metrics are exported every 30 seconds. Gauges and the WebSocket and write-queue totals are read at export time, and a gauge only reports what exists at that moment: a closed session or connection disappears from the next export. With telemetry off nothing is measured. ### Logs diff --git a/packages/browserhive/src/composition/adapters/metrics.test.ts b/packages/browserhive/src/composition/adapters/metrics.test.ts index 5e0efb1..08027a4 100644 --- a/packages/browserhive/src/composition/adapters/metrics.test.ts +++ b/packages/browserhive/src/composition/adapters/metrics.test.ts @@ -253,7 +253,11 @@ describe('wireMetrics: observable instruments', () => { ); await Promise.resolve(); expect(collect('dbWriteQueueDepth')).toEqual([{ value: 3, attrs: {} }]); - expect(collect('dbDroppedWrites')).toEqual([{ value: 2, attrs: { table: 'tool_calls' } }]); + // Every recorder table from the start, at 0 until it loses a write. + const dropped = collect('dbDroppedWrites'); + expect(dropped).toHaveLength(8); + expect(dropped).toContainEqual({ value: 2, attrs: { table: 'tool_calls' } }); + expect(dropped).toContainEqual({ value: 0, attrs: { table: 'logs' } }); expect(collect('dbSizeBytes')).toEqual([{ value: 4096, attrs: {} }]); expect(collect('browserRssBytes')).toEqual([{ value: 123_456, attrs: { session_id: 's-1' } }]); expect(collect('processEventLoopLag')).toEqual([{ value: 12.5, attrs: {} }]); diff --git a/packages/browserhive/src/composition/adapters/metrics.ts b/packages/browserhive/src/composition/adapters/metrics.ts index 32450bd..3ecbdeb 100644 --- a/packages/browserhive/src/composition/adapters/metrics.ts +++ b/packages/browserhive/src/composition/adapters/metrics.ts @@ -10,6 +10,7 @@ import type { WriteQueue, } from '@browserhive/core/runtime'; import { + DB_WRITE_TABLES, readProcessMemory, startEventLoopLagMonitor, WS_BUFFERED_BYTES_MAX_SERIES, @@ -181,8 +182,10 @@ export function wireMetrics( [ instruments.dbDroppedWrites, (r) => { - for (const [table, count] of sources.queue.droppedWritesByTable) { - r.observe(count, { table }); + const dropped = sources.queue.droppedWritesByTable; + for (const table of DB_WRITE_TABLES) r.observe(dropped.get(table) ?? 0, { table }); + for (const [table, count] of dropped) { + if (!(DB_WRITE_TABLES as readonly string[]).includes(table)) r.observe(count, { table }); } }, ], diff --git a/packages/browserhive/test/composition/metrics-guard.test.ts b/packages/browserhive/test/composition/metrics-guard.test.ts index 0242090..fb7cfa6 100644 --- a/packages/browserhive/test/composition/metrics-guard.test.ts +++ b/packages/browserhive/test/composition/metrics-guard.test.ts @@ -27,6 +27,7 @@ import { createRedactor, createSystemClock, createTelemetry, + DB_WRITE_TABLES, DegradationService, type DomainEvents, type Logger, @@ -108,6 +109,13 @@ describe('documented == defined', () => { expect(SPEC).toEqual(CATALOGUE); }); + it('spec 10 §7 names the recorder tables the dropped-writes counter reports from the start', () => { + const text = readFileSync(join(ROOT, 'specs/10-error-handling-and-telemetry.md'), 'utf8'); + const row = text.split('\n').find((l) => l.startsWith('| `browserhive.db.dropped_writes`')); + const listed = /\(the recorder table: ([^)]*)\)/.exec(row ?? '')?.[1] ?? ''; + expect(ticked(listed)).toEqual([...DB_WRITE_TABLES]); + }); + it('the telemetry guide lists exactly the catalogue (types without the observable note)', () => { expect(DOCS).toEqual(CATALOGUE.map((d) => ({ ...d, observable: false }))); }); @@ -321,8 +329,8 @@ const SOURCES: readonly Source[] = [ name: 'vault broker and blocklist (vault.access, blocklist.hit)', drives: ['browserhive.vault.fills', 'browserhive.blocklist.hits'], async run(w) { - publish(w, 'vault.access', { type: 'vault.access', row: { result: 'filled' } }); - publish(w, 'blocklist.hit', { type: 'blocklist.hit', row: { source: 'config' } }); + publish(w, 'vault.access', { type: 'vault.access', row: { result: 'success' } }); + publish(w, 'blocklist.hit', { type: 'blocklist.hit', row: { source: 'request' } }); }, }, { @@ -409,7 +417,7 @@ const SOURCES: readonly Source[] = [ drives: ['browserhive.notifications.actions'], async run(w) { await w.ops.actions.press({ - token: 'a'.repeat(22), + token: 'a'.repeat(11), origin: null, actor: { platform: 'ntfy', id: null, name: null }, }); @@ -435,6 +443,7 @@ describe('documented == recorded', () => { let server: ReturnType | undefined; let telemetry: Telemetry | undefined; let received = new Map(); + let first = new Map(); const stops: (() => unknown)[] = []; beforeAll(async () => { @@ -605,9 +614,14 @@ describe('documented == recorded', () => { schedule: () => () => undefined, every: () => () => undefined, }); - // A "browser" whose process tree is this test process: the real reader, the real sampler. + // "Browsers" whose process trees are this test process and its parent: the real reader, the + // real sampler. The second one closes before the last export. + const browsers = [ + { id: SESSION, browserPid: async () => process.pid }, + { id: 'gone-00000001', browserPid: async () => process.ppid }, + ]; const sampler = startBrowserMemorySampler({ - sessions: () => [{ id: SESSION, browserPid: async () => process.pid }], + sessions: () => browsers, reader: createProcessTreeReader(), repeat: () => () => undefined, }); @@ -623,8 +637,12 @@ describe('documented == recorded', () => { ); const world: World = { bus, queue, realtime, ops, retention, sampler, channelId, hooks }; for (const source of SOURCES) await source.run(world); - // The size gauge reports the value read at the previous collection. await telemetry.forceFlush(); + first = receivedMetrics(bodies); + // A session closes; the next export no longer reports its browser. (The size gauge also + // reports the value read at the previous collection.) + browsers.pop(); + await sampler.sample(); await telemetry.forceFlush(); received = receivedMetrics(bodies); }); @@ -658,6 +676,13 @@ describe('documented == recorded', () => { }); } + it('stops reporting a gauge series that is no longer observed (a closed session)', () => { + const sessions = (m: Map) => + (m.get('browserhive.browser.rss_bytes')?.points ?? []).map((p) => p.attrs['session_id']); + expect(sessions(first).sort()).toEqual(['gone-00000001', SESSION]); + expect(sessions(received)).toEqual([SESSION]); + }); + it('carries the values of the driven sources', () => { const points = (name: MetricName) => received.get(name)?.points ?? []; expect(points('browserhive.tool_calls')).toContainEqual({ @@ -679,9 +704,14 @@ describe('documented == recorded', () => { expect(points('browserhive.retention.pruned_rows')).toEqual([ { value: 4, attrs: { table: 'tool_calls' } }, ]); - expect(points('browserhive.db.dropped_writes')).toEqual([ - { value: 1, attrs: { table: 'tool_calls' } }, - ]); + expect(points('browserhive.db.dropped_writes')).toContainEqual({ + value: 1, + attrs: { table: 'tool_calls' }, + }); + expect(points('browserhive.db.dropped_writes')).toContainEqual({ + value: 0, + attrs: { table: 'logs' }, + }); expect(points('browserhive.ws.connections')).toContainEqual({ value: 1, attrs: {} }); expect(points('browserhive.ws.frames_dropped')).toContainEqual({ value: 1, diff --git a/packages/core/src/infra/telemetry/metrics.ts b/packages/core/src/infra/telemetry/metrics.ts index ccf11bc..44d3cbc 100644 --- a/packages/core/src/infra/telemetry/metrics.ts +++ b/packages/core/src/infra/telemetry/metrics.ts @@ -43,6 +43,21 @@ export type MetricName = (typeof METRIC)[keyof typeof METRIC]; /** Cap on `connection_id` series for the buffered-bytes gauge (spec 10 §7). */ export const WS_BUFFERED_BYTES_MAX_SERIES = 50; +/** + * The recorder tables of the write queue (spec 10 §7 `browserhive.db.dropped_writes{table}`): each + * is reported from the start, at 0 until it loses a write, so a rate over the series works. + */ +export const DB_WRITE_TABLES = [ + 'sessions', + 'tool_calls', + 'pages', + 'screenshots', + 'blocked_requests', + 'vault_access', + 'events', + 'logs', +] as const; + /** How often the browser-memory sampler reads the process trees (spec 10 §7). */ export const BROWSER_RSS_SAMPLE_INTERVAL_MS = 10_000; diff --git a/packages/core/src/infra/telemetry/telemetry.ts b/packages/core/src/infra/telemetry/telemetry.ts index 89630da..bf86559 100644 --- a/packages/core/src/infra/telemetry/telemetry.ts +++ b/packages/core/src/infra/telemetry/telemetry.ts @@ -9,7 +9,7 @@ import { } from '@opentelemetry/api'; import { logs, type Logger as OtelLogger } from '@opentelemetry/api-logs'; import type { LogRecordExporter } from '@opentelemetry/sdk-logs'; -import type { PushMetricExporter } from '@opentelemetry/sdk-metrics'; +import type { InstrumentType, PushMetricExporter } from '@opentelemetry/sdk-metrics'; import type { SpanExporter } from '@opentelemetry/sdk-trace-base'; import { createInstruments, type Instruments } from './metrics.ts'; import { TRACER_NAME } from './spans.ts'; @@ -174,9 +174,9 @@ async function enabledTelemetry(options: TelemetryOptions): Promise { protocol === 'http/json' ? await import('@opentelemetry/exporter-metrics-otlp-http') : await import('@opentelemetry/exporter-metrics-otlp-proto'); - const exporter = trackMetricExporter( - new OTLPMetricExporter(exporterConfig('/v1/metrics')), - tracker, + const exporter = reportOnlyObservedGauges( + trackMetricExporter(new OTLPMetricExporter(exporterConfig('/v1/metrics')), tracker), + sdk, ); const provider = new sdk.MeterProvider({ resource, @@ -315,6 +315,28 @@ function trackMetricExporter( }; } +/** + * Observable gauges are collected with delta temporality, which OTLP gauges do not carry, so each + * export holds only the series observed at that collection: a closed session's browser memory or a + * closed WebSocket's buffered bytes stops being reported instead of repeating its last value. + * Every other instrument keeps the exporter's temporality (cumulative by default). + */ +function reportOnlyObservedGauges( + exporter: PushMetricExporter, + sdk: Pick< + typeof import('@opentelemetry/sdk-metrics'), + 'AggregationTemporality' | 'InstrumentType' + >, +): PushMetricExporter { + const inner = exporter.selectAggregationTemporality?.bind(exporter); + return Object.assign(exporter, { + selectAggregationTemporality: (type: InstrumentType) => + type === sdk.InstrumentType.OBSERVABLE_GAUGE || type === sdk.InstrumentType.GAUGE + ? sdk.AggregationTemporality.DELTA + : (inner?.(type) ?? sdk.AggregationTemporality.CUMULATIVE), + }); +} + function trackLogExporter(inner: LogRecordExporter, tracker: FailureTracker): LogRecordExporter { return { export: (items, cb) => diff --git a/packages/core/src/public/runtime.ts b/packages/core/src/public/runtime.ts index 2faf8cf..c283e87 100644 --- a/packages/core/src/public/runtime.ts +++ b/packages/core/src/public/runtime.ts @@ -45,6 +45,7 @@ export type { LogSink } from '../infra/logging/sinks.ts'; export { createBunProcessRunner } from '../infra/process/bun-process-runner.ts'; export { BROWSER_RSS_SAMPLE_INTERVAL_MS, + DB_WRITE_TABLES, type Instruments, METRIC, METRIC_DEFINITIONS, diff --git a/specs/10-error-handling-and-telemetry.md b/specs/10-error-handling-and-telemetry.md index 703f137..56a4495 100644 --- a/specs/10-error-handling-and-telemetry.md +++ b/specs/10-error-handling-and-telemetry.md @@ -339,7 +339,7 @@ All spans use `@opentelemetry/api`'s tracer `browserhive`; attributes use the `b ## 7. Metrics catalogue -This table is the single source of truth: `infra/telemetry/metrics.ts` defines exactly these instruments, the guard test (09 §3.2) fails when a row here, a row of `docs/guide/telemetry.md` and the catalogue disagree on name, type, unit or attributes, or when an instrument is never recorded. Attribute values in parentheses are the closed set that can occur. Types are the OTLP data a collector receives: *counter* is a monotonic sum, *up-down counter* a non-monotonic sum, *gauge* a gauge; "observable" instruments are read by a callback at each export (every 30 s) instead of being written on each event. +This table is the single source of truth: `infra/telemetry/metrics.ts` defines exactly these instruments, the guard test (09 §3.2) fails when a row here, a row of `docs/guide/telemetry.md` and the catalogue disagree on name, type, unit or attributes, or when an instrument is never recorded. Attribute values in parentheses are the closed set that can occur. Types are the OTLP data a collector receives: *counter* is a monotonic sum, *up-down counter* a non-monotonic sum, *gauge* a gauge; "observable" instruments are read by a callback at each export (every 30 s) instead of being written on each event. An observable gauge reports only the series observed at that export (it is collected with delta temporality, which an OTLP gauge does not carry), so a closed session or WebSocket stops being reported instead of repeating its last value. | Instrument | Type | Unit | Attributes | Recorded | |---|---|---|---|---| @@ -352,9 +352,9 @@ This table is the single source of truth: `infra/telemetry/metrics.ts` defines e | `browserhive.ws.buffered_bytes` | gauge (observable) | `By` | `connection_id` (the 50 connections with the most buffered bytes) | bytes the socket has not flushed yet, per connection (http transport only) | | `browserhive.ws.frames_dropped` | counter (observable) | `{frame}` | `channel` (`screencast`, `logs`, `feed`) | the hub's running totals: a screencast frame replaced by a newer one or refused by a congested socket; a log record skipped for a congested `logs` subscriber; a feed frame the socket refused (the connection closes with 1013 and replays from its cursor) | | `browserhive.db.write_queue.depth` | gauge (observable) | `{write}` | — | recorder writes waiting in the queue | -| `browserhive.db.dropped_writes` | counter (observable) | `{write}` | `table` (the recorder table: `sessions`, `tool_calls`, `pages`, `screenshots`, `blocked_requests`, `vault_access`, `events`, `logs`) | the write queue's running totals: writes refused (queue full or closed) or failed inside their drain | +| `browserhive.db.dropped_writes` | counter (observable) | `{write}` | `table` (the recorder table: `sessions`, `tool_calls`, `pages`, `screenshots`, `blocked_requests`, `vault_access`, `events`, `logs`) | the write queue's running totals: writes refused (queue full or closed) or failed inside their drain; every recorder table is reported from the start, at 0 until it loses a write | | `browserhive.db.size_bytes` | gauge (observable) | `By` | — | database file size, read in the background and reported at the next export | -| `browserhive.browser.rss_bytes` | gauge (observable) | `By` | `session_id` | resident memory of each live session's browser process tree (the browser process and every descendant: renderers, GPU, utilities), sampled every 10 s; Linux and macOS only, no data points on Windows | +| `browserhive.browser.rss_bytes` | gauge (observable) | `By` | `session_id` | resident memory of each live session's browser process tree: the sum of the RSS of the browser process and every descendant (renderers, GPU, utilities; memory they share counts once per process), sampled every 10 s from `/proc` or `ps` with the browser pid read once over a browser-level DevTools session; Linux and macOS only, no data points on Windows | | `browserhive.attention.open` | up-down counter | `{request}` | `kind` (`attention`, `vault_confirm`) | +1 when an operator request opens, −1 when it settles | | `browserhive.attention.wait` | histogram | `ms` | `status` (`resolved`, `rejected`, `timeout`, `cancelled`) | once per settled attention request, its `waited_ms` | | `browserhive.vault.fills` | counter | `{fill}` | `result` | once per `vault.access` | From da4e5300bce4a8b5d639b9f600dc5338002230e6 Mon Sep 17 00:00:00 2001 From: Amir Ghorbani Date: Tue, 29 Sep 2026 12:48:16 -0400 Subject: [PATCH 5/6] fix(observability): wire the metric sources only when metrics are exported With --otelSignals traces,logs the meter is a no-op, so the bus consumers, the browser-memory sampler and the event-loop monitor are no longer started for nothing. --- .../browserhive/src/composition/phases/wire-observers.ts | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/packages/browserhive/src/composition/phases/wire-observers.ts b/packages/browserhive/src/composition/phases/wire-observers.ts index f945334..9c94ac1 100644 --- a/packages/browserhive/src/composition/phases/wire-observers.ts +++ b/packages/browserhive/src/composition/phases/wire-observers.ts @@ -26,7 +26,7 @@ export function patchrightVersion(): string | null { /** * The OTel metrics consumers and the browser-memory sampler (spec 10 §7); only called with - * telemetry on, so with `--otel` off nothing is subscribed, sampled or timed. + * `--otel` on and the `metrics` signal exported. */ function wireMetricsFor(ctx: BootContext, logger: Logger): () => void { const { telemetry } = part(ctx.observability, 'observability'); @@ -155,7 +155,9 @@ export async function wireObserversPhase(ctx: BootContext): Promise .refreshLastBackup() .catch((err: unknown) => logger.warn('backup scan failed', { err: serializeError(err) })); domain.sweeper.start(); - const unwireMetrics = telemetry.enabled ? wireMetricsFor(ctx, logger) : () => undefined; + // Only with the metrics signal exported: otherwise nothing is subscribed, sampled or timed. + const metricsOn = telemetry.enabled && config.otelSignals.includes('metrics'); + const unwireMetrics = metricsOn ? wireMetricsFor(ctx, logger) : () => undefined; ctx.observers = { status }; return { From e20eafc1aeaf44adf6bd75a035bf758e73c9f9a6 Mon Sep 17 00:00:00 2001 From: Amir Ghorbani Date: Tue, 29 Sep 2026 12:53:55 -0400 Subject: [PATCH 6/6] docs(observability): say the metric adapters run only when metrics are exported --- packages/browserhive/src/composition/adapters/browser-memory.ts | 2 +- packages/browserhive/src/composition/adapters/metrics.ts | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/browserhive/src/composition/adapters/browser-memory.ts b/packages/browserhive/src/composition/adapters/browser-memory.ts index f5d04d0..584f710 100644 --- a/packages/browserhive/src/composition/adapters/browser-memory.ts +++ b/packages/browserhive/src/composition/adapters/browser-memory.ts @@ -1,4 +1,4 @@ -/** @module composition/adapters/browser-memory — samples the resident memory of each live session's browser process tree every 10 s for `browserhive.browser.rss_bytes` (spec 10 §7). Only started when telemetry is on. */ +/** @module composition/adapters/browser-memory — samples the resident memory of each live session's browser process tree every 10 s for `browserhive.browser.rss_bytes` (spec 10 §7). Only started when metrics are exported. */ import type { ProcessTreeReader } from '@browserhive/core/runtime'; import { BROWSER_RSS_SAMPLE_INTERVAL_MS } from '@browserhive/core/runtime'; diff --git a/packages/browserhive/src/composition/adapters/metrics.ts b/packages/browserhive/src/composition/adapters/metrics.ts index 3ecbdeb..72236b6 100644 --- a/packages/browserhive/src/composition/adapters/metrics.ts +++ b/packages/browserhive/src/composition/adapters/metrics.ts @@ -1,4 +1,4 @@ -/** @module composition/adapters/metrics — OTel metrics consumers (spec 10 §7): bus events → counters/histograms; observable instruments over the write queue, the database size, the realtime hub, the browser-memory sampler and this process. Only wired when telemetry is on. */ +/** @module composition/adapters/metrics — OTel metrics consumers (spec 10 §7): bus events → counters/histograms; observable instruments over the write queue, the database size, the realtime hub, the browser-memory sampler and this process. Only wired when metrics are exported. */ import { metricHarness } from '@browserhive/contracts/harness'; import type {