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

331 lines
11 KiB
TypeScript

import { describe, expect, test } from "bun:test";
import { BuildEventStreamHub } from "../../server/build-event-stream";
import type { BuildStatus } from "../../shared/build-protocol";
const base: BuildStatus = {
version: 1,
id: "a",
state: "running",
createdAt: "now",
};
const tick = () => Bun.sleep(10);
describe("build SSE hub", () => {
test("reports an explicit gap when the resume cursor predates retained logs", async () => {
const abort = new AbortController();
const hub = new BuildEventStreamHub(
{
getBuildStatus: async () => base,
reconcileBuild: async () => base,
getBuildEvents: async (_id, after) => [
{ type: "status" as const, status: base },
...[4, 5]
.filter((sequence) => sequence > (after ?? 0))
.map((sequence) => ({
type: "log" as const,
id: "a",
sequence,
message: `log ${sequence}\n`,
})),
],
},
{ pollMs: 100, heartbeatMs: 1_000 },
);
const response = await hub.open(
"a",
new Request("https://test/builds/a/events", {
headers: { "last-event-id": "1" },
signal: abort.signal,
}),
);
const reader = response.body!.getReader();
const decoder = new TextDecoder();
const first = decoder.decode((await reader.read()).value);
expect(first).toContain("event: gap\n");
expect(first).toContain('"missing":2');
expect(first).toContain("id: 3\n");
expect(decoder.decode((await reader.read()).value)).toContain("id: 4\nevent: log");
abort.abort();
await reader.cancel();
});
test("coalesces reconciliation, resumes logs, dedupes status and drains terminal", async () => {
let calls = 0;
let status = base;
const hub = new BuildEventStreamHub(
{
getBuildStatus: async () => status,
reconcileBuild: async () => {
calls++;
return status;
},
getBuildEvents: async (_id, after) => [
{ type: "status" as const, status: base },
...[1, 2]
.filter((sequence) => sequence > (after ?? 0))
.map((sequence) => ({
type: "log" as const,
id: "a",
sequence,
message: `log ${sequence}\n`,
})),
],
},
{ pollMs: 20, heartbeatMs: 40 },
);
const a = new AbortController();
const b = new AbortController();
const request = (signal: AbortSignal, after: string) =>
new Request(`https://test/builds/a/events?after=${after}`, { signal });
const first = await hub.open("a", request(a.signal, "1"));
const second = await hub.open("a", request(b.signal, "2"));
expect(first.headers.get("content-type")).toContain("text/event-stream");
const reader1 = first.body!.getReader();
const reader2 = second.body!.getReader();
const text = new TextDecoder();
expect(text.decode((await reader1.read()).value)).toContain(
"id: 2\nevent: log\n",
);
expect(text.decode((await reader1.read()).value)).toContain(
"event: status\n",
);
expect(text.decode((await reader2.read()).value)).toContain(
"event: status\n",
);
await tick();
expect(calls).toBe(1);
status = { ...base, state: "succeeded" };
const remaining = async (reader: typeof reader1) => {
let output = "";
for (;;) {
const { value, done } = await reader.read();
if (done) return output;
output += text.decode(value);
}
};
expect(await remaining(reader1)).toContain('"state":"succeeded"');
expect(await remaining(reader2)).toContain('"state":"succeeded"');
a.abort();
b.abort();
});
test("validates cursor, stops aborted subscribers and emits heartbeat comments", async () => {
const hub = new BuildEventStreamHub(
{
getBuildStatus: async () => base,
reconcileBuild: async () => base,
getBuildEvents: async () => [],
},
{ pollMs: 5, heartbeatMs: 5 },
);
await expect(
hub.open("a", new Request("https://test/events?after=-1")),
).rejects.toThrow(RangeError);
const abort = new AbortController();
const response = await hub.open(
"a",
new Request("https://test/events", { signal: abort.signal }),
);
const reader = response.body!.getReader();
expect(new TextDecoder().decode((await reader.read()).value)).toContain(
"event: status",
);
expect(new TextDecoder().decode((await reader.read()).value)).toBe(
": heartbeat\n\n",
);
abort.abort();
expect((await reader.read()).done).toBe(true);
});
test("reads persisted progress when another replica owns reconciliation", async () => {
let status = base;
let sequence = 0;
let reconciles = 0;
const hub = new BuildEventStreamHub(
{
getBuildStatus: async () => status,
reconcileBuild: async () => {
reconciles++;
throw new Error("build lease held by another replica");
},
getBuildEvents: async (_id, after) =>
sequence > (after ?? 0)
? [
{
type: "log" as const,
id: "a",
sequence,
message: "remote log\n",
},
]
: [],
},
{ pollMs: 5, heartbeatMs: 50 },
);
const abort = new AbortController();
const response = await hub.open(
"a",
new Request("https://test/events", { signal: abort.signal }),
);
const reader = response.body!.getReader();
const decoder = new TextDecoder();
expect(decoder.decode((await reader.read()).value)).toContain(
'"state":"running"',
);
sequence = 1;
expect(decoder.decode((await reader.read()).value)).toContain(
"id: 1\nevent: log",
);
status = { ...base, state: "succeeded" };
let terminalFrame = "";
for (;;) {
const { value, done } = await reader.read();
if (done) break;
terminalFrame += decoder.decode(value);
}
expect(terminalFrame).toContain('"state":"succeeded"');
expect(reconciles).toBeGreaterThan(0);
abort.abort();
});
test("heartbeats during slow reconciliation and tears down on abort", async () => {
const blocked = new Promise<BuildStatus>(() => {});
const signals: AbortSignal[] = [];
let observations = 0;
const hub = new BuildEventStreamHub(
{
getBuildStatus: async () => base,
reconcileBuild: async (_id, options) => {
if (options?.signal) signals.push(options.signal);
return blocked;
},
getBuildEvents: async () => {
observations++;
return [];
},
},
{ pollMs: 5, heartbeatMs: 5 },
);
const abort = new AbortController();
const response = await hub.open(
"a",
new Request("https://test/events", { signal: abort.signal }),
);
const reader = response.body!.getReader();
expect(new TextDecoder().decode((await reader.read()).value)).toContain(
"event: status",
);
expect(new TextDecoder().decode((await reader.read()).value)).toBe(
": heartbeat\n\n",
);
await Bun.sleep(25);
expect(observations).toBeGreaterThan(1);
expect(signals).toHaveLength(1);
abort.abort();
for (;;) {
if ((await reader.read()).done) break;
}
expect(signals[0]?.aborted).toBe(true);
const stoppedAt = observations;
await Bun.sleep(20);
expect(observations).toBe(stoppedAt);
const next = new AbortController();
const reopened = await hub.open(
"a",
new Request("https://test/events", { signal: next.signal }),
);
const nextReader = reopened.body!.getReader();
expect(new TextDecoder().decode((await nextReader.read()).value)).toContain(
"event: status",
);
expect(signals).toHaveLength(2);
expect(signals[1]?.aborted).toBe(false);
next.abort();
await nextReader.cancel();
});
test("bounds shared reconciliation writes while observing two subscribers independently", async () => {
let reconciles = 0;
let observations = 0;
const hub = new BuildEventStreamHub(
{
getBuildStatus: async () => base,
reconcileBuild: async () => {
reconciles++;
return base;
},
getBuildEvents: async () => {
observations++;
return [];
},
},
{ pollMs: 10, reconcileMs: 80, heartbeatMs: 1_000 },
);
const a = new AbortController();
const b = new AbortController();
const first = await hub.open(
"a",
new Request("https://test/events", { signal: a.signal }),
);
const second = await hub.open(
"a",
new Request("https://test/events", { signal: b.signal }),
);
expect(
new TextDecoder().decode((await first.body!.getReader().read()).value),
).toContain("event: status");
expect(
new TextDecoder().decode((await second.body!.getReader().read()).value),
).toContain("event: status");
await Bun.sleep(45);
expect(reconciles).toBe(1); // One lease/write attempt, not one per poll or subscriber.
expect(observations).toBeGreaterThanOrEqual(3);
for (let i = 0; i < 15 && reconciles < 2; i++) await tick();
expect(reconciles).toBe(2);
a.abort();
b.abort();
});
test("closes a slow subscriber at the bounded queue limit for resumable logs", async () => {
const hub = new BuildEventStreamHub(
{
getBuildStatus: async () => base,
reconcileBuild: async () => base,
getBuildEvents: async (_id, after) =>
Array.from({ length: 40 }, (_, index) => index + 1)
.filter((sequence) => sequence > (after ?? 0))
.map((sequence) => ({
type: "log" as const,
id: "a",
sequence,
message: `log ${sequence}\n`,
})),
},
{ pollMs: 5, heartbeatMs: 50 },
);
const response = await hub.open("a", new Request("https://test/events"));
await tick(); // Allow the producer to fill its queue before consumption.
const reader = response.body!.getReader();
const frames: string[] = [];
for (;;) {
const { value, done } = await reader.read();
if (done) break;
frames.push(new TextDecoder().decode(value));
}
expect(frames).toHaveLength(16);
expect(frames[0]).toContain("id: 1\nevent: log");
expect(frames[15]).toContain("id: 16\nevent: log");
const resumed = await hub.open(
"a",
new Request("https://test/events", {
headers: { "last-event-id": "16" },
}),
);
const resumedReader = resumed.body!.getReader();
expect(
new TextDecoder().decode((await resumedReader.read()).value),
).toContain("id: 17\nevent: log");
await resumedReader.cancel();
});
});