import type { BuildController } from "./build-controller"; import type { BuildStatus } from "../shared/build-protocol"; const terminal = (status: BuildStatus) => status.state === "succeeded" || status.state === "failed"; function cursor(value: string | null): number { if (value === null) return 0; if (!/^(0|[1-9]\d*)$/.test(value) || !Number.isSafeInteger(Number(value))) throw new RangeError("Event cursor must be a non-negative safe integer"); return Number(value); } type Subscriber = { sequence: number; status?: string; lastActivity: number; controller: ReadableStreamDefaultController; stop: () => void; }; type Group = { members: Set; abort: AbortController; nextReconcileAt: number; reconciling: boolean; }; /** One hub per app/controller instance; call only after the route's capability and workspace checks. * The authenticated GET /api/v2/builds/:id/events route should branch on * Accept: text/event-stream, then return hub.open(id, request); retain JSON otherwise. * A hub coalesces reconciliation and persisted event reads per physical build on this * replica. BuildController's lease coordinates reconciliation across replicas; its * persisted event store supplies logs even when a different replica owns the job. */ export class BuildEventStreamHub { private readonly builds = new Map(); constructor( private readonly controller: Pick< BuildController, "getBuildStatus" | "getBuildEvents" | "reconcileBuild" >, private readonly options: { pollMs?: number; heartbeatMs?: number; reconcileMs?: number; } = {}, ) {} async open(id: string, request: Request): Promise { const header = request.headers.get("last-event-id"); const query = new URL(request.url).searchParams.get("after"); const after = cursor(header ?? query); // Validate existence before committing response headers. await this.controller.getBuildStatus(id); let subscriber: Subscriber; const stream = new ReadableStream( { start: (controller) => { let group: Group; const stop = () => { request.signal.removeEventListener("abort", stop); group.members.delete(subscriber); if (group.members.size === 0 && this.builds.get(id) === group) { this.builds.delete(id); group.abort.abort(); } try { controller.close(); } catch { /* Already closed or cancelled. */ } }; subscriber = { sequence: after, lastActivity: Date.now(), controller, stop, }; const existing = this.builds.get(id); if (existing) { group = existing; group.members.add(subscriber); } else { group = { members: new Set([subscriber]), abort: new AbortController(), nextReconcileAt: 0, reconciling: false, }; this.builds.set(id, group); void this.run(id, group); } request.signal.addEventListener("abort", stop, { once: true }); if (request.signal.aborted) stop(); }, cancel: () => subscriber.stop(), }, { highWaterMark: 16 }, ); return new Response(stream, { headers: { "content-type": "text/event-stream; charset=utf-8", "cache-control": "no-cache, no-transform", connection: "keep-alive", "x-accel-buffering": "no", }, }); } private async run(id: string, group: Group): Promise { const encoder = new TextEncoder(); const pollMs = this.options.pollMs ?? 1_000; const reconcileMs = this.options.reconcileMs ?? 7_000; const { members } = group; // Bun's default HTTP idle timeout is 10 seconds; keep bytes flowing well // inside it even while Kubernetes reconciliation makes no progress. const heartbeatMs = this.options.heartbeatMs ?? 5_000; const heartbeat = setInterval(() => { if (group.abort.signal.aborted) return; for (const member of members) { if (Date.now() - member.lastActivity < heartbeatMs) continue; if ((member.controller.desiredSize ?? 0) <= 0) member.stop(); else { member.controller.enqueue(encoder.encode(": heartbeat\n\n")); member.lastActivity = Date.now(); } } }, heartbeatMs); heartbeat.unref?.(); try { while (!group.abort.signal.aborted) { try { const initial = await this.controller.getBuildStatus(id); if (group.abort.signal.aborted) break; // Reconciliation is shared per build and never blocks persisted progress. // The controller's lease still coordinates attempts across replicas. if ( !terminal(initial) && !group.reconciling && Date.now() >= group.nextReconcileAt ) { group.reconciling = true; group.nextReconcileAt = Date.now() + reconcileMs; const signal = group.abort.signal; let onAbort = () => {}; const aborted = new Promise((resolve) => { onAbort = resolve; signal.addEventListener("abort", onAbort, { once: true }); }); const reconciliation = this.controller .reconcileBuild(id, { signal }) .catch(() => { // Another replica may hold the lease; keep observing its progress. }); // Drop this group's references even if an external API ignores abort // and leaves the reconciliation promise unresolved indefinitely. void Promise.race([reconciliation, aborted]).finally(() => { signal.removeEventListener("abort", onAbort); group.reconciling = false; }); } const status = await this.controller.getBuildStatus(id); if (group.abort.signal.aborted) break; const oldest = Math.min( ...[...members].map((member) => member.sequence), ); const events = await this.controller.getBuildEvents(id, oldest); if (group.abort.signal.aborted) break; const firstLogSequence = events.find( (event) => event.type === "log", )?.sequence; for (const member of members) { const send = (frame: string) => { // A slow consumer reconnects from its last delivered sequence instead // of retaining an unbounded in-memory backlog on this replica. if ((member.controller.desiredSize ?? 0) <= 0) { member.stop(); return false; } member.controller.enqueue(encoder.encode(frame)); member.lastActivity = Date.now(); return true; }; if ( firstLogSequence !== undefined && firstLogSequence > member.sequence + 1 ) { const gap = { type: "gap", after: member.sequence, before: firstLogSequence, missing: firstLogSequence - member.sequence - 1, message: `Build log history was trimmed; ${firstLogSequence - member.sequence - 1} log event(s) before sequence ${firstLogSequence} are unavailable.`, }; if ( !send( `id: ${firstLogSequence - 1}\nevent: gap\ndata: ${JSON.stringify(gap)}\n\n`, ) ) continue; // Advance over the unavailable range so reconnects do not repeat // the marker. The cursor now refers to the last unavailable // sequence; delivered log records retain their normal IDs. member.sequence = firstLogSequence - 1; } for (const event of events) { if (event.type !== "log" || event.sequence <= member.sequence) continue; if ( !send( `id: ${event.sequence}\nevent: log\ndata: ${JSON.stringify(event)}\n\n`, ) ) break; member.sequence = event.sequence; } if (!members.has(member)) continue; const fingerprint = JSON.stringify(status); if (member.status !== fingerprint) { if ( !send( `event: status\ndata: ${JSON.stringify({ type: "status", status })}\n\n`, ) ) continue; member.status = fingerprint; } if (terminal(status)) member.stop(); } } catch { // A transient API/store failure closes streams so clients reconnect with // their last sequence; no silent terminal state is fabricated. for (const member of members) member.stop(); } if (group.abort.signal.aborted) break; await new Promise((resolve) => { const stop = () => { clearTimeout(timer); resolve(); }; const timer = setTimeout(() => { group.abort.signal.removeEventListener("abort", stop); resolve(); }, pollMs); group.abort.signal.addEventListener("abort", stop, { once: true }); }); } } finally { clearInterval(heartbeat); group.abort.abort(); } } }