132 lines
4.9 KiB
TypeScript
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}`;
|
|
}
|