feat: prepare 2.6.1-rc5 shared databases and build SSE
This commit is contained in:
@@ -0,0 +1,330 @@
|
||||
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();
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user