import { describe, expect, test } from "bun:test"; import type { WorkbenchApplication, WorkbenchTemporalGateway } from "./contracts.ts"; import { createWorkbenchActivities } from "./activities.ts"; import { createWorkbenchDispatcher } from "./dispatcher.ts"; import { createCloudWorkbenchApplication } from "./cloud-application.ts"; import { createWorkbenchHttpApp } from "./http.ts"; import { createNativeTestTemporalGateway } from "./native-temporal.ts"; import { workbenchRuntime } from "./runtime.ts"; import { runWorkbenchCli } from "../../tools/src/workbench-cli.ts"; test("temporal API runtime does not require worker Kafka publisher configuration", async () => { const runtime = workbenchRuntime({ WORKBENCH_MODE: "temporal", WORKBENCH_CLOUD_API_URL: "http://cloud.internal", WORKBENCH_CLOUD_API_AUTHORIZATION: "internal-key", WORKBENCH_TEMPORAL_ADDRESS: "127.0.0.1:7233", WORKBENCH_TEMPORAL_NAMESPACE: "unidesk", WORKBENCH_TEMPORAL_TASK_QUEUE: "hwlab-v03-workbench", WORKBENCH_ACTIVITY_START_TO_CLOSE_TIMEOUT_MS: "30000", WORKBENCH_ACTIVITY_RETRY_INITIAL_INTERVAL_MS: "1000", WORKBENCH_ACTIVITY_RETRY_MAXIMUM_INTERVAL_MS: "30000", WORKBENCH_ACTIVITY_RETRY_MAXIMUM_ATTEMPTS: "5" }, { eventPublisher: null }); expect(runtime.eventPublisher).toBeNull(); await runtime.close(); }); describe("Workbench dispatcher", () => { test("submits through Temporal without synthesizing terminal state", async () => { const calls: string[] = []; let temporalInput: Record | null = null; const application: WorkbenchApplication = { async health() { return { status: "ready" }; }, async createSession() { return { session: { sessionId: "ses_test" } }; }, async admitTurn(input) { calls.push("admit"); return input; }, async dispatchTurn() { throw new Error("API dispatcher must not dispatch directly"); }, async cancelTurn() { throw new Error("API dispatcher must not cancel directly"); } }; const temporal: WorkbenchTemporalGateway = { async submit(input) { calls.push("temporal"); temporalInput = input; return { workflowId: `workbench-turn-${input.traceId}`, workflowRunId: "run-1", reused: false }; }, async cancel() { return { workflowId: "cancel", workflowRunId: "run-2", reused: false }; }, async close() {} }; const dispatch = createWorkbenchDispatcher({ application, temporal, mode: "temporal" }); const result = await dispatch({ operation: "turn.submit", actor: { id: "usr_test" }, traceId: "trc_test", params: { sessionId: "ses_test", message: "hi", model: "gpt-5.3-codex-spark", reasoningEffort: "low" } }); expect(calls).toEqual(["admit", "temporal"]); expect(temporalInput?.params).toMatchObject({ model: "gpt-5.3-codex-spark", reasoningEffort: "low" }); expect(result).toMatchObject({ ok: true, mode: "temporal", data: { status: "running", traceId: "trc_test", terminalAuthority: "hwlab.event.v1", resultSynthesized: false } }); }); test("cancel is a durable Temporal command", async () => { const application = stubApplication(); const temporal: WorkbenchTemporalGateway = { async submit() { throw new Error("not used"); }, async cancel(input) { return { workflowId: `workbench-cancel-${input.traceId}`, workflowRunId: "run-cancel", reused: false }; }, async close() {} }; const result = await createWorkbenchDispatcher({ application, temporal, mode: "temporal" })({ operation: "turn.cancel", actor: { id: "usr_test" }, traceId: "trc_test" }); expect(result).toMatchObject({ ok: true, data: { status: "cancel_requested", resultSynthesized: false } }); }); test("steer is accepted through a distinct durable Temporal workflow", async () => { let temporalInput: Record | null = null; const application = { ...stubApplication(), async steerTurn() { throw new Error("API dispatcher must not steer directly"); } }; const temporal = { async submit() { throw new Error("not used"); }, async cancel() { throw new Error("not used"); }, async steer(input: Record) { temporalInput = input; return { workflowId: `workbench-steer-${input.steerTraceId}`, workflowRunId: "run-steer", reused: false, traceId: "trc_actual_followup_turn", requestedDelivery: "steer", delivery: "turn", agentRun: { runId: "run_reused", commandId: "cmd_followup" } }; }, async close() {} } as WorkbenchTemporalGateway; const result = await createWorkbenchDispatcher({ application, temporal, mode: "temporal" })({ operation: "turn.steer", actor: { id: "usr_test" }, traceId: "trc_target", params: { sessionId: "ses_test", message: "adjust course", targetTraceId: "trc_target", steerTraceId: "trc_steer_control" } } as any); expect(temporalInput).toMatchObject({ sessionId: "ses_test", targetTraceId: "trc_target", steerTraceId: "trc_steer_control", params: { message: "adjust course" } }); expect(result).toMatchObject({ ok: true, operation: "turn.steer", data: { accepted: true, status: "running", traceId: "trc_actual_followup_turn", targetTraceId: "trc_target", steerTraceId: "trc_steer_control", sessionId: "ses_test", requestedDelivery: "steer", delivery: "turn", agentRun: { runId: "run_reused", commandId: "cmd_followup" }, terminalAuthority: "hwlab.event.v1", resultSynthesized: false } }); }); test("native Temporal gateway preserves the actual AgentRun steer decision", async () => { const application = { ...stubApplication(), async steerTurn() { return { traceId: "trc_actual_native_turn", requestedDelivery: "steer", delivery: "turn", agentRun: { runId: "run_native_reused", commandId: "cmd_native_followup" } }; } }; const temporal = createNativeTestTemporalGateway(application); const result = await temporal.steer({ actor: { id: "usr_test" }, sessionId: "ses_test", targetTraceId: "trc_target", steerTraceId: "trc_requested_steer", params: { message: "continue" } }); expect(result).toMatchObject({ traceId: "trc_actual_native_turn", requestedDelivery: "steer", delivery: "turn", agentRun: { runId: "run_native_reused", commandId: "cmd_native_followup" }, workflowId: "native-workbench-steer-trc_requested_steer", reused: false }); }); }); describe("Workbench Cloud application adapter", () => { test("persists admission before dispatch through distinct internal routes", async () => { const calls: Array<{ path: string; authorization: string | null }> = []; const application = createCloudWorkbenchApplication({ baseUrl: "http://cloud.internal", authorization: "unified-internal-key", fetchImpl: (async (url: string | URL | Request, init?: RequestInit) => { calls.push({ path: new URL(String(url)).pathname, authorization: new Headers(init?.headers).get("authorization") }); return Response.json({ ok: true, accepted: true, status: "admitted" }); }) as typeof fetch }); const input = { actor: { id: "usr_test" }, traceId: "trc_test", sessionId: "ses_test", params: { sessionId: "ses_test", message: "hi" } }; await application.admitTurn(input); await application.dispatchTurn(input); expect(calls).toEqual([ { path: "/v1/internal/workbench/admit", authorization: "Bearer unified-internal-key" }, { path: "/v1/agent/chat", authorization: "Bearer unified-internal-key" } ]); }); test("preserves an existing Bearer authorization scheme", async () => { let authorization: string | null = null; const application = createCloudWorkbenchApplication({ baseUrl: "http://cloud.internal", authorization: "Bearer unified-internal-key", fetchImpl: (async (_url: string | URL | Request, init?: RequestInit) => { authorization = new Headers(init?.headers).get("authorization"); return Response.json({ ok: true, session: { sessionId: "ses_test" } }); }) as typeof fetch }); await application.createSession({ actor: { id: "usr_test" }, params: {} }); expect(authorization).toBe("Bearer unified-internal-key"); }); test("forwards steer from the Temporal activity to the existing Cloud API route", async () => { let request: { path: string; body: Record } | null = null; const application = createCloudWorkbenchApplication({ baseUrl: "http://cloud.internal", authorization: "unified-internal-key", fetchImpl: (async (url: string | URL | Request, init?: RequestInit) => { request = { path: new URL(String(url)).pathname, body: JSON.parse(String(init?.body ?? "{}")) }; return Response.json({ ok: true, accepted: true, status: "running" }); }) as typeof fetch }); await application.steerTurn({ actor: { id: "usr_test" }, sessionId: "ses_cloud_steer", targetTraceId: "trc_cloud_target", steerTraceId: "trc_cloud_steer", params: { message: "adjust course" } }); expect(request).toEqual({ path: "/v1/agent/chat/steer", body: { message: "adjust course", traceId: "trc_cloud_target", targetTraceId: "trc_cloud_target", steerTraceId: "trc_cloud_steer", sessionId: "ses_cloud_steer", shortConnection: true } }); }); test("preserves Cloud API retryability and billing failure semantics", async () => { const application = createCloudWorkbenchApplication({ baseUrl: "http://cloud.internal", authorization: "unified-internal-key", fetchImpl: (async () => Response.json({ ok: false, error: { code: "insufficient_credits", layer: "billing", retryable: false, message: "credits exhausted", userMessage: "余额不足" } }, { status: 402 })) as typeof fetch }); const input = { actor: { id: "usr_test" }, traceId: "trc_billing", sessionId: "ses_billing", params: { message: "hi" } }; await expect(application.dispatchTurn(input)).rejects.toMatchObject({ code: "insufficient_credits", layer: "billing", retryable: false, status: 402, userMessage: "余额不足" }); }); }); describe("Workbench Temporal activities", () => { const activity = { startToCloseTimeoutMs: 1000, retryInitialIntervalMs: 100, retryMaximumIntervalMs: 500, retryMaximumAttempts: 3 }; const input = { actor: { id: "usr_test" }, traceId: "trc_activity", sessionId: "ses_activity", params: { message: "hi" } }; test("publishes semantic Kafka failure and stops retrying non-retryable billing errors", async () => { const published: any[] = []; const application = stubApplication(); application.dispatchTurn = async () => { throw Object.assign(new Error("credits exhausted"), { code: "insufficient_credits", layer: "billing", retryable: false, userMessage: "余额不足" }); }; const activities = createWorkbenchActivities(application, { activity, attempt: () => 1, eventPublisher: { async publish(event) { published.push(event); return { published: true }; } } }); await expect(activities.dispatchWorkbenchTurn(input)).rejects.toMatchObject({ type: "insufficient_credits", nonRetryable: true }); expect(published).toHaveLength(2); expect(published[0]).toMatchObject({ event: { type: "error", failureDomain: "billing", code: "insufficient_credits", retryable: false, willRetry: false, attempt: 1, maxAttempts: 3, terminal: false } }); expect(published[1]).toMatchObject({ event: { type: "terminal_status", terminalStatus: "failed", errorCode: "insufficient_credits", message: "余额不足", terminal: true } }); }); test("publishes the YAML-controlled exponential retry without sealing the turn early", async () => { const published: any[] = []; const application = stubApplication(); application.dispatchTurn = async () => { throw Object.assign(new Error("billing unavailable"), { code: "user_billing_preflight_failed", layer: "billing", retryable: true }); }; const activities = createWorkbenchActivities(application, { activity, attempt: () => 2, eventPublisher: { async publish(event) { published.push(event); return { published: true }; } } }); await expect(activities.dispatchWorkbenchTurn(input)).rejects.toMatchObject({ type: "user_billing_preflight_failed", nonRetryable: false }); expect(published).toHaveLength(1); expect(published[0]).toMatchObject({ event: { type: "backend_status", retryPhase: "retryScheduled", attempt: 2, maxAttempts: 3, backoffMs: 200, terminal: false } }); }); }); describe("Workbench native HTTP adapter", () => { test("native session list exposes the first user preview for summary-only rail loading", async () => { const app = createWorkbenchHttpApp({ async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: { ses_title: { sessionId: "ses_title", updatedAt: "2026-07-20T01:00:00.000Z", messages: [{ role: "user", content: "persisted native title", createdAt: "2026-07-20T00:59:00.000Z" }] }, ses_empty: { sessionId: "ses_empty", updatedAt: "2026-07-20T00:00:00.000Z", messages: [] } }, turns: {} }; } }); const response = await app.fetch(new Request("http://native.test/v1/workbench/sessions")); expect(response.status).toBe(200); expect(await response.json()).toMatchObject({ sessions: [ { sessionId: "ses_title", firstUserMessagePreview: "persisted native title" }, { sessionId: "ses_empty", firstUserMessagePreview: null } ] }); }); test("native session list bounds the loaded realtime scope with the existing cursor contract", async () => { const sessions = Object.fromEntries(Array.from({ length: 25 }, (_, index) => { const rank = String(25 - index).padStart(2, "0"); const sessionId = `ses_page_${rank}`; return [sessionId, { sessionId, updatedAt: `2026-07-20T00:${rank}:00.000Z`, messages: [] }]; })); const app = createWorkbenchHttpApp({ async dispatch() { return { ok: true }; }, async snapshot() { return { sessions, turns: {} }; } }); const first = await app.fetch(new Request("http://native.test/v1/workbench/sessions?limit=20")); expect(first.status).toBe(200); const firstBody = await first.json(); expect(firstBody).toMatchObject({ count: 20, total: 25, hasMore: true, nextCursor: "idx:20" }); expect(firstBody.sessions).toHaveLength(20); const second = await app.fetch(new Request("http://native.test/v1/workbench/sessions?limit=20&cursor=idx%3A20")); expect(second.status).toBe(200); const secondBody = await second.json(); expect(secondBody).toMatchObject({ count: 5, total: 25, cursor: "idx:20", hasMore: false, nextCursor: null }); expect(secondBody.sessions).toHaveLength(5); }); test("AgentRun native session list merges Kafka summaries with persisted empty sessions", async () => { const app = createWorkbenchHttpApp({ mode: "agentrun-native", kafkaEventBridge: { workbenchSessionSummaries() { return { ok: true, sessions: [{ sessionId: "ses_title", firstUserMessagePreview: "Kafka title", lastUserMessageAt: "2026-07-20T01:00:00.000Z", lastTraceId: "trc_target", status: "completed" }] }; } }, async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: { ses_title: { sessionId: "ses_title", updatedAt: "2026-07-20T00:00:00.000Z", lastTraceId: "trc_steer", status: "running", messages: [{ role: "user", content: "admission snapshot title" }] }, ses_admission_only: { sessionId: "ses_admission_only", createdAt: "2026-07-20T02:00:00.000Z", updatedAt: "2026-07-20T02:00:00.000Z", status: "idle", messages: [] } }, turns: {} }; } }); const response = await app.fetch(new Request("http://native.test/v1/workbench/sessions")); const body = await response.json(); expect(body).toMatchObject({ count: 2, total: 2 }); expect(body.sessions).toEqual([ expect.objectContaining({ sessionId: "ses_admission_only", createdAt: "2026-07-20T02:00:00.000Z", updatedAt: "2026-07-20T02:00:00.000Z", firstUserMessagePreview: null }), expect.objectContaining({ sessionId: "ses_title", firstUserMessagePreview: "Kafka title", lastUserMessageAt: "2026-07-20T01:00:00.000Z", lastTraceId: "trc_target", status: "completed" }) ]); }); test("authorizes only the Workbench navigation entry", async () => { const app = createWorkbenchHttpApp({ async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: {}, turns: {} }; } }); for (const path of ["/auth/session", "/auth/bootstrap"]) { const response = await app.fetch(new Request(`http://native.test${path}`)); expect(response.status).toBe(200); expect(await response.json()).toMatchObject({ authenticated: true, access: { nav: { profileId: "workbench-native", allowedIds: ["workbench.code"] } } }); } }); test("serves read-only empty states for Workbench auxiliary panels", async () => { const app = createWorkbenchHttpApp({ async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: {}, turns: {} }; } }); for (const path of ["/v1/caserun/cases", "/v1/hwpod/specs?probe=1", "/v1/hwpod-node-ops"]) { const response = await app.fetch(new Request(`http://native.test${path}`)); expect(response.status).toBe(200); expect(await response.json()).toMatchObject({ ok: true, status: "unavailable", mode: "native-test" }); } const performance = await app.fetch(new Request("http://native.test/v1/web-performance", { method: "POST", body: "{}" })); expect(performance.status).toBe(202); expect(await performance.json()).toMatchObject({ ok: true, status: "accepted", persisted: false, mode: "native-test" }); }); test("L0 native HTTP maps the Web steer route to the shared Workbench command", async () => { let command: Record | null = null; const app = createWorkbenchHttpApp({ mode: "native-test", async dispatch(input) { command = input; return { ok: true, operation: input.operation, mode: "native-test", data: { accepted: true, status: "running", traceId: "trc_target", targetTraceId: "trc_target", steerTraceId: "trc_steer_http", sessionId: "ses_steer_http" } }; }, async snapshot() { return { sessions: {}, turns: {} }; } }); const response = await app.fetch(new Request("http://localhost/v1/agent/chat/steer", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ sessionId: "ses_steer_http", traceId: "trc_target", targetTraceId: "trc_target", steerTraceId: "trc_steer_http", message: "adjust course" }) })); expect(response.status).toBe(202); expect(command).toMatchObject({ operation: "turn.steer", actor: { id: "usr_native" }, traceId: "trc_target", params: { sessionId: "ses_steer_http", targetTraceId: "trc_target", steerTraceId: "trc_steer_http", message: "adjust course" } }); }); test("L0 CLI submits steer over the shared command transport", async () => { let command: Record | null = null; const server = Bun.serve({ port: 0, async fetch(request) { command = await request.json(); return Response.json({ ok: true, operation: "turn.steer", mode: "native-test", data: { accepted: true, status: "running", traceId: "trc_cli_target", targetTraceId: "trc_cli_target", steerTraceId: "trc_cli_steer", sessionId: "ses_cli_steer" } }, { status: 202 }); } }); try { const result = await runWorkbenchCli([ "turn", "steer", "--actor-id", "usr_native", "--session-id", "ses_cli_steer", "--trace-id", "trc_cli_target", "--steer-trace-id", "trc_cli_steer", "--message", "adjust course", "--over-api", "--api-url", `http://127.0.0.1:${server.port}` ], {}); expect(command).toMatchObject({ operation: "turn.steer", traceId: "trc_cli_target", params: { sessionId: "ses_cli_steer", targetTraceId: "trc_cli_target", steerTraceId: "trc_cli_steer", message: "adjust course" } }); expect(result).toMatchObject({ ok: true, transport: "api", httpStatus: 202, identity: { sessionId: "ses_cli_steer", traceId: "trc_cli_target", steerTraceId: "trc_cli_steer" } }); } finally { server.stop(true); } }); test("L0 CLI observes AgentRun-derived events only through the Kafka SSE route", async () => { const subscribers = new Set<(envelope: any, transport?: any) => void>(); const bridge = { capabilities: { liveKafkaSse: true }, ready: Promise.resolve(), subscribeLiveHwlabEvents(listener: (envelope: any, transport?: any) => void) { subscribers.add(listener); return () => subscribers.delete(listener); } }; const app = createWorkbenchHttpApp({ mode: "agentrun-native", kafkaEventBridge: bridge, async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: {}, turns: {} }; } }); let authorization: string | null = null; const server = Bun.serve({ port: 0, fetch(request) { authorization = request.headers.get("authorization"); return app.fetch(request); } }); try { const pending = runWorkbenchCli([ "events", "inspect", "--session-id", "ses_l0_sse", "--trace-id", "trc_l0_sse", "--over-api", "--api-url", `http://127.0.0.1:${server.port}`, "--timeout-ms", "1000" ], { HWLAB_API_KEY: "l0-sse-api-key" }); for (let index = 0; index < 20 && subscribers.size === 0; index += 1) await Bun.sleep(5); for (const listener of subscribers) listener({ schema: "hwlab.event.v1", sessionId: "ses_l0_sse", traceId: "trc_l0_sse", event: { type: "backend", eventType: "backend", label: "agentrun:backend:command-created" } }, { topic: "hwlab.event.v1", partition: 0, offset: "7" }); expect(await pending).toMatchObject({ ok: true, operation: "events.inspect", connected: true, connectedContract: { deliverySemantics: "live-only", realtimeSource: "hwlab.event.v1", replay: false, liveOnly: true, refreshReplay: null }, kafkaTransport: { eventCount: 1, identifiedCount: 1, missingCount: 0, invalidCount: 0, partitions: [{ topic: "hwlab.event.v1", partition: 0, firstOffset: "7", lastOffset: "7", eventCount: 1, nonMonotonicCount: 0 }] }, eventCount: 1, eventTypes: ["backend"] }); expect(authorization).toBe("Bearer l0-sse-api-key"); } finally { server.stop(true); await app.close(); } }); test("L0 CLI reports the Workbench SSE HTTP status when authentication fails", async () => { const server = Bun.serve({ port: 0, fetch: () => new Response("unauthorized", { status: 401 }) }); try { await expect(runWorkbenchCli([ "events", "inspect", "--trace-id", "trc_l0_auth_failure", "--over-api", "--api-url", `http://127.0.0.1:${server.port}`, "--timeout-ms", "1000" ], {})).rejects.toMatchObject({ code: "workbench_sse_read_failed", message: expect.stringContaining("Workbench SSE returned HTTP 401") }); } finally { server.stop(true); } }); test("L0 CLI replays retained Kafka events before the same live SSE stream", async () => { const sessionId = "ses_l0_refresh"; const traceId = "trc_l0_refresh"; const bridge = { capabilities: { liveKafkaSse: true, kafkaRefreshReplay: true }, refreshReplay: { groupIdPrefix: "hwlab-l0-refresh", timeoutMs: 1000, scanLimit: 100, matchedEventLimit: 20, liveBufferLimit: 20 }, ...nativeRailProjectionStub(), ready: Promise.resolve(), subscribeLiveHwlabEvents() { return () => {}; }, async queryHwlabEventRetention(options: any) { expect(options).toMatchObject({ sessionId, traceId, partitionKey: sessionId, fromBeginning: true }); return { topic: "hwlab.event.v1", events: Array.from({ length: 100 }, (_, index) => ({ topic: "hwlab.event.v1", partition: 0, offset: String(index), value: { schema: "hwlab.event.v1", eventType: "hwlab.trace.event.projected", eventId: `hwlab:evt_l0_refresh_${index}`, sourceEventId: `evt_l0_refresh_${index}`, sessionId, hwlabSessionId: sessionId, traceId, event: { type: "backend", eventType: "backend", sessionId, traceId, sourceEventId: `evt_l0_refresh_${index}` } } })), completionReason: "end-offset", completion: { reason: "end-offset", complete: true, barrierReached: true, retentionStartVerified: true }, reachedEndOffsets: true, endOffsetsAvailable: true, endOffsets: [{ partition: 0, startOffset: "0", endOffset: "100" }], scannedCount: 100, parsedCount: 100, matchedCount: 100, filterRejectedCount: 0, keyRejectedCount: 0, partitionKeyScoped: true, targetPartition: 0, timing: { endOffsetSnapshotMs: 0, groupOffsetsMs: 0, consumerConnectMs: 0, consumerSubscribeMs: 0, scanMs: 0, cleanupMs: 0, indexLookupMs: 1, totalMs: 1 }, index: { hit: true, source: "process-session-index", indexedEventCount: 300, indexedSessionCount: 3, incompleteSessionCount: 0, valuesRedacted: true } }; } }; const app = createWorkbenchHttpApp({ mode: "agentrun-native", kafkaEventBridge: bridge, async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: {}, turns: {} }; } }); const server = Bun.serve({ port: 0, fetch: (request) => app.fetch(request) }); try { const result = await runWorkbenchCli([ "events", "inspect", "--session-id", sessionId, "--trace-id", traceId, "--over-api", "--api-url", `http://127.0.0.1:${server.port}`, "--timeout-ms", "1000" ], {}); expect(result).toMatchObject({ ok: true, connected: true, connectedContract: { deliverySemantics: "kafka-retention-then-live", realtimeSource: "hwlab.event.v1", replay: true, liveOnly: false, refreshReplay: { phase: "live", counts: { replayed: 100 }, query: { complete: true, scannedCount: 100, parsedCount: 100, matchedCount: 100, filterRejectedCount: 0, keyRejectedCount: 0, partitionKeyScoped: true, targetPartition: 0, index: { hit: true, source: "process-session-index", indexedEventCount: 300, indexedSessionCount: 3, incompleteSessionCount: 0, blocking: false }, timing: { scanMs: 0, indexLookupMs: 1, totalMs: 1 } } } }, eventCount: 100, kafkaTransport: { eventCount: 100, identifiedCount: 100, missingCount: 0, invalidCount: 0, partitions: [{ topic: "hwlab.event.v1", partition: 0, firstOffset: "0", lastOffset: "99", eventCount: 100, nonMonotonicCount: 0 }] }, eventTypes: ["backend"], eventTypeCounts: { backend: 100 } }); expect(JSON.stringify(result).length).toBeLessThan(10_000); } finally { server.stop(true); await app.close(); } }); test("L0 native SSE projects the loaded rail scope and replays only the active session", async () => { const sessionIds = Array.from({ length: 20 }, (_, index) => `ses_l0_scope_${String(index).padStart(2, "0")}`); const queryCalls: any[] = []; const bridge = { capabilities: { liveKafkaSse: true, kafkaRefreshReplay: true }, refreshReplay: { groupIdPrefix: "hwlab-l0-scope", timeoutMs: 1000, scanLimit: 100, matchedEventLimit: 20, liveBufferLimit: 20 }, ...nativeRailProjectionStub(), ready: Promise.resolve(), subscribeLiveHwlabEvents() { return () => {}; }, async queryHwlabEventRetention(options: any) { queryCalls.push(options); return { topic: "hwlab.event.v1", events: [{ topic: "hwlab.event.v1", partition: 0, offset: "0", value: { schema: "hwlab.event.v1", eventType: "hwlab.trace.event.projected", eventId: "hwlab:evt_l0_scope_active", sourceEventId: "evt_l0_scope_active", sessionId: sessionIds[0], hwlabSessionId: sessionIds[0], traceId: "trc_l0_scope_active", event: { type: "backend", eventType: "backend", sessionId: sessionIds[0], traceId: "trc_l0_scope_active", sourceEventId: "evt_l0_scope_active" } } }], completionReason: "end-offset", completion: { reason: "end-offset", complete: true, barrierReached: true, retentionStartVerified: true }, reachedEndOffsets: true, endOffsetsAvailable: true, endOffsets: [{ partition: 0, startOffset: "0", endOffset: "1" }], scannedCount: 1, parsedCount: 1, matchedCount: 1 }; } }; const app = createWorkbenchHttpApp({ mode: "agentrun-native", kafkaEventBridge: bridge, async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: {}, turns: {} }; } }); const server = Bun.serve({ port: 0, fetch: (request) => app.fetch(request) }); const controller = new AbortController(); try { const response = await fetch( `http://127.0.0.1:${server.port}/v1/workbench/events?${sessionIds.map((sessionId) => `sessionId=${sessionId}`).join("&")}&sessionId=${sessionIds[0]}&activeSessionId=${sessionIds[0]}`, { signal: controller.signal } ); expect(response.status).toBe(200); const reader = response.body!.getReader(); const decoder = new TextDecoder(); let body = ""; for (let index = 0; index < 20 && !body.includes("event: workbench.connected"); index += 1) { const result = await reader.read(); if (result.done) break; body += decoder.decode(result.value, { stream: true }); } await reader.cancel(); expect(queryCalls).toHaveLength(1); expect(queryCalls[0]).toMatchObject({ sessionIds: [sessionIds[0]], sessionId: sessionIds[0], partitionKey: sessionIds[0], fromBeginning: true }); expect(body.match(/event: workbench\.session\.rail/gu)).toHaveLength(20); expect(body.match(/event: hwlab\.event\.v1/gu)).toHaveLength(1); expect(body).toContain("id: kafka|hwlab.event.v1|0|0"); expect(body).toContain(`"sessionIds":[${sessionIds.map((sessionId) => `"${sessionId}"`).join(",")}]`); expect(body).toContain(`"activeSessionId":"${sessionIds[0]}"`); expect(body).not.toContain('"traceId":"trc_l0_scope_1"'); expect(body.indexOf("event: hwlab.event.v1")).toBeLessThan(body.indexOf("event: workbench.connected")); } finally { controller.abort(); server.stop(true); await app.close(); } }); test("L0 native SSE keeps active replay and live delivery when the session rail fails", async () => { const sessionId = "ses_l0_rail_degraded"; const traceId = "trc_l0_rail_degraded"; const subscribers = new Set<(envelope: any, transport?: any) => void>(); let rejectRailReady: (error: Error) => void = () => {}; const railReady = new Promise((_resolve, reject) => { rejectRailReady = reject; }); let retentionQueries = 0; const retainedEnvelope = { schema: "hwlab.event.v1", eventType: "hwlab.trace.event.projected", eventId: "hwlab:evt_l0_rail_retained", sourceEventId: "evt_l0_rail_retained", sessionId, hwlabSessionId: sessionId, traceId, event: { type: "backend", eventType: "backend", sessionId, traceId, sourceEventId: "evt_l0_rail_retained" } }; const liveEnvelope = { schema: "hwlab.event.v1", eventType: "hwlab.trace.event.projected", eventId: "hwlab:evt_l0_rail_live", sourceEventId: "evt_l0_rail_live", sessionId, hwlabSessionId: sessionId, traceId, event: { type: "assistant", eventType: "assistant", sessionId, traceId, sourceEventId: "evt_l0_rail_live" } }; const app = createWorkbenchHttpApp({ mode: "agentrun-native", kafkaEventBridge: { capabilities: { liveKafkaSse: true, kafkaRefreshReplay: true }, refreshReplay: { groupIdPrefix: "hwlab-l0-rail-degraded", timeoutMs: 1000, scanLimit: 20, matchedEventLimit: 10, liveBufferLimit: 10 }, ready: Promise.resolve(), waitWorkbenchSessionRailReady() { return railReady; }, workbenchSessionSummaries() { throw new Error("rail summaries must not run after readiness fails"); }, subscribeLiveHwlabEvents(listener: (envelope: any, transport?: any) => void) { subscribers.add(listener); return () => subscribers.delete(listener); }, async queryHwlabEventRetention() { retentionQueries += 1; return { topic: "hwlab.event.v1", events: [{ topic: "hwlab.event.v1", partition: 0, offset: "10", value: retainedEnvelope }], completionReason: "end-offset", completion: { reason: "end-offset", complete: true, barrierReached: true, retentionStartVerified: true }, reachedEndOffsets: true, endOffsetsAvailable: true, endOffsets: [{ partition: 0, startOffset: "10", endOffset: "11" }], scannedCount: 1, parsedCount: 1, matchedCount: 1 }; } }, async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: {}, turns: {} }; } }); const response = await app.fetch(new Request(`http://localhost/v1/workbench/events?sessionId=${sessionId}&activeSessionId=${sessionId}&traceId=${traceId}`)); expect(response.status).toBe(200); const reader = response.body!.getReader(); const decoder = new TextDecoder(); let body = ""; let liveSent = false; try { rejectRailReady(Object.assign(new Error("session rail bootstrap unavailable"), { code: "workbench_kafka_session_rail_bootstrap_failed" })); for (let index = 0; index < 20; index += 1) { const result = await reader.read(); if (result.done) break; body += decoder.decode(result.value, { stream: true }); if (!liveSent && body.includes("event: workbench.connected")) { liveSent = true; for (const listener of subscribers) listener(liveEnvelope, { topic: "hwlab.event.v1", partition: 0, offset: "11" }); } if (body.includes("evt_l0_rail_live") && body.includes("workbench_kafka_session_rail_bootstrap_failed")) break; } } finally { await reader.cancel(); await app.close(); } expect(retentionQueries).toBe(1); expect(body).toContain("event: workbench.error"); expect(body).toContain('"status":"degraded"'); expect(body).toContain('"phase":"session-rail"'); expect(body).toContain('"blocking":false'); expect(body).toContain("evt_l0_rail_retained"); expect(body).toContain("id: kafka|hwlab.event.v1|0|10"); expect(body).toContain("event: workbench.connected"); expect(body).toContain("evt_l0_rail_live"); expect(body).toContain("id: kafka|hwlab.event.v1|0|11"); }); test("L0 CLI waits for explicit Kafka event semantics in one SSE window", async () => { const subscribers = new Set<(envelope: any) => void>(); const bridge = { capabilities: { liveKafkaSse: true }, ready: Promise.resolve(), subscribeLiveHwlabEvents(listener: (envelope: any) => void) { subscribers.add(listener); return () => subscribers.delete(listener); } }; const app = createWorkbenchHttpApp({ mode: "agentrun-native", kafkaEventBridge: bridge, async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: {}, turns: {} }; } }); const server = Bun.serve({ port: 0, fetch: (request) => app.fetch(request) }); try { const pending = runWorkbenchCli([ "events", "inspect", "--session-id", "ses_l0_semantics", "--trace-id", "trc_l0_semantics", "--over-api", "--api-url", `http://127.0.0.1:${server.port}`, "--timeout-ms", "1000", "--wait-for", "user,backend,assistant,terminal,final" ], {}); for (let index = 0; index < 20 && subscribers.size === 0; index += 1) await Bun.sleep(5); const events = [ { type: "user", eventType: "user", createdAt: "2026-07-19T16:00:00.000Z" }, { type: "backend", eventType: "backend", createdAt: "2026-07-19T16:00:02.000Z" }, { type: "assistant", eventType: "assistant", createdAt: "2026-07-19T16:01:02.000Z", final: true, replyAuthority: true, finalResponse: { text: "done" } }, { type: "result", eventType: "terminal", createdAt: "2026-07-19T16:01:03.000Z", terminal: true, terminalStatus: "completed", status: "completed" } ]; for (const event of events) { for (const listener of subscribers) listener({ schema: "hwlab.event.v1", sessionId: "ses_l0_semantics", traceId: "trc_l0_semantics", event }); } expect(await pending).toMatchObject({ ok: true, eventCount: 4, eventTypes: ["user", "backend", "assistant", "terminal"], waitFor: ["user", "backend", "assistant", "terminal", "final"], observedSemantics: ["user", "backend", "assistant", "terminal", "final"], missingSemantics: [], semanticCounts: { user: 1, backend: 1, assistant: 1, terminal: 1, final: 1 }, terminalStatuses: ["completed"], timing: { clock: "agentrun-source-event-created-at", phasesMs: { firstEventToBackend: 2000, admissionAndDispatch: null, runnerProvisioning: null, runtimePreparation: null, providerExecution: null, backendToAssistant: 60000, assistantToTerminal: 1000, firstEventToTerminal: 63000, firstEventToFinal: 62000 }, slowestGap: { durationMs: 60000, toIndex: 2, toEventType: "assistant" }, valuesPrinted: false } }); } finally { server.stop(true); await app.close(); } }); test("L0 CLI reports missing Kafka event semantics on timeout", async () => { const subscribers = new Set<(envelope: any) => void>(); const app = createWorkbenchHttpApp({ mode: "agentrun-native", kafkaEventBridge: { capabilities: { liveKafkaSse: true }, ready: Promise.resolve(), subscribeLiveHwlabEvents(listener: (envelope: any) => void) { subscribers.add(listener); return () => subscribers.delete(listener); } }, async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: {}, turns: {} }; } }); const server = Bun.serve({ port: 0, fetch: (request) => app.fetch(request) }); try { const pending = runWorkbenchCli([ "events", "inspect", "--trace-id", "trc_l0_missing_terminal", "--over-api", "--api-url", `http://127.0.0.1:${server.port}`, "--timeout-ms", "250", "--wait-for", "assistant,terminal" ], {}); for (let index = 0; index < 20 && subscribers.size === 0; index += 1) await Bun.sleep(2); for (const listener of subscribers) listener({ schema: "hwlab.event.v1", traceId: "trc_l0_missing_terminal", event: { type: "assistant", eventType: "assistant" } }); expect(await pending).toMatchObject({ ok: false, status: "event-timeout", observedSemantics: ["assistant"], missingSemantics: ["terminal"], terminalStatuses: [], error: { code: "workbench_kafka_sse_event_missing" } }); } finally { server.stop(true); await app.close(); } }); test("L0 SSE exposes Kafka startup failure without waiting for a read model", async () => { const startupError = Object.assign(new Error("Kafka broker unavailable"), { code: "kafka_broker_unavailable" }); const response = await createWorkbenchHttpApp({ mode: "agentrun-native", kafkaEventBridge: { capabilities: { liveKafkaSse: true }, liveReady: Promise.reject(startupError), subscribeLiveHwlabEvents() { return () => {}; } }, async dispatch() { return { ok: true }; }, async snapshot() { return { sessions: {}, turns: {} }; } }).fetch(new Request("http://localhost/v1/workbench/events?traceId=trc_l0_kafka_error")); expect(response.status).toBe(200); const body = await response.text(); expect(body).toContain("event: workbench.connected"); expect(body).toContain("event: workbench.error"); expect(body).toContain("kafka_broker_unavailable"); }); }); function stubApplication(): WorkbenchApplication { return { async health() { return {}; }, async createSession() { return {}; }, async admitTurn(input) { return input; }, async dispatchTurn() { return {}; }, async steerTurn() { return {}; }, async cancelTurn() { return {}; } }; } function nativeRailProjectionStub() { return { async waitWorkbenchSessionRailReady() { return { ready: true }; }, workbenchSessionSummaries(options: { sessionIds?: string[] } = {}) { return { ok: true, sessions: (options.sessionIds ?? []).map((sessionId, index) => ({ sessionId, firstUserMessagePreview: `title ${sessionId}`, lastUserMessageAt: "2026-07-20T00:00:00.000Z", lastTraceId: null, status: "idle", revision: { topic: "hwlab.event.v1", partition: 0, offset: String(index + 1), valuesRedacted: true }, valuesRedacted: true })) }; } }; }