diff --git a/EVENT_QUEUE_DESIGN.md b/EVENT_QUEUE_DESIGN.md index 070208f..fc50552 100644 --- a/EVENT_QUEUE_DESIGN.md +++ b/EVENT_QUEUE_DESIGN.md @@ -218,7 +218,8 @@ const stats = queue.getStats(); // total: 42, // Total events in queue // pending: 10, // Awaiting processing // processing: 2, // Currently being handled -// processed: 30, // Successfully handled +// retrying: 3, // Failed with retries remaining +// processed: 27, // Successfully handled // failed: 0, // Exhausted retries // oldestEventAge: 1500 // ms since oldest event enqueued // } @@ -320,6 +321,7 @@ const stats = propagator.getQueueStats(); const isHealthy = stats.processing <= 1 && // Not backed up + stats.retrying === 0 && // No retry backlog stats.failed === 0 && // No exhausted retries stats.oldestEventAge < 60_000; // No old events (< 60s) diff --git a/EVENT_QUEUE_IMPLEMENTATION_SUMMARY.md b/EVENT_QUEUE_IMPLEMENTATION_SUMMARY.md index 4d2fd91..b1e2429 100644 --- a/EVENT_QUEUE_IMPLEMENTATION_SUMMARY.md +++ b/EVENT_QUEUE_IMPLEMENTATION_SUMMARY.md @@ -200,7 +200,9 @@ stmt.run(event.id, JSON.stringify(event), "pending", 0, Date.now()); **Verification**: - Test: `reports accurate queue statistics` ✅ - Enqueue 5 events in different states - - Verify stats: total=5, pending=3, processing=1, processed=1 + - Verify stats: total=5, pending=3, retrying=1, processed=1, processing=0 + +- Test: `separates retrying failures from processing and exhausted failed events` ✅ - Test: `calculates oldest event age` ✅ - Enqueue event @@ -219,8 +221,9 @@ return { total: sum of all counts, pending: WHERE status = 'pending', processing: WHERE status = 'processing', + retrying: WHERE status = 'failed' AND attempts < maxRetries, processed: WHERE status = 'processed', - failed: WHERE status = 'failed', + failed: WHERE status = 'failed' AND attempts >= maxRetries, oldestEventAge: Date.now() - earliest enqueueTime }; ``` @@ -353,6 +356,7 @@ EventQueue ✓ does not recover processed events Queue Statistics ✓ reports accurate queue statistics + ✓ separates retrying failures from processing and exhausted failed events ✓ calculates oldest event age Cleanup ✓ removes old processed events diff --git a/engine-bridge/src/__tests__/event-queue.test.ts b/engine-bridge/src/__tests__/event-queue.test.ts index e6d3551..532da4b 100644 --- a/engine-bridge/src/__tests__/event-queue.test.ts +++ b/engine-bridge/src/__tests__/event-queue.test.ts @@ -257,8 +257,37 @@ describe("EventQueue", () => { expect(stats.total).toBe(5); expect(stats.processed).toBe(1); expect(stats.pending).toBe(3); - expect(stats.processing).toBe(1); // evt-2 is still in processing state - expect(stats.failed).toBe(0); // evt-2 was retried as pending, not failed + expect(stats.processing).toBe(0); + expect(stats.retrying).toBe(1); // evt-2 failed with retries remaining + expect(stats.failed).toBe(0); + + cleanupTestQueue(queue, path); + }); + + it("separates retrying failures from processing and exhausted failed events", () => { + const { queue, path } = createTestQueue(); + + queue.enqueue(createTestEvent("evt-exhausted", "test")); + for (let i = 0; i < 3; i++) { + expect(queue.dequeue()).not.toBeNull(); + queue.markFailed("evt-exhausted", new Error(`attempt ${i + 1}`)); + } + + queue.enqueue(createTestEvent("evt-processing", "test")); + expect(queue.dequeue()!.id).toBe("evt-processing"); + + queue.enqueue(createTestEvent("evt-retrying", "test")); + expect(queue.dequeue()!.id).toBe("evt-retrying"); + queue.markFailed("evt-retrying", new Error("transient")); + + const stats = queue.getStats(); + + expect(stats.total).toBe(3); + expect(stats.processing).toBe(1); + expect(stats.retrying).toBe(1); + expect(stats.failed).toBe(1); + expect(stats.pending).toBe(0); + expect(stats.processed).toBe(0); cleanupTestQueue(queue, path); }); diff --git a/engine-bridge/src/__tests__/heartbeat.test.ts b/engine-bridge/src/__tests__/heartbeat.test.ts index 2365795..ac8c37e 100644 --- a/engine-bridge/src/__tests__/heartbeat.test.ts +++ b/engine-bridge/src/__tests__/heartbeat.test.ts @@ -72,6 +72,26 @@ describe("HeartbeatMonitor", () => { propagator.stop(); }); + it("surfaces retrying queue stats separately from processing", async () => { + rpc.call = jest.fn().mockResolvedValue({}); + + monitor.start(); + await new Promise(resolve => setTimeout(resolve, 10)); + + const pulse = logSpy.mock.calls.find( + (call) => typeof call[1] === "string" && call[1].includes("eventQueue") + ); + expect(pulse).toBeDefined(); + const report = JSON.parse(pulse![1] as string); + expect(report.eventQueue).toEqual( + expect.objectContaining({ + processing: 0, + retrying: 0, + failed: 0, + }) + ); + }); + it("logs at intervals", async () => { rpc.call = jest.fn().mockResolvedValue({}); diff --git a/engine-bridge/src/event-propagator.ts b/engine-bridge/src/event-propagator.ts index debaa02..c9eb21b 100644 --- a/engine-bridge/src/event-propagator.ts +++ b/engine-bridge/src/event-propagator.ts @@ -8,7 +8,7 @@ */ import { RpcClient } from "./rpc-client"; -import { EventQueue } from "./event-queue"; +import { EventQueue, type QueueStats } from "./event-queue"; import { logger } from "./logger"; export interface EngineEvent { @@ -153,8 +153,8 @@ export class EventPropagator { } } - /** Get queue statistics. */ - getQueueStats() { + /** Get queue statistics, including the retrying backlog. */ + getQueueStats(): QueueStats { return this.queue.getStats(); } diff --git a/engine-bridge/src/event-queue.ts b/engine-bridge/src/event-queue.ts index be15429..0873eae 100644 --- a/engine-bridge/src/event-queue.ts +++ b/engine-bridge/src/event-queue.ts @@ -33,6 +33,27 @@ export interface QueuedEvent { nextAttempt?: number; } +export interface QueueStats { + total: number; + pending: number; + processing: number; + /** Failed events that still have retries remaining — not healthy in-flight work. */ + retrying: number; + processed: number; + failed: number; + oldestEventAge: number | null; +} + +const EMPTY_STATS: QueueStats = { + total: 0, + pending: 0, + processing: 0, + retrying: 0, + processed: 0, + failed: 0, + oldestEventAge: null, +}; + export class EventQueue { private db: Database.Database; private readonly dbPath: string; @@ -239,15 +260,12 @@ export class EventQueue { /** * Get queue statistics (size, oldest event, error rate). + * + * Failed events that still have retries remaining are reported as + * `retrying`, not `processing`, so an outage-driven retry backlog is + * visible instead of looking like healthy in-flight work. */ - getStats(): { - total: number; - pending: number; - processing: number; - processed: number; - failed: number; - oldestEventAge: number | null; - } { + getStats(): QueueStats { try { const allStmt = this.db.prepare("SELECT status, attempts FROM events"); const rows = allStmt.all() as any[]; @@ -256,6 +274,7 @@ export class EventQueue { total: rows.length, pending: 0, processing: 0, + retrying: 0, processed: 0, failed: 0, }; @@ -269,7 +288,7 @@ export class EventQueue { stats.processing++; } else if (row.status === "failed") { if (row.attempts < this.maxRetries) { - stats.processing++; + stats.retrying++; } else { stats.failed++; } @@ -286,14 +305,7 @@ export class EventQueue { }; } catch (err) { logger.error("[EventQueue] Get stats failed:", err); - return { - total: 0, - pending: 0, - processing: 0, - processed: 0, - failed: 0, - oldestEventAge: null, - }; + return { ...EMPTY_STATS }; } } diff --git a/engine-bridge/src/heartbeat-monitor.ts b/engine-bridge/src/heartbeat-monitor.ts index d31dce8..f72fb29 100644 --- a/engine-bridge/src/heartbeat-monitor.ts +++ b/engine-bridge/src/heartbeat-monitor.ts @@ -71,6 +71,7 @@ export class HeartbeatMonitor { running: this.propagator.isRunning(), cursor: this.propagator.getCursor() ?? "none", }, + eventQueue: this.propagator.getQueueStats(), system: { memoryRssMb: Math.round(process.memoryUsage().rss / 1024 / 1024), uptimeSec: Math.round(process.uptime()), diff --git a/engine-bridge/src/index.ts b/engine-bridge/src/index.ts index 414b745..3e87093 100644 --- a/engine-bridge/src/index.ts +++ b/engine-bridge/src/index.ts @@ -1,3 +1,5 @@ +export { EventQueue } from "./event-queue"; +export type { QueuedEvent, QueueStats } from "./event-queue"; export { RpcClient } from "./rpc-client"; export { NonceManager } from "./nonce-manager"; export { EventPropagator } from "./event-propagator";