Files
pikasTech-HWLAB/internal/cloud/workbench-debug-fake-sse.test.ts
2026-07-09 15:46:28 +02:00

132 lines
4.9 KiB
TypeScript

import assert from "node:assert/strict";
import { createServer } from "node:http";
import { test } from "bun:test";
import { handleWorkbenchDebugFakeSseHttp } from "./workbench-debug-fake-sse.ts";
test("workbench debug fake SSE API resets, steps, appends, runs, and streams typed events", async () => {
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchDebugFakeSseHttp(request, response, url, { logger: null });
});
await listen(server);
const baseUrl = serverUrl(server);
const queueId = "api-test";
const abort = new AbortController();
try {
const initial = await jsonFetch(`${baseUrl}/v1/workbench/debug/fake-sse?queueId=${queueId}`);
assert.equal(initial.ok, true);
assert.equal(initial.queue.queueId, queueId);
assert.equal(initial.queue.cursor, 0);
assert.equal(initial.queue.remaining, 6);
const stream = await fetch(`${baseUrl}/v1/workbench/debug/fake-sse/events?queueId=${queueId}`, { signal: abort.signal });
assert.equal(stream.status, 200);
assert.match(stream.headers.get("content-type") ?? "", /text\/event-stream/u);
const reader = stream.body?.getReader();
assert.ok(reader, "SSE response must expose a readable stream");
const connected = await readUntil(reader, "workbench.connected");
assert.match(connected, /workbench\.connected/u);
const reset = await jsonFetch(`${baseUrl}/v1/workbench/debug/fake-sse/reset`, {
method: "POST",
body: JSON.stringify({ queueId, sequenceId: "trace-card-basic" })
});
assert.equal(reset.ok, true);
assert.equal(reset.queue.remaining, 6);
const next = await jsonFetch(`${baseUrl}/v1/workbench/debug/fake-sse/next`, {
method: "POST",
body: JSON.stringify({ queueId })
});
assert.equal(next.ok, true);
assert.equal(next.delivered.eventName, "workbench.message.snapshot");
assert.equal(next.delivered.deliveredTo, 1);
assert.equal(next.queue.cursor, 1);
const firstEvent = await readUntil(reader, "workbench.message.snapshot");
assert.match(firstEvent, /workbench-realtime-authority-v2/u);
assert.match(firstEvent, /msg_debug_agent/u);
const append = await jsonFetch(`${baseUrl}/v1/workbench/debug/fake-sse/append`, {
method: "POST",
body: JSON.stringify({
queueId,
events: [{
type: "trace.event",
sessionId: "ses_debug_fake_sse",
threadId: "thr_debug_fake_sse",
traceId: "trc_debug_fake_sse",
event: {
traceId: "trc_debug_fake_sse",
projectedSeq: 99,
createdAt: "2026-07-09T08:00:09.000Z",
label: "assistant:message",
type: "assistant_message",
status: "running",
message: "appended from api test"
}
}]
})
});
assert.equal(append.ok, true);
assert.equal(append.appended, 1);
assert.equal(append.queue.eventCount, 7);
assert.equal(append.queue.remaining, 6);
const run = await jsonFetch(`${baseUrl}/v1/workbench/debug/fake-sse/run`, {
method: "POST",
body: JSON.stringify({ queueId })
});
assert.equal(run.ok, true);
assert.equal(run.delivered.count, 6);
assert.equal(run.queue.remaining, 0);
const runStream = await readUntil(reader, "appended from api test");
assert.match(runStream, /workbench\.trace\.event/u);
} finally {
abort.abort();
await close(server);
}
});
async function jsonFetch(url: string, init: RequestInit = {}) {
const response = await fetch(url, {
...init,
headers: {
accept: "application/json",
...(init.body ? { "content-type": "application/json" } : {}),
...(init.headers ?? {})
}
});
const body = await response.json();
assert.equal(response.ok, true, JSON.stringify(body));
return body as any;
}
async function readUntil(reader: ReadableStreamDefaultReader<Uint8Array>, pattern: string): Promise<string> {
const decoder = new TextDecoder();
let text = "";
const deadline = Date.now() + 3000;
while (!text.includes(pattern)) {
if (Date.now() > deadline) throw new Error(`SSE stream did not include ${pattern}: ${text}`);
const next = await reader.read();
if (next.done) break;
text += decoder.decode(next.value, { stream: true });
}
return text;
}
async function listen(server: ReturnType<typeof createServer>) {
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
}
async function close(server: ReturnType<typeof createServer>) {
await new Promise<void>((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
}
function serverUrl(server: ReturnType<typeof createServer>): string {
const address = server.address();
assert.equal(typeof address, "object");
assert.ok(address && typeof address.port === "number");
return `http://127.0.0.1:${address.port}`;
}