Files
pikasTech-HWLAB/internal/cloud/workbench-kafka-sse-debug.test.ts
2026-07-10 15:58:59 +02:00

537 lines
26 KiB
TypeScript

import assert from "node:assert/strict";
import { createServer } from "node:http";
import { test } from "bun:test";
import { workbenchKafkaDebugCapability } from "./workbench-kafka-debug-capability.ts";
import { handleWorkbenchKafkaSseDebugHttp } from "./workbench-kafka-sse-debug.ts";
test("workbench Kafka SSE debug endpoint streams filtered raw HWLAB Kafka events", async () => {
const fakeKafka = createFakeKafkaFactory();
const env = {
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092",
HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api",
HWLAB_KAFKA_EVENT_TOPIC: "hwlab.event.v1"
};
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, accessController: fakeAccessController(), logger: null });
});
await listen(server);
const abort = new AbortController();
try {
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab&sessionId=ses_kafka_sse_debug`, { signal: abort.signal });
assert.equal(response.status, 200);
assert.match(response.headers.get("content-type") ?? "", /text\/event-stream/u);
const reader = response.body?.getReader();
assert.ok(reader, "SSE response must expose a readable stream");
const connected = await readUntil(reader, "hwlab.kafka.connected");
assert.match(connected, /hwlab\.event\.v1/u);
assert.match(connected, /"consumerReady":true/u);
assert.match(connected, new RegExp(String(fakeKafka.groupId).replace(/[.*+?^${}()|[\]\\]/gu, "\\$&"), "u"));
await fakeKafka.emit({ eventType: "hwlab.trace.event.projected", sessionId: "ses_other", event: { label: "ignored" } });
await fakeKafka.emit({ eventType: "hwlab.trace.event.projected", sessionId: "ses_kafka_sse_debug", traceId: "trc_kafka_sse_debug", context: { runId: "run_kafka_sse_debug", commandId: "cmd_kafka_sse_debug" }, event: { label: "agentrun:event:debug", status: "running" } });
const streamed = await readUntil(reader, "trc_kafka_sse_debug");
assert.match(streamed, /hwlab\.kafka\.event/u);
assert.match(streamed, /ses_kafka_sse_debug/u);
assert.doesNotMatch(streamed, /ses_other/u);
} finally {
abort.abort();
await close(server);
}
});
test("isolated Workbench Kafka replay uses its YAML-owned topic, unique group, and bound", async () => {
const fakeKafka = createFakeKafkaFactory();
const env = {
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092",
HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true",
HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.debug.v1",
HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-test-workbench-isolated-debug",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "10",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "3000"
};
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, accessController: fakeAccessController(), logger: null });
});
await listen(server);
const abort = new AbortController();
try {
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_isolated_debug&fromBeginning=true`, { signal: abort.signal });
assert.equal(response.status, 200);
const reader = response.body?.getReader();
assert.ok(reader);
const connected = await readUntil(reader, "hwlab.kafka.connected");
assert.match(connected, /hwlab\.event\.debug\.v1/u);
assert.match(connected, /"debugIsolation":true/u);
assert.match(connected, /"deliverySemantics":"debug-replay"/u);
assert.match(connected, /"liveOnly":false/u);
assert.match(connected, /"replay":true/u);
assert.match(fakeKafka.groupId ?? "", /^hwlab-test-workbench-isolated-debug-/u);
assert.equal(fakeKafka.subscribedTopic, "hwlab.event.debug.v1");
assert.equal(fakeKafka.fromBeginning, true);
await fakeKafka.emit({ schema: "hwlab.event.debug.v1", eventType: "assistant", traceId: "trc_other", event: { traceId: "trc_other", type: "assistant", text: "ignored" } });
await fakeKafka.emit({ schema: "hwlab.event.debug.v1", eventType: "assistant", traceId: "trc_isolated_debug", event: { traceId: "trc_isolated_debug", type: "assistant", text: "isolated-visible" } });
await fakeKafka.emit({ schema: "hwlab.event.debug.v1", eventType: "terminal", traceId: "trc_isolated_debug", event: { traceId: "trc_isolated_debug", eventType: "terminal", type: "result", terminal: true, status: "completed" } });
const streamed = await readUntil(reader, "hwlab.kafka.replay-complete");
assert.match(streamed, /isolated-visible/u);
assert.match(streamed, /"reason":"terminal"/u);
assert.match(streamed, /"count":2/u);
assert.match(streamed, /"terminalObserved":true/u);
assert.doesNotMatch(streamed, /trc_other/u);
} finally {
abort.abort();
await close(server);
}
});
test("isolated Workbench Kafka replay applies replayId and offset barriers with bounded counts", async () => {
const fakeKafka = createFakeKafkaFactory();
const env = {
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092",
HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true",
HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.debug.v1",
HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-test-workbench-isolated-debug",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "10",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "3000"
};
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, accessController: fakeAccessController(), logger: null });
});
await listen(server);
try {
const query = "stream=hwlab-debug&traceId=trc_barrier&replayId=rpl_barrier&fromBeginning=true&partition=0&firstOffset=10&lastOffset=11";
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?${query}`);
const reader = response.body?.getReader();
assert.ok(reader);
const connected = await readUntil(reader, "hwlab.kafka.connected");
assert.match(connected, /"replayId":"rpl_barrier"/u);
assert.match(connected, /"firstOffset":"10"/u);
assert.match(connected, /"replayMode":"correlated-v2"/u);
assert.match(connected, /"seekApplied":true/u);
assert.equal(fakeKafka.fromBeginning, false);
assert.deepEqual(fakeKafka.seekCalls, [{ topic: "hwlab.event.debug.v1", partition: 0, offset: "10" }]);
await fakeKafka.emit({ schema: "hwlab.event.debug.v1", traceId: "trc_barrier", replayId: "rpl_barrier", debugLineage: { inputTraceId: "trc_barrier", inputSessionId: "ses_agentrun_barrier", sourceSeq: 1 }, event: { type: "assistant" } }, { offset: "10" });
await fakeKafka.emit({ schema: "hwlab.event.debug.v1", traceId: "trc_barrier", replayId: "rpl_barrier", debugLineage: { inputTraceId: "trc_barrier", inputSessionId: "ses_agentrun_barrier", sourceSeq: 2 }, eventType: "terminal", event: { type: "result", terminal: true } }, { offset: "11" });
const replay = await readUntil(reader, "hwlab.kafka.replay-complete");
assert.match(replay, /"code":"terminal_complete"/u);
assert.match(replay, /"reason":"barrier"/u);
assert.match(replay, /"scanned":2/u);
assert.match(replay, /"barrierScanned":2/u);
assert.match(replay, /"parsed":2/u);
assert.match(replay, /"matched":2/u);
assert.match(replay, /"delivered":2/u);
assert.match(replay, /"firstOffset":"10","lastOffset":"11","count":2/u);
assert.match(replay, /"inputSessionId":"ses_agentrun_barrier"/u);
assert.match(replay, /"barrierComplete":true/u);
assert.match(replay, /"clientCountsAvailable":false/u);
await waitUntil(() => fakeKafka.stopCalls > 0);
} finally {
await close(server);
}
});
test("isolated Workbench Kafka capability rejects product topic and group identities", () => {
const env = {
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true",
HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.debug.v1",
HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-test-workbench-isolated-debug",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "10",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "3000",
HWLAB_KAFKA_HWLAB_EVENT_TOPIC: "hwlab.event.v1",
HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID: "hwlab-v03-workbench-live-sse"
};
assert.throws(
() => workbenchKafkaDebugCapability({ ...env, HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.v1" }),
(error) => error?.code === "workbench_kafka_debug_capability_invalid"
&& error?.envName === "HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC"
);
assert.throws(
() => workbenchKafkaDebugCapability({ ...env, HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-v03-workbench-live-sse" }),
(error) => error?.code === "workbench_kafka_debug_capability_invalid"
&& error?.envName === "HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX"
);
});
test("isolated Workbench Kafka debug fails closed when disabled and rejects replay", async () => {
const fakeKafka = createFakeKafkaFactory();
const env = {
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "false",
HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.debug.v1",
HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-test-workbench-isolated-debug",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "10",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "3000"
};
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, accessController: fakeAccessController(), logger: null });
});
await listen(server);
try {
const disabled = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_disabled&fromBeginning=true`);
assert.equal(disabled.status, 503);
assert.match(await disabled.text(), /workbench_kafka_debug_disabled/u);
assert.equal(fakeKafka.consumerCreated, false);
env.HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED = "true";
const replay = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay`);
assert.equal(replay.status, 400);
assert.match(await replay.text(), /workbench_kafka_debug_replay_required/u);
assert.equal(fakeKafka.consumerCreated, false);
const invalidRange = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&partition=0&firstOffset=20`);
assert.equal(invalidRange.status, 400);
assert.match(await invalidRange.text(), /workbench_kafka_debug_offset_range_invalid/u);
assert.equal(fakeKafka.consumerCreated, false);
const invalidReplayId = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&replayId=not-a-replay`);
assert.equal(invalidReplayId.status, 400);
assert.match(await invalidReplayId.text(), /workbench_kafka_debug_replay_id_invalid/u);
assert.equal(fakeKafka.consumerCreated, false);
const replayOnly = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&replayId=rpl_incomplete`);
assert.equal(replayOnly.status, 400);
assert.match(await replayOnly.text(), /workbench_kafka_debug_v2_correlation_incomplete/u);
assert.equal(fakeKafka.consumerCreated, false);
const rangeOnly = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&partition=0&firstOffset=10&lastOffset=11`);
assert.equal(rangeOnly.status, 400);
assert.match(await rangeOnly.text(), /workbench_kafka_debug_v2_correlation_incomplete/u);
assert.equal(fakeKafka.consumerCreated, false);
const producerNotInvoked = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&replayId=rpl_not_invoked&partition=0&firstOffset=10&lastOffset=11&producerInvoked=false`);
assert.equal(producerNotInvoked.status, 409);
assert.match(await producerNotInvoked.text(), /workbench_kafka_debug_producer_not_invoked/u);
assert.equal(fakeKafka.consumerCreated, false);
const cardinalityMismatch = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay&fromBeginning=true&replayId=rpl_bad_count&partition=0&firstOffset=10&lastOffset=11&publishedCount=3`);
assert.equal(cardinalityMismatch.status, 400);
assert.match(await cardinalityMismatch.text(), /workbench_kafka_debug_producer_cardinality_mismatch/u);
assert.equal(fakeKafka.consumerCreated, false);
} finally {
await close(server);
}
});
test("isolated Workbench Kafka replay reports matched events without terminal", async () => {
const fakeKafka = createFakeKafkaFactory();
const env = {
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true",
HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.debug.v1",
HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-test-workbench-isolated-debug",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "10",
HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "100"
};
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, accessController: fakeAccessController(), logger: null });
});
await listen(server);
try {
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_incomplete&fromBeginning=true`);
assert.equal(response.status, 200);
const reader = response.body?.getReader();
assert.ok(reader);
await readUntil(reader, "hwlab.kafka.connected");
await fakeKafka.emit({ schema: "hwlab.event.debug.v1", traceId: "trc_incomplete", eventType: "assistant", event: { traceId: "trc_incomplete", type: "assistant" } }, { offset: "5" });
const replay = await readUntil(reader, "hwlab.kafka.replay-complete");
assert.match(replay, /"reason":"timeout"/u);
assert.match(replay, /"count":1/u);
assert.match(replay, /"code":"terminal_missing"/u);
assert.match(replay, /"terminalObserved":false/u);
await waitUntil(() => fakeKafka.stopCalls > 0);
} finally {
await close(server);
}
});
test("workbench Kafka SSE debug endpoint streams filtered codex stdio Kafka events", async () => {
const fakeKafka = createFakeKafkaFactory();
const env = {
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092",
HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api",
HWLAB_KAFKA_STDIO_TOPIC: "codex-stdio.raw.v1"
};
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, accessController: fakeAccessController(), logger: null });
});
await listen(server);
const abort = new AbortController();
try {
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=stdio&traceId=trc_stdio_sse_debug&fromBeginning=true`, { signal: abort.signal });
assert.equal(response.status, 200);
const reader = response.body?.getReader();
assert.ok(reader, "SSE response must expose a readable stream");
const connected = await readUntil(reader, "hwlab.kafka.connected");
assert.match(connected, /codex-stdio\.raw\.v1/u);
await fakeKafka.emit({ traceId: "trc_other", stream: "stdout", line: "ignored" });
await fakeKafka.emit({ trace_id: "trc_stdio_sse_debug", metadata: { session_id: "ses_stdio_sse_debug", run_id: "run_stdio_sse_debug" }, stream: "stdout", line: "visible stdio frame" });
const streamed = await readUntil(reader, "visible stdio frame");
assert.match(streamed, /hwlab\.kafka\.event/u);
assert.match(streamed, /codex-stdio\.raw\.v1/u);
assert.match(streamed, /ses_stdio_sse_debug/u);
assert.doesNotMatch(streamed, /trc_other/u);
} finally {
abort.abort();
await close(server);
}
});
test("workbench Kafka SSE debug endpoint resolves traceId to AgentRun Kafka ids", async () => {
const fakeKafka = createFakeKafkaFactory();
const traceStore = {
snapshot: () => ({
traceId: "trc_resolved_kafka_debug",
events: [
{ type: "backend", runId: "run_resolved_kafka_debug", commandId: "cmd_resolved_kafka_debug", sessionId: "ses_agentrun_resolved_kafka_debug" }
],
lastEvent: null
})
};
const env = {
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092",
HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api",
HWLAB_KAFKA_EVENT_TOPIC: "hwlab.event.v1"
};
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, traceStore, accessController: fakeAccessController(), logger: null });
});
await listen(server);
const abort = new AbortController();
try {
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab&traceId=trc_resolved_kafka_debug`, { signal: abort.signal });
assert.equal(response.status, 200);
const reader = response.body?.getReader();
assert.ok(reader, "SSE response must expose a readable stream");
const connected = await readUntil(reader, "resolvedFilters");
assert.match(connected, /trc_resolved_kafka_debug/u);
assert.match(connected, /run_resolved_kafka_debug/u);
await fakeKafka.emit({ eventType: "hwlab.trace.event.projected", traceId: null, sessionId: null, context: { runId: "run_other", commandId: "cmd_other" }, event: { label: "ignored" } });
await fakeKafka.emit({ eventType: "hwlab.trace.event.projected", traceId: null, sessionId: null, context: { runId: "run_resolved_kafka_debug", commandId: "cmd_resolved_kafka_debug" }, event: { label: "resolved", status: "running" } });
const streamed = await readUntil(reader, "run_resolved_kafka_debug");
assert.match(streamed, /hwlab\.kafka\.event/u);
assert.match(streamed, /resolved/u);
assert.doesNotMatch(streamed, /run_other/u);
} finally {
abort.abort();
await close(server);
}
});
test("workbench Kafka SSE debug endpoint is admin-only", async () => {
const fakeKafka = createFakeKafkaFactory();
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, {
env: { HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092" },
kafkaFactory: fakeKafka.factory,
accessController: fakeAccessController("user"),
logger: null
});
});
await listen(server);
try {
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab`);
assert.equal(response.status, 403);
assert.match(await response.text(), /admin_required/u);
assert.equal(fakeKafka.consumerCreated, false);
} finally {
await close(server);
}
});
test("workbench Kafka SSE debug endpoint reports connected only after consumer run is ready", async () => {
const fakeKafka = createFakeKafkaFactory({ deferRun: true });
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, {
env: {
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092",
HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api",
HWLAB_KAFKA_EVENT_TOPIC: "hwlab.event.v1"
},
kafkaFactory: fakeKafka.factory,
accessController: fakeAccessController(),
logger: null
});
});
await listen(server);
const abort = new AbortController();
try {
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab`, { signal: abort.signal });
const reader = response.body?.getReader();
assert.ok(reader);
const firstRead = reader.read();
const early = await Promise.race([firstRead.then(() => "event"), delay(25).then(() => "pending")]);
assert.equal(early, "pending");
fakeKafka.releaseRun();
const first = await firstRead;
assert.equal(first.done, false);
const connected = new TextDecoder().decode(first.value);
assert.match(connected, /hwlab\.kafka\.connected/u);
assert.match(connected, /"consumerReady":true/u);
assert.match(connected, /"groupId":"test-hwlab-cloud-api-debug-sse-/u);
} finally {
abort.abort();
fakeKafka.releaseRun();
await close(server);
}
});
test("workbench Kafka SSE debug endpoint stops a consumer that becomes ready after client close", async () => {
const fakeKafka = createFakeKafkaFactory({ deferRun: true });
const server = createServer((request, response) => {
const url = new URL(request.url ?? "/", "http://127.0.0.1");
void handleWorkbenchKafkaSseDebugHttp(request, response, url, {
env: { HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092", HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api", HWLAB_KAFKA_EVENT_TOPIC: "hwlab.event.v1" },
kafkaFactory: fakeKafka.factory,
accessController: fakeAccessController(),
logger: null
});
});
await listen(server);
const abort = new AbortController();
try {
const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab`, { signal: abort.signal });
const reader = response.body?.getReader();
assert.ok(reader);
const firstRead = reader.read().catch(() => null);
await delay(20);
abort.abort();
fakeKafka.releaseRun();
await waitUntil(() => fakeKafka.stopCalls > 0);
const first = await firstRead;
const text = first && !first.done && first.value ? new TextDecoder().decode(first.value) : "";
assert.doesNotMatch(text, /hwlab\.kafka\.connected/u);
assert.equal(fakeKafka.disconnectCalls > 0, true);
} finally {
abort.abort();
fakeKafka.releaseRun();
await close(server);
}
});
function createFakeKafkaFactory(options: { deferRun?: boolean } = {}) {
let eachMessage: ((input: any) => Promise<void>) | null = null;
let subscribedTopic = "hwlab.event.v1";
let fromBeginning = false;
let releaseRun: (() => void) | null = null;
const runGate = options.deferRun ? new Promise<void>((resolve) => { releaseRun = resolve; }) : null;
let groupId: string | null = null;
let stopCalls = 0;
let disconnectCalls = 0;
let consumerCreated = false;
let runStarted = false;
const seekCalls: Array<{ topic: string; partition: number; offset: string }> = [];
const consumer = {
connect: async () => undefined,
subscribe: async (input: { topic?: string; fromBeginning?: boolean }) => {
subscribedTopic = String(input.topic ?? subscribedTopic);
fromBeginning = input.fromBeginning === true;
},
run: async (input: any) => {
if (runGate) await runGate;
runStarted = true;
eachMessage = input.eachMessage;
},
seek: async (input: { topic: string; partition: number; offset: string }) => {
assert.equal(runStarted, true, "seek must happen after consumer.run assignment");
seekCalls.push(input);
},
stop: async () => { stopCalls += 1; },
disconnect: async () => { disconnectCalls += 1; }
};
return {
factory: () => ({
consumer: (input: { groupId?: string }) => {
consumerCreated = true;
groupId = String(input.groupId ?? "");
return consumer;
}
}),
get groupId() { return groupId; },
get subscribedTopic() { return subscribedTopic; },
get fromBeginning() { return fromBeginning; },
get stopCalls() { return stopCalls; },
get disconnectCalls() { return disconnectCalls; },
get consumerCreated() { return consumerCreated; },
get seekCalls() { return seekCalls; },
releaseRun: () => releaseRun?.(),
emit: async (value: Record<string, unknown>, transport: { offset?: string; partition?: number } = {}) => {
assert.ok(eachMessage, "consumer.run must be called before emitting fake Kafka events");
await eachMessage({
topic: subscribedTopic,
partition: transport.partition ?? 0,
message: {
offset: transport.offset ?? String(value.sessionId === "ses_other" ? 1 : 2),
key: Buffer.from(String(value.sessionId ?? "unknown")),
timestamp: "2026-07-09T18:30:00.000Z",
value: Buffer.from(JSON.stringify(value))
}
});
}
};
}
function fakeAccessController(role: "admin" | "user" = "admin") {
return {
async ensureBootstrap() {},
async authenticate() {
return { ok: true, status: 200, actor: { id: role === "admin" ? "usr_debug_admin" : "usr_debug_user", role } };
}
};
}
async function delay(ms: number) {
await new Promise<void>((resolve) => setTimeout(resolve, ms));
}
async function waitUntil(predicate: () => boolean, timeoutMs = 1000) {
const deadline = Date.now() + timeoutMs;
while (!predicate()) {
if (Date.now() >= deadline) throw new Error("condition did not become true before timeout");
await delay(10);
}
}
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}`;
}