256 lines
9.5 KiB
TypeScript
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();
|
|
}
|
|
}
|
|
}
|