diff --git a/internal/cloud/server-workbench-realtime-http.test.ts b/internal/cloud/server-workbench-realtime-http.test.ts index 17811e7d..782f1f18 100644 --- a/internal/cloud/server-workbench-realtime-http.test.ts +++ b/internal/cloud/server-workbench-realtime-http.test.ts @@ -45,6 +45,24 @@ const TRANSACTIONAL_REALTIME_ENV = Object.freeze({ HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "true" }); +function projectionRealtimeBridge(capabilities = {}) { + return { + started: true, + capabilities: { + directPublish: false, + liveKafkaSse: false, + transactionalProjector: false, + projectionOutboxRelay: false, + projectionRealtime: true, + ...capabilities + }, + ready: Promise.resolve(), + subscribeLiveHwlabEvents() { return () => {}; }, + subscribeProjectionCommits() { return () => {}; }, + async stop() {} + }; +} + test("projection outbox emits immutable assistant versions in row order", () => { const traceId = "trc_outbox_versions"; const sessionId = "ses_outbox_versions"; @@ -345,6 +363,7 @@ test("workbench realtime initial connection emits current snapshot without repla const server = createCloudApiServer({ accessController: realtimeAccessController(), workbenchRuntime: runtime, + kafkaEventBridge: projectionRealtimeBridge(), env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" } }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); @@ -384,6 +403,7 @@ test("live and projection realtime remain independently reachable when both capa const server = createCloudApiServer({ accessController: realtimeAccessController(), workbenchRuntime: runtime, + kafkaEventBridge: projectionRealtimeBridge({ directPublish: true, liveKafkaSse: true }), env: { HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true", HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", @@ -676,7 +696,7 @@ test("workbench realtime resets a cursor ahead of scoped cutoff with one fresh s return { facts: {}, events: [], cutoffOutboxSeq: 5, cursorOutboxSeq: 5, hasMore: false }; } }; - const server = createCloudApiServer({ accessController: realtimeAccessController(), workbenchRuntime: runtime, env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" } }); + const server = createCloudApiServer({ accessController: realtimeAccessController(), workbenchRuntime: runtime, kafkaEventBridge: projectionRealtimeBridge(), env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" } }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); try { const events = await getSseEvents(server.address().port, `/v1/workbench/events?sessionId=${sessionId}&afterSeq=42`, 2); @@ -726,6 +746,7 @@ test("workbench realtime stream surfaces facts blocker instead of legacy trace f traceStore, codeAgentChatResults: results, workbenchRuntime, + kafkaEventBridge: projectionRealtimeBridge(), env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", @@ -793,6 +814,7 @@ test("workbench session realtime follows durable outbox beyond stale lastTraceId const serverWithKafka = createCloudApiServer({ accessController, workbenchRuntime, + kafkaEventBridge: projectionRealtimeBridge(), env: { ...TRANSACTIONAL_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", diff --git a/internal/cloud/workbench-realtime-authority-contract.test.ts b/internal/cloud/workbench-realtime-authority-contract.test.ts index f1a6e689..a98aac3d 100644 --- a/internal/cloud/workbench-realtime-authority-contract.test.ts +++ b/internal/cloud/workbench-realtime-authority-contract.test.ts @@ -22,6 +22,22 @@ const LIVE_REALTIME_ENV = Object.freeze({ HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" }); +function projectionRealtimeBridge() { + return { + started: true, + capabilities: { + directPublish: false, + liveKafkaSse: false, + transactionalProjector: false, + projectionOutboxRelay: false, + projectionRealtime: true + }, + ready: Promise.resolve(), + subscribeProjectionCommits() { return () => {}; }, + async stop() {} + }; +} + test("workbench sync fails closed before projection storage in live Kafka capability", async () => { let syncReads = 0; const server = createCloudApiServer({ @@ -104,6 +120,7 @@ test("workbench sync and SSE read the same atomic projection outbox", async () = const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, + kafkaEventBridge: projectionRealtimeBridge(), env: { ...PROJECTION_REALTIME_ENV, HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", @@ -152,7 +169,7 @@ test("workbench sync delta includes all authority families and marks trace rows }), outboxRows: [] }); - const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, env: { ...PROJECTION_REALTIME_ENV } }); + const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, kafkaEventBridge: projectionRealtimeBridge(), env: { ...PROJECTION_REALTIME_ENV } }); await listen(server); try { @@ -182,7 +199,7 @@ test("workbench sync preserves session and trace scope for narrow replay", async const traceId = "trc_realtime_authority_narrow"; const outboxQueries = []; const runtime = createRuntime({ facts: durableFacts({ sessionId, traceId }), outboxQueries }); - const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, env: { ...PROJECTION_REALTIME_ENV } }); + const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, kafkaEventBridge: projectionRealtimeBridge(), env: { ...PROJECTION_REALTIME_ENV } }); await listen(server); try { @@ -205,7 +222,7 @@ test("workbench trace event detail route declares detail-only authority", async facts: durableFacts({ sessionId, traceId, finalText: "sealed final", traceEventText: "detail row" }), outboxRows: [] }); - const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, env: { ...PROJECTION_REALTIME_ENV } }); + const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: runtime, kafkaEventBridge: projectionRealtimeBridge(), env: { ...PROJECTION_REALTIME_ENV } }); await listen(server); try { @@ -224,7 +241,7 @@ test("workbench trace event detail route declares detail-only authority", async }); test("workbench sync rejects unscoped automatic repair requests", async () => { - const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: createRuntime({ facts: durableFacts({}) }), env: { ...PROJECTION_REALTIME_ENV } }); + const server = createCloudApiServer({ accessController: createAccessController(), workbenchRuntime: createRuntime({ facts: durableFacts({}) }), kafkaEventBridge: projectionRealtimeBridge(), env: { ...PROJECTION_REALTIME_ENV } }); await listen(server); try {