Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion EVENT_QUEUE_DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
// }
Expand Down Expand Up @@ -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)

Expand Down
8 changes: 6 additions & 2 deletions EVENT_QUEUE_IMPLEMENTATION_SUMMARY.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
};
```
Expand Down Expand Up @@ -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
Expand Down
33 changes: 31 additions & 2 deletions engine-bridge/src/__tests__/event-queue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
});
Expand Down
20 changes: 20 additions & 0 deletions engine-bridge/src/__tests__/heartbeat.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({});

Expand Down
6 changes: 3 additions & 3 deletions engine-bridge/src/event-propagator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -153,8 +153,8 @@ export class EventPropagator {
}
}

/** Get queue statistics. */
getQueueStats() {
/** Get queue statistics, including the retrying backlog. */
getQueueStats(): QueueStats {
return this.queue.getStats();
}

Expand Down
46 changes: 29 additions & 17 deletions engine-bridge/src/event-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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[];
Expand All @@ -256,6 +274,7 @@ export class EventQueue {
total: rows.length,
pending: 0,
processing: 0,
retrying: 0,
processed: 0,
failed: 0,
};
Expand All @@ -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++;
}
Expand All @@ -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 };
}
}

Expand Down
1 change: 1 addition & 0 deletions engine-bridge/src/heartbeat-monitor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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()),
Expand Down
2 changes: 2 additions & 0 deletions engine-bridge/src/index.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down
Loading