Files
kuber/server/build-event-stream.ts

256 lines
9.5 KiB
TypeScript

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<Uint8Array>;
stop: () => void;
};
type Group = {
members: Set<Subscriber>;
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<string, Group>();
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<Response> {
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<Uint8Array>(
{
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<void> {
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<void>((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<void>((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();
}
}
}