From 8a8d2fdabfae2b51ed92744d974da6ef5a920601 Mon Sep 17 00:00:00 2001 From: Marc Liu Date: Wed, 22 Jul 2026 05:02:43 -0400 Subject: [PATCH 1/3] feat: instrument pending external-event queue health Add an in-memory metrics collector for the pending external-event delivery queue, plus a delivery-terminal hook on ExternalEventStore so every delivered and terminal-failed transition is counted once from a single observation point. SpaceRuntime records enqueues (by source + target run/node state), flush attempts, claim conflicts, stale-session skips, and structured logs at the previously-silent skip seams, and exposes getQueueHealthSnapshot() merging cumulative counters with live gauges (depth, event-age p95, in-flight, digest backlog, persisted pending). Surface it via the space.externalEvents.queueHealth RPC. Service wires the shared metrics to the store hook. --- .../external-events/external-event-store.ts | 49 ++- .../external-events/queue-health-metrics.ts | 295 ++++++++++++++++++ packages/daemon/src/lib/rpc-handlers/index.ts | 6 + .../space/runtime/space-runtime-service.ts | 34 ++ .../src/lib/space/runtime/space-runtime.ts | 116 ++++++- .../storage/external-event-store.test.ts | 81 +++++ .../storage/queue-health-metrics.test.ts | 188 +++++++++++ .../space-runtime-external-events.test.ts | 155 ++++++++- 8 files changed, 916 insertions(+), 8 deletions(-) create mode 100644 packages/daemon/src/lib/external-events/queue-health-metrics.ts create mode 100644 packages/daemon/tests/unit/4-space-storage/storage/queue-health-metrics.test.ts diff --git a/packages/daemon/src/lib/external-events/external-event-store.ts b/packages/daemon/src/lib/external-events/external-event-store.ts index e1c19d4a55..341dc25d75 100644 --- a/packages/daemon/src/lib/external-events/external-event-store.ts +++ b/packages/daemon/src/lib/external-events/external-event-store.ts @@ -31,6 +31,7 @@ import { TERMINAL_DELIVERY_STATES, TERMINAL_EVENT_STATES, } from './types'; +import type { DeliveryTerminalEvent } from './queue-health-metrics'; import { validateLiteralTopic, validateSource } from './topic-validator'; interface ExternalEventRow { @@ -90,6 +91,20 @@ export class ExternalEventValidationError extends Error { export class ExternalEventStore { constructor(private readonly db: BunDatabase) {} + /** + * Optional hook fired when a delivery row transitions to a terminal state + * (`delivered`, or `failed` via `failure.terminal=true`). Set by the space + * runtime so queue-health metrics can count every delivery outcome from a + * single observation point, regardless of which call path reached the + * transition. Only fired on an actual transition (`changes > 0`). + */ + private deliveryTerminalHook?: (event: DeliveryTerminalEvent) => void; + + /** Install the delivery-terminal observation hook. */ + setDeliveryTerminalHook(hook: (event: DeliveryTerminalEvent) => void): void { + this.deliveryTerminalHook = hook; + } + // --------------------------------------------------------------------------- // Source event lifecycle // --------------------------------------------------------------------------- @@ -430,7 +445,7 @@ export class ExternalEventStore { /** Mark the delivery row terminal `delivered`. No-op if already terminal. */ markDeliveryDelivered(eventId: string, deliveryKey: string): void { const now = Date.now(); - this.db + const result = this.db .prepare( `UPDATE space_external_event_deliveries SET state = 'delivered', failure_reason = NULL, delivered_at = ?, updated_at = ? @@ -438,6 +453,9 @@ export class ExternalEventStore { AND state NOT IN ('delivered', 'failed')` ) .run(now, now, eventId, deliveryKey); + if (result.changes > 0 && this.deliveryTerminalHook) { + this.deliveryTerminalHook({ eventId, deliveryKey, outcome: 'delivered', reason: null }); + } } /** @@ -452,7 +470,7 @@ export class ExternalEventStore { markDeliveryFailed(eventId: string, deliveryKey: string, failure: DeliveryFailure): void { const now = Date.now(); const newState: ExternalEventDeliveryState = failure.terminal ? 'failed' : 'pending'; - this.db + const result = this.db .prepare( `UPDATE space_external_event_deliveries SET state = ?, failure_reason = ?, updated_at = ? @@ -460,6 +478,14 @@ export class ExternalEventStore { AND state NOT IN ('delivered', 'failed')` ) .run(newState, failure.reason, now, eventId, deliveryKey); + if (failure.terminal && result.changes > 0 && this.deliveryTerminalHook) { + this.deliveryTerminalHook({ + eventId, + deliveryKey, + outcome: 'failed', + reason: failure.reason, + }); + } } /** List delivery rows for an event (for diagnostics and tests). */ @@ -555,6 +581,25 @@ export class ExternalEventStore { return rows.map(deliveryRowToRecord); } + /** + * Return the event-age (`now - event.created_at`, in ms) of every + * DB-persisted `pending` delivery row, in a single joined query. Used by the + * queue-health snapshot's persisted-pending age gauge. The anchor is the + * source event's ingestion time (`space_external_events.created_at`), matching + * the runtime's event-age TTL semantics — see `EXTERNAL_EVENT_QUEUE_TTL_MS`. + */ + getPendingDeliveryAges(now: number = Date.now()): number[] { + const rows = this.db + .prepare( + `SELECT (? - e.created_at) AS age + FROM space_external_event_deliveries d + INNER JOIN space_external_events e ON e.id = d.event_id + WHERE d.state = 'pending'` + ) + .all(now) as { age: number }[]; + return rows.map((row) => row.age); + } + getDelivery(eventId: string, deliveryKey: string): ExternalEventDeliveryRecord | null { const row = this.db .prepare( diff --git a/packages/daemon/src/lib/external-events/queue-health-metrics.ts b/packages/daemon/src/lib/external-events/queue-health-metrics.ts new file mode 100644 index 0000000000..5ef4ce1d8f --- /dev/null +++ b/packages/daemon/src/lib/external-events/queue-health-metrics.ts @@ -0,0 +1,295 @@ +/** + * In-memory health metrics for the pending external-event delivery queue. + * + * Pending external-event delivery is the runtime's core reliability mechanism + * for getting a GitHub (or other source) event to a workflow node agent that + * is not yet ready to receive it. These counters + live gauges give operators + * a single view of queue health: how much is being enqueued and from where, + * which target states force queuing, how often deliveries are skipped + * (claim conflicts, stale sessions, deliverability guards), how often the + * queue evicts (cap/TTL), and the terminal outcome of every delivery. + * + * Cumulative counters are process-lifetime (since the runtime started) and are + * NOT persisted — they reset on daemon restart. Live gauges (depth, age, + * in-flight) are computed at read time by {@link SpaceRuntime.getQueueHealthSnapshot} + * from the runtime's in-memory maps and the durable store, then merged with + * these counters into a {@link QueueHealthSnapshot}. + */ + +/** + * A delivery row that just transitioned to a terminal state. Emitted by + * `ExternalEventStore` via its delivery-terminal hook — the single source of + * truth for `delivered` and `finalFailuresByReason`, so every terminal + * transition is counted exactly once regardless of which call path reached it. + */ +export interface DeliveryTerminalEvent { + eventId: string; + deliveryKey: string; + /** `delivered` for success, `failed` for a terminal failure. */ + outcome: 'delivered' | 'failed'; + /** Free-form failure reason (`null` when delivered). */ + reason: string | null; +} + +/** Min/max/avg/p95 of a set of millisecond ages (event-age of queued items). */ +export interface QueueAgeStats { + count: number; + minMs: number; + maxMs: number; + avgMs: number; + p95Ms: number; +} + +/** Cumulative counters since the runtime started counting. */ +export interface QueueHealthCounters { + /** Epoch ms when counting started (runtime start). */ + since: number; + /** Total items enqueued into the in-memory pending queue. */ + enqueue: number; + /** Enqueues broken down by event source (e.g. `github`). */ + enqueueBySource: Record; + /** + * Enqueues broken down by the target's run+node state at enqueue time + * (e.g. `run=in_progress;node=pending`). Surfaces which states force + * events to be queued rather than delivered immediately. + */ + enqueueByTargetState: Record; + /** Number of pending-queue flush attempts (target activation / retry). */ + flushAttempts: number; + /** Total items handed off to dispatch across all flush attempts. */ + flushItemsDispatched: number; + /** Deliveries that reached terminal `delivered` (success). */ + delivered: number; + /** Terminal failures broken down by persisted `failure_reason`. */ + finalFailuresByReason: Record; + /** + * Non-terminal skips: a delivery was ready to dispatch but another path + * already had it in flight (`externalEventDeliveriesInFlight`). The delivery + * stays pending; this counts how often concurrent dispatch races occurred. + */ + claimConflicts: number; + /** + * Non-terminal skips: the target's worker session was no longer live (or its + * space was paused) at injection time, so the delivery was deferred/requeued + * rather than injected. + */ + staleSessionSkips: number; +} + +/** Live gauges computed at read time from in-memory + DB state. */ +export interface QueueHealthGauges { + /** In-memory pending items across all target queues. */ + queueDepth: number; + /** Distinct target queues currently holding pending items. */ + queueKeys: number; + /** Delivery keys currently mid-dispatch (`externalEventDeliveriesInFlight`). */ + inFlight: number; + /** Items buffered in rate-limit digests awaiting the next digest flush. */ + digestBacklog: number; + /** Active bounded retry timers. */ + retryTimers: number; + /** DB-persisted `pending` delivery rows (global, across all runs). */ + persistedPending: number; + /** Event-age stats for in-memory queued items, or `null` when empty. */ + queueAgeMs: QueueAgeStats | null; + /** Event-age stats for DB-persisted pending deliveries, or `null` when none. */ + persistedAgeMs: QueueAgeStats | null; +} + +/** Aggregate queue-health snapshot surfaced to operators/debug views. */ +export interface QueueHealthSnapshot { + /** Epoch ms the snapshot was collected. */ + collectedAt: number; + counters: QueueHealthCounters; + /** + * Terminal failures grouped into operator-meaningful categories derived from + * {@link QueueHealthCounters.finalFailuresByReason}. Each terminal failure is + * counted in exactly one category — this is a view over the same underlying + * counters, not a second tally. + */ + failuresByCategory: Record; + gauges: QueueHealthGauges; +} + +export type FailureCategory = + | 'ttl_expired' + | 'cap_eviction' + | 'deliverability' + | 'retry_exhausted' + | 'injection_error' + | 'other'; + +const FAILURE_CATEGORY_PREFIXES: Array<{ + category: FailureCategory; + test: (reason: string) => boolean; +}> = [ + { category: 'ttl_expired', test: (r) => r === 'ttl_expired' }, + { category: 'cap_eviction', test: (r) => r === 'pending_node_queue_overflow' }, + { + category: 'deliverability', + test: (r) => + r === 'run_not_externally_deliverable' || + r === 'target_task_terminal' || + r === 'subscription_no_longer_active' || + r === 'blocked_run_gate_not_opened', + }, + // Retry-exhaustion terminal failures carry the underlying activation/delivery + // reason (e.g. `node_execution_not_active`, `activation_failed; ...`). + { + category: 'retry_exhausted', + test: (r) => + r === 'node_execution_not_active' || + r === 'node_execution_pending' || + r.startsWith('activation_failed'), + }, + // Injection errors surface as a `deliveryMode:; ` reason. + { + category: 'injection_error', + test: (r) => r.startsWith('deliveryMode:'), + }, +]; + +/** + * Map a persisted terminal failure reason to an operator category. Reasons that + * do not match a known prefix fall back to `other`. + */ +export function categorizeFailureReason(reason: string): FailureCategory { + for (const { category, test } of FAILURE_CATEGORY_PREFIXES) { + if (test(reason)) return category; + } + return 'other'; +} + +const FAILURE_CATEGORIES: readonly FailureCategory[] = [ + 'ttl_expired', + 'cap_eviction', + 'deliverability', + 'retry_exhausted', + 'injection_error', + 'other', +]; + +function computeAgeStats(ages: readonly number[]): QueueAgeStats | null { + if (ages.length === 0) return null; + let min = Infinity; + let max = -Infinity; + let sum = 0; + for (const age of ages) { + if (age < min) min = age; + if (age > max) max = age; + sum += age; + } + const sorted = [...ages].sort((a, b) => a - b); + // Nearest-rank p95 (no interpolation): the ceil(0.95 * n)th value, clamped. + const p95Index = Math.min(sorted.length - 1, Math.ceil(sorted.length * 0.95) - 1); + return { + count: ages.length, + minMs: min, + maxMs: max, + avgMs: Math.round(sum / ages.length), + p95Ms: sorted[p95Index]!, + }; +} + +/** + * Lightweight in-memory counter store for pending external-event queue health. + * + * Not thread-safe in any concurrent sense — it relies on the single-threaded + * event loop of {@link SpaceRuntime}, which owns all queue mutations. + */ +export class ExternalEventQueueMetrics { + private readonly since: number; + private enqueue = 0; + private readonly enqueueBySource = new Map(); + private readonly enqueueByTargetState = new Map(); + private flushAttempts = 0; + private flushItemsDispatched = 0; + private delivered = 0; + private readonly finalFailuresByReason = new Map(); + private claimConflicts = 0; + private staleSessionSkips = 0; + + constructor(now: number = Date.now()) { + this.since = now; + } + + /** Record an enqueue into the pending queue, attributed to a source + target state. */ + recordEnqueue(source: string, targetState: string): void { + this.enqueue += 1; + this.enqueueBySource.set(source, (this.enqueueBySource.get(source) ?? 0) + 1); + this.enqueueByTargetState.set( + targetState, + (this.enqueueByTargetState.get(targetState) ?? 0) + 1 + ); + } + + /** Record a flush attempt and how many items it handed off to dispatch. */ + recordFlushAttempt(itemsDispatched: number): void { + this.flushAttempts += 1; + this.flushItemsDispatched += itemsDispatched; + } + + /** Record a non-terminal skip caused by the delivery already being in flight. */ + recordClaimConflict(): void { + this.claimConflicts += 1; + } + + /** Record a non-terminal skip caused by the target session not being live. */ + recordStaleSessionSkip(): void { + this.staleSessionSkips += 1; + } + + /** + * Record a delivery terminal transition. Called from the store's + * delivery-terminal hook — the single point that observes every `delivered` + * and terminal `failed` transition. + */ + recordDeliveryTerminal(event: DeliveryTerminalEvent): void { + if (event.outcome === 'delivered') { + this.delivered += 1; + return; + } + const reason = event.reason ?? 'unknown'; + this.finalFailuresByReason.set(reason, (this.finalFailuresByReason.get(reason) ?? 0) + 1); + } + + /** Read-only copy of the cumulative counters. */ + getCounters(): QueueHealthCounters { + return { + since: this.since, + enqueue: this.enqueue, + enqueueBySource: recordToObject(this.enqueueBySource), + enqueueByTargetState: recordToObject(this.enqueueByTargetState), + flushAttempts: this.flushAttempts, + flushItemsDispatched: this.flushItemsDispatched, + delivered: this.delivered, + finalFailuresByReason: recordToObject(this.finalFailuresByReason), + claimConflicts: this.claimConflicts, + staleSessionSkips: this.staleSessionSkips, + }; + } + + /** + * Build the operator snapshot by merging the cumulative counters with + * live gauges (computed by the caller from in-memory + DB state). + */ + snapshot(gauges: QueueHealthGauges, now: number = Date.now()): QueueHealthSnapshot { + const counters = this.getCounters(); + const failuresByCategory = {} as Record; + for (const category of FAILURE_CATEGORIES) failuresByCategory[category] = 0; + for (const [reason, count] of Object.entries(counters.finalFailuresByReason)) { + const category = categorizeFailureReason(reason); + failuresByCategory[category] += count; + } + return { collectedAt: now, counters, failuresByCategory, gauges }; + } +} + +/** Exposed for age-stat computation in {@link SpaceRuntime.getQueueHealthSnapshot}. */ +export { computeAgeStats as computeQueueAgeStats }; + +function recordToObject(map: Map): Record { + const obj: Record = {}; + for (const [key, value] of map) obj[key] = value; + return obj; +} diff --git a/packages/daemon/src/lib/rpc-handlers/index.ts b/packages/daemon/src/lib/rpc-handlers/index.ts index 4f8bd03a35..13de64b2ae 100644 --- a/packages/daemon/src/lib/rpc-handlers/index.ts +++ b/packages/daemon/src/lib/rpc-handlers/index.ts @@ -916,6 +916,12 @@ export function setupRPCHandlers(deps: RPCHandlerDependencies): RPCHandlerSetupR { longHorizonAgentRepo } ); + // Operator/debug view of pending external-event queue health (daemon-wide). + // Returns cumulative counters + live gauges; read-only, no side effects. + deps.messageHub.onRequest('space.externalEvents.queueHealth', async () => { + return spaceRuntimeService.getQueueHealthSnapshot(); + }); + // Space Worktree Manager — one worktree per task, shared by all node agents. const spaceWorktreeManager = new SpaceWorktreeManager(deps.db.getDatabase()); diff --git a/packages/daemon/src/lib/space/runtime/space-runtime-service.ts b/packages/daemon/src/lib/space/runtime/space-runtime-service.ts index 9780795839..debe7dea13 100644 --- a/packages/daemon/src/lib/space/runtime/space-runtime-service.ts +++ b/packages/daemon/src/lib/space/runtime/space-runtime-service.ts @@ -72,6 +72,10 @@ import { import { encodeActorIdComponent, longTermAgentSessionId } from '../long-term-agent-session'; import type { DaemonCommandMap, InternalCommandBus } from '../../internal-command-bus'; import type { ExternalEventStore } from '../../external-events/external-event-store'; +import { + type QueueHealthSnapshot, + ExternalEventQueueMetrics, +} from '../../external-events/queue-health-metrics'; import type { ExternalEventService } from '../../external-events/external-event-service'; import type { AgentMemoryRepository } from '../../../storage/repositories/agent-memory-repository'; import type { SDKUserMessage } from '@hyperneo/shared/sdk'; @@ -164,6 +168,13 @@ export interface SpaceRuntimeServiceConfig { internalEventBus?: InternalEventBus; commandBus?: InternalCommandBus; externalEventStore?: ExternalEventStore; + /** + * Optional queue-health metrics collector shared with the runtime. When + * provided, the service wires it to the store's delivery-terminal hook so + * terminal outcomes are counted from a single observation point. Defaults to + * a new in-memory instance. + */ + queueHealthMetrics?: ExternalEventQueueMetrics; /** External event publisher, available for runtime-owned direct publications if needed. */ externalEventService?: ExternalEventService; /** @@ -197,6 +208,12 @@ export interface SpaceRuntimeServiceConfig { export class SpaceRuntimeService { private readonly runtime: SpaceRuntime; + /** + * Queue-health metrics shared with the runtime and the store's + * delivery-terminal hook. Held on the service so the hook and the runtime + * observe the same counters. + */ + private readonly queueHealthMetrics: ExternalEventQueueMetrics; private started = false; /** Unsubscribe handles for InternalEventBus event subscriptions (daemon-lifetime). */ private readonly unsubscribers: Array<() => void> = []; @@ -257,9 +274,17 @@ export class SpaceRuntimeService { ? new SpaceActorRegistryAdapter(config.actorRegistryRepos) : null; this.auditLogRepo = new McpAuditLogRepository(this.config.db); + this.queueHealthMetrics = config.queueHealthMetrics ?? new ExternalEventQueueMetrics(); + // Observe every terminal delivery transition from a single point so + // delivered/failure-by-reason counters stay accurate regardless of which + // runtime call path reached the transition. + config.externalEventStore?.setDeliveryTerminalHook((event) => + this.queueHealthMetrics.recordDeliveryTerminal(event) + ); this.runtime = new SpaceRuntime({ ...config, nodeExecutionRepo: this.nodeExecutionRepo, + queueHealthMetrics: this.queueHealthMetrics, selectWorkflowWithLlm: config.selectWorkflowWithLlm ?? selectWorkflowWithLlmDefault, internalEventBus: config.internalEventBus, onTaskUpdated: async ({ spaceId, task, archiveSource }) => { @@ -1758,6 +1783,15 @@ export class SpaceRuntimeService { return this.runtime; } + /** + * Aggregate health snapshot for the pending external-event delivery queue. + * Daemon-wide (the runtime is a shared singleton handling all spaces). + * Surfaced to operators/debug views via the `space.externalEvents.queueHealth` RPC. + */ + getQueueHealthSnapshot(): QueueHealthSnapshot { + return this.runtime.getQueueHealthSnapshot(); + } + refreshLongHorizonAgentSubscriptions( spaceId: string, agentId: string diff --git a/packages/daemon/src/lib/space/runtime/space-runtime.ts b/packages/daemon/src/lib/space/runtime/space-runtime.ts index 8feefd20ad..0e46f68f5a 100644 --- a/packages/daemon/src/lib/space/runtime/space-runtime.ts +++ b/packages/daemon/src/lib/space/runtime/space-runtime.ts @@ -45,6 +45,12 @@ import { import type { ReactiveDatabase } from '../../../storage/reactive-database'; import type { ExternalEventPublishedPayload } from '../../external-events/external-event-service'; import type { ExternalEventStore } from '../../external-events/external-event-store'; +import { + type QueueHealthGauges, + type QueueHealthSnapshot, + ExternalEventQueueMetrics, + computeQueueAgeStats, +} from '../../external-events/queue-health-metrics'; import type { ExternalEvent } from '../../external-events/types'; import { KNOWN_SOURCES, validateGlobPattern } from '../../external-events/topic-validator'; import { ChannelCycleRepository } from '../../../storage/repositories/channel-cycle-repository'; @@ -192,6 +198,13 @@ export interface SpaceRuntimeConfig { commandBus?: InternalCommandBus; /** Persistent external-event delivery state store. */ externalEventStore?: ExternalEventStore; + /** + * Optional queue-health metrics collector for the pending external-event + * delivery queue. Defaults to a new in-memory instance. The same instance is + * wired to the store's delivery-terminal hook by the service so terminal + * outcomes are counted from a single observation point. + */ + queueHealthMetrics?: ExternalEventQueueMetrics; /** * Completion detector — inspects the canonical `SpaceTask` to decide whether * a workflow run is complete or ready for runtime resolution. @@ -1040,6 +1053,12 @@ export class SpaceRuntime { private readonly cancelledLongHorizonDeliveries = new Set(); private readonly longHorizonSubscriptionPatterns = new Map(); private readonly externalEventRateLimits = new Map(); + /** + * Pending external-event queue health counters. Defaults to a fresh + * in-memory instance; the service wires the store's delivery-terminal hook + * to the same instance (when shared) so terminal outcomes are counted once. + */ + private readonly queueHealthMetrics: ExternalEventQueueMetrics; private unsubscribeExternalEventPublished?: () => void; private unsubscribeSdkToolUseCreated?: () => void; private unsubscribeSdkToolUseConsumed?: () => void; @@ -1099,6 +1118,7 @@ export class SpaceRuntime { if (hasSqlExec(config.db)) { this.toolContinuationRepo.ensureSchema(); } + this.queueHealthMetrics = config.queueHealthMetrics ?? new ExternalEventQueueMetrics(); this.subscribeExternalEventPublished(); this.subscribeSdkToolUseCreated(); this.unsubscribeSpaceResumed = this.config.spaceManager.onSpaceResumedRegister?.((spaceId) => @@ -1830,6 +1850,7 @@ export class SpaceRuntime { if (!prepared) return; const { targetWithExecution, dispatchable } = prepared; + let dispatched = 0; for (const item of dispatchable) { // A concurrent cleanup may have marked this delivery terminal while the // batch was being prepared. Skip dispatch rather than injecting into a @@ -1856,6 +1877,7 @@ export class SpaceRuntime { continue; } this.clearExternalEventRetry(item.deliveryKey); + dispatched += 1; void this.enqueueDeliverableExternalEvent( targetWithExecution, item.event, @@ -1865,6 +1887,7 @@ export class SpaceRuntime { true ); } + this.queueHealthMetrics.recordFlushAttempt(dispatched); } private async flushPendingNodeQueueAsync( @@ -1884,6 +1907,7 @@ export class SpaceRuntime { dispatchable.sort((a, b) => a.createdAt - b.createdAt); } + let dispatched = 0; for (const item of dispatchable) { // A concurrent cleanup (unregisterExecution, subscription removal, or // run terminalization) may have marked this delivery terminal while the @@ -1908,6 +1932,7 @@ export class SpaceRuntime { continue; } this.clearExternalEventRetry(item.deliveryKey); + dispatched += 1; await this.enqueueDeliverableExternalEvent( targetWithExecution, item.event, @@ -1917,6 +1942,7 @@ export class SpaceRuntime { true ); } + this.queueHealthMetrics.recordFlushAttempt(dispatched); } private preparePendingNodeQueueDispatchable( @@ -2030,7 +2056,14 @@ export class SpaceRuntime { for (const { delivery, eventRecord } of deliveries) { if (store.isDeliveryTerminal(delivery.eventId, delivery.deliveryKey)) continue; - if (this.externalEventDeliveriesInFlight.has(delivery.deliveryKey)) continue; + if (this.externalEventDeliveriesInFlight.has(delivery.deliveryKey)) { + this.queueHealthMetrics.recordClaimConflict(); + log.debug('SpaceRuntime: external event delivery already in flight; skipped flush', { + runId: delivery.workflowRunId, + deliveryKey: delivery.deliveryKey, + }); + continue; + } if (!runDeliverable) { store.markDeliveryFailed(delivery.eventId, delivery.deliveryKey, { terminal: true, @@ -2038,6 +2071,10 @@ export class SpaceRuntime { }); store.markEventFailedIfAllDeliveriesTerminal(delivery.eventId); this.clearExternalEventRetry(delivery.deliveryKey); + log.debug('SpaceRuntime: external event delivery skipped — run not deliverable', { + runId: delivery.workflowRunId, + deliveryKey: delivery.deliveryKey, + }); continue; } if (this.isTargetTaskTerminal(target.taskId)) { @@ -2047,6 +2084,10 @@ export class SpaceRuntime { }); store.markEventFailedIfAllDeliveriesTerminal(delivery.eventId); this.clearExternalEventRetry(delivery.deliveryKey); + log.debug('SpaceRuntime: external event delivery skipped — target task terminal', { + runId: delivery.workflowRunId, + deliveryKey: delivery.deliveryKey, + }); continue; } if (!this.isTargetStillSubscribed(target, eventRecord.event.topic)) { @@ -2056,6 +2097,10 @@ export class SpaceRuntime { }); store.markEventFailedIfAllDeliveriesTerminal(delivery.eventId); this.clearExternalEventRetry(delivery.deliveryKey); + log.debug('SpaceRuntime: external event delivery skipped — subscription removed', { + runId: delivery.workflowRunId, + deliveryKey: delivery.deliveryKey, + }); continue; } @@ -2269,10 +2314,15 @@ export class SpaceRuntime { // since the delivery was registered is picked up. const resolved = this.resolveSubscriptionTarget(target); try { - if ( - store.isDeliveryTerminal(payload.eventId, deliveryKey) || - this.externalEventDeliveriesInFlight.has(deliveryKey) - ) { + if (store.isDeliveryTerminal(payload.eventId, deliveryKey)) { + return; + } + if (this.externalEventDeliveriesInFlight.has(deliveryKey)) { + this.queueHealthMetrics.recordClaimConflict(); + log.debug('SpaceRuntime: external event delivery already in flight; skipped dispatch', { + runId: resolved.workflowRunId, + deliveryKey, + }); return; } @@ -2714,8 +2764,10 @@ export class SpaceRuntime { : null; const spacePaused = !!(pausedRun && this.pausedSpaceIds.has(pausedRun.spaceId)); if (!target.sessionId || spacePaused) { + const sessionLoss = !target.sessionId; const reason = spacePaused ? 'space_paused' : 'session loss'; for (const item of items) { + if (sessionLoss) this.queueHealthMetrics.recordStaleSessionSkip(); store.markDeliveryFailed(item.event.eventId, item.deliveryKey, { terminal: false, reason: `deliveryMode:${item.deliveryMode}; digest requeued after ${reason}`, @@ -2732,6 +2784,12 @@ export class SpaceRuntime { // queue and will be re-claimed on the next flush. this.externalEventDeliveriesInFlight.delete(item.deliveryKey); } + if (sessionLoss) { + log.debug('SpaceRuntime: external event digest requeued — target session not live', { + runId: target.workflowRunId, + count: items.length, + }); + } return; } const deliveryKeys = items.map((item) => item.deliveryKey); @@ -3220,6 +3278,54 @@ export class SpaceRuntime { } queue.push({ event, deliveryKey, deliveryMode, createdAt }); this.pendingExternalEventQueue.set(key, queue); + this.queueHealthMetrics.recordEnqueue(event.source, this.describeEnqueueTargetState(target)); + } + + /** + * Describe the target's run + node-execution state at enqueue time, for the + * queue-health `enqueueByTargetState` breakdown. Surfaces which states force + * an event to be queued rather than delivered immediately (e.g. an + * `in_progress` run whose node is still `pending`, or a `blocked` run). + */ + private describeEnqueueTargetState(target: WorkflowSubscriptionTarget): string { + const run = this.config.workflowRunRepo.getRun(target.workflowRunId); + const nodeStatus = this.getCurrentQueueableOrActiveExecution(target)?.status ?? 'none'; + return `run=${run?.status ?? 'unknown'};node=${nodeStatus}`; + } + + /** + * Aggregate health snapshot for the pending external-event delivery queue, + * surfaced to operators/debug views. Merges cumulative counters (enqueue, + * flush, skips, delivered, failures by reason) with live gauges (depth, age, + * in-flight, digest backlog) computed from this runtime's in-memory state and + * the durable store. Counters are process-lifetime and reset on restart. + */ + getQueueHealthSnapshot(): QueueHealthSnapshot { + const now = Date.now(); + let queueDepth = 0; + const inMemoryAges: number[] = []; + for (const queue of this.pendingExternalEventQueue.values()) { + queueDepth += queue.length; + for (const item of queue) inMemoryAges.push(now - item.createdAt); + } + let digestBacklog = 0; + for (const state of this.externalEventRateLimits.values()) { + digestBacklog += state.pendingDigest.length; + } + const store = this.config.externalEventStore; + const persistedPending = store ? store.listPendingDeliveries().length : 0; + const persistedAges = store ? store.getPendingDeliveryAges(now) : []; + const gauges: QueueHealthGauges = { + queueDepth, + queueKeys: this.pendingExternalEventQueue.size, + inFlight: this.externalEventDeliveriesInFlight.size, + digestBacklog, + retryTimers: this.externalEventRetryTimers.size, + persistedPending, + queueAgeMs: computeQueueAgeStats(inMemoryAges), + persistedAgeMs: computeQueueAgeStats(persistedAges), + }; + return this.queueHealthMetrics.snapshot(gauges, now); } /** diff --git a/packages/daemon/tests/unit/4-space-storage/storage/external-event-store.test.ts b/packages/daemon/tests/unit/4-space-storage/storage/external-event-store.test.ts index 9cd5003da4..cb900ca2c6 100644 --- a/packages/daemon/tests/unit/4-space-storage/storage/external-event-store.test.ts +++ b/packages/daemon/tests/unit/4-space-storage/storage/external-event-store.test.ts @@ -886,3 +886,84 @@ describe('cross-event isolation', () => { expect(store.getDelivery('evt-b', 'dk-b')!.state).toBe('pending'); }); }); + +// Delivery-terminal hook + pending-age query (queue-health instrumentation) +describe('delivery-terminal hook', () => { + function registerPending(deliveryKey = 'dk-1'): void { + store.store(EVENT_A); + store.registerExpectedDelivery('evt-a', deliveryKey, { + workflowRunId: 'run-1', + taskId: 'task-1', + nodeId: 'node-1', + agentName: 'coder', + }); + } + + test('fires delivered on a real terminal transition', () => { + const events: Array<{ outcome: string; reason: string | null }> = []; + store.setDeliveryTerminalHook((event) => + events.push({ outcome: event.outcome, reason: event.reason }) + ); + registerPending(); + store.markDeliveryDelivered('evt-a', 'dk-1'); + + expect(events).toEqual([{ outcome: 'delivered', reason: null }]); + }); + + test('fires failed only for terminal failures, with reason', () => { + const events: Array<{ outcome: string; reason: string | null }> = []; + store.setDeliveryTerminalHook((event) => + events.push({ outcome: event.outcome, reason: event.reason }) + ); + registerPending(); + + // Non-terminal (retryable) failure must NOT fire the hook. + store.markDeliveryFailed('evt-a', 'dk-1', { + terminal: false, + reason: 'node_execution_not_active', + }); + expect(events).toHaveLength(0); + + // Terminal failure fires with the reason. + store.markDeliveryFailed('evt-a', 'dk-1', { terminal: true, reason: 'ttl_expired' }); + expect(events).toEqual([{ outcome: 'failed', reason: 'ttl_expired' }]); + }); + + test('does not fire when the row is already terminal (no double-count)', () => { + const events: string[] = []; + store.setDeliveryTerminalHook((event) => events.push(event.outcome)); + registerPending(); + + store.markDeliveryDelivered('evt-a', 'dk-1'); + store.markDeliveryDelivered('evt-a', 'dk-1'); // no-op, already delivered + store.markDeliveryFailed('evt-a', 'dk-1', { terminal: true, reason: 'late' }); // no-op + + expect(events).toEqual(['delivered']); + }); +}); + +describe('getPendingDeliveryAges', () => { + test('returns event-age for pending deliveries, empty when none', () => { + const now = Date.now(); + expect(store.getPendingDeliveryAges(now)).toEqual([]); + + store.store(EVENT_A); + store.registerExpectedDelivery('evt-a', 'dk-1', { + workflowRunId: 'run-1', + taskId: 'task-1', + nodeId: 'node-1', + agentName: 'coder', + }); + // created_at is ingestion time (the TTL anchor), set by store() to the + // current time. Backdate it 60s so the age is deterministic. + db.prepare(`UPDATE space_external_events SET created_at = ? WHERE id = ?`).run( + now - 60_000, + 'evt-a' + ); + + const ages = store.getPendingDeliveryAges(now); + expect(ages).toHaveLength(1); + expect(ages[0]).toBeGreaterThanOrEqual(59_000); + expect(ages[0]).toBeLessThanOrEqual(61_000); + }); +}); diff --git a/packages/daemon/tests/unit/4-space-storage/storage/queue-health-metrics.test.ts b/packages/daemon/tests/unit/4-space-storage/storage/queue-health-metrics.test.ts new file mode 100644 index 0000000000..ee775cc33f --- /dev/null +++ b/packages/daemon/tests/unit/4-space-storage/storage/queue-health-metrics.test.ts @@ -0,0 +1,188 @@ +/** + * ExternalEventQueueMetrics unit tests. + * + * Pure in-memory counter behavior — no DB. Covers enqueue attribution, flush, + * skip counters, terminal-outcome recording, and failure categorization in the + * snapshot. + */ + +import { describe, expect, test } from 'bun:test'; +import { + ExternalEventQueueMetrics, + categorizeFailureReason, + computeQueueAgeStats, + type QueueHealthGauges, +} from '../../../../src/lib/external-events/queue-health-metrics'; + +const EMPTY_GAUGES: QueueHealthGauges = { + queueDepth: 0, + queueKeys: 0, + inFlight: 0, + digestBacklog: 0, + retryTimers: 0, + persistedPending: 0, + queueAgeMs: null, + persistedAgeMs: null, +}; + +describe('ExternalEventQueueMetrics — enqueue attribution', () => { + test('counts total enqueues and breaks them down by source + target state', () => { + const metrics = new ExternalEventQueueMetrics(1000); + metrics.recordEnqueue('github', 'run=in_progress;node=pending'); + metrics.recordEnqueue('github', 'run=in_progress;node=pending'); + metrics.recordEnqueue('slack', 'run=blocked;node=in_progress'); + + const counters = metrics.getCounters(); + expect(counters.enqueue).toBe(3); + expect(counters.enqueueBySource).toEqual({ github: 2, slack: 1 }); + expect(counters.enqueueByTargetState).toEqual({ + 'run=in_progress;node=pending': 2, + 'run=blocked;node=in_progress': 1, + }); + }); +}); + +describe('ExternalEventQueueMetrics — flush + skip counters', () => { + test('accumulates flush attempts and dispatched items', () => { + const metrics = new ExternalEventQueueMetrics(1000); + metrics.recordFlushAttempt(3); + metrics.recordFlushAttempt(0); + metrics.recordFlushAttempt(2); + + const counters = metrics.getCounters(); + expect(counters.flushAttempts).toBe(3); + expect(counters.flushItemsDispatched).toBe(5); + }); + + test('counts claim conflicts and stale-session skips separately', () => { + const metrics = new ExternalEventQueueMetrics(1000); + metrics.recordClaimConflict(); + metrics.recordClaimConflict(); + metrics.recordStaleSessionSkip(); + + const counters = metrics.getCounters(); + expect(counters.claimConflicts).toBe(2); + expect(counters.staleSessionSkips).toBe(1); + }); +}); + +describe('ExternalEventQueueMetrics — terminal outcomes', () => { + test('counts delivered and terminal failures by reason', () => { + const metrics = new ExternalEventQueueMetrics(1000); + metrics.recordDeliveryTerminal({ + eventId: 'e1', + deliveryKey: 'd1', + outcome: 'delivered', + reason: null, + }); + metrics.recordDeliveryTerminal({ + eventId: 'e2', + deliveryKey: 'd2', + outcome: 'delivered', + reason: null, + }); + metrics.recordDeliveryTerminal({ + eventId: 'e3', + deliveryKey: 'd3', + outcome: 'failed', + reason: 'ttl_expired', + }); + metrics.recordDeliveryTerminal({ + eventId: 'e4', + deliveryKey: 'd4', + outcome: 'failed', + reason: 'run_not_externally_deliverable', + }); + + const counters = metrics.getCounters(); + expect(counters.delivered).toBe(2); + expect(counters.finalFailuresByReason).toEqual({ + ttl_expired: 1, + run_not_externally_deliverable: 1, + }); + }); +}); + +describe('ExternalEventQueueMetrics — snapshot', () => { + test('merges counters with gauges and categorizes failures', () => { + const metrics = new ExternalEventQueueMetrics(1000); + metrics.recordEnqueue('github', 'run=in_progress;node=pending'); + metrics.recordDeliveryTerminal({ + eventId: 'e1', + deliveryKey: 'd1', + outcome: 'delivered', + reason: null, + }); + metrics.recordDeliveryTerminal({ + eventId: 'e2', + deliveryKey: 'd2', + outcome: 'failed', + reason: 'pending_node_queue_overflow', + }); + metrics.recordDeliveryTerminal({ + eventId: 'e3', + deliveryKey: 'd3', + outcome: 'failed', + reason: 'subscription_no_longer_active', + }); + metrics.recordDeliveryTerminal({ + eventId: 'e4', + deliveryKey: 'd4', + outcome: 'failed', + reason: 'deliveryMode:immediate; boom', + }); + + const snapshot = metrics.snapshot({ ...EMPTY_GAUGES, queueDepth: 4, inFlight: 1 }, 5000); + + expect(snapshot.collectedAt).toBe(5000); + expect(snapshot.counters.since).toBe(1000); + expect(snapshot.counters.enqueue).toBe(1); + expect(snapshot.counters.delivered).toBe(1); + expect(snapshot.gauges.queueDepth).toBe(4); + expect(snapshot.gauges.inFlight).toBe(1); + // Each terminal failure lands in exactly one category. + expect(snapshot.failuresByCategory).toEqual({ + ttl_expired: 0, + cap_eviction: 1, + deliverability: 1, + retry_exhausted: 0, + injection_error: 1, + other: 0, + }); + }); +}); + +describe('categorizeFailureReason', () => { + test.each([ + ['ttl_expired', 'ttl_expired'], + ['pending_node_queue_overflow', 'cap_eviction'], + ['run_not_externally_deliverable', 'deliverability'], + ['target_task_terminal', 'deliverability'], + ['subscription_no_longer_active', 'deliverability'], + ['blocked_run_gate_not_opened', 'deliverability'], + ['node_execution_not_active', 'retry_exhausted'], + ['node_execution_pending', 'retry_exhausted'], + ['activation_failed; timeout', 'retry_exhausted'], + ['deliveryMode:immediate; inject failed', 'injection_error'], + ['something_unexpected', 'other'], + ])('categorizes %s -> %s', (reason, expected) => { + expect(categorizeFailureReason(reason)).toBe(expected); + }); +}); + +describe('computeQueueAgeStats', () => { + test('returns null for an empty set', () => { + expect(computeQueueAgeStats([])).toBeNull(); + }); + + test('computes min/max/avg/p95 for a small set', () => { + const stats = computeQueueAgeStats([100, 200, 300, 400, 500]); + expect(stats).not.toBeNull(); + expect(stats!.count).toBe(5); + expect(stats!.minMs).toBe(100); + expect(stats!.maxMs).toBe(500); + expect(stats!.avgMs).toBe(300); + // p95 nearest-rank of 5 values -> ceil(0.95*5)-1 = 4th index (500). + expect(stats!.p95Ms).toBe(500); + }); +}); diff --git a/packages/daemon/tests/unit/5-space/runtime/space-runtime-external-events.test.ts b/packages/daemon/tests/unit/5-space/runtime/space-runtime-external-events.test.ts index 2cbca2fb19..c5f9e21188 100644 --- a/packages/daemon/tests/unit/5-space/runtime/space-runtime-external-events.test.ts +++ b/packages/daemon/tests/unit/5-space/runtime/space-runtime-external-events.test.ts @@ -1,8 +1,9 @@ -import { beforeEach, describe, expect, setDefaultTimeout, test } from 'bun:test'; +import { afterEach, beforeEach, describe, expect, setDefaultTimeout, test } from 'bun:test'; import { Database } from 'bun:sqlite'; import type { SpaceTask, SpaceWorkflow } from '@hyperneo/shared'; import { ExternalEventService } from '../../../../src/lib/external-events/external-event-service'; import { ExternalEventStore } from '../../../../src/lib/external-events/external-event-store'; +import { ExternalEventQueueMetrics } from '../../../../src/lib/external-events/queue-health-metrics'; import type { ExternalEvent } from '../../../../src/lib/external-events/types'; import { createInternalCommandBus } from '../../../../src/lib/internal-command-bus'; import { createDaemonInternalEventBus } from '../../../../src/lib/internal-event-bus'; @@ -7014,3 +7015,155 @@ describe('SpaceRuntime event-driven gate evaluation', () => { expect(injected).toHaveLength(1); }); }); + +describe('SpaceRuntime queue-health snapshot', () => { + let db: Database; + let workflowRunRepo: SpaceWorkflowRunRepository; + let taskRepo: SpaceTaskRepository; + let nodeExecutionRepo: NodeExecutionRepository; + let workflowManager: SpaceWorkflowManager; + let runtime: SpaceRuntime; + let eventStore: ExternalEventStore; + let queueHealthMetrics: ExternalEventQueueMetrics; + let eventService: ExternalEventService; + let injected: Array<{ sessionId: string; message: string; deliveryMode?: string }>; + let tam: MockTaskAgentManager; + let bus: ReturnType; + + beforeEach(() => { + db = makeDb(); + workflowRunRepo = new SpaceWorkflowRunRepository(db); + taskRepo = new SpaceTaskRepository(db); + nodeExecutionRepo = new NodeExecutionRepository(db); + workflowManager = new SpaceWorkflowManager(new SpaceWorkflowRepository(db)); + bus = createDaemonInternalEventBus(); + const commandBus = createInternalCommandBus(); + eventStore = new ExternalEventStore(db); + queueHealthMetrics = new ExternalEventQueueMetrics(); + // Mirror SpaceRuntimeService wiring: observe terminal transitions from a + // single point so delivered/failure counters stay accurate. + eventStore.setDeliveryTerminalHook((event) => queueHealthMetrics.recordDeliveryTerminal(event)); + eventService = new ExternalEventService(eventStore, bus); + injected = []; + commandBus.register('agent.message.inject', async (command) => { + injected.push({ + sessionId: command.sessionId, + message: command.message, + deliveryMode: command.deliveryMode, + }); + return { ok: true }; + }); + tam = new MockTaskAgentManager(); + runtime = new SpaceRuntime({ + db, + spaceManager: new SpaceManager(db), + spaceAgentManager: new SpaceAgentManager(new SpaceAgentRepository(db)), + spaceWorkflowManager: workflowManager, + workflowRunRepo, + taskRepo, + nodeExecutionRepo, + internalEventBus: bus, + commandBus, + externalEventStore: eventStore, + queueHealthMetrics, + taskAgentManager: tam as never, + }); + }); + + afterEach(() => { + void runtime.stop(); + }); + + function createWorkflow(): SpaceWorkflow { + return workflowManager.createWorkflow({ + spaceId: SPACE_ID, + name: `Workflow ${Math.random()}`, + description: '', + nodes: [ + { + id: 'code', + name: 'Code', + agents: [{ agentId: AGENT_ID, name: 'coder' }], + }, + ], + transitions: [], + startNodeId: 'code', + rules: [], + tags: [], + }); + } + + async function startRun(): Promise<{ runId: string; taskId: string }> { + const workflow = createWorkflow(); + const { run, tasks } = await runtime.startWorkflowRun(SPACE_ID, workflow.id, 'Run'); + const task = tasks[0]!; + runtime.registerSubscription(run.id, task.id, 'code', 'coder', DEFAULT_TOPIC); + return { runId: run.id, taskId: task.id }; + } + + test('reports zero gauges and counters before any event', () => { + const snapshot = runtime.getQueueHealthSnapshot(); + expect(snapshot.gauges.queueDepth).toBe(0); + expect(snapshot.gauges.queueKeys).toBe(0); + expect(snapshot.gauges.inFlight).toBe(0); + expect(snapshot.gauges.persistedPending).toBe(0); + expect(snapshot.counters.enqueue).toBe(0); + expect(snapshot.counters.delivered).toBe(0); + expect(snapshot.counters.flushAttempts).toBe(0); + expect(snapshot.gauges.queueAgeMs).toBeNull(); + }); + + test('enqueues pending events and reports depth, source, target state, and age', async () => { + const { runId } = await startRun(); + // Node execution starts `pending` with no session, so the event is queued. + expect(nodeExecutionRepo.listByNode(runId, 'code')[0]!.status).toBe('pending'); + await eventService.publish(makeEvent()); + + const snapshot = runtime.getQueueHealthSnapshot(); + expect(snapshot.counters.enqueue).toBe(1); + expect(snapshot.counters.enqueueBySource).toEqual({ github: 1 }); + expect(snapshot.counters.enqueueByTargetState).toEqual({ + 'run=in_progress;node=pending': 1, + }); + expect(snapshot.gauges.queueDepth).toBe(1); + expect(snapshot.gauges.queueKeys).toBe(1); + expect(snapshot.gauges.queueAgeMs).not.toBeNull(); + expect(snapshot.gauges.queueAgeMs!.count).toBe(1); + }); + + test('counts cap-eviction terminal failures when a target queue overflows', async () => { + await startRun(); + // 50 items fill the per-target queue; the 51st evicts the oldest, which the + // store hook records as a terminal pending_node_queue_overflow failure. + for (let i = 0; i < 51; i++) { + await eventService.publish(makeEvent()); + } + + const snapshot = runtime.getQueueHealthSnapshot(); + expect(snapshot.counters.enqueue).toBe(51); + expect(snapshot.counters.finalFailuresByReason['pending_node_queue_overflow']).toBe(1); + expect(snapshot.failuresByCategory.cap_eviction).toBe(1); + // Queue is capped at 50. + expect(snapshot.gauges.queueDepth).toBe(50); + }); + + test('counts delivered after a successful injection into a live session', async () => { + const { runId } = await startRun(); + const execution = nodeExecutionRepo.listByNode(runId, 'code')[0]!; + nodeExecutionRepo.update(execution.id, { + status: 'in_progress', + agentSessionId: 'session-live', + startedAt: Date.now(), + }); + tam.alive.add('session-live'); + + await eventService.publish(makeEvent()); + + expect(injected).toHaveLength(1); + const snapshot = runtime.getQueueHealthSnapshot(); + expect(snapshot.counters.delivered).toBe(1); + expect(snapshot.counters.enqueue).toBe(0); + expect(snapshot.gauges.queueDepth).toBe(0); + expect(snapshot.counters.flushAttempts).toBeGreaterThanOrEqual(1); + }); +}); From 608e9dd745544a6877789b67b41ba2309dfa40dc Mon Sep 17 00:00:00 2001 From: Marc Liu Date: Wed, 22 Jul 2026 05:03:38 -0400 Subject: [PATCH 2/3] feat(web): add pending external-event queue-health operator view Add a QueueHealthSummary panel to the external-events settings that fetches the daemon-wide queue-health snapshot via the new space.externalEvents.queueHealth RPC and renders cumulative counters (enqueue by source/target state, flush, delivered, failures by reason/category, skips) and live gauges (depth, age p95, in-flight, digest backlog). Add typed store helper + QueueHealthSnapshot types. --- .../components/space/QueueHealthSummary.tsx | 219 ++++++++++++++++++ .../space/SpaceExternalEventsSettings.tsx | 3 + .../__tests__/QueueHealthSummary.test.tsx | 101 ++++++++ .../SpaceExternalEventsSettings.test.tsx | 6 + packages/web/src/lib/space-store.ts | 62 +++++ 5 files changed, 391 insertions(+) create mode 100644 packages/web/src/components/space/QueueHealthSummary.tsx create mode 100644 packages/web/src/components/space/__tests__/QueueHealthSummary.test.tsx diff --git a/packages/web/src/components/space/QueueHealthSummary.tsx b/packages/web/src/components/space/QueueHealthSummary.tsx new file mode 100644 index 0000000000..ebaa1f5944 --- /dev/null +++ b/packages/web/src/components/space/QueueHealthSummary.tsx @@ -0,0 +1,219 @@ +/** + * QueueHealthSummary — operator/debug view of pending external-event queue health. + * + * Daemon-wide aggregate (the runtime is a shared singleton). Shows cumulative + * counters (enqueue by source/target state, flush, delivered, failures by + * reason/category, skips) and live gauges (depth, age, in-flight, digest + * backlog). Counters are process-lifetime and reset on daemon restart. + */ + +import { useCallback, useEffect, useState } from 'preact/hooks'; +import { spaceStore, type QueueAgeStats, type QueueHealthSnapshot } from '../../lib/space-store.ts'; +import { Button } from '../ui/Button.tsx'; + +function formatAge(ms: number): string { + if (ms < 1000) return `${Math.round(ms)} ms`; + if (ms < 60_000) return `${(ms / 1000).toFixed(1)} s`; + return `${(ms / 60_000).toFixed(1)} m`; +} + +function formatRelative(epochMs: number): string { + const seconds = Math.max(0, Math.round((Date.now() - epochMs) / 1000)); + if (seconds < 60) return `${seconds}s ago`; + const minutes = Math.floor(seconds / 60); + if (minutes < 60) return `${minutes}m ago`; + const hours = Math.floor(minutes / 60); + if (hours < 24) return `${hours}h ago`; + return `${Math.floor(hours / 24)}d ago`; +} + +function ageStats(stats: QueueAgeStats | null): string { + if (!stats) return '—'; + return `p95 ${formatAge(stats.p95Ms)} · max ${formatAge(stats.maxMs)} · ${stats.count} item${stats.count === 1 ? '' : 's'}`; +} + +function entryList(record: Record): Array<{ key: string; value: number }> { + return Object.entries(record) + .filter(([, value]) => value > 0) + .map(([key, value]) => ({ key, value })) + .sort((a, b) => b.value - a.value); +} + +function Metric({ + label, + value, + hint, +}: { + label: string; + value: string; + hint?: string; +}): preact.JSX.Element { + return ( +
+
{label}
+
{value}
+ {hint ?
{hint}
: null} +
+ ); +} + +function Breakdown({ + title, + entries, + emptyHint, +}: { + title: string; + entries: Array<{ key: string; value: number }>; + emptyHint: string; +}): preact.JSX.Element { + return ( +
+
{title}
+ {entries.length === 0 ? ( +
{emptyHint}
+ ) : ( +
    + {entries.map((entry) => ( +
  • + + {entry.key} + + {entry.value} +
  • + ))} +
+ )} +
+ ); +} + +export function QueueHealthSummary(): preact.JSX.Element { + const [snapshot, setSnapshot] = useState(null); + const [loading, setLoading] = useState(false); + const [error, setError] = useState(null); + + const refresh = useCallback(async () => { + setLoading(true); + setError(null); + try { + const result = await spaceStore.getExternalEventQueueHealth(); + setSnapshot(result); + } catch (err) { + setError(err instanceof Error ? err.message : String(err)); + } finally { + setLoading(false); + } + }, []); + + useEffect(() => { + void refresh(); + }, [refresh]); + + const counters = snapshot?.counters; + const gauges = snapshot?.gauges; + const totalFailures = counters + ? Object.values(counters.finalFailuresByReason).reduce((sum, value) => sum + value, 0) + : 0; + const settled = counters ? counters.delivered + totalFailures : 0; + const successRate = + settled > 0 && counters ? Math.round((counters.delivered / settled) * 100) : null; + + return ( +
+
+
+
Queue health
+

+ Daemon-wide pending external-event delivery queue.{' '} + {snapshot ? `Counting since ${formatRelative(snapshot.counters.since)}` : ''} + {snapshot ? ` · updated ${formatRelative(snapshot.collectedAt)}.` : ''} +

+
+ +
+ + {error ? ( +
+ Failed to load queue health: {error} +
+ ) : null} + + {!snapshot && !error ? ( +
+ {loading ? 'Loading…' : 'No data yet. Click Refresh.'} +
+ ) : null} + + {snapshot && counters && gauges ? ( +
+
+ + + + 0 ? formatAge(gauges.queueAgeMs?.p95Ms ?? 0) : '—'} + hint={ageStats(gauges.queueAgeMs)} + /> + + + + +
+ +
+ + + + +
+ +
+ Persisted pending age: {ageStats(gauges.persistedAgeMs)} +
+
+ ) : null} +
+ ); +} diff --git a/packages/web/src/components/space/SpaceExternalEventsSettings.tsx b/packages/web/src/components/space/SpaceExternalEventsSettings.tsx index 1bcd07f1b1..4001f0bc5f 100644 --- a/packages/web/src/components/space/SpaceExternalEventsSettings.tsx +++ b/packages/web/src/components/space/SpaceExternalEventsSettings.tsx @@ -14,6 +14,7 @@ import { cn } from '../../lib/utils.ts'; import { Button } from '../ui/Button.tsx'; import { CopyButton } from '../ui/CopyButton.tsx'; import { Spinner } from '../ui/Spinner.tsx'; +import { QueueHealthSummary } from './QueueHealthSummary.tsx'; interface SpaceExternalEventsSettingsProps { spaceId: string; @@ -805,6 +806,8 @@ export function SpaceExternalEventsSettings({ onRefresh={refreshDeliveries} onSelect={setSelectedDelivery} /> + + )} diff --git a/packages/web/src/components/space/__tests__/QueueHealthSummary.test.tsx b/packages/web/src/components/space/__tests__/QueueHealthSummary.test.tsx new file mode 100644 index 0000000000..52f4e62bd0 --- /dev/null +++ b/packages/web/src/components/space/__tests__/QueueHealthSummary.test.tsx @@ -0,0 +1,101 @@ +// @ts-nocheck +import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'; +import { render, fireEvent, waitFor, cleanup } from '@testing-library/preact'; + +const mockGetExternalEventQueueHealth = vi.fn(); + +vi.mock('../../../lib/space-store', () => ({ + spaceStore: { + get getExternalEventQueueHealth() { + return mockGetExternalEventQueueHealth; + }, + }, +})); + +vi.mock('../../ui/Button', () => ({ + Button: ({ children, onClick, type, loading }) => ( + + ), +})); + +import { QueueHealthSummary } from '../QueueHealthSummary'; + +const sampleSnapshot = { + collectedAt: Date.now(), + counters: { + since: Date.now() - 60_000, + enqueue: 5, + enqueueBySource: { github: 4, slack: 1 }, + enqueueByTargetState: { 'run=in_progress;node=pending': 5 }, + flushAttempts: 3, + flushItemsDispatched: 4, + delivered: 2, + finalFailuresByReason: { ttl_expired: 1, pending_node_queue_overflow: 1 }, + claimConflicts: 1, + staleSessionSkips: 2, + }, + failuresByCategory: { + ttl_expired: 1, + cap_eviction: 1, + deliverability: 0, + retry_exhausted: 0, + injection_error: 0, + other: 0, + }, + gauges: { + queueDepth: 2, + queueKeys: 1, + inFlight: 1, + digestBacklog: 0, + retryTimers: 1, + persistedPending: 0, + queueAgeMs: { count: 2, minMs: 1000, maxMs: 5000, avgMs: 3000, p95Ms: 5000 }, + persistedAgeMs: null, + }, +}; + +describe('QueueHealthSummary', () => { + beforeEach(() => { + cleanup(); + mockGetExternalEventQueueHealth.mockReset(); + }); + + afterEach(() => cleanup()); + + it('renders counters, gauges, and breakdowns from the snapshot', async () => { + mockGetExternalEventQueueHealth.mockResolvedValue(sampleSnapshot); + const { getByText, getAllByText, findByText, getByTestId } = render(); + + await findByText('Queue health'); + expect(getByTestId('queue-health-summary')).toBeTruthy(); + // Breakdown entries render the source keys and failure reasons. + await waitFor(() => { + expect(getByText('github')).toBeTruthy(); + expect(getByText('slack')).toBeTruthy(); + // ttl_expired appears in both the category and reason breakdowns. + expect(getAllByText('ttl_expired')).toHaveLength(2); + expect(getByText('pending_node_queue_overflow')).toBeTruthy(); + }); + expect(mockGetExternalEventQueueHealth).toHaveBeenCalledTimes(1); + }); + + it('shows an error banner when the fetch rejects', async () => { + mockGetExternalEventQueueHealth.mockRejectedValue(new Error('boom')); + const { findByText } = render(); + + expect(await findByText(/Failed to load queue health: boom/)).toBeTruthy(); + }); + + it('re-fetches the snapshot when Refresh is clicked', async () => { + mockGetExternalEventQueueHealth.mockResolvedValue(sampleSnapshot); + const { findByText, getByText: getByTextSync } = render(); + + await findByText('Queue health'); + fireEvent.click(getByTextSync('Refresh')); + await waitFor(() => { + expect(mockGetExternalEventQueueHealth).toHaveBeenCalledTimes(2); + }); + }); +}); diff --git a/packages/web/src/components/space/__tests__/SpaceExternalEventsSettings.test.tsx b/packages/web/src/components/space/__tests__/SpaceExternalEventsSettings.test.tsx index f5b66b4ad9..a9d3f1d1b9 100644 --- a/packages/web/src/components/space/__tests__/SpaceExternalEventsSettings.test.tsx +++ b/packages/web/src/components/space/__tests__/SpaceExternalEventsSettings.test.tsx @@ -7,6 +7,7 @@ const mockGetHubIfConnected = vi.fn(); const mockToastSuccess = vi.fn(); const mockToastError = vi.fn(); const mockListExternalEventDeliveries = vi.fn(); +const mockGetExternalEventQueueHealth = vi.fn(); vi.mock('../../../lib/connection-manager', () => ({ connectionManager: { @@ -32,6 +33,9 @@ vi.mock('../../../lib/space-store', () => ({ get listExternalEventDeliveries() { return mockListExternalEventDeliveries; }, + get getExternalEventQueueHealth() { + return mockGetExternalEventQueueHealth; + }, }, })); @@ -207,6 +211,8 @@ describe('SpaceExternalEventsSettings', () => { mockToastError.mockReset(); mockListExternalEventDeliveries.mockReset(); mockListExternalEventDeliveries.mockResolvedValue([]); + mockGetExternalEventQueueHealth.mockReset(); + mockGetExternalEventQueueHealth.mockResolvedValue(null); }); afterEach(() => cleanup()); diff --git a/packages/web/src/lib/space-store.ts b/packages/web/src/lib/space-store.ts index a1660db173..f4e7e7f11d 100644 --- a/packages/web/src/lib/space-store.ts +++ b/packages/web/src/lib/space-store.ts @@ -89,6 +89,57 @@ export interface SpaceWithTasks extends Space { export type ExternalEventDeliveryStatus = 'pending' | 'delivered' | 'failed'; +/** Min/max/avg/p95 of a set of millisecond ages. */ +export interface QueueAgeStats { + count: number; + minMs: number; + maxMs: number; + avgMs: number; + p95Ms: number; +} + +/** Cumulative pending-external-event queue counters (process-lifetime). */ +export interface QueueHealthCounters { + since: number; + enqueue: number; + enqueueBySource: Record; + enqueueByTargetState: Record; + flushAttempts: number; + flushItemsDispatched: number; + delivered: number; + finalFailuresByReason: Record; + claimConflicts: number; + staleSessionSkips: number; +} + +/** Live queue gauges computed at read time. */ +export interface QueueHealthGauges { + queueDepth: number; + queueKeys: number; + inFlight: number; + digestBacklog: number; + retryTimers: number; + persistedPending: number; + queueAgeMs: QueueAgeStats | null; + persistedAgeMs: QueueAgeStats | null; +} + +export type QueueHealthFailureCategory = + | 'ttl_expired' + | 'cap_eviction' + | 'deliverability' + | 'retry_exhausted' + | 'injection_error' + | 'other'; + +/** Daemon-wide aggregate queue-health snapshot. */ +export interface QueueHealthSnapshot { + collectedAt: number; + counters: QueueHealthCounters; + failuresByCategory: Record; + gauges: QueueHealthGauges; +} + export interface SpaceExternalEventDeliveryLogRecord { eventId: string; deliveryKey: string; @@ -2456,6 +2507,17 @@ class SpaceStore { return result?.deliveries ?? []; } + /** + * Fetch the daemon-wide pending external-event queue-health snapshot + * (counters + live gauges). Not space-scoped — the runtime is shared. + */ + async getExternalEventQueueHealth(): Promise { + const hub = connectionManager.getHubIfConnected(); + if (!hub) throw new Error('Not connected'); + const result = await hub.request('space.externalEvents.queueHealth', {}); + return result ?? null; + } + /** * List all artifacts for a workflow run. */ From 521ae4a0689f61f67f914e7aa4ddcb10ac1a4bff Mon Sep 17 00:00:00 2001 From: Marc Liu Date: Sun, 26 Jul 2026 15:01:16 -0400 Subject: [PATCH 3/3] fix: address queue-health review (perf, skips, category mapping) Address PR #2257 review feedback: - Compute persisted-pending health via SQL aggregates (summarizePendingDeliveries: COUNT/MIN/MAX/AVG + LIMIT/OFFSET p95) instead of materializing every pending row + a full in-memory sort on each snapshot read. - Count paused-space deferrals as a separate pausedSpaceSkips counter (recorded at the digest-requeue, fresh-delivery, and retry-injection pause sites), so staleSessionSkips stays a true session-loss signal and the pause-induced backlog scenario is no longer reported as zero skips. - Instrument preservePendingDigestItem with recordEnqueue so rehoused digest items keep the cumulative enqueue counter consistent with queueDepth. - Fix failure-category mapping: drop unreachable blocked_run_gate_not_opened (event-level markEventFailed never fires the delivery hook) and route auto_pr_subscription_cleared / node_execution_cancelled / run_interests_rebuilt / run_terminal_cleanup to deliverability. --- .../external-events/external-event-store.ts | 57 ++++++++++++++----- .../external-events/queue-health-metrics.ts | 33 +++++++++-- .../src/lib/space/runtime/space-runtime.ts | 23 ++++++-- .../storage/external-event-store.test.ts | 49 ++++++++++++---- .../storage/queue-health-metrics.test.ts | 15 ++++- .../components/space/QueueHealthSummary.tsx | 4 +- .../__tests__/QueueHealthSummary.test.tsx | 1 + packages/web/src/lib/space-store.ts | 1 + 8 files changed, 146 insertions(+), 37 deletions(-) diff --git a/packages/daemon/src/lib/external-events/external-event-store.ts b/packages/daemon/src/lib/external-events/external-event-store.ts index 341dc25d75..c7b631b1f5 100644 --- a/packages/daemon/src/lib/external-events/external-event-store.ts +++ b/packages/daemon/src/lib/external-events/external-event-store.ts @@ -31,7 +31,7 @@ import { TERMINAL_DELIVERY_STATES, TERMINAL_EVENT_STATES, } from './types'; -import type { DeliveryTerminalEvent } from './queue-health-metrics'; +import type { DeliveryTerminalEvent, QueueAgeStats } from './queue-health-metrics'; import { validateLiteralTopic, validateSource } from './topic-validator'; interface ExternalEventRow { @@ -582,22 +582,53 @@ export class ExternalEventStore { } /** - * Return the event-age (`now - event.created_at`, in ms) of every - * DB-persisted `pending` delivery row, in a single joined query. Used by the - * queue-health snapshot's persisted-pending age gauge. The anchor is the - * source event's ingestion time (`space_external_events.created_at`), matching - * the runtime's event-age TTL semantics — see `EXTERNAL_EVENT_QUEUE_TTL_MS`. + * Summarize DB-persisted `pending` deliveries for the queue-health snapshot + * without materializing every row: count + min/max/avg via SQL aggregates, and + * p95 via a single `LIMIT 1 OFFSET k` lookup over a sorted scan (no per-row JS + * allocation). The anchor is the source event's ingestion time + * (`space_external_events.created_at`), matching the runtime's event-age TTL + * semantics — see `EXTERNAL_EVENT_QUEUE_TTL_MS`. Returns `null` when there are + * no pending deliveries. */ - getPendingDeliveryAges(now: number = Date.now()): number[] { - const rows = this.db + summarizePendingDeliveries(now: number = Date.now()): QueueAgeStats | null { + const agg = this.db + .prepare( + `SELECT + COUNT(*) AS count, + MIN(? - e.created_at) AS minMs, + MAX(? - e.created_at) AS maxMs, + AVG(? - e.created_at) AS avgMs + FROM space_external_event_deliveries d + INNER JOIN space_external_events e ON e.id = d.event_id + WHERE d.state = 'pending'` + ) + .get(now, now, now) as { + count: number; + minMs: number | null; + maxMs: number | null; + avgMs: number | null; + }; + const count = agg?.count ?? 0; + if (count === 0) return null; + // Nearest-rank p95 (no interpolation): the ceil(0.95 * n)th value, clamped. + const p95Offset = Math.min(count - 1, Math.ceil(count * 0.95) - 1); + const p95 = this.db .prepare( `SELECT (? - e.created_at) AS age - FROM space_external_event_deliveries d - INNER JOIN space_external_events e ON e.id = d.event_id - WHERE d.state = 'pending'` + FROM space_external_event_deliveries d + INNER JOIN space_external_events e ON e.id = d.event_id + WHERE d.state = 'pending' + ORDER BY age + LIMIT 1 OFFSET ?` ) - .all(now) as { age: number }[]; - return rows.map((row) => row.age); + .get(now, p95Offset) as { age: number } | undefined; + return { + count, + minMs: agg.minMs ?? 0, + maxMs: agg.maxMs ?? 0, + avgMs: Math.round(agg.avgMs ?? 0), + p95Ms: p95?.age ?? agg.maxMs ?? 0, + }; } getDelivery(eventId: string, deliveryKey: string): ExternalEventDeliveryRecord | null { diff --git a/packages/daemon/src/lib/external-events/queue-health-metrics.ts b/packages/daemon/src/lib/external-events/queue-health-metrics.ts index 5ef4ce1d8f..f970c200e3 100644 --- a/packages/daemon/src/lib/external-events/queue-health-metrics.ts +++ b/packages/daemon/src/lib/external-events/queue-health-metrics.ts @@ -69,11 +69,18 @@ export interface QueueHealthCounters { */ claimConflicts: number; /** - * Non-terminal skips: the target's worker session was no longer live (or its - * space was paused) at injection time, so the delivery was deferred/requeued - * rather than injected. + * Non-terminal skips: the target's worker session was no longer live at + * injection time, so the delivery was deferred/requeued rather than + * injected. A session-loss signal (worker crashed or was superseded). */ staleSessionSkips: number; + /** + * Non-terminal skips: a delivery was deferred because the target's space was + * paused/stopped at injection time. Distinct from `staleSessionSkips` (a + * reliability signal) since pausing a space is intentional — the delivery + * stays pending and is requeued by `onSpaceResumed`. + */ + pausedSpaceSkips: number; } /** Live gauges computed at read time from in-memory + DB state. */ @@ -125,13 +132,24 @@ const FAILURE_CATEGORY_PREFIXES: Array<{ }> = [ { category: 'ttl_expired', test: (r) => r === 'ttl_expired' }, { category: 'cap_eviction', test: (r) => r === 'pending_node_queue_overflow' }, + // Deliverability: the target/run is no longer a valid delivery destination + // (run not deliverable, task terminal, subscription removed/cleared, node + // cancelled, or queued deliveries swept during run-interests rebuild / + // terminal cleanup). These are all terminal `markDeliveryFailed` reasons + // that reach the delivery-terminal hook. + // NOTE: `blocked_run_gate_not_opened` is intentionally absent — it is set via + // event-level `markEventFailed`, which never fires the delivery hook, so it + // can never appear in `finalFailuresByReason`. { category: 'deliverability', test: (r) => r === 'run_not_externally_deliverable' || r === 'target_task_terminal' || r === 'subscription_no_longer_active' || - r === 'blocked_run_gate_not_opened', + r === 'auto_pr_subscription_cleared' || + r === 'node_execution_cancelled' || + r === 'run_interests_rebuilt' || + r === 'run_terminal_cleanup', }, // Retry-exhaustion terminal failures carry the underlying activation/delivery // reason (e.g. `node_execution_not_active`, `activation_failed; ...`). @@ -208,6 +226,7 @@ export class ExternalEventQueueMetrics { private readonly finalFailuresByReason = new Map(); private claimConflicts = 0; private staleSessionSkips = 0; + private pausedSpaceSkips = 0; constructor(now: number = Date.now()) { this.since = now; @@ -239,6 +258,11 @@ export class ExternalEventQueueMetrics { this.staleSessionSkips += 1; } + /** Record a non-terminal skip caused by the target's space being paused/stopped. */ + recordPausedSpaceSkip(): void { + this.pausedSpaceSkips += 1; + } + /** * Record a delivery terminal transition. Called from the store's * delivery-terminal hook — the single point that observes every `delivered` @@ -266,6 +290,7 @@ export class ExternalEventQueueMetrics { finalFailuresByReason: recordToObject(this.finalFailuresByReason), claimConflicts: this.claimConflicts, staleSessionSkips: this.staleSessionSkips, + pausedSpaceSkips: this.pausedSpaceSkips, }; } diff --git a/packages/daemon/src/lib/space/runtime/space-runtime.ts b/packages/daemon/src/lib/space/runtime/space-runtime.ts index 0e46f68f5a..d1fca1a1cc 100644 --- a/packages/daemon/src/lib/space/runtime/space-runtime.ts +++ b/packages/daemon/src/lib/space/runtime/space-runtime.ts @@ -2352,6 +2352,7 @@ export class SpaceRuntime { // auto-subscription, which now matches PR events during pause.) const targetRun = this.config.workflowRunRepo.getRun(resolved.workflowRunId); if (targetRun && this.pausedSpaceIds.has(targetRun.spaceId)) { + this.queueHealthMetrics.recordPausedSpaceSkip(); return; } const eventRecord = store.getById(payload.eventId); @@ -2768,6 +2769,7 @@ export class SpaceRuntime { const reason = spacePaused ? 'space_paused' : 'session loss'; for (const item of items) { if (sessionLoss) this.queueHealthMetrics.recordStaleSessionSkip(); + else this.queueHealthMetrics.recordPausedSpaceSkip(); store.markDeliveryFailed(item.event.eventId, item.deliveryKey, { terminal: false, reason: `deliveryMode:${item.deliveryMode}; digest requeued after ${reason}`, @@ -2901,7 +2903,10 @@ export class SpaceRuntime { // pause that bypass the fresh-delivery guard in deliverExternalEventToWorkflowTarget; // leaving the persisted delivery pending lets onSpaceResumed requeue it. const pausedRun = this.config.workflowRunRepo.getRun(target.workflowRunId); - if (pausedRun && this.pausedSpaceIds.has(pausedRun.spaceId)) return; + if (pausedRun && this.pausedSpaceIds.has(pausedRun.spaceId)) { + this.queueHealthMetrics.recordPausedSpaceSkip(); + return; + } this.externalEventDeliveriesInFlight.add(deliveryKey); try { if (!this.config.commandBus) { @@ -3244,6 +3249,13 @@ export class SpaceRuntime { createdAt: item.createdAt, }); this.pendingExternalEventQueue.set(key, queue); + // These items are rehoused from the rate-limit digest (they bypassed + // queueForPendingNode), so count the enqueue here to keep the cumulative + // counter consistent with the live queueDepth gauge. + this.queueHealthMetrics.recordEnqueue( + item.event.source, + this.describeEnqueueTargetState(item.target) + ); } private queueForPendingNode( @@ -3312,18 +3324,19 @@ export class SpaceRuntime { for (const state of this.externalEventRateLimits.values()) { digestBacklog += state.pendingDigest.length; } + // Persisted-pending count + age via SQL aggregates — avoids materializing + // every pending row (and a full in-memory sort) on each snapshot read. const store = this.config.externalEventStore; - const persistedPending = store ? store.listPendingDeliveries().length : 0; - const persistedAges = store ? store.getPendingDeliveryAges(now) : []; + const persisted = store ? store.summarizePendingDeliveries(now) : null; const gauges: QueueHealthGauges = { queueDepth, queueKeys: this.pendingExternalEventQueue.size, inFlight: this.externalEventDeliveriesInFlight.size, digestBacklog, retryTimers: this.externalEventRetryTimers.size, - persistedPending, + persistedPending: persisted?.count ?? 0, queueAgeMs: computeQueueAgeStats(inMemoryAges), - persistedAgeMs: computeQueueAgeStats(persistedAges), + persistedAgeMs: persisted, }; return this.queueHealthMetrics.snapshot(gauges, now); } diff --git a/packages/daemon/tests/unit/4-space-storage/storage/external-event-store.test.ts b/packages/daemon/tests/unit/4-space-storage/storage/external-event-store.test.ts index cb900ca2c6..6f82ce0d21 100644 --- a/packages/daemon/tests/unit/4-space-storage/storage/external-event-store.test.ts +++ b/packages/daemon/tests/unit/4-space-storage/storage/external-event-store.test.ts @@ -942,28 +942,55 @@ describe('delivery-terminal hook', () => { }); }); -describe('getPendingDeliveryAges', () => { - test('returns event-age for pending deliveries, empty when none', () => { - const now = Date.now(); - expect(store.getPendingDeliveryAges(now)).toEqual([]); +describe('summarizePendingDeliveries', () => { + test('returns null when there are no pending deliveries', () => { + expect(store.summarizePendingDeliveries(Date.now())).toBeNull(); + }); + test('returns count + min/max/avg/p95 age without materializing rows', () => { + const now = Date.now(); store.store(EVENT_A); - store.registerExpectedDelivery('evt-a', 'dk-1', { + store.store(EVENT_B); + store.registerExpectedDelivery('evt-a', 'dk-a', { + workflowRunId: 'run-1', + taskId: 'task-1', + nodeId: 'node-1', + agentName: 'coder', + }); + store.registerExpectedDelivery('evt-b', 'dk-b', { workflowRunId: 'run-1', taskId: 'task-1', nodeId: 'node-1', agentName: 'coder', }); - // created_at is ingestion time (the TTL anchor), set by store() to the - // current time. Backdate it 60s so the age is deterministic. + // created_at is ingestion time (the TTL anchor). Backdate evt-a 60s and + // evt-b 30s so min/max/avg/p95 are deterministic. db.prepare(`UPDATE space_external_events SET created_at = ? WHERE id = ?`).run( now - 60_000, 'evt-a' ); + db.prepare(`UPDATE space_external_events SET created_at = ? WHERE id = ?`).run( + now - 30_000, + 'evt-b' + ); - const ages = store.getPendingDeliveryAges(now); - expect(ages).toHaveLength(1); - expect(ages[0]).toBeGreaterThanOrEqual(59_000); - expect(ages[0]).toBeLessThanOrEqual(61_000); + const summary = store.summarizePendingDeliveries(now); + expect(summary).not.toBeNull(); + expect(summary!.count).toBe(2); + expect(summary!.minMs).toBeGreaterThanOrEqual(29_000); + expect(summary!.minMs).toBeLessThanOrEqual(31_000); + expect(summary!.maxMs).toBeGreaterThanOrEqual(59_000); + expect(summary!.maxMs).toBeLessThanOrEqual(61_000); + expect(summary!.avgMs).toBeGreaterThanOrEqual(44_000); + expect(summary!.avgMs).toBeLessThanOrEqual(46_000); + // With 2 values, nearest-rank p95 = the max. + expect(summary!.p95Ms).toBe(summary!.maxMs); + + // Delivering one drops the count to 1 and recomputes the single-value stats. + store.markDeliveryDelivered('evt-b', 'dk-b'); + const afterDeliver = store.summarizePendingDeliveries(now); + expect(afterDeliver!.count).toBe(1); + expect(afterDeliver!.minMs).toBe(afterDeliver!.maxMs); + expect(afterDeliver!.p95Ms).toBe(afterDeliver!.maxMs); }); }); diff --git a/packages/daemon/tests/unit/4-space-storage/storage/queue-health-metrics.test.ts b/packages/daemon/tests/unit/4-space-storage/storage/queue-health-metrics.test.ts index ee775cc33f..f3229f4783 100644 --- a/packages/daemon/tests/unit/4-space-storage/storage/queue-health-metrics.test.ts +++ b/packages/daemon/tests/unit/4-space-storage/storage/queue-health-metrics.test.ts @@ -54,15 +54,19 @@ describe('ExternalEventQueueMetrics — flush + skip counters', () => { expect(counters.flushItemsDispatched).toBe(5); }); - test('counts claim conflicts and stale-session skips separately', () => { + test('counts claim conflicts, stale-session skips, and paused-space skips separately', () => { const metrics = new ExternalEventQueueMetrics(1000); metrics.recordClaimConflict(); metrics.recordClaimConflict(); metrics.recordStaleSessionSkip(); + metrics.recordPausedSpaceSkip(); + metrics.recordPausedSpaceSkip(); + metrics.recordPausedSpaceSkip(); const counters = metrics.getCounters(); expect(counters.claimConflicts).toBe(2); expect(counters.staleSessionSkips).toBe(1); + expect(counters.pausedSpaceSkips).toBe(3); }); }); @@ -159,11 +163,18 @@ describe('categorizeFailureReason', () => { ['run_not_externally_deliverable', 'deliverability'], ['target_task_terminal', 'deliverability'], ['subscription_no_longer_active', 'deliverability'], - ['blocked_run_gate_not_opened', 'deliverability'], + ['auto_pr_subscription_cleared', 'deliverability'], + ['node_execution_cancelled', 'deliverability'], + ['run_interests_rebuilt', 'deliverability'], + ['run_terminal_cleanup', 'deliverability'], ['node_execution_not_active', 'retry_exhausted'], ['node_execution_pending', 'retry_exhausted'], ['activation_failed; timeout', 'retry_exhausted'], ['deliveryMode:immediate; inject failed', 'injection_error'], + // blocked_run_gate_not_opened is event-level (markEventFailed) and never + // reaches the delivery hook, so it cannot appear in finalFailuresByReason; + // if it ever did, it would fall through to `other`. + ['blocked_run_gate_not_opened', 'other'], ['something_unexpected', 'other'], ])('categorizes %s -> %s', (reason, expected) => { expect(categorizeFailureReason(reason)).toBe(expected); diff --git a/packages/web/src/components/space/QueueHealthSummary.tsx b/packages/web/src/components/space/QueueHealthSummary.tsx index ebaa1f5944..94b8691576 100644 --- a/packages/web/src/components/space/QueueHealthSummary.tsx +++ b/packages/web/src/components/space/QueueHealthSummary.tsx @@ -181,8 +181,8 @@ export function QueueHealthSummary(): preact.JSX.Element { diff --git a/packages/web/src/components/space/__tests__/QueueHealthSummary.test.tsx b/packages/web/src/components/space/__tests__/QueueHealthSummary.test.tsx index 52f4e62bd0..1857ac5522 100644 --- a/packages/web/src/components/space/__tests__/QueueHealthSummary.test.tsx +++ b/packages/web/src/components/space/__tests__/QueueHealthSummary.test.tsx @@ -35,6 +35,7 @@ const sampleSnapshot = { finalFailuresByReason: { ttl_expired: 1, pending_node_queue_overflow: 1 }, claimConflicts: 1, staleSessionSkips: 2, + pausedSpaceSkips: 0, }, failuresByCategory: { ttl_expired: 1, diff --git a/packages/web/src/lib/space-store.ts b/packages/web/src/lib/space-store.ts index f4e7e7f11d..4795ee8587 100644 --- a/packages/web/src/lib/space-store.ts +++ b/packages/web/src/lib/space-store.ts @@ -110,6 +110,7 @@ export interface QueueHealthCounters { finalFailuresByReason: Record; claimConflicts: number; staleSessionSkips: number; + pausedSpaceSkips: number; } /** Live queue gauges computed at read time. */