181 lines
6.7 KiB
TypeScript
181 lines
6.7 KiB
TypeScript
import { test, expect } from "vitest";
|
|
import { buildApp } from "../src/app.js";
|
|
import { loadConfig } from "../src/config.js";
|
|
import { SseHub } from "../src/sse/sse-hub.js";
|
|
|
|
async function readUntil(
|
|
reader: ReadableStreamDefaultReader<Uint8Array>,
|
|
predicate: (text: string) => boolean,
|
|
): Promise<string> {
|
|
const decoder = new TextDecoder();
|
|
let text = "";
|
|
const deadline = Date.now() + 1_000;
|
|
while (Date.now() < deadline) {
|
|
const result = await Promise.race([
|
|
reader.read(),
|
|
new Promise<{ timeout: true }>((resolve) =>
|
|
setTimeout(() => resolve({ timeout: true }), Math.max(0, deadline - Date.now())),
|
|
),
|
|
]);
|
|
if ("timeout" in result || result.done) break;
|
|
text += decoder.decode(result.value);
|
|
if (predicate(text)) break;
|
|
}
|
|
return text;
|
|
}
|
|
|
|
async function captureReplay(options: { query?: string; lastEventId?: string }): Promise<string> {
|
|
const hub = new SseHub();
|
|
hub.publish("s1", "info", { type: "info", text: "one" });
|
|
hub.publish("s1", "info", { type: "info", text: "two" });
|
|
hub.publish("s1", "info", { type: "info", text: "three" });
|
|
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), {
|
|
hub,
|
|
thtRunner: {} as any,
|
|
});
|
|
await app.listen({ port: 0, host: "127.0.0.1" });
|
|
const port = (app.server.address() as { port: number }).port;
|
|
const controller = new AbortController();
|
|
|
|
try {
|
|
const response = await fetch(
|
|
`http://127.0.0.1:${port}/sessions/s1/events${options.query ?? ""}`,
|
|
{
|
|
headers: options.lastEventId ? { "Last-Event-ID": options.lastEventId } : undefined,
|
|
signal: controller.signal,
|
|
},
|
|
);
|
|
const reader = response.body!.getReader();
|
|
const body = await readUntil(reader, (text) => text.includes('"three"'));
|
|
await reader.cancel();
|
|
return body;
|
|
} finally {
|
|
controller.abort();
|
|
await app.close();
|
|
}
|
|
}
|
|
|
|
async function expectStreamEnd(reader: ReadableStreamDefaultReader<Uint8Array>): Promise<void> {
|
|
const result = await Promise.race([
|
|
reader.read(),
|
|
new Promise<{ timeout: true }>((resolve) =>
|
|
setTimeout(() => resolve({ timeout: true }), 1_000),
|
|
),
|
|
]);
|
|
expect(result).not.toEqual({ timeout: true });
|
|
expect("done" in result && result.done).toBe(true);
|
|
}
|
|
|
|
test("SSE emits ids and honors the native Last-Event-ID replay cursor", async () => {
|
|
const body = await captureReplay({ lastEventId: "1" });
|
|
|
|
expect(body).not.toContain('"one"');
|
|
expect(body).toContain("id: 2\nevent: info\n");
|
|
expect(body).toContain('data: {"type":"info","text":"two"}');
|
|
expect(body).toContain("id: 3\nevent: info\n");
|
|
expect(body).toContain('data: {"type":"info","text":"three"}');
|
|
});
|
|
|
|
test("SSE honors the manual lastEventId query cursor", async () => {
|
|
const body = await captureReplay({ query: "?lastEventId=2" });
|
|
|
|
expect(body).not.toContain('"one"');
|
|
expect(body).not.toContain('"two"');
|
|
expect(body).toContain("id: 3\nevent: info\n");
|
|
expect(body).toContain('data: {"type":"info","text":"three"}');
|
|
});
|
|
|
|
test("SSE uses the newer valid cursor when header and query are both present", async () => {
|
|
const body = await captureReplay({ query: "?lastEventId=2", lastEventId: "1" });
|
|
|
|
expect(body).not.toContain('"two"');
|
|
expect(body).toContain("id: 3\nevent: info\n");
|
|
});
|
|
|
|
test("SSE resets a cursor from an older process and emits a buffered pending gate once", async () => {
|
|
const hub = new SseHub();
|
|
const descriptor = { id: "gate-fresh", widget: "select", title: "Choose" };
|
|
hub.publish("s1", "info", { type: "info", text: "fresh generation" });
|
|
hub.publish("s1", "ui_request", { type: "ui_request", ui_request: descriptor });
|
|
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), {
|
|
hub,
|
|
thtRunner: {} as any,
|
|
mgr: {
|
|
get: () => ({ bridge: { pendingWidget: () => ({ ...descriptor }) } }),
|
|
} as any,
|
|
});
|
|
await app.listen({ port: 0, host: "127.0.0.1" });
|
|
const port = (app.server.address() as { port: number }).port;
|
|
const controller = new AbortController();
|
|
|
|
try {
|
|
const response = await fetch(`http://127.0.0.1:${port}/sessions/s1/events`, {
|
|
headers: { "Last-Event-ID": "900" },
|
|
signal: controller.signal,
|
|
});
|
|
const reader = response.body!.getReader();
|
|
const body = await readUntil(reader, (text) => text.includes('"gate-fresh"'));
|
|
await reader.cancel();
|
|
|
|
expect(body).toContain("id: 1\nevent: info\n");
|
|
expect(body).toContain('data: {"type":"info","text":"fresh generation"}');
|
|
expect(body.match(/event: ui_request/g)).toHaveLength(1);
|
|
expect(body).toContain("id: 2\nevent: ui_request\n");
|
|
expect(body).toContain('"ui_request":{"id":"gate-fresh","widget":"select","title":"Choose"}');
|
|
} finally {
|
|
controller.abort();
|
|
await app.close();
|
|
}
|
|
});
|
|
|
|
test("clear ends every live SSE response and a cursor reconnect replays post-clear events", async () => {
|
|
const hub = new SseHub();
|
|
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), {
|
|
hub,
|
|
thtRunner: {} as any,
|
|
});
|
|
await app.listen({ port: 0, host: "127.0.0.1" });
|
|
const port = (app.server.address() as { port: number }).port;
|
|
const controllers = [new AbortController(), new AbortController(), new AbortController()];
|
|
|
|
try {
|
|
const responses = await Promise.all(controllers.slice(0, 2).map((controller) => fetch(
|
|
`http://127.0.0.1:${port}/sessions/s1/events`,
|
|
{ signal: controller.signal },
|
|
)));
|
|
const readers = responses.map((response) => response.body!.getReader());
|
|
|
|
expect(hub.publish("s1", "info", { type: "info", text: "before" })).toBe(1);
|
|
const initial = await Promise.all(readers.map((reader) =>
|
|
readUntil(reader, (text) => text.includes('"before"'))));
|
|
expect(initial.every((body) => body.includes("id: 1\nevent: info\n"))).toBe(true);
|
|
|
|
hub.clear("s1");
|
|
await Promise.all(readers.map(expectStreamEnd));
|
|
|
|
expect(hub.publish("s1", "info", { type: "info", text: "after" })).toBe(2);
|
|
expect(hub.publish("s1", "ui_request", {
|
|
type: "ui_request",
|
|
ui_request: { id: "gate-1", widget: "select" },
|
|
})).toBe(3);
|
|
|
|
const reconnected = await fetch(
|
|
`http://127.0.0.1:${port}/sessions/s1/events`,
|
|
{
|
|
headers: { "Last-Event-ID": "1" },
|
|
signal: controllers[2].signal,
|
|
},
|
|
);
|
|
const reconnectReader = reconnected.body!.getReader();
|
|
const replay = await readUntil(reconnectReader, (text) => text.includes('"gate-1"'));
|
|
expect(replay).toContain("id: 2\nevent: info\n");
|
|
expect(replay).toContain('data: {"type":"info","text":"after"}');
|
|
expect(replay).toContain("id: 3\nevent: ui_request\n");
|
|
expect(replay).toContain('"ui_request":{"id":"gate-1","widget":"select"}');
|
|
await reconnectReader.cancel();
|
|
} finally {
|
|
for (const controller of controllers) controller.abort();
|
|
await app.close();
|
|
}
|
|
});
|