From 3cee88b272b095afa8f934caf8ca22f48289388d Mon Sep 17 00:00:00 2001 From: Ronaldo Martins Date: Mon, 10 Aug 2026 17:56:27 -0300 Subject: [PATCH 1/2] perf(api): replace SSE polling with shared EventEmitter (ENG-1669) - New core/task-event-bus.ts: shared EventEmitter broadcasting persisted task events in-process (setMaxListeners(0) for many SSE connections) - db.createTaskEvent now emits each event on the bus (with a light task-status lookup, only when listeners exist) preserving RML-716 - GET /api/logs/stream: removed the 2s-per-connection DB polling loop; now push-based via bus subscription, with a one-time catch-up query honoring Last-Event-ID/cursor and dedupe guards for the overlap window - keep-alive ping and abort cleanup preserved; unit tests added --- packages/api/src/core/task-event-bus.test.ts | 54 ++++++++++ packages/api/src/core/task-event-bus.ts | 37 +++++++ packages/api/src/integrations/db.ts | 23 ++++- packages/api/src/router.ts | 102 ++++++++++++------- 4 files changed, 178 insertions(+), 38 deletions(-) create mode 100644 packages/api/src/core/task-event-bus.test.ts create mode 100644 packages/api/src/core/task-event-bus.ts diff --git a/packages/api/src/core/task-event-bus.test.ts b/packages/api/src/core/task-event-bus.test.ts new file mode 100644 index 0000000..d89bdc7 --- /dev/null +++ b/packages/api/src/core/task-event-bus.test.ts @@ -0,0 +1,54 @@ +import { describe, it, expect } from "bun:test"; +import { taskEventBus, type BroadcastTaskEvent } from "./task-event-bus"; + +function makeEvent(overrides: Partial = {}): BroadcastTaskEvent { + return { + id: "11111111-1111-1111-1111-111111111111", + taskId: "22222222-2222-2222-2222-222222222222", + eventType: "AGENT_COMPLETED", + agent: "coder", + inputSummary: null, + outputSummary: "done", + tokensUsed: 10, + durationMs: 100, + metadata: null, + createdAt: new Date(), + ...overrides, + } as BroadcastTaskEvent; +} + +describe("taskEventBus (ENG-1669)", () => { + it("delivers emitted events to subscribers", () => { + const received: BroadcastTaskEvent[] = []; + const unsubscribe = taskEventBus.onTaskEvent((e) => received.push(e)); + + const event = makeEvent({ taskStatus: "IN_PROGRESS" }); + taskEventBus.emitTaskEvent(event); + + expect(received).toHaveLength(1); + expect(received[0]!.id).toBe(event.id); + expect(received[0]!.taskStatus).toBe("IN_PROGRESS"); + unsubscribe(); + }); + + it("unsubscribe stops delivery and updates listener count", () => { + const before = taskEventBus.listenerCountTaskEvent; + const received: BroadcastTaskEvent[] = []; + const unsubscribe = taskEventBus.onTaskEvent((e) => received.push(e)); + expect(taskEventBus.listenerCountTaskEvent).toBe(before + 1); + + unsubscribe(); + expect(taskEventBus.listenerCountTaskEvent).toBe(before); + + taskEventBus.emitTaskEvent(makeEvent()); + expect(received).toHaveLength(0); + }); + + it("supports many concurrent subscribers (no MaxListeners warning cap)", () => { + const unsubs = Array.from({ length: 50 }, () => + taskEventBus.onTaskEvent(() => {}), + ); + expect(taskEventBus.listenerCountTaskEvent).toBeGreaterThanOrEqual(50); + for (const u of unsubs) u(); + }); +}); diff --git a/packages/api/src/core/task-event-bus.ts b/packages/api/src/core/task-event-bus.ts new file mode 100644 index 0000000..b759a1a --- /dev/null +++ b/packages/api/src/core/task-event-bus.ts @@ -0,0 +1,37 @@ +/** + * Shared EventEmitter for task events (ENG-1669). + * + * Replaces the per-connection DB polling in the SSE endpoint + * (GET /api/logs/stream). `db.createTaskEvent` emits every persisted + * event here; SSE connections subscribe instead of polling. + * + * Note: in-process only. If the API is ever scaled to multiple + * processes, this must be replaced by Postgres LISTEN/NOTIFY or Redis + * pub/sub. + */ +import { EventEmitter } from "node:events"; +import type { TaskEvent } from "./types"; + +export const TASK_EVENT = "task-event"; + +/** TaskEvent enriched with the task's current status (RML-716). */ +export type BroadcastTaskEvent = TaskEvent & { taskStatus?: string }; + +class TaskEventBus extends EventEmitter { + emitTaskEvent(event: BroadcastTaskEvent): void { + this.emit(TASK_EVENT, event); + } + + onTaskEvent(listener: (event: BroadcastTaskEvent) => void): () => void { + this.on(TASK_EVENT, listener); + return () => this.off(TASK_EVENT, listener); + } + + get listenerCountTaskEvent(): number { + return this.listenerCount(TASK_EVENT); + } +} + +export const taskEventBus = new TaskEventBus(); +// One listener per SSE connection; do not cap at the default of 10. +taskEventBus.setMaxListeners(0); diff --git a/packages/api/src/integrations/db.ts b/packages/api/src/integrations/db.ts index 418f92d..b5b7304 100644 --- a/packages/api/src/integrations/db.ts +++ b/packages/api/src/integrations/db.ts @@ -6,6 +6,7 @@ import { OrchestrationState, OrchestrationStateSchema, } from "../core/types"; +import { taskEventBus } from "../core/task-event-bus"; const connectionString = process.env.DATABASE_URL; @@ -300,7 +301,27 @@ export const db = { ) RETURNING * `; - return this.mapTaskEvent(result); + const mapped = this.mapTaskEvent(result); + + // ENG-1669: broadcast to in-process subscribers (SSE) instead of + // having each connection poll the DB. Only pay for the task-status + // lookup (RML-716) when someone is actually listening. + if (taskEventBus.listenerCountTaskEvent > 0) { + try { + const [task] = await sql` + SELECT status FROM tasks WHERE id = ${mapped.taskId} + `; + taskEventBus.emitTaskEvent({ + ...mapped, + taskStatus: task?.status as string | undefined, + }); + } catch (err) { + // Broadcasting is best-effort; never fail the write because of it. + console.error("[task-event-bus] Failed to broadcast event:", err); + } + } + + return mapped; }, async getTaskEvents(taskId: string): Promise { diff --git a/packages/api/src/router.ts b/packages/api/src/router.ts index 57060e7..4145b55 100644 --- a/packages/api/src/router.ts +++ b/packages/api/src/router.ts @@ -3,9 +3,11 @@ import { GitHubCheckRunEvent, GitHubPullRequestReviewEvent, Task, + TaskEvent, defaultConfig, JobStatus, } from "./core/types"; +import { taskEventBus } from "./core/task-event-bus"; import { Orchestrator } from "./core/orchestrator"; import { TaskRunner } from "./core/task-runner"; import { db } from "./integrations/db"; @@ -4296,48 +4298,74 @@ route("GET", "/api/logs/stream", async (req) => { ), ); - // Poll for new events every 2 seconds - const pollInterval = setInterval(async () => { - if (!isActive) { - clearInterval(pollInterval); + function sendEvent(event: TaskEvent & { taskStatus?: string }): void { + const cursor = formatCursor({ + createdAt: event.createdAt, + id: event.id, + }); + + controller.enqueue(encoder.encode(`id: ${cursor}\n`)); + controller.enqueue( + encoder.encode( + `data: ${JSON.stringify({ + type: "event", + id: event.id, + taskId: event.taskId, + eventType: event.eventType, + agent: event.agent, + message: event.outputSummary || event.eventType, + timestamp: event.createdAt, + level: getLogLevel(event.eventType), + tokensUsed: event.tokensUsed, + durationMs: event.durationMs, + // Include current task status for real-time UI updates (RML-716) + taskStatus: event.taskStatus, + })}\n\n`, + ), + ); + + lastCursor = { createdAt: event.createdAt, id: event.id }; + } + + // ENG-1669: push-based delivery via shared EventEmitter instead of + // polling the DB every 2s per connection. + const unsubscribe = taskEventBus.onTaskEvent((event) => { + if (!isActive) return; + if (taskId && event.taskId !== taskId) return; + // Skip events already delivered by the catch-up query below. + if ( + event.createdAt.getTime() < lastCursor.createdAt.getTime() || + (event.createdAt.getTime() === lastCursor.createdAt.getTime() && + event.id <= lastCursor.id) + ) { return; } - try { - const events = await db.getRecentTaskEvents(lastCursor, taskId); - - for (const event of events) { - const cursor = formatCursor({ - createdAt: event.createdAt, - id: event.id, - }); - - controller.enqueue(encoder.encode(`id: ${cursor}\n`)); - controller.enqueue( - encoder.encode( - `data: ${JSON.stringify({ - type: "event", - id: event.id, - taskId: event.taskId, - eventType: event.eventType, - agent: event.agent, - message: event.outputSummary || event.eventType, - timestamp: event.createdAt, - level: getLogLevel(event.eventType), - tokensUsed: event.tokensUsed, - durationMs: event.durationMs, - // Include current task status for real-time UI updates (RML-716) - taskStatus: (event as any).taskStatus, - })}\n\n`, - ), - ); + sendEvent(event); + } catch (err) { + console.error("[SSE] Error sending event:", err); + } + }); - lastCursor = { createdAt: event.createdAt, id: event.id }; + // One-time catch-up for events missed since Last-Event-ID/cursor + // (also covers events written between connect and subscribe). + try { + const missed = await db.getRecentTaskEvents(lastCursor, taskId); + for (const event of missed) { + if (!isActive) break; + // Skip anything already delivered live while the query ran. + if ( + event.createdAt.getTime() < lastCursor.createdAt.getTime() || + (event.createdAt.getTime() === lastCursor.createdAt.getTime() && + event.id <= lastCursor.id) + ) { + continue; } - } catch (err) { - console.error("[SSE] Error fetching events:", err); + sendEvent(event); } - }, 2000); + } catch (err) { + console.error("[SSE] Error fetching catch-up events:", err); + } // Send keep-alive ping every 30 seconds const keepAlive = setInterval(() => { @@ -4351,7 +4379,7 @@ route("GET", "/api/logs/stream", async (req) => { // Cleanup on close (handled by abort signal) req.signal?.addEventListener("abort", () => { isActive = false; - clearInterval(pollInterval); + unsubscribe(); clearInterval(keepAlive); controller.close(); }); From fbe318d0dad3dbc4e93f808f7174023eb40aa399 Mon Sep 17 00:00:00 2001 From: Ronaldo Martins Date: Mon, 10 Aug 2026 21:47:34 -0300 Subject: [PATCH 2/2] fix(ENG-1669): close SSE catch-up/live race, pagination, listener leak, backpressure - Extract SSE catch-up/live-merge logic into core/sse-log-stream.ts with a buffer-then-drain design: subscribe to the live bus first (buffering events instead of sending), paginate the catch-up query to full drain (no truncating LIMIT), then flush the buffer deduped against catch-up before switching to direct live delivery. Dedup and cursor advancement now happen at a single point keyed by event id, removing the two independently-advancing cursors that raced on timestamp. - router.ts: register the abort/cancel cleanup (unsubscribe + keep-alive clear + controller close) immediately in start(), before any await, so a disconnect during catch-up can no longer leak the live-bus listener. - router.ts: check controller.desiredSize before every enqueue; close the connection when a client falls behind instead of growing the buffer unbounded. - db.ts: split the task-status SELECT from the event emit into separate try/catch blocks so a transient SELECT failure degrades to broadcasting without enrichment instead of suppressing the broadcast entirely. - Add sse-log-stream.test.ts covering the race, >50-event pagination, listener-leak-on-disconnect, backpressure cutoff, taskId filtering, and live/catch-up dedup (8 tests, 159 assertions, all passing) --- packages/api/src/core/sse-log-stream.test.ts | 288 +++++++++++++++++++ packages/api/src/core/sse-log-stream.ts | 177 ++++++++++++ packages/api/src/integrations/db.ts | 16 +- packages/api/src/router.ts | 220 +++++++------- 4 files changed, 601 insertions(+), 100 deletions(-) create mode 100644 packages/api/src/core/sse-log-stream.test.ts create mode 100644 packages/api/src/core/sse-log-stream.ts diff --git a/packages/api/src/core/sse-log-stream.test.ts b/packages/api/src/core/sse-log-stream.test.ts new file mode 100644 index 0000000..fd74d5c --- /dev/null +++ b/packages/api/src/core/sse-log-stream.test.ts @@ -0,0 +1,288 @@ +import { describe, it, expect, beforeEach } from "bun:test"; +import { + runSseLogStream, + parseCursor, + formatCursor, + DEFAULT_CURSOR, + type TaskEventBusLike, + type TaskEventSourceLike, + type SseSink, +} from "./sse-log-stream"; +import type { BroadcastTaskEvent } from "./task-event-bus"; +import type { TaskEvent } from "./types"; + +type Ev = TaskEvent & { taskStatus?: string }; + +function makeEvent(id: string, createdAt: Date, taskId = "task-1"): Ev { + return { + id, + taskId, + eventType: "CODED", + agent: "coder", + outputSummary: `event ${id}`, + tokensUsed: 1, + durationMs: 1, + createdAt, + }; +} + +/** Fake bus that lets the test push live events on demand. */ +class FakeBus implements TaskEventBusLike { + private listeners: Array<(event: BroadcastTaskEvent) => void> = []; + + onTaskEvent(listener: (event: BroadcastTaskEvent) => void): () => void { + this.listeners.push(listener); + return () => { + this.listeners = this.listeners.filter((l) => l !== listener); + }; + } + + emit(event: BroadcastTaskEvent): void { + for (const l of [...this.listeners]) l(event); + } + + get listenerCount(): number { + return this.listeners.length; + } +} + +/** Fake paginated DB source backed by an in-memory sorted array. */ +class FakeSource implements TaskEventSourceLike { + constructor(private rows: Ev[]) {} + + async getRecentTaskEvents( + since: { createdAt: Date; id: string }, + taskId: string | undefined, + limit: number, + ): Promise { + const filtered = this.rows + .filter((e) => !taskId || e.taskId === taskId) + .filter( + (e) => + e.createdAt.getTime() > since.createdAt.getTime() || + (e.createdAt.getTime() === since.createdAt.getTime() && + e.id > since.id), + ) + .sort( + (a, b) => + a.createdAt.getTime() - b.createdAt.getTime() || + a.id.localeCompare(b.id), + ); + return filtered.slice(0, limit); + } +} + +/** Fake sink that records everything sent and can simulate a slow client. */ +class FakeSink implements SseSink { + sent: BroadcastTaskEvent[] = []; + active = true; + desiredSize: number | null = 10; + + isActive(): boolean { + return this.active; + } + + send(event: BroadcastTaskEvent): void { + this.sent.push(event); + } +} + +describe("sse-log-stream cursor helpers", () => { + it("parseCursor returns DEFAULT_CURSOR for null/invalid input", () => { + expect(parseCursor(null)).toEqual(DEFAULT_CURSOR); + expect(parseCursor("garbage")).toEqual(DEFAULT_CURSOR); + expect(parseCursor("not-a-date|abc")).toEqual(DEFAULT_CURSOR); + }); + + it("formatCursor/parseCursor round-trip", () => { + const d = new Date("2026-01-01T00:00:00.000Z"); + const formatted = formatCursor({ createdAt: d, id: "abc-123" }); + const parsed = parseCursor(formatted); + expect(parsed.createdAt.getTime()).toBe(d.getTime()); + expect(parsed.id).toBe("abc-123"); + }); +}); + +describe("runSseLogStream (ENG-1669 rework)", () => { + let bus: FakeBus; + let sink: FakeSink; + + beforeEach(() => { + bus = new FakeBus(); + sink = new FakeSink(); + }); + + it("BLOCKER 1: a live event racing ahead of catch-up does not cause earlier backlog to be skipped", async () => { + const t0 = new Date("2026-01-01T00:00:00.000Z"); + const backlog = Array.from({ length: 5 }, (_, i) => + makeEvent(`backlog-${i}`, new Date(t0.getTime() + i * 1000)), + ); + const source = new FakeSource(backlog); + + // Fire a "live" event with a LATER timestamp than the whole backlog, + // synchronously during the catch-up query (simulated by emitting it + // right after subscribe but before we await catch-up completion). + const liveEvent = makeEvent( + "live-1", + new Date(t0.getTime() + 100_000), // far ahead of backlog + ); + + const originalGet = source.getRecentTaskEvents.bind(source); + let firstCall = true; + source.getRecentTaskEvents = async (since, taskId, limit) => { + if (firstCall) { + firstCall = false; + // Simulate the live event arriving while the catch-up query is in flight. + bus.emit(liveEvent); + } + return originalGet(since, taskId, limit); + }; + + await runSseLogStream({ + taskId: "task-1", + initialCursor: DEFAULT_CURSOR, + bus, + source, + sink, + pageSize: 50, + }); + + const sentIds = sink.sent.map((e) => e.id); + // All backlog events must still be delivered, in order, despite the + // live event racing ahead on timestamp. + for (const b of backlog) { + expect(sentIds).toContain(b.id); + } + expect(sentIds).toContain("live-1"); + // Backlog must come before the live event in delivery order. + const lastBacklogIdx = Math.max( + ...backlog.map((b) => sentIds.indexOf(b.id)), + ); + expect(sentIds.indexOf("live-1")).toBeGreaterThan(lastBacklogIdx); + // No duplicates. + expect(new Set(sentIds).size).toBe(sentIds.length); + }); + + it("BLOCKER 2: backlog larger than a single page (>50) is delivered in full via pagination", async () => { + const t0 = new Date("2026-01-01T00:00:00.000Z"); + const backlog = Array.from({ length: 137 }, (_, i) => + makeEvent(`ev-${String(i).padStart(4, "0")}`, new Date(t0.getTime() + i * 10)), + ); + const source = new FakeSource(backlog); + + await runSseLogStream({ + taskId: "task-1", + initialCursor: DEFAULT_CURSOR, + bus, + source, + sink, + pageSize: 50, // matches prior default LIMIT + }); + + expect(sink.sent).toHaveLength(137); + const sentIds = sink.sent.map((e) => e.id); + for (const b of backlog) { + expect(sentIds).toContain(b.id); + } + // Delivered in ascending createdAt order. + const times = sink.sent.map((e) => e.createdAt.getTime()); + expect(times).toEqual([...times].sort((a, b) => a - b)); + }); + + it("HIGH 2: listener is unsubscribed via stop() and does not leak on early disconnect", async () => { + const source = new FakeSource([]); + expect(bus.listenerCount).toBe(0); + + const handlePromise = runSseLogStream({ + taskId: "task-1", + initialCursor: DEFAULT_CURSOR, + bus, + source, + sink, + }); + + // Listener should be registered synchronously before any await resolves. + expect(bus.listenerCount).toBe(1); + + const handle = await handlePromise; + handle.stop(); + + expect(bus.listenerCount).toBe(0); + + // Idempotent stop() should not throw or double-count. + handle.stop(); + expect(bus.listenerCount).toBe(0); + }); + + it("MED: stops sending once the sink reports backpressure/inactive", async () => { + const t0 = new Date("2026-01-01T00:00:00.000Z"); + const backlog = Array.from({ length: 10 }, (_, i) => + makeEvent(`ev-${i}`, new Date(t0.getTime() + i * 10)), + ); + const source = new FakeSource(backlog); + + let deliveries = 0; + sink.send = (event) => { + deliveries++; + if (deliveries === 3) { + sink.active = false; // simulate connection closing mid-flight + } + }; + + await runSseLogStream({ + taskId: "task-1", + initialCursor: DEFAULT_CURSOR, + bus, + source, + sink, + pageSize: 50, + }); + + expect(deliveries).toBe(3); + }); + + it("filters events by taskId", async () => { + const t0 = new Date("2026-01-01T00:00:00.000Z"); + const source = new FakeSource([ + makeEvent("a", t0, "task-1"), + makeEvent("b", new Date(t0.getTime() + 10), "task-2"), + ]); + + await runSseLogStream({ + taskId: "task-1", + initialCursor: DEFAULT_CURSOR, + bus, + source, + sink, + }); + + expect(sink.sent.map((e) => e.id)).toEqual(["a"]); + }); + + it("dedups an event that appears in both catch-up and the live buffer", async () => { + const t0 = new Date("2026-01-01T00:00:00.000Z"); + const shared = makeEvent("shared-1", t0); + const source = new FakeSource([shared]); + + const originalGet = source.getRecentTaskEvents.bind(source); + let firstCall = true; + source.getRecentTaskEvents = async (since, taskId, limit) => { + if (firstCall) { + firstCall = false; + // Same event also delivered live before catch-up resolves. + bus.emit(shared); + } + return originalGet(since, taskId, limit); + }; + + await runSseLogStream({ + taskId: "task-1", + initialCursor: DEFAULT_CURSOR, + bus, + source, + sink, + }); + + expect(sink.sent.filter((e) => e.id === "shared-1")).toHaveLength(1); + }); +}); diff --git a/packages/api/src/core/sse-log-stream.ts b/packages/api/src/core/sse-log-stream.ts new file mode 100644 index 0000000..0b267ac --- /dev/null +++ b/packages/api/src/core/sse-log-stream.ts @@ -0,0 +1,177 @@ +/** + * SSE log stream (ENG-1669 rework). + * + * Extracted from router.ts so the buffer-then-drain catch-up/live-merge + * logic can be unit tested without booting the full router/HTTP stack. + * + * Design (buffer-then-drain): + * 1. Subscribe to the live event bus FIRST, buffering everything into a + * local array (never sending directly) while catch-up runs. + * 2. Run the catch-up query in a loop, paginating until fully drained + * (no truncating LIMIT), sending each catch-up event and recording its + * id in a `sentIds` set. + * 3. Drain the live buffer collected during step 2, deduplicating against + * `sentIds`, sending anything not already delivered by catch-up. + * 4. Only after the buffer is drained do live events get sent directly. + * + * There is a single point that advances the cursor/dedup state: the + * `sentIds` set plus `lastSent` (the last emitted event's createdAt/id), + * updated exclusively inside `sendEvent`. Live events and catch-up events + * never race to independently advance a shared cursor. + */ +import type { BroadcastTaskEvent } from "./task-event-bus"; +import type { TaskEvent } from "./types"; + +export interface CursorPosition { + createdAt: Date; + id: string; +} + +export const DEFAULT_CURSOR: CursorPosition = { + createdAt: new Date(0), + id: "00000000-0000-0000-0000-000000000000", +}; + +export function parseCursor(cursor: string | null | undefined): CursorPosition { + if (!cursor) return DEFAULT_CURSOR; + const [createdAtStr, id] = cursor.split("|"); + const createdAt = new Date(createdAtStr ?? ""); + if (!id || Number.isNaN(createdAt.getTime())) return DEFAULT_CURSOR; + return { createdAt, id }; +} + +export function formatCursor(event: CursorPosition): string { + return `${event.createdAt.toISOString()}|${event.id}`; +} + +/** Minimal dependency surface needed from the task event bus. */ +export interface TaskEventBusLike { + onTaskEvent(listener: (event: BroadcastTaskEvent) => void): () => void; +} + +/** Minimal dependency surface needed from the DB layer. */ +export interface TaskEventSourceLike { + /** Fetch a single page of events strictly after `since`, oldest first. */ + getRecentTaskEvents( + since: CursorPosition, + taskId: string | undefined, + limit: number, + ): Promise<(TaskEvent & { taskStatus?: string })[]>; +} + +/** Abstraction over the SSE wire so tests don't need a real ReadableStream. */ +export interface SseSink { + /** Send one formatted SSE frame (id + data lines). */ + send(event: BroadcastTaskEvent): void; + /** True once the connection has been aborted/closed. */ + isActive(): boolean; +} + +export interface RunSseLogStreamOptions { + taskId?: string; + initialCursor: CursorPosition; + bus: TaskEventBusLike; + source: TaskEventSourceLike; + sink: SseSink; + /** Page size for catch-up pagination (default 50, matches prior LIMIT). */ + pageSize?: number; + /** Safety valve against pathological catch-up backlogs (default 200 pages). */ + maxPages?: number; +} + +export interface SseLogStreamHandle { + /** Unsubscribe from the live bus. Idempotent. */ + stop(): void; +} + +function matchesFilter(event: BroadcastTaskEvent, taskId?: string): boolean { + return !taskId || event.taskId === taskId; +} + +function eventKey(event: { id: string }): string { + return event.id; +} + +/** + * Starts the buffer-then-drain SSE delivery flow. Resolves once catch-up + * and buffer-drain are complete and the stream has switched to direct live + * delivery. Callers must call `stop()` on abort/cancel regardless of + * whether this promise has resolved. + */ +export async function runSseLogStream( + options: RunSseLogStreamOptions, +): Promise { + const { + taskId, + initialCursor, + bus, + source, + sink, + pageSize = 50, + maxPages = 200, + } = options; + + const sentIds = new Set(); + let liveMode = false; + let liveBuffer: BroadcastTaskEvent[] = []; + + function deliver(event: BroadcastTaskEvent): void { + if (sentIds.has(eventKey(event))) return; + sentIds.add(eventKey(event)); + sink.send(event); + } + + // Step 1: subscribe FIRST. While catch-up is running, buffer instead of + // sending directly so nothing is lost or duplicated relative to catch-up. + const unsubscribe = bus.onTaskEvent((event) => { + if (!sink.isActive()) return; + if (!matchesFilter(event, taskId)) return; + if (liveMode) { + deliver(event); + } else { + liveBuffer.push(event); + } + }); + + const handle: SseLogStreamHandle = { + stop() { + unsubscribe(); + }, + }; + + try { + // Step 2: paginate catch-up until fully drained (no truncating LIMIT). + let cursor = initialCursor; + for (let page = 0; page < maxPages; page++) { + if (!sink.isActive()) return handle; + const batch = await source.getRecentTaskEvents(cursor, taskId, pageSize); + if (batch.length === 0) break; + + for (const event of batch) { + if (!sink.isActive()) return handle; + deliver(event); + } + + const last = batch[batch.length - 1]!; + cursor = { createdAt: last.createdAt, id: last.id }; + + if (batch.length < pageSize) break; // fully drained this pass + } + } catch (err) { + console.error("[SSE] Error fetching catch-up events:", err); + } + + // Step 3: drain whatever arrived live while catch-up ran, deduped against + // everything catch-up already sent. + const bufferedDuringCatchUp = liveBuffer; + liveBuffer = []; + for (const event of bufferedDuringCatchUp) { + if (!sink.isActive()) return handle; + deliver(event); + } + + // Step 4: switch to direct live delivery. + liveMode = true; + + return handle; +} diff --git a/packages/api/src/integrations/db.ts b/packages/api/src/integrations/db.ts index b5b7304..9e3243c 100644 --- a/packages/api/src/integrations/db.ts +++ b/packages/api/src/integrations/db.ts @@ -307,13 +307,27 @@ export const db = { // having each connection poll the DB. Only pay for the task-status // lookup (RML-716) when someone is actually listening. if (taskEventBus.listenerCountTaskEvent > 0) { + // The status lookup is a best-effort enrichment: if it fails we still + // want to broadcast the event (without taskStatus) rather than drop + // the broadcast entirely. Keep this SELECT in its own try/catch so a + // transient failure here can never suppress the emit below. + let taskStatus: string | undefined; try { const [task] = await sql` SELECT status FROM tasks WHERE id = ${mapped.taskId} `; + taskStatus = task?.status as string | undefined; + } catch (err) { + console.error( + "[task-event-bus] Failed to look up task status for broadcast (degrading without enrichment):", + err, + ); + } + + try { taskEventBus.emitTaskEvent({ ...mapped, - taskStatus: task?.status as string | undefined, + taskStatus, }); } catch (err) { // Broadcasting is best-effort; never fail the write because of it. diff --git a/packages/api/src/router.ts b/packages/api/src/router.ts index 4145b55..5ce4a01 100644 --- a/packages/api/src/router.ts +++ b/packages/api/src/router.ts @@ -8,6 +8,13 @@ import { JobStatus, } from "./core/types"; import { taskEventBus } from "./core/task-event-bus"; +import { + runSseLogStream, + parseCursor, + formatCursor, + type SseLogStreamHandle, + type SseSink, +} from "./core/sse-log-stream"; import { Orchestrator } from "./core/orchestrator"; import { TaskRunner } from "./core/task-runner"; import { db } from "./integrations/db"; @@ -4254,6 +4261,14 @@ function getLogLevel(eventType: string): "INFO" | "SUCCESS" | "WARN" | "ERROR" { return "INFO"; } +/** + * Backpressure threshold: if the underlying sink reports it is this far (or + * more) behind, the client is not draining fast enough. We drop the + * connection rather than let the buffer grow unbounded (MED finding, + * ENG-1669 rework). + */ +const SSE_BACKPRESSURE_THRESHOLD_BYTES = 0; + /** * GET /api/logs/stream - SSE endpoint for real-time task events * Query params: @@ -4263,34 +4278,36 @@ route("GET", "/api/logs/stream", async (req) => { const url = new URL(req.url); const taskId = url.searchParams.get("taskId") || undefined; - const DEFAULT_CURSOR = { - createdAt: new Date(0), - id: "00000000-0000-0000-0000-000000000000", - }; - - function parseCursor(cursor: string | null): { createdAt: Date; id: string } { - if (!cursor) return DEFAULT_CURSOR; - const [createdAtStr, id] = cursor.split("|"); - const createdAt = new Date(createdAtStr); - if (!id || Number.isNaN(createdAt.getTime())) return DEFAULT_CURSOR; - return { createdAt, id }; - } - - function formatCursor(event: { createdAt: Date; id: string }): string { - return `${event.createdAt.toISOString()}|${event.id}`; - } - - // Track last cursor sent (SSE "Last-Event-ID" compatible) - const initialCursor = + const initialCursorParam = url.searchParams.get("cursor") || req.headers.get("last-event-id"); - let lastCursor = parseCursor(initialCursor); + const initialCursor = parseCursor(initialCursorParam); + let isActive = true; + let streamHandle: SseLogStreamHandle | null = null; + let keepAlive: ReturnType | null = null; // Create a readable stream for SSE const stream = new ReadableStream({ - async start(controller) { + start(controller) { const encoder = new TextEncoder(); + function cleanup(): void { + isActive = false; + streamHandle?.stop(); + if (keepAlive) clearInterval(keepAlive); + try { + controller.close(); + } catch { + // Already closed (e.g. client disconnected mid-write) - ignore. + } + } + + // HIGH 2 fix: register the abort cleanup IMMEDIATELY, before any + // await, so a disconnect during catch-up can never leak the live + // listener or the keep-alive timer. try/finally below also runs + // cleanup on any unexpected throw from the async flow. + req.signal?.addEventListener("abort", cleanup); + // Send initial connection message controller.enqueue( encoder.encode( @@ -4298,91 +4315,96 @@ route("GET", "/api/logs/stream", async (req) => { ), ); - function sendEvent(event: TaskEvent & { taskStatus?: string }): void { - const cursor = formatCursor({ - createdAt: event.createdAt, - id: event.id, - }); + const sink: SseSink = { + isActive: () => isActive, + send(event) { + // MED fix: backpressure check before enqueue. `desiredSize` is + // negative/near-zero when the client isn't draining fast enough. + if ( + controller.desiredSize !== null && + controller.desiredSize <= SSE_BACKPRESSURE_THRESHOLD_BYTES + ) { + console.error( + "[SSE] Backpressure threshold exceeded, closing slow consumer", + ); + cleanup(); + return; + } - controller.enqueue(encoder.encode(`id: ${cursor}\n`)); - controller.enqueue( - encoder.encode( - `data: ${JSON.stringify({ - type: "event", - id: event.id, - taskId: event.taskId, - eventType: event.eventType, - agent: event.agent, - message: event.outputSummary || event.eventType, - timestamp: event.createdAt, - level: getLogLevel(event.eventType), - tokensUsed: event.tokensUsed, - durationMs: event.durationMs, - // Include current task status for real-time UI updates (RML-716) - taskStatus: event.taskStatus, - })}\n\n`, - ), - ); + const cursor = formatCursor({ + createdAt: event.createdAt, + id: event.id, + }); - lastCursor = { createdAt: event.createdAt, id: event.id }; - } + try { + controller.enqueue(encoder.encode(`id: ${cursor}\n`)); + controller.enqueue( + encoder.encode( + `data: ${JSON.stringify({ + type: "event", + id: event.id, + taskId: event.taskId, + eventType: event.eventType, + agent: event.agent, + message: event.outputSummary || event.eventType, + timestamp: event.createdAt, + level: getLogLevel(event.eventType), + tokensUsed: event.tokensUsed, + durationMs: event.durationMs, + // Include current task status for real-time UI updates (RML-716) + taskStatus: event.taskStatus, + })}\n\n`, + ), + ); + } catch (err) { + console.error("[SSE] Error sending event:", err); + } + }, + }; - // ENG-1669: push-based delivery via shared EventEmitter instead of - // polling the DB every 2s per connection. - const unsubscribe = taskEventBus.onTaskEvent((event) => { - if (!isActive) return; - if (taskId && event.taskId !== taskId) return; - // Skip events already delivered by the catch-up query below. - if ( - event.createdAt.getTime() < lastCursor.createdAt.getTime() || - (event.createdAt.getTime() === lastCursor.createdAt.getTime() && - event.id <= lastCursor.id) - ) { - return; - } + // ENG-1669 rework: buffer-then-drain catch-up (paginated, no + // truncating LIMIT) merged with live delivery via a single dedup set, + // instead of two independently-advancing cursors racing on + // timestamp. See core/sse-log-stream.ts for the full design. + void (async () => { try { - sendEvent(event); + streamHandle = await runSseLogStream({ + taskId, + initialCursor, + bus: taskEventBus, + source: db, + sink, + }); } catch (err) { - console.error("[SSE] Error sending event:", err); + console.error("[SSE] Fatal error running log stream:", err); + } finally { + if (!isActive) return; // already cleaned up via abort + // Start keep-alive only once catch-up/drain has completed and + // we're live (or failed fatally) - matches prior behavior. + keepAlive = setInterval(() => { + if (!isActive) { + if (keepAlive) clearInterval(keepAlive); + return; + } + try { + controller.enqueue(encoder.encode(`: keep-alive\n\n`)); + } catch (err) { + console.error("[SSE] Error sending keep-alive:", err); + } + }, 30000); } - }); + })(); - // One-time catch-up for events missed since Last-Event-ID/cursor - // (also covers events written between connect and subscribe). - try { - const missed = await db.getRecentTaskEvents(lastCursor, taskId); - for (const event of missed) { - if (!isActive) break; - // Skip anything already delivered live while the query ran. - if ( - event.createdAt.getTime() < lastCursor.createdAt.getTime() || - (event.createdAt.getTime() === lastCursor.createdAt.getTime() && - event.id <= lastCursor.id) - ) { - continue; - } - sendEvent(event); - } - } catch (err) { - console.error("[SSE] Error fetching catch-up events:", err); - } - - // Send keep-alive ping every 30 seconds - const keepAlive = setInterval(() => { - if (!isActive) { - clearInterval(keepAlive); - return; - } - controller.enqueue(encoder.encode(`: keep-alive\n\n`)); - }, 30000); - - // Cleanup on close (handled by abort signal) - req.signal?.addEventListener("abort", () => { - isActive = false; - unsubscribe(); - clearInterval(keepAlive); - controller.close(); - }); + // ReadableStream.cancel() is invoked by the platform when the + // consumer stops reading without an abort signal firing (e.g. some + // proxies/runtimes). Route it through the same cleanup so the + // listener/timer never leak. + // (Bun/undici call `cancel` on the underlying source, wired below.) + }, + cancel() { + isActive = false; + streamHandle?.stop(); + if (keepAlive) clearInterval(keepAlive); }, });