From f9f9f42c4ac7bdbf0b27287a5f5b52660d39f58b Mon Sep 17 00:00:00 2001 From: AgentRun-Codex Date: Tue, 14 Jul 2026 12:43:46 +0200 Subject: [PATCH] fix: restore Kafka workbench authority --- deploy/deploy.yaml | 8 + internal/cloud/kafka-event-bridge.test.ts | 78 ++-- internal/cloud/kafka-event-bridge.ts | 393 ++++++++++++++-- .../cloud/server-code-agent-admission-http.ts | 75 ++-- .../server-workbench-realtime-http.test.ts | 41 +- .../cloud/server-workbench-realtime-http.ts | 421 +++++++++++++++++- internal/cloud/server.ts | 3 + .../cloud/workbench-realtime-capabilities.ts | 35 ++ internal/dev-entrypoint/cloud-web-runtime.mjs | 5 + scripts/gitops-render.test.ts | 24 +- ...time-realtime-capabilities-config.test.mjs | 29 +- ...h-fixed-transactional-architecture.test.ts | 110 ++--- .../workbench-realtime-runtime.test.ts | 14 +- .../src/api/workbench-events.test.ts | 2 + .../src/api/workbench-events.ts | 16 +- web/hwlab-cloud-web/src/config/runtime.ts | 22 + .../src/stores/workbench-event-reducer.ts | 8 + .../stores/workbench-isolated-kafka-debug.ts | 17 +- .../src/stores/workbench-live-kafka-event.ts | 4 +- web/hwlab-cloud-web/src/stores/workbench.ts | 96 +++- web/hwlab-cloud-web/src/types/global.d.ts | 5 + .../src/utils/workbench-realtime-runtime.ts | 2 +- .../src/utils/workbench-stream-transport.ts | 16 +- 23 files changed, 1170 insertions(+), 254 deletions(-) create mode 100644 internal/cloud/workbench-realtime-capabilities.ts diff --git a/deploy/deploy.yaml b/deploy/deploy.yaml index 2ae60ebb..5ac749c0 100644 --- a/deploy/deploy.yaml +++ b/deploy/deploy.yaml @@ -621,6 +621,9 @@ lanes: HWLAB_CLOUD_WEB_DISPLAY_TIME_ZONE: Asia/Shanghai HWLAB_CLOUD_WEB_DISPLAY_TIME_LOCALE: zh-CN HWLAB_CLOUD_WEB_DISPLAY_TIME_LABEL: 北京时间 + HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true" + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "true" + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true" HWLAB_CLOUD_WEB_OPENCODE_UPSTREAM_URL: http://opencode-server.hwlab-v03.svc.cluster.local:4096 HWLAB_CLOUD_WEB_OPENCODE_USERNAME: secretRef:hwlab-opencode-server-auth/username @@ -847,6 +850,11 @@ lanes: HWLAB_CODE_AGENT_AGENTRUN_REPO_URL: http://gitea-http.devops-infra.svc.cluster.local:3000/mirrors/pikasTech-HWLAB.git HWLAB_KAFKA_ENABLED: "true" HWLAB_KAFKA_AGENTRUN_EVENT_CONSUME_ENABLED: "true" + HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true" + HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true" + HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false" + HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false" + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true" HWLAB_KAFKA_BOOTSTRAP_SERVERS: platform-infra-kafka-kafka-bootstrap.platform-infra.svc.cluster.local:9092 HWLAB_KAFKA_STDIO_TOPIC: codex-stdio.raw.v1 diff --git a/internal/cloud/kafka-event-bridge.test.ts b/internal/cloud/kafka-event-bridge.test.ts index 9ed5f38b..68b5ecd2 100644 --- a/internal/cloud/kafka-event-bridge.test.ts +++ b/internal/cloud/kafka-event-bridge.test.ts @@ -31,6 +31,23 @@ const PROJECTOR_ENV = Object.freeze({ HWLAB_KAFKA_OUTBOX_RELAY_RETRY_BACKOFF_MS: "1000" }); +const LIVE_ENV = Object.freeze({ + HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true", + HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", + HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", + HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false", + HWLAB_KAFKA_ENABLED: "true", + HWLAB_KAFKA_AGENTRUN_EVENT_CONSUME_ENABLED: "true", + HWLAB_KAFKA_BOOTSTRAP_SERVERS: "kafka.test:9092", + HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC: "agentrun.event.v1", + HWLAB_KAFKA_EVENT_TOPIC: "hwlab.event.v1", + HWLAB_KAFKA_CLIENT_ID: "hwlab-live-test", + HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID: "hwlab-live-bridge-fixed", + HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID: "hwlab-live-fanout-fixed" +}); + test("projects AgentRun assistant_message Kafka event into HWLAB trace event", () => { const projected = projectAgentRunKafkaEventToHwlabEvent({ schema: "agentrun.event.v1", @@ -305,40 +322,39 @@ test("projects AgentRun terminal_status Kafka event into terminal HWLAB event", assert.equal(projected.event.message, "AgentRun completed"); }); -test("Workbench bridge ignores obsolete architecture env and fixes transactional authority", () => { - const config = kafkaEventBridgeConfig({ - ...PROJECTOR_ENV, - HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true", - HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", - HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "true", - HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", - HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" +test("Workbench bridge fixes direct publish and live Kafka authority", () => { + const config = kafkaEventBridgeConfig(LIVE_ENV); + assert.deepEqual(config?.capabilities, { + directPublish: true, + liveKafkaSse: true, + kafkaRefreshReplay: false, + transactionalProjector: false, + projectionOutboxRelay: false, + projectionRealtime: false }); - - assert.equal(config?.projectorGroupId, "hwlab-projector-test-v1"); - assert.deepEqual(config?.relay, { - intervalMs: 60000, - batchSize: 1, - leaseMs: 30000, - sendTimeoutMs: 10000, - retryBackoffMs: 1000 - }); - assert.deepEqual(Object.keys(config ?? {}).sort(), ["agentRunTopic", "brokers", "clientId", "hwlabTopic", "projectorGroupId", "projectorHeartbeatIntervalMs", "relay"]); + assert.equal(config?.directPublishGroupId, "hwlab-live-bridge-fixed"); + assert.equal(config?.hwlabEventGroupId, "hwlab-live-fanout-fixed"); + assert.equal(config?.projectorHeartbeatIntervalMs, null); + assert.equal(config?.relay, null); }); -test("Workbench bridge does not require the six obsolete architecture env", () => { - const env = { ...PROJECTOR_ENV }; - for (const name of [ - "HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED", - "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED", - "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED", - "HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED", - "HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED", - "HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED" - ]) delete env[name]; - - assert.doesNotThrow(() => kafkaEventBridgeConfig(env)); +test("live Kafka startup does not touch transactional v8 readiness", async () => { + let runtimeReads = 0; + const runtimeStore = new Proxy({}, { get() { runtimeReads += 1; throw new Error("live Kafka must not require transactional v8 schema"); } }); + const consumers = []; + const kafkaFactory = () => ({ + consumer() { + const consumer = { async connect() {}, async subscribe() {}, async run() {}, async stop() {}, async disconnect() {} }; + consumers.push(consumer); + return consumer; + }, + producer() { return { async connect() {}, async send() {}, async disconnect() {} }; } + }); + const bridge = startHwlabKafkaEventBridge({ env: LIVE_ENV, runtimeStore, kafkaFactory, logger: null }); + await bridge.ready; + assert.equal(runtimeReads, 0); + assert.equal(consumers.length, 2); + await bridge.stop(); }); test("projects AgentRun run-created envelope even when no source event is present", () => { diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index 6462721e..3f2670ac 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -4,7 +4,12 @@ import { createHash, randomUUID } from "node:crypto"; import { Kafka, logLevel } from "kafkajs"; +import { emitCodeAgentOtelSpan } from "./otel-trace.ts"; import { buildWorkbenchProjectionEventFacts } from "./workbench-projection-writer.ts"; +import { + WORKBENCH_REALTIME_CAPABILITY_ENVS, + workbenchRealtimeCapabilities +} from "./workbench-realtime-capabilities.ts"; const TRUE_VALUES = new Set(["1", "true", "yes", "on"]); export const DEFAULT_AGENTRUN_EVENT_TOPIC = "agentrun.event.v1"; @@ -17,6 +22,7 @@ export const AGENTRUN_STDIO_RECONSTRUCTION_KIND = "stdio-derived partial reconst const DEFAULT_CLIENT_ID = "hwlab-v03-cloud-api"; const DEFAULT_QUERY_TIMEOUT_MS = 5000; const DEFAULT_QUERY_LIMIT = 50; +const LIVE_KAFKA_COMMAND_OUTPUT_OTEL_SAMPLE_MODULUS = 32; const REQUIRED_KAFKA_ENV = Object.freeze([ "HWLAB_KAFKA_BOOTSTRAP_SERVERS", "HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC", @@ -25,19 +31,35 @@ const REQUIRED_KAFKA_ENV = Object.freeze([ ]); export function kafkaEventBridgeConfig(env = process.env) { + const featureEnvPresent = Object.values(WORKBENCH_REALTIME_CAPABILITY_ENVS).some((name) => stringValue(env[name])); const legacyKafkaEnabled = truthy(env.HWLAB_KAFKA_AGENTRUN_EVENT_CONSUME_ENABLED ?? env.HWLAB_KAFKA_AGENTRUN_CONSUME_ENABLED ?? env.HWLAB_KAFKA_ENABLED); - const kafkaConfigured = REQUIRED_KAFKA_ENV.every((name) => stringValue(env[name])); - if (!legacyKafkaEnabled && !kafkaConfigured) return null; - const required = [ - ...REQUIRED_KAFKA_ENV, - "HWLAB_KAFKA_PROJECTOR_GROUP_ID", - "HWLAB_KAFKA_PROJECTOR_HEARTBEAT_INTERVAL_MS", + if (!featureEnvPresent && !legacyKafkaEnabled) return null; + const capabilities = workbenchRealtimeCapabilities(env); + if (capabilities.kafkaRefreshReplay && !capabilities.liveKafkaSse) { + throw contractError("hwlab_kafka_refresh_live_dependency_missing", "Kafka refresh replay requires the independently enabled live Kafka SSE capability for handoff."); + } + const kafkaEnabled = capabilities.directPublish || capabilities.liveKafkaSse || capabilities.kafkaRefreshReplay || capabilities.transactionalProjector || capabilities.projectionOutboxRelay; + const transactionalEnabled = capabilities.transactionalProjector || capabilities.projectionOutboxRelay || capabilities.projectionRealtime; + if (!kafkaEnabled && !transactionalEnabled) return null; + const required = [...REQUIRED_KAFKA_ENV]; + if (!kafkaEnabled) required.length = 0; + if (capabilities.directPublish) required.push("HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID"); + if (capabilities.liveKafkaSse) required.push("HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID"); + if (capabilities.kafkaRefreshReplay) required.push( + "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_GROUP_PREFIX", + "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_TIMEOUT_MS", + "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_SCAN_LIMIT", + "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_MATCHED_EVENT_LIMIT", + "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_LIVE_BUFFER_LIMIT" + ); + if (capabilities.transactionalProjector) required.push("HWLAB_KAFKA_PROJECTOR_GROUP_ID", "HWLAB_KAFKA_PROJECTOR_HEARTBEAT_INTERVAL_MS"); + if (capabilities.projectionOutboxRelay) required.push( "HWLAB_KAFKA_OUTBOX_RELAY_INTERVAL_MS", "HWLAB_KAFKA_OUTBOX_RELAY_BATCH_SIZE", "HWLAB_KAFKA_OUTBOX_RELAY_LEASE_MS", "HWLAB_KAFKA_OUTBOX_RELAY_SEND_TIMEOUT_MS", "HWLAB_KAFKA_OUTBOX_RELAY_RETRY_BACKOFF_MS" - ]; + ); const missing = required.filter((name) => !stringValue(env[name])); if (missing.length > 0) { const error = new Error(`HWLAB Kafka capability configuration is incomplete: ${missing.join(", ")}`); @@ -49,37 +71,66 @@ export function kafkaEventBridgeConfig(env = process.env) { const agentRunTopic = stringValue(env.HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC); const hwlabTopic = stringValue(env.HWLAB_KAFKA_EVENT_TOPIC); const clientId = stringValue(env.HWLAB_KAFKA_CLIENT_ID); - const relay = { + const relay = capabilities.projectionOutboxRelay ? { intervalMs: strictPositiveInteger(env.HWLAB_KAFKA_OUTBOX_RELAY_INTERVAL_MS, "HWLAB_KAFKA_OUTBOX_RELAY_INTERVAL_MS"), batchSize: strictPositiveInteger(env.HWLAB_KAFKA_OUTBOX_RELAY_BATCH_SIZE, "HWLAB_KAFKA_OUTBOX_RELAY_BATCH_SIZE"), leaseMs: strictPositiveInteger(env.HWLAB_KAFKA_OUTBOX_RELAY_LEASE_MS, "HWLAB_KAFKA_OUTBOX_RELAY_LEASE_MS"), sendTimeoutMs: strictPositiveInteger(env.HWLAB_KAFKA_OUTBOX_RELAY_SEND_TIMEOUT_MS, "HWLAB_KAFKA_OUTBOX_RELAY_SEND_TIMEOUT_MS"), retryBackoffMs: strictPositiveInteger(env.HWLAB_KAFKA_OUTBOX_RELAY_RETRY_BACKOFF_MS, "HWLAB_KAFKA_OUTBOX_RELAY_RETRY_BACKOFF_MS") - }; + } : null; + const refreshReplay = capabilities.kafkaRefreshReplay ? { + groupIdPrefix: stringValue(env.HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_GROUP_PREFIX), + timeoutMs: strictPositiveInteger(env.HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_TIMEOUT_MS, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_TIMEOUT_MS"), + scanLimit: strictPositiveInteger(env.HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_SCAN_LIMIT, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_SCAN_LIMIT"), + matchedEventLimit: strictPositiveInteger(env.HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_MATCHED_EVENT_LIMIT, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_MATCHED_EVENT_LIMIT"), + liveBufferLimit: strictPositiveInteger(env.HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_LIVE_BUFFER_LIMIT, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_LIVE_BUFFER_LIMIT") + } : null; if (relay && relay.sendTimeoutMs + 5000 >= relay.leaseMs) { const error = new Error("HWLAB Kafka relay lease must exceed producer send timeout by at least 5000ms."); error.code = "hwlab_kafka_relay_lease_invalid"; throw error; } + const enabledConsumerGroups = [ + capabilities.directPublish ? stringValue(env.HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID) : null, + capabilities.transactionalProjector ? stringValue(env.HWLAB_KAFKA_PROJECTOR_GROUP_ID) : null, + capabilities.liveKafkaSse ? stringValue(env.HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID) : null + ].filter(Boolean); + if (new Set(enabledConsumerGroups).size !== enabledConsumerGroups.length) { + const error = new Error("HWLAB Kafka direct publisher, projector, and live SSE consumer groups must be distinct when enabled together."); + error.code = "hwlab_kafka_consumer_groups_not_distinct"; + error.groupIds = enabledConsumerGroups; + throw error; + } return { + capabilities, brokers, agentRunTopic, hwlabTopic, clientId, + directPublishGroupId: stringValue(env.HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID), projectorGroupId: stringValue(env.HWLAB_KAFKA_PROJECTOR_GROUP_ID), - projectorHeartbeatIntervalMs: strictPositiveInteger(env.HWLAB_KAFKA_PROJECTOR_HEARTBEAT_INTERVAL_MS, "HWLAB_KAFKA_PROJECTOR_HEARTBEAT_INTERVAL_MS"), + hwlabEventGroupId: stringValue(env.HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID), + refreshReplay, + projectorHeartbeatIntervalMs: capabilities.transactionalProjector ? strictPositiveInteger(env.HWLAB_KAFKA_PROJECTOR_HEARTBEAT_INTERVAL_MS, "HWLAB_KAFKA_PROJECTOR_HEARTBEAT_INTERVAL_MS") : null, relay }; } -export function startHwlabKafkaEventBridge({ env = process.env, logger = console, kafkaFactory = defaultKafkaFactory, runtimeStore = null, now = () => new Date().toISOString() } = {}) { +export function startHwlabKafkaEventBridge({ env = process.env, logger = console, kafkaFactory = defaultKafkaFactory, runtimeStore = null, now = () => new Date().toISOString(), otelSpanEmitter = emitCodeAgentOtelSpan } = {}) { const config = kafkaEventBridgeConfig(env); if (!config) return { started: false, reason: "disabled_or_unconfigured", stop() {}, subscribeProjectionCommits() { return () => {}; }, valuesPrinted: false }; - return startTransactionalHwlabKafkaEventBridge({ config, logger, kafkaFactory, runtimeStore, now }); + const components = []; + if (config.capabilities.directPublish || config.capabilities.liveKafkaSse) { + components.push(startLiveHwlabKafkaEventBridge({ config, env, logger, kafkaFactory, otelSpanEmitter })); + } + if (config.capabilities.transactionalProjector || config.capabilities.projectionOutboxRelay || config.capabilities.projectionRealtime) { + components.push(startTransactionalHwlabKafkaEventBridge({ config, logger, kafkaFactory, runtimeStore, now })); + } + return components.length === 1 ? components[0] : combineKafkaEventBridgeComponents(config, components); } function startTransactionalHwlabKafkaEventBridge({ config, logger = console, kafkaFactory = defaultKafkaFactory, runtimeStore = null, now = () => new Date().toISOString() } = {}) { - requireKafkaProjectorStore(runtimeStore); + requireKafkaProjectorStore(runtimeStore, config.capabilities); let stopped = false; let consumer = null; @@ -98,14 +149,16 @@ function startTransactionalHwlabKafkaEventBridge({ config, logger = console, kaf let startupReady = false; const ready = (async () => { await runtimeStore.assertWorkbenchTransactionalRealtimeReady(); - const kafka = kafkaFactory(config); - consumer = kafka.consumer({ groupId: config.projectorGroupId, allowAutoTopicCreation: false }); - producer = kafka.producer({ allowAutoTopicCreation: false }); - await producer.connect(); - await consumer.connect(); - stopProjectionNotifications = await runtimeStore.subscribeWorkbenchProjectionCommits(publishProjectionCommit); - await consumer.subscribe({ topic: config.agentRunTopic, fromBeginning: false }); - await consumer.run({ + const kafka = config.capabilities.transactionalProjector || config.capabilities.projectionOutboxRelay ? kafkaFactory(config) : null; + if (config.capabilities.transactionalProjector) consumer = kafka.consumer({ groupId: config.projectorGroupId, allowAutoTopicCreation: false }); + if (config.capabilities.projectionOutboxRelay) producer = kafka.producer({ allowAutoTopicCreation: false }); + if (producer) await producer.connect(); + if (consumer) await consumer.connect(); + if (config.capabilities.projectionRealtime && typeof runtimeStore.subscribeWorkbenchProjectionCommits === "function") { + stopProjectionNotifications = await runtimeStore.subscribeWorkbenchProjectionCommits(publishProjectionCommit); + } + if (consumer) await consumer.subscribe({ topic: config.agentRunTopic, fromBeginning: false }); + if (consumer) await consumer.run({ autoCommit: false, eachBatchAutoResolve: false, eachBatch: async ({ batch, resolveOffset, heartbeat, isRunning, isStale }) => { @@ -132,6 +185,7 @@ function startTransactionalHwlabKafkaEventBridge({ config, logger = console, kaf } }); const runRelay = async () => { + if (!config.capabilities.projectionOutboxRelay) return; if (stopped || relayRunning) return; relayRunning = true; try { @@ -140,15 +194,18 @@ function startTransactionalHwlabKafkaEventBridge({ config, logger = console, kaf relayRunning = false; } }; - relayTimer = setInterval(() => { void runRelay().catch((error) => logWarn(logger, "hwlab-kafka-outbox-relay-failed", { message: errorMessage(error), valuesPrinted: false })); }, config.relay.intervalMs); - relayTimer.unref?.(); - await runRelay(); + if (config.capabilities.projectionOutboxRelay) { + relayTimer = setInterval(() => { void runRelay().catch((error) => logWarn(logger, "hwlab-kafka-outbox-relay-failed", { message: errorMessage(error), valuesPrinted: false })); }, config.relay.intervalMs); + relayTimer.unref?.(); + await runRelay(); + } startupReady = true; logInfo(logger, "hwlab-kafka-event-bridge-started", { agentRunTopic: config.agentRunTopic, hwlabTopic: config.hwlabTopic, clientId: config.clientId, projectorGroupId: config.projectorGroupId, + capabilities: config.capabilities, valuesPrinted: false }); })(); @@ -181,8 +238,10 @@ function startTransactionalHwlabKafkaEventBridge({ config, logger = console, kaf const transactionalReadiness = typeof runtimeStore.workbenchTransactionalRealtimeReadiness === "function" ? await runtimeStore.workbenchTransactionalRealtimeReadiness() : { ready: true, valuesRedacted: true }; - const projectorStatus = await runtimeStore.hwlabKafkaProjectorStatus(); - return { status: transactionalReadiness.ready === false ? "blocked" : "ready", transactionalReadiness, ...projectorStatus, valuesRedacted: true }; + const projectorStatus = config.capabilities.transactionalProjector || config.capabilities.projectionOutboxRelay + ? await runtimeStore.hwlabKafkaProjectorStatus() + : {}; + return { status: transactionalReadiness.ready === false ? "blocked" : "ready", capabilities: config.capabilities, transactionalReadiness, ...projectorStatus, valuesRedacted: true }; }, startupStatus() { if (startupError) return { status: "blocked", errorCode: startupError.code ?? "hwlab_kafka_projector_start_failed", message: errorMessage(startupError), valuesRedacted: true }; @@ -197,6 +256,269 @@ function startTransactionalHwlabKafkaEventBridge({ config, logger = console, kaf }; } +function startLiveHwlabKafkaEventBridge({ config, env = process.env, logger = console, kafkaFactory = defaultKafkaFactory, otelSpanEmitter = emitCodeAgentOtelSpan } = {}) { + let stopped = false; + let bridgeConsumer = null; + let fanoutConsumer = null; + let producer = null; + let startupError = null; + let startupReady = false; + const subscribers = new Set(); + const lagByPartition = new Map(); + const publishLiveEnvelope = (envelope, transport) => { + for (const listener of [...subscribers]) { + try { + listener(envelope, transport); + } catch (error) { + logWarn(logger, "hwlab-live-kafka-subscriber-failed", { message: errorMessage(error), valuesPrinted: false }); + } + } + }; + const ready = (async () => { + const kafka = kafkaFactory(config); + if (config.capabilities.directPublish) { + bridgeConsumer = kafka.consumer({ groupId: config.directPublishGroupId, allowAutoTopicCreation: false }); + producer = kafka.producer({ allowAutoTopicCreation: false }); + } + if (config.capabilities.liveKafkaSse) fanoutConsumer = kafka.consumer({ groupId: config.hwlabEventGroupId, allowAutoTopicCreation: false }); + if (producer) await producer.connect(); + if (bridgeConsumer) await bridgeConsumer.connect(); + if (fanoutConsumer) await fanoutConsumer.connect(); + if (bridgeConsumer) await bridgeConsumer.subscribe({ topic: config.agentRunTopic, fromBeginning: false }); + if (fanoutConsumer) await fanoutConsumer.subscribe({ topic: config.hwlabTopic, fromBeginning: false }); + if (fanoutConsumer) await fanoutConsumer.run({ + eachBatchAutoResolve: false, + eachBatch: async ({ batch, resolveOffset, heartbeat, isRunning, isStale }) => { + let lastOffset = null; + for (const message of batch.messages) { + if (stopped || !isRunning() || isStale()) break; + const envelope = parseJson(message?.value ? Buffer.from(message.value).toString("utf8") : ""); + if (!envelope || envelope.schema !== "hwlab.event.v1") { + logWarn(logger, "hwlab-live-kafka-envelope-ignored", { reason: "schema-invalid", valuesPrinted: false }); + } else { + const transport = { topic: batch.topic, partition: batch.partition, offset: message.offset }; + emitLiveKafkaOtelSpan("hwlab.kafka.live.fanout_receive", envelope, transport, { env, otelSpanEmitter }); + publishLiveEnvelope(envelope, transport); + } + resolveOffset(message.offset); + lastOffset = message.offset; + await heartbeat(); + } + if (lastOffset !== null) lagByPartition.set(batch.partition, kafkaPartitionLag(batch.highWatermark, lastOffset)); + } + }); + if (bridgeConsumer) await bridgeConsumer.run({ + eachMessage: async ({ topic, partition, message }) => { + if (stopped) return; + await publishAgentRunKafkaMessageLive({ topic, partition, message }, { producer, config, env, logger, otelSpanEmitter }); + } + }); + startupReady = true; + logInfo(logger, "hwlab-live-kafka-event-bridge-started", { + capabilities: config.capabilities, + agentRunTopic: config.agentRunTopic, + hwlabTopic: config.hwlabTopic, + clientId: config.clientId, + directPublishGroupId: config.directPublishGroupId, + hwlabEventGroupId: config.hwlabEventGroupId, + valuesPrinted: false + }); + })(); + void ready.catch((error) => { + startupError = error; + logWarn(logger, "hwlab-live-kafka-event-bridge-start-failed", { + message: errorMessage(error), + agentRunTopic: config.agentRunTopic, + hwlabTopic: config.hwlabTopic, + valuesPrinted: false + }); + }); + return { + started: true, + ...config, + ready, + liveReady: ready, + async stop() { + stopped = true; + subscribers.clear(); + await Promise.allSettled([bridgeConsumer?.stop?.(), fanoutConsumer?.stop?.()]); + await Promise.allSettled([bridgeConsumer?.disconnect?.(), fanoutConsumer?.disconnect?.(), producer?.disconnect?.()]); + }, + async status() { + if (startupError) return { status: "blocked", capabilities: config.capabilities, errorCode: startupError.code ?? "hwlab_live_kafka_start_failed", message: errorMessage(startupError), valuesRedacted: true }; + return { status: startupReady ? "ready" : "initializing", capabilities: config.capabilities, subscriberCount: subscribers.size, consumerLag: liveKafkaLagSummary(lagByPartition), lossPossible: true, valuesRedacted: true }; + }, + startupStatus() { + if (startupError) return { status: "blocked", capabilities: config.capabilities, errorCode: startupError.code ?? "hwlab_live_kafka_start_failed", message: errorMessage(startupError), valuesRedacted: true }; + return { status: startupReady ? "ready" : "initializing", capabilities: config.capabilities, valuesRedacted: true }; + }, + subscribeProjectionCommits() { return () => {}; }, + subscribeLiveHwlabEvents(listener) { + if (typeof listener !== "function") return () => {}; + subscribers.add(listener); + return () => subscribers.delete(listener); + }, + queryHwlabEventRetention(params = {}) { + return queryKafkaEventStream({ ...params, env, stream: "hwlab", topic: config.hwlabTopic, kafkaFactory }); + }, + liveSubscriberCount() { return subscribers.size; }, + valuesPrinted: false + }; +} + +function kafkaPartitionLag(highWatermark, lastOffset) { + try { + const lag = BigInt(String(highWatermark ?? "0")) - (BigInt(String(lastOffset ?? "0")) + 1n); + return Number(lag > 0n ? lag : 0n); + } catch { + return null; + } +} + +function liveKafkaLagSummary(lagByPartition) { + const partitions = [...lagByPartition.entries()].map(([partition, lag]) => ({ partition, lag })).sort((left, right) => left.partition - right.partition); + const known = partitions.map((item) => item.lag).filter((lag) => Number.isFinite(lag)); + return { + known: known.length > 0, + total: known.length > 0 ? known.reduce((sum, lag) => sum + lag, 0) : null, + partitions, + valuesRedacted: true + }; +} + +function combineKafkaEventBridgeComponents(config, components) { + const ready = Promise.all(components.map((component) => component.ready)); + const liveOwner = components.find((component) => typeof component.subscribeLiveHwlabEvents === "function"); + const componentStartup = () => components.map((component) => component.startupStatus?.() ?? { status: "initializing" }); + return { + started: true, + ...config, + ready, + liveReady: liveOwner?.liveReady ?? liveOwner?.ready ?? ready, + async stop() { await Promise.allSettled(components.map((component) => component.stop())); }, + async status() { + const statuses = await Promise.all(components.map((component) => component.status())); + const blocked = statuses.find((status) => status.status === "blocked"); + const initializing = statuses.some((status) => status.status !== "ready"); + return { + status: blocked ? "blocked" : initializing ? "initializing" : "ready", + capabilities: config.capabilities, + components: statuses, + errorCode: blocked?.errorCode ?? null, + message: blocked?.message ?? null, + valuesRedacted: true + }; + }, + startupStatus() { + const statuses = componentStartup(); + const blocked = statuses.find((status) => status.status === "blocked"); + return { + status: blocked ? "blocked" : statuses.every((status) => status.status === "ready") ? "ready" : "initializing", + capabilities: config.capabilities, + components: statuses, + errorCode: blocked?.errorCode ?? null, + message: blocked?.message ?? null, + valuesRedacted: true + }; + }, + subscribeProjectionCommits(listener) { + const stops = components.map((component) => component.subscribeProjectionCommits?.(listener)).filter((stop) => typeof stop === "function"); + return () => stops.forEach((stop) => stop()); + }, + subscribeLiveHwlabEvents(listener) { + return liveOwner?.subscribeLiveHwlabEvents(listener) ?? (() => {}); + }, + queryHwlabEventRetention(params = {}) { + if (typeof liveOwner?.queryHwlabEventRetention !== "function") throw contractError("hwlab_kafka_refresh_query_unconfigured", "Kafka refresh retention query is not configured."); + return liveOwner.queryHwlabEventRetention(params); + }, + liveSubscriberCount() { + const owner = components.find((component) => typeof component.liveSubscriberCount === "function"); + return owner?.liveSubscriberCount() ?? 0; + }, + valuesPrinted: false + }; +} + +export async function publishAgentRunKafkaMessageLive(kafkaMessage = {}, { producer, config, env = process.env, logger = console, otelSpanEmitter = emitCodeAgentOtelSpan } = {}) { + const transport = kafkaTransport(kafkaMessage); + const decoded = decodeCanonicalAgentRunKafkaMessage(kafkaMessage); + if (!decoded.ok) { + logWarn(logger, "hwlab-live-kafka-message-invalid", { topic: transport.sourceTopic, partition: transport.sourcePartition, offset: transport.sourceOffset, errorCode: decoded.error.code, valuesPrinted: false }); + return { handled: true, invalid: true, published: false, valuesPrinted: false }; + } + if (!stringValue(decoded.event.traceId) || !stringValue(decoded.event.hwlabSessionId)) { + return { handled: true, ignored: true, published: false, valuesPrinted: false }; + } + const projected = projectAgentRunKafkaEventToHwlabEvent(decoded.event, { + source: config.clientId, + sourceTopic: transport.sourceTopic, + sourcePartition: transport.sourcePartition, + sourceOffset: transport.sourceOffset, + sourceKey: transport.sourceKey, + inputSha256: transport.inputSha256 + }); + const publishMetadata = await producer.send({ topic: config.hwlabTopic, messages: [buildHwlabKafkaProducerMessage(projected)] }); + const publishedRecord = Array.isArray(publishMetadata) ? publishMetadata[0] : null; + emitLiveKafkaOtelSpan("hwlab.kafka.live.direct_publish", projected, { + topic: config.hwlabTopic, + partition: publishedRecord?.partition, + offset: publishedRecord?.baseOffset ?? publishedRecord?.offset, + sourceTopic: transport.sourceTopic, + sourcePartition: transport.sourcePartition, + sourceOffset: transport.sourceOffset + }, { env, otelSpanEmitter }); + return { handled: true, published: true, envelope: projected, sessionId: projected.sessionId, traceId: projected.traceId, valuesPrinted: false }; +} + +export function shouldEmitLiveKafkaOtelSpan(envelope = {}) { + const event = objectValue(envelope.event); + const context = objectValue(envelope.context); + const eventType = firstText(context.agentRunEventType, event.eventType, event.type, envelope.eventType)?.toLowerCase(); + if (event.terminal === true || eventType === "terminal" || eventType === "terminal_status") return true; + if (["assistant", "assistant_progress", "assistant_message", "tool", "tool_call"].includes(eventType)) return true; + if (eventType !== "command_output" && event.type !== "output") return true; + const sourceSeq = integerValue(context.sourceSeq ?? event.sourceSeq); + if (Number.isInteger(sourceSeq) && sourceSeq > 0) return sourceSeq % LIVE_KAFKA_COMMAND_OUTPUT_OTEL_SAMPLE_MODULUS === 0; + const identity = firstText(envelope.sourceEventId, envelope.eventId, envelope.traceId, envelope.hwlabSessionId, envelope.sessionId); + if (!identity) return false; + return createHash("sha256").update(identity).digest().readUInt32BE(0) % LIVE_KAFKA_COMMAND_OUTPUT_OTEL_SAMPLE_MODULUS === 0; +} + +export function liveKafkaOtelSpanAttributes(envelope = {}, transport = {}) { + const event = objectValue(envelope.event); + const context = objectValue(envelope.context); + const sourceEvent = objectValue(envelope.sourceEvent); + return { + businessTraceId: firstText(envelope.traceId, event.traceId), + hwlabSessionId: firstText(envelope.hwlabSessionId, envelope.sessionId, event.sessionId), + runId: firstText(envelope.runId, context.runId, event.runId), + commandId: firstText(envelope.commandId, context.commandId, event.commandId), + topic: firstText(transport.topic, sourceEvent.topic), + partition: integerValue(transport.partition ?? sourceEvent.partition), + offset: firstText(transport.offset, sourceEvent.offset), + sourceTopic: firstText(transport.sourceTopic), + sourcePartition: integerValue(transport.sourcePartition), + sourceOffset: firstText(transport.sourceOffset), + eventType: firstText(context.agentRunEventType, event.eventType, event.type, envelope.eventType), + terminal: event.terminal === true, + valuesRedacted: true + }; +} + +export function emitLiveKafkaOtelSpan(name, envelope = {}, transport = {}, { env = process.env, otelSpanEmitter = emitCodeAgentOtelSpan } = {}) { + if (typeof otelSpanEmitter !== "function" || !shouldEmitLiveKafkaOtelSpan(envelope)) return { emitted: false, valuesRedacted: true }; + const attributes = liveKafkaOtelSpanAttributes(envelope, transport); + if (!attributes.businessTraceId) return { emitted: false, valuesRedacted: true }; + try { + const pending = otelSpanEmitter(name, attributes.businessTraceId, env, { attributes }); + void Promise.resolve(pending).catch(() => undefined); + } catch { + // OTel is best effort and must never change the live Kafka delivery path. + } + return { emitted: true, valuesRedacted: true }; +} + export async function projectAgentRunKafkaMessage(kafkaMessage = {}, { runtimeStore, config, logger = console } = {}) { const transport = kafkaTransport(kafkaMessage); const decoded = decodeCanonicalAgentRunKafkaMessage(kafkaMessage); @@ -620,18 +942,13 @@ function kafkaTransport(kafkaMessage = {}) { }; } -function requireKafkaProjectorStore(runtimeStore) { - const required = [ - "commitAgentRunKafkaProjection", - "recordFailedAgentRunKafkaMessage", - "recordIgnoredAgentRunKafkaMessage", - "claimHwlabKafkaOutbox", - "completeHwlabKafkaOutbox", - "retryHwlabKafkaOutbox", - "subscribeWorkbenchProjectionCommits", - "hwlabKafkaProjectorStatus", - "assertWorkbenchTransactionalRealtimeReady" - ]; +function requireKafkaProjectorStore(runtimeStore, capabilities = {}) { + const required = []; + if (capabilities.transactionalProjector) required.push("commitAgentRunKafkaProjection", "recordFailedAgentRunKafkaMessage", "recordIgnoredAgentRunKafkaMessage"); + if (capabilities.projectionOutboxRelay) required.push("claimHwlabKafkaOutbox", "completeHwlabKafkaOutbox", "retryHwlabKafkaOutbox"); + if (capabilities.projectionRealtime) required.push("subscribeWorkbenchProjectionCommits"); + if (capabilities.transactionalProjector || capabilities.projectionOutboxRelay) required.push("hwlabKafkaProjectorStatus"); + if (capabilities.transactionalProjector || capabilities.projectionOutboxRelay || capabilities.projectionRealtime) required.push("assertWorkbenchTransactionalRealtimeReady"); const missing = required.filter((name) => typeof runtimeStore?.[name] !== "function"); if (missing.length > 0) throw contractError("hwlab_kafka_projector_store_invalid", `Kafka durable capabilities require a runtime store: ${missing.join(", ")}`); } diff --git a/internal/cloud/server-code-agent-admission-http.ts b/internal/cloud/server-code-agent-admission-http.ts index 029837b9..725433b3 100644 --- a/internal/cloud/server-code-agent-admission-http.ts +++ b/internal/cloud/server-code-agent-admission-http.ts @@ -20,6 +20,7 @@ import { codeAgentSessionLifecycleSummary } from "./code-agent-session-lifecycle import { messageAuthorityTextValue } from "./code-agent-agentrun-prompt.ts"; import { promoteWorkbenchTurnAdmission as persistWorkbenchTurnAdmissionPromotion, writeWorkbenchSessionAdmissionFact } from "./workbench-projection-writer.ts"; import { createWorkbenchReadModel } from "./workbench-read-model.ts"; +import { workbenchRealtimeCapabilities } from "./workbench-realtime-capabilities.ts"; import { codeAgentOtelTraceFields, emitCodeAgentOtelSpan } from "./otel-trace.ts"; import { firstHeaderValue, @@ -1042,6 +1043,8 @@ async function submitCodeAgentChatTurn({ params, options, traceId }) { } }); if (adapterEnabled) { + const realtimeCapabilities = workbenchRealtimeCapabilities(runtimeEnv); + const transactionalAdmission = realtimeCapabilities.transactionalProjector; const pendingPromotion = results.get(traceId); if (pendingPromotion?.admissionState === "promotion-pending" && pendingPromotion?.agentRun?.durableDispatch === true) { const retryParams = { ...params, userBillingReservation: pendingPromotion.userBillingReservation ?? null }; @@ -1056,13 +1059,15 @@ async function submitCodeAgentChatTurn({ params, options, traceId }) { } const lifecycle = codeAgentTurnLifecycleFields(traceId, params); - await persistCodeAgentAdmissionInputState({ - params, - options, - traceId, - status: "admitting", - messageId: lifecycle.userMessageId - }); + if (transactionalAdmission) { + await persistCodeAgentAdmissionInputState({ + params, + options, + traceId, + status: "admitting", + messageId: lifecycle.userMessageId + }); + } traceStore.append(traceId, { type: "request", @@ -1083,14 +1088,16 @@ async function submitCodeAgentChatTurn({ params, options, traceId }) { error: billingPreflight?.blocked ? billingPreflight?.payload?.error?.message ?? billingPreflight?.blocker?.summary ?? "billing preflight blocked" : null }); if (billingPreflight?.blocked) { - await persistCodeAgentAdmissionInputState({ - params, - options, - traceId, - status: "blocked", - messageId: lifecycle.userMessageId, - error: billingPreflight.payload?.error ?? billingPreflight.blocker ?? { code: "billing_preflight_blocked" } - }); + if (transactionalAdmission) { + await persistCodeAgentAdmissionInputState({ + params, + options, + traceId, + status: "blocked", + messageId: lifecycle.userMessageId, + error: billingPreflight.payload?.error ?? billingPreflight.blocker ?? { code: "billing_preflight_blocked" } + }); + } traceStore.append(traceId, { type: "billing", status: "blocked", @@ -1116,14 +1123,16 @@ async function submitCodeAgentChatTurn({ params, options, traceId }) { const failure = codeAgentDispatchFailure(error); const rejected = codeAgentSubmissionUnavailablePayload(failure, { params: dispatchParams, options: executionOptions, traceId }); await finalizeCodeAgentBillingUsage({ payload: rejected, params: dispatchParams, options: executionOptions }); - await persistCodeAgentAdmissionInputState({ - params: dispatchParams, - options: executionOptions, - traceId, - status: "failed", - messageId: lifecycle.userMessageId, - error: failure.payload?.error ?? error - }); + if (transactionalAdmission) { + await persistCodeAgentAdmissionInputState({ + params: dispatchParams, + options: executionOptions, + traceId, + status: "failed", + messageId: lifecycle.userMessageId, + error: failure.payload?.error ?? error + }); + } void emitCodeAgentOtelSpan("agentrun_dispatch", traceId, executionOptions.env ?? process.env, { startTimeMs: dispatchStartedAt, status: "error", @@ -1144,15 +1153,17 @@ async function submitCodeAgentChatTurn({ params, options, traceId }) { }, dispatchParams); results.set(traceId, pending); try { - await persistCodeAgentAdmissionInputState({ - params: dispatchParams, - options: executionOptions, - traceId, - status: "admitting", - messageId: lifecycle.userMessageId, - commandId: payload.agentRun?.commandId ?? null, - agentRun: payload.agentRun ?? null - }); + if (transactionalAdmission) { + await persistCodeAgentAdmissionInputState({ + params: dispatchParams, + options: executionOptions, + traceId, + status: "admitting", + messageId: lifecycle.userMessageId, + commandId: payload.agentRun?.commandId ?? null, + agentRun: payload.agentRun ?? null + }); + } await promoteCodeAgentTurnAdmission({ payload: pending, params: dispatchParams, options: executionOptions, traceId }); } catch (error) { throw codeAgentPromotionUnavailableError(error, { payload: pending, params: dispatchParams, options: executionOptions, traceId }); diff --git a/internal/cloud/server-workbench-realtime-http.test.ts b/internal/cloud/server-workbench-realtime-http.test.ts index 1de352a2..835cceef 100644 --- a/internal/cloud/server-workbench-realtime-http.test.ts +++ b/internal/cloud/server-workbench-realtime-http.test.ts @@ -111,47 +111,36 @@ test("workbench realtime initial connection emits current snapshot without repla } }); -test("product realtime uses the projection outbox authority", async () => { +test("product realtime uses live Kafka SSE without projection reads", async () => { const sessionId = "ses_composable_realtime"; const traceId = "trc_composable_realtime"; + const subscribers = new Set(); let projectionReads = 0; - const runtime = { - async readAtomicWorkbenchProjectionSync() { - projectionReads += 1; - return { - facts: { - sessions: [{ sessionId, lastTraceId: traceId }], - messages: [{ messageId: "msg_composable_realtime", sessionId, traceId, role: "agent", text: "projection snapshot", projectedSeq: 1 }], - turns: [] - }, - events: [], - cutoffOutboxSeq: 1, - cursorOutboxSeq: 1, - hasMore: false - }; - } - }; + const runtime = new Proxy({}, { get() { projectionReads += 1; throw new Error("live Kafka realtime must not read projection facts"); } }); + const capabilities = { directPublish: true, liveKafkaSse: true, kafkaRefreshReplay: false, transactionalProjector: false, projectionOutboxRelay: false, projectionRealtime: false }; const server = createCloudApiServer({ - accessController: realtimeAccessController(), + accessController: realtimeAccessController({ sessions: [{ id: sessionId, ownerUserId: ACTOR.id, lastTraceId: traceId }] }), workbenchRuntime: runtime, - kafkaEventBridge: projectionRealtimeBridge(), + kafkaEventBridge: { capabilities, ready: Promise.resolve(), subscribeLiveHwlabEvents(listener) { subscribers.add(listener); return () => subscribers.delete(listener); }, async stop() {} }, env: { HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true", HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "false", HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "true", - HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000", - HWLAB_WORKBENCH_OUTBOX_TAIL_BATCH_SIZE: "100" + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" } }); 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}`, 2); - assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.message.snapshot"]); - assert.equal(events[0].data.realtimeSource, "projection-outbox"); - assert.equal(projectionReads, 1); + const eventsPromise = getSseEvents(server.address().port, `/v1/workbench/events?sessionId=${sessionId}`, 2); + await waitForCondition(() => subscribers.size === 1); + for (const listener of subscribers) listener({ schema: "hwlab.event.v1", eventId: "evt_live", traceId, hwlabSessionId: sessionId, sessionId, event: { type: "assistant", traceId, sessionId, text: "live" } }, { topic: "hwlab.event.v1", partition: 0, offset: "1" }); + const events = await eventsPromise; + assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "hwlab.event.v1"]); + assert.deepEqual(events[0].data.capabilities, capabilities); + assert.equal(events[1].id, null); + assert.equal(projectionReads, 0); } finally { await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); } diff --git a/internal/cloud/server-workbench-realtime-http.ts b/internal/cloud/server-workbench-realtime-http.ts index 4f6a7115..0c5bc0c1 100644 --- a/internal/cloud/server-workbench-realtime-http.ts +++ b/internal/cloud/server-workbench-realtime-http.ts @@ -14,10 +14,15 @@ import { } from "./server-http-utils.ts"; import { createWorkbenchReadModel } from "./workbench-read-model.ts"; import { createWorkbenchRuntimeClient } from "./workbench-runtime-client.ts"; +import { emitLiveKafkaOtelSpan } from "./kafka-event-bridge.ts"; +import { createWorkbenchKafkaRefreshHandoff, workbenchKafkaRefreshErrorPayload } from "./workbench-kafka-refresh-handoff.ts"; import { buildWorkbenchSessionDetail, compactLaunchContext, includeMessagesForSessionDetail } from "./workbench-session-detail-response.ts"; import { projectionOutboxRealtimeEvents } from "./workbench-projection-outbox-events.ts"; import { durableTraceStatus, RUNNING_STATUSES, TERMINAL_STATUSES } from "./workbench-turn-projection.ts"; import { emitCodeAgentOtelSpan, emitHttpServerRequestSpan } from "./otel-trace.ts"; +import { + workbenchRealtimeCapabilities +} from "./workbench-realtime-capabilities.ts"; import * as workbenchFacts from "./server-workbench-facts.ts"; const DEFAULT_PAGE_LIMIT = 50; @@ -177,16 +182,21 @@ export async function drainWorkbenchRealtimeConnections(options = {}) { export async function handleWorkbenchRealtimeHttp(request, response, url, options = {}) { try { if (request.method !== "GET") return methodNotAllowed(response, "GET"); + const realtimeCapabilities = workbenchRealtimeCapabilities(options.env ?? process.env); const projectionOnly = url.pathname === "/v1/workbench/projection-events"; - if (projectionOnly) { - sendJson(response, 410, workbenchError("workbench_projection_events_route_removed", "Use /v1/workbench/events for the single transactional projection authority.")); + if (!projectionOnly && realtimeCapabilities.liveKafkaSse) { + await handleLiveKafkaWorkbenchRealtimeHttp(request, response, url, options); + return; + } + if (!realtimeCapabilities.projectionRealtime) { + sendJson(response, 503, workbenchError("workbench_realtime_disabled", "No Workbench realtime SSE capability is enabled.")); return; } const requestedSessionId = safeSessionId(url.searchParams.get("sessionId") ?? url.searchParams.get("includeSessionId")); const requestedTraceId = safeTraceId(url.searchParams.get("traceId")); const heartbeatMs = parsePositiveInteger(options.env?.HWLAB_WORKBENCH_SSE_HEARTBEAT_MS, DEFAULT_WORKBENCH_SSE_HEARTBEAT_MS); attachWorkbenchRealtimeOtelContext(request, { - route: "/v1/workbench/events", + route: projectionOnly ? "/v1/workbench/projection-events" : "/v1/workbench/events", sessionId: requestedSessionId, traceId: requestedTraceId, heartbeatMs, @@ -442,6 +452,405 @@ export async function handleWorkbenchRealtimeHttp(request, response, url, option } } +async function handleLiveKafkaWorkbenchRealtimeHttp(request, response, url, options = {}) { + const requestedSessionId = safeSessionId(url.searchParams.get("sessionId") ?? url.searchParams.get("includeSessionId")); + const requestedTraceId = safeTraceId(url.searchParams.get("traceId")); + const realtimeCapabilities = workbenchRealtimeCapabilities(options.env ?? process.env); + const heartbeatMs = parsePositiveInteger(options.env?.HWLAB_WORKBENCH_SSE_HEARTBEAT_MS, DEFAULT_WORKBENCH_SSE_HEARTBEAT_MS); + attachWorkbenchRealtimeOtelContext(request, { + sessionId: requestedSessionId, + traceId: requestedTraceId, + heartbeatMs, + realtimeSource: "live-kafka-sse", + realtimeCapabilities + }); + const perf = options.backendPerformance; + const auth = perf ? await perf.measure("workbench_auth", () => authenticateWorkbenchRead(request, response, options)) : await authenticateWorkbenchRead(request, response, options); + if (!auth) return; + if (url.searchParams.has("projectId") || url.searchParams.has("workspaceId")) { + sendJson(response, 400, workbenchError("workbench_authority_removed", "Workbench realtime is keyed by sessionId/traceId only.")); + return; + } + if (!requestedSessionId && !requestedTraceId) { + sendJson(response, 400, workbenchError("workbench_realtime_scope_required", "Workbench realtime requires sessionId or traceId.")); + return; + } + const authorization = await authorizeLiveKafkaWorkbenchRealtimeScope(options, auth.actor, { + sessionId: requestedSessionId, + traceId: requestedTraceId + }); + if (!authorization.ok) { + sendJson(response, authorization.status, authorization.body); + return; + } + const authorizedSessionId = authorization.sessionId; + const bridge = options.kafkaEventBridge; + if (bridge?.capabilities?.liveKafkaSse !== true || typeof bridge?.subscribeLiveHwlabEvents !== "function") { + sendJson(response, 503, workbenchError("workbench_live_kafka_unconfigured", "Workbench live Kafka fanout is not configured.")); + return; + } + if (realtimeCapabilities.kafkaRefreshReplay) { + await handleKafkaRefreshReplayWorkbenchRealtimeHttp(request, response, options, { + bridge, + heartbeatMs, + realtimeCapabilities, + requestedSessionId: authorizedSessionId, + requestedTraceId + }); + return; + } + + let closed = false; + let realtimeCloseReason = "client_close"; + let realtimeCloseSignal = null; + const realtimeStartedAtMs = Date.now(); + const cleanup = []; + const isActive = () => !closed && !response.destroyed && !response.writableEnded; + const closeConnection = () => { + if (closed) return; + closed = true; + emitWorkbenchRealtimeClosedOtelSpan(request, options.env, { + reason: realtimeCloseReason, + signal: realtimeCloseSignal, + startedAtMs: realtimeStartedAtMs, + sessionId: authorizedSessionId, + threadId: null, + traceId: requestedTraceId, + activeConnectionCount: activeWorkbenchRealtimeConnections.size + }); + for (const item of cleanup.splice(0)) item(); + }; + + response.once("close", closeConnection); + request.once?.("aborted", closeConnection); + request.socket?.once?.("close", closeConnection); + cleanup.push(() => request.off?.("aborted", closeConnection)); + cleanup.push(() => request.socket?.off?.("close", closeConnection)); + + const bufferedEnvelopes = []; + let liveDeliveryStarted = false; + let enqueueEvent = null; + const unsubscribe = bridge.subscribeLiveHwlabEvents((envelope, transport) => { + if (!liveKafkaEnvelopeMatches(envelope, authorizedSessionId, requestedTraceId)) return; + if (!liveDeliveryStarted || typeof enqueueEvent !== "function") { + bufferedEnvelopes.push({ envelope, transport }); + return; + } + void enqueueEvent("hwlab.event.v1", envelope, transport); + }); + if (typeof unsubscribe === "function") cleanup.push(unsubscribe); + + try { + await (bridge.liveReady ?? bridge.ready); + } catch (error) { + closeConnection(); + throw error; + } + if (!isActive()) { + closeConnection(); + return; + } + + response.writeHead(200, { + "content-type": "text/event-stream; charset=utf-8", + "cache-control": "no-store, no-transform", + connection: "keep-alive", + "x-accel-buffering": "no", + "x-content-type-options": "nosniff" + }); + if (typeof response.flushHeaders === "function") response.flushHeaders(); + + const writeEvent = async (name, payload, transport = null) => { + if (!isActive()) return false; + const block = `event: ${name}\ndata: ${JSON.stringify(payload)}\n\n`; + const writable = response.write(block); + if (!writable) await waitForResponseDrain(response, isActive); + const active = isActive(); + if (active && name === "hwlab.event.v1") { + emitLiveKafkaOtelSpan("hwlab.workbench.live_sse.business_event_write", payload, transport, { + env: options.env, + otelSpanEmitter: options.otelSpanEmitter ?? emitCodeAgentOtelSpan + }); + } + return active; + }; + let writeChain = Promise.resolve(true); + enqueueEvent = (name, payload, transport = null) => { + writeChain = writeChain.then(() => writeEvent(name, payload, transport)); + return writeChain; + }; + + const realtimeConnection = { + close(fields = {}) { + if (!isActive()) return false; + realtimeCloseReason = textValue(fields.reason) || "server_shutdown"; + realtimeCloseSignal = textValue(fields.signal) || null; + void enqueueEvent("workbench.server_draining", { + type: "server.draining", + status: "closing", + reason: realtimeCloseReason, + signal: realtimeCloseSignal, + capabilities: realtimeCapabilities, + lossPossible: true + }).finally(() => { + if (!response.writableEnded) response.end(); + }); + return true; + }, + isActive + }; + activeWorkbenchRealtimeConnections.add(realtimeConnection); + cleanup.push(() => activeWorkbenchRealtimeConnections.delete(realtimeConnection)); + + await enqueueEvent("workbench.connected", { + type: "connected", + status: "connected", + capabilities: realtimeCapabilities, + realtimeSource: "hwlab.event.v1", + deliverySemantics: "live-only", + liveOnly: true, + replay: false, + replaySupported: false, + lossPossible: true, + filters: { sessionId: authorizedSessionId, traceId: requestedTraceId } + }); + liveDeliveryStarted = true; + const readyBufferedEnvelopes = bufferedEnvelopes.splice(0); + for (const { envelope, transport } of readyBufferedEnvelopes) void enqueueEvent("hwlab.event.v1", envelope, transport); + await writeChain; + if (!isActive()) { + closeConnection(); + return; + } + const heartbeatTimer = setInterval(() => { + void enqueueEvent("workbench.heartbeat", { + type: "heartbeat", + status: "connected", + realtimeSource: "hwlab.event.v1", + liveOnly: true, + replay: false, + lossPossible: true, + serverSentAt: new Date().toISOString(), + valuesPrinted: false + }); + }, heartbeatMs); + heartbeatTimer.unref?.(); + cleanup.push(() => clearInterval(heartbeatTimer)); + emitWorkbenchRealtimeAcceptedOtelSpan(request, options.env); +} + +async function handleKafkaRefreshReplayWorkbenchRealtimeHttp(request, response, options, input) { + const { + bridge, + heartbeatMs, + realtimeCapabilities, + requestedSessionId, + requestedTraceId + } = input; + const refreshReplay = bridge?.refreshReplay; + if (!refreshReplay || typeof bridge?.queryHwlabEventRetention !== "function") { + sendJson(response, 503, workbenchError("workbench_kafka_refresh_unconfigured", "Workbench Kafka refresh replay is enabled without its retention query runtime.")); + return; + } + let closed = false; + let realtimeCloseReason = "client_close"; + let realtimeCloseSignal = null; + let failureClosing = false; + const realtimeStartedAtMs = Date.now(); + const cleanup = []; + const isActive = () => !closed && !response.destroyed && !response.writableEnded; + const closeConnection = () => { + if (closed) return; + closed = true; + emitWorkbenchRealtimeClosedOtelSpan(request, options.env, { + reason: realtimeCloseReason, + signal: realtimeCloseSignal, + startedAtMs: realtimeStartedAtMs, + sessionId: requestedSessionId, + threadId: null, + traceId: requestedTraceId, + activeConnectionCount: activeWorkbenchRealtimeConnections.size + }); + for (const item of cleanup.splice(0)) item(); + }; + + response.once("close", closeConnection); + request.once?.("aborted", closeConnection); + request.socket?.once?.("close", closeConnection); + cleanup.push(() => request.off?.("aborted", closeConnection)); + cleanup.push(() => request.socket?.off?.("close", closeConnection)); + + try { + await (bridge.liveReady ?? bridge.ready); + } catch (error) { + for (const item of cleanup.splice(0)) item(); + throw error; + } + if (!isActive()) { + closeConnection(); + return; + } + + response.writeHead(200, { + "content-type": "text/event-stream; charset=utf-8", + "cache-control": "no-store, no-transform", + connection: "keep-alive", + "x-accel-buffering": "no", + "x-content-type-options": "nosniff" + }); + if (typeof response.flushHeaders === "function") response.flushHeaders(); + + const writeEvent = async (name, payload, transport = null) => { + if (!isActive()) return false; + const block = `event: ${name}\ndata: ${JSON.stringify(payload)}\n\n`; + const writable = response.write(block); + if (!writable) await waitForResponseDrain(response, isActive); + const active = isActive(); + if (active && name === "hwlab.event.v1") { + emitLiveKafkaOtelSpan("hwlab.workbench.live_sse.business_event_write", payload, transport, { + env: options.env, + otelSpanEmitter: options.otelSpanEmitter ?? emitCodeAgentOtelSpan + }); + } + return active; + }; + let writeChain = Promise.resolve(true); + const enqueueEvent = (name, payload, transport = null) => { + writeChain = writeChain.then(() => writeEvent(name, payload, transport)); + return writeChain; + }; + const closeWithRefreshFailure = async (error) => { + if (failureClosing || !isActive()) return; + failureClosing = true; + await enqueueEvent("workbench.error", workbenchKafkaRefreshErrorPayload(error, { + sessionId: requestedSessionId, + traceId: requestedTraceId + })); + await writeChain; + if (!response.writableEnded) response.end(); + }; + + const realtimeConnection = { + close(fields = {}) { + if (!isActive()) return false; + realtimeCloseReason = textValue(fields.reason) || "server_shutdown"; + realtimeCloseSignal = textValue(fields.signal) || null; + void enqueueEvent("workbench.server_draining", { + type: "server.draining", + status: "closing", + reason: realtimeCloseReason, + signal: realtimeCloseSignal, + capabilities: realtimeCapabilities, + lossPossible: false + }).finally(() => { + if (!response.writableEnded) response.end(); + }); + return true; + }, + isActive + }; + activeWorkbenchRealtimeConnections.add(realtimeConnection); + cleanup.push(() => activeWorkbenchRealtimeConnections.delete(realtimeConnection)); + + const handoff = createWorkbenchKafkaRefreshHandoff({ + sessionId: requestedSessionId, + traceId: requestedTraceId, + liveBufferLimit: refreshReplay.liveBufferLimit, + subscribeLive: (listener) => bridge.subscribeLiveHwlabEvents(listener), + queryRetention: ({ signal }) => bridge.queryHwlabEventRetention({ + sessionId: requestedSessionId, + traceId: requestedTraceId, + limit: refreshReplay.matchedEventLimit, + scanLimit: refreshReplay.scanLimit, + timeoutMs: refreshReplay.timeoutMs, + groupIdPrefix: refreshReplay.groupIdPrefix, + fromBeginning: true, + signal + }), + deliverEvent: (envelope, transport) => enqueueEvent("hwlab.event.v1", envelope, transport), + deliverConnected: (summary) => enqueueEvent("workbench.connected", { + type: "connected", + status: "connected", + capabilities: realtimeCapabilities, + realtimeSource: "hwlab.event.v1", + deliverySemantics: "kafka-retention-then-live", + liveOnly: false, + replay: true, + replaySupported: true, + lossPossible: false, + filters: { sessionId: requestedSessionId, traceId: requestedTraceId }, + refreshReplay: summary + }), + onFailure: closeWithRefreshFailure + }); + cleanup.push(() => handoff.stop("connection-closed")); + + try { + await handoff.start(); + } catch (error) { + if (isActive()) await closeWithRefreshFailure(error); + return; + } + if (!isActive()) return; + + const heartbeatTimer = setInterval(() => { + void enqueueEvent("workbench.heartbeat", { + type: "heartbeat", + status: "connected", + realtimeSource: "hwlab.event.v1", + deliverySemantics: "kafka-retention-then-live", + liveOnly: false, + replay: true, + lossPossible: false, + serverSentAt: new Date().toISOString(), + valuesPrinted: false + }); + }, heartbeatMs); + heartbeatTimer.unref?.(); + cleanup.push(() => clearInterval(heartbeatTimer)); + emitWorkbenchRealtimeAcceptedOtelSpan(request, options.env); +} + +async function authorizeLiveKafkaWorkbenchRealtimeScope(options, actor, { sessionId, traceId }) { + const access = options.accessController; + if ((sessionId && typeof access?.getAgentSession !== "function") || (traceId && typeof access?.getAgentSessionByTraceId !== "function")) { + return { + ok: false, + status: 503, + body: workbenchError("workbench_realtime_authorization_unconfigured", "Workbench live realtime ownership authorization is not configured.") + }; + } + + const [sessionById, sessionByTrace] = await Promise.all([ + sessionId ? access.getAgentSession(sessionId) : null, + traceId ? access.getAgentSessionByTraceId(traceId) : null + ]); + const sessionByIdId = sessionById ? safeSessionId(sessionById.id ?? sessionById.sessionId) : null; + const sessionByTraceId = sessionByTrace ? safeSessionId(sessionByTrace.id ?? sessionByTrace.sessionId) : null; + const scopeMissing = (sessionId && sessionByIdId !== sessionId) + || (traceId && !sessionByTraceId) + || (sessionId && traceId && sessionByIdId !== sessionByTraceId); + const resolvedSessions = [sessionById, sessionByTrace].filter(Boolean); + const ownedByActor = actor?.role === "admin" || (Boolean(actor?.id) + && resolvedSessions.length > 0 + && resolvedSessions.every((session) => textValue(session.ownerUserId) === actor.id)); + if (scopeMissing || !ownedByActor) { + return { + ok: false, + status: 404, + body: workbenchError("workbench_realtime_scope_not_found", "Workbench realtime scope is not visible to the current actor.") + }; + } + return { ok: true, sessionId: sessionByIdId ?? sessionByTraceId }; +} + +function liveKafkaEnvelopeMatches(envelope, sessionId, traceId) { + if (!envelope || envelope.schema !== "hwlab.event.v1") return false; + if (sessionId && safeSessionId(envelope.hwlabSessionId ?? envelope.sessionId) !== sessionId) return false; + if (traceId && safeTraceId(envelope.traceId) !== traceId) return false; + return true; +} + export function attachWorkbenchRealtimeOtelContext(request, fields = {}) { const context = request?.hwlabHttpRequestContext; if (!context) return; @@ -455,9 +864,9 @@ export function attachWorkbenchRealtimeOtelContext(request, fields = {}) { "workbench.thread_id": fields.threadId ?? null, "workbench.sse.heartbeat_ms": fields.heartbeatMs ?? null, "workbench.sse.realtime_source": fields.realtimeSource ?? null, - "workbench.sse.live_kafka_enabled": false, - "workbench.sse.projection_realtime_enabled": true, - "workbench.sse.outbox_mode": true + "workbench.sse.live_kafka_enabled": fields.realtimeCapabilities?.liveKafkaSse ?? null, + "workbench.sse.projection_realtime_enabled": fields.realtimeCapabilities?.projectionRealtime ?? null, + "workbench.sse.outbox_mode": fields.realtimeCapabilities?.liveKafkaSse === true ? false : true }; } diff --git a/internal/cloud/server.ts b/internal/cloud/server.ts index 945259b4..148283d8 100644 --- a/internal/cloud/server.ts +++ b/internal/cloud/server.ts @@ -372,7 +372,10 @@ export async function buildHealthPayload(options = {}) { async function kafkaProjectorHealth(projector, { liveProbe = false } = {}) { if (!projector?.started) return { started: false, reason: projector?.reason ?? "disabled", valuesRedacted: true }; const identity = { + capabilities: projector.capabilities ?? null, + directPublishGroupId: projector.directPublishGroupId ?? null, projectorGroupId: projector.projectorGroupId ?? null, + hwlabEventGroupId: projector.hwlabEventGroupId ?? null, agentRunTopic: projector.agentRunTopic, hwlabTopic: projector.hwlabTopic }; diff --git a/internal/cloud/workbench-realtime-capabilities.ts b/internal/cloud/workbench-realtime-capabilities.ts new file mode 100644 index 00000000..8d87d200 --- /dev/null +++ b/internal/cloud/workbench-realtime-capabilities.ts @@ -0,0 +1,35 @@ +// SPEC: PJ2026-0104010803 Workbench composable realtime capabilities. +// Responsibility: validate YAML-owned capability switches without code defaults or implied mode selection. + +export const WORKBENCH_REALTIME_CAPABILITY_ENVS = Object.freeze({ + directPublish: "HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED", + liveKafkaSse: "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED", + kafkaRefreshReplay: "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED", + transactionalProjector: "HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED", + projectionOutboxRelay: "HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED", + projectionRealtime: "HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED" +}); + +export function workbenchRealtimeCapabilities(env = process.env) { + return Object.fromEntries(Object.entries(WORKBENCH_REALTIME_CAPABILITY_ENVS).map(([key, name]) => [key, requiredBooleanEnv(env, name)])); +} + +export function requiredBooleanEnv(env, name) { + const value = String(env?.[name] ?? "").trim().toLowerCase(); + if (value === "true" || value === "1") return true; + if (value === "false" || value === "0") return false; + throw capabilityError( + "hwlab_workbench_realtime_capability_invalid", + `${name} is required and must be explicitly true or false in the owning YAML.`, + name + ); +} + +function capabilityError(code, message, envName) { + return Object.assign(new Error(message), { + code, + envName, + statusCode: 503, + valuesRedacted: true + }); +} diff --git a/internal/dev-entrypoint/cloud-web-runtime.mjs b/internal/dev-entrypoint/cloud-web-runtime.mjs index 096f6cda..25f1e254 100644 --- a/internal/dev-entrypoint/cloud-web-runtime.mjs +++ b/internal/dev-entrypoint/cloud-web-runtime.mjs @@ -400,6 +400,11 @@ function validateDisplayLocale(locale) { function workbenchRuntimeConfigFromEnv() { const traceTimeline = {}; const result = { + realtimeFeatures: { + liveKafkaSse: requiredWorkbenchRealtimeFeature(process.env.HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED, "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED"), + kafkaRefreshReplay: requiredWorkbenchRealtimeFeature(process.env.HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED"), + projectionRealtime: requiredWorkbenchRealtimeFeature(process.env.HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED, "HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED") + }, debugCapabilities: { isolatedKafka: requiredWorkbenchRealtimeFeature(process.env.HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED, "HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED"), rawHwlabEventWindow: { diff --git a/scripts/gitops-render.test.ts b/scripts/gitops-render.test.ts index eb9d7663..531707e8 100644 --- a/scripts/gitops-render.test.ts +++ b/scripts/gitops-render.test.ts @@ -687,14 +687,11 @@ test("v03 render keeps node identity as data instead of generated structure", as assert.equal(cloudApiEnv.get("HWLAB_CLOUD_DB_SSL_MODE"), "require"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_ENABLED"), "true"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_AGENTRUN_EVENT_CONSUME_ENABLED"), "true"); - for (const name of [ - "HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED", - "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED", - "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED", - "HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED", - "HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED", - "HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED" - ]) assert.equal(cloudApiEnv.has(name), false); + assert.equal(cloudApiEnv.get("HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED"), "true"); + assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED"), "true"); + assert.equal(cloudApiEnv.get("HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED"), "false"); + assert.equal(cloudApiEnv.get("HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED"), "false"); + assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED"), "false"); assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED"), "true"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_BOOTSTRAP_SERVERS"), "platform-infra-kafka-kafka-bootstrap.platform-infra.svc.cluster.local:9092"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_STDIO_TOPIC"), "codex-stdio.raw.v1"); @@ -712,14 +709,9 @@ test("v03 render keeps node identity as data instead of generated structure", as const cloudWeb = (workloadsJson.items ?? []).find((item) => item.kind === "Deployment" && item.metadata?.name === "hwlab-cloud-web"); const cloudWebContainer = collectContainersFromItem(cloudWeb).find((container) => container.name === "hwlab-cloud-web"); const cloudWebEnv = new Map((cloudWebContainer?.env ?? []).map((entry) => [entry.name, entry.value])); - for (const name of [ - "HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED", - "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED", - "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED", - "HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED", - "HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED", - "HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED" - ]) assert.equal(cloudWebEnv.has(name), false); + assert.equal(cloudWebEnv.get("HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED"), "true"); + assert.equal(cloudWebEnv.get("HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED"), "true"); + assert.equal(cloudWebEnv.get("HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED"), "false"); assert.equal(cloudWebEnv.get("HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED"), "true"); assert.equal(cloudWebEnv.get("HWLAB_WORKBENCH_RAW_HWLAB_EVENT_WINDOW_ENABLED"), "1"); assert.equal(cloudWebEnv.get("HWLAB_WORKBENCH_RAW_HWLAB_EVENT_WINDOW_MAX_ENTRIES"), "200"); diff --git a/scripts/runtime-realtime-capabilities-config.test.mjs b/scripts/runtime-realtime-capabilities-config.test.mjs index a037d55d..b24f08b1 100644 --- a/scripts/runtime-realtime-capabilities-config.test.mjs +++ b/scripts/runtime-realtime-capabilities-config.test.mjs @@ -3,14 +3,19 @@ import test from "node:test"; import { readStructuredFile } from "./src/structured-config.mjs"; -const OBSOLETE_ARCHITECTURE_ENV = Object.freeze([ - "HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED", - "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED", - "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED", - "HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED", - "HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED", - "HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED" -]); +const CLOUD_API_ARCHITECTURE_ENV = Object.freeze({ + HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "true", + HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", + HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: "false", + HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: "false", + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" +}); + +const CLOUD_WEB_ARCHITECTURE_ENV = Object.freeze({ + HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", + HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: "true", + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" +}); const LIVE_KAFKA_CONTRACT = Object.freeze({ HWLAB_KAFKA_STDIO_TOPIC: "codex-stdio.raw.v1", @@ -37,17 +42,15 @@ const ISOLATED_DEBUG_CONTRACT = Object.freeze({ HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "15000" }); -test("v03 YAML fixes one transactional projection architecture without architecture env", async () => { +test("v03 YAML fixes direct publish, live Kafka SSE and retention replay authority", async () => { const deploy = await readStructuredFile(process.cwd(), "deploy/deploy.yaml"); const cloudApi = deploy.lanes?.v03?.services?.find((service) => service.serviceId === "hwlab-cloud-api"); const cloudWeb = deploy.lanes?.v03?.serviceDeclarations?.["hwlab-cloud-web"]; assert.ok(cloudApi, "deploy.lanes.v03.services must declare hwlab-cloud-api"); assert.ok(cloudWeb, "deploy.lanes.v03.serviceDeclarations must declare hwlab-cloud-web"); - for (const name of OBSOLETE_ARCHITECTURE_ENV) { - assert.equal(cloudApi.env?.[name], undefined, `${name} must not be rendered into cloud-api`); - assert.equal(cloudWeb.env?.[name], undefined, `${name} must not be rendered into cloud-web`); - } + for (const [name, expected] of Object.entries(CLOUD_API_ARCHITECTURE_ENV)) assert.equal(cloudApi.env?.[name], expected, `${name} must fix the cloud-api Kafka authority`); + for (const [name, expected] of Object.entries(CLOUD_WEB_ARCHITECTURE_ENV)) assert.equal(cloudWeb.env?.[name], expected, `${name} must fix the cloud-web Kafka authority`); for (const [name, expected] of Object.entries(LIVE_KAFKA_CONTRACT)) { assert.equal(cloudApi.env?.[name], expected, `${name} must preserve the pure Kafka live contract`); } diff --git a/web/hwlab-cloud-web/scripts/workbench-fixed-transactional-architecture.test.ts b/web/hwlab-cloud-web/scripts/workbench-fixed-transactional-architecture.test.ts index a46acec0..cdfaec2b 100644 --- a/web/hwlab-cloud-web/scripts/workbench-fixed-transactional-architecture.test.ts +++ b/web/hwlab-cloud-web/scripts/workbench-fixed-transactional-architecture.test.ts @@ -4,90 +4,72 @@ import test from "node:test"; import { fileURLToPath } from "node:url"; const repositoryRoot = fileURLToPath(new URL("../../..", import.meta.url)); -const productionFiles = [ - "internal/cloud/kafka-event-bridge.ts", - "internal/cloud/server-workbench-realtime-http.ts", - "internal/cloud/server-workbench-facts.ts", - "internal/cloud/server-code-agent-admission-http.ts", - "internal/cloud/server.ts", - "web/hwlab-cloud-web/src/config/runtime.ts", - "web/hwlab-cloud-web/src/config/workbench-runtime-policy.ts", - "web/hwlab-cloud-web/src/types/global.d.ts", - "web/hwlab-cloud-web/src/api/workbench-events.ts", - "web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts", - "web/hwlab-cloud-web/src/utils/workbench-realtime-runtime.ts", - "web/hwlab-cloud-web/src/stores/workbench.ts", - "web/hwlab-cloud-web/src/stores/workbench-colada-queries.ts", - "web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts", - "web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts", - "web/hwlab-cloud-web/src/stores/workbench-session.ts", - "web/hwlab-cloud-web/src/stores/workbench-session-messages-read-budget.ts" -]; -const forbiddenArchitectureIdentifiers = [ - "directPublish", - "liveKafkaSse", - "kafkaRefreshReplay", - "transactionalProjector", - "projectionOutboxRelay", - "projectionRealtime", - "startLiveHwlabKafkaEventBridge", - "handleLiveKafkaWorkbenchRealtimeHttp", - "handleKafkaRefreshReplayWorkbenchRealtimeHttp", - "subscribeLiveHwlabEvents", - "queryHwlabEventRetention", - "WorkbenchRealtimeCapabilities", - "workbenchRealtimeCapabilities", - "workbenchHistoryAuthorityPolicy", - "workbenchRealtimeTransportEnabled", - "workbenchRealtimeTraceIdForCapabilities", - "workbenchProjectionEventStreamPath", +const requiredProductionIdentifiers = [ + ["internal/cloud/kafka-event-bridge.ts", "startLiveHwlabKafkaEventBridge"], + ["internal/cloud/kafka-event-bridge.ts", "queryHwlabEventRetention"], + ["internal/cloud/server-workbench-realtime-http.ts", "handleLiveKafkaWorkbenchRealtimeHttp"], + ["internal/cloud/server-workbench-realtime-http.ts", "handleKafkaRefreshReplayWorkbenchRealtimeHttp"], + ["web/hwlab-cloud-web/src/stores/workbench.ts", "applyLiveKafkaBusinessEvent"], + ["web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts", "hwlab.event.v1"], + ["web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts", "workbenchRealtimeTraceIdForCapabilities"] +] as const; +const forbiddenRepairIdentifiers = [ + "handleWorkbenchSyncHttp", + "/v1/workbench/sync", "scheduleSessionListRefresh", "invalidateSessionList", "fetchWorkbenchTurnStatus", "refreshTurnStatusByTraceId", "refreshMessageProjectionForTrace", - "selectActiveTurnStatusRefreshTraceIds", - "sessionDetailAutoReadDecision", - "projectionMergeCommitSummary", "forceRead", "sealRestoredActiveTurnMessages", "messageNeedsRestoredTurnSeal", "refreshSessionStatusAuthority", - "readTerminalTraceDetailGaps", - "workbenchTraceMessagesWindowLimit", - "handleWorkbenchSyncHttp", - "/v1/workbench/sync" + "readTerminalTraceDetailGaps" ]; -const obsoleteArchitectureEnv = [ - "HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED", - "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED", - "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED", - "HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED", - "HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED", - "HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED" +const productionFiles = [ + "internal/cloud/server-workbench-facts.ts", + "internal/cloud/server-workbench-realtime-http.ts", + "web/hwlab-cloud-web/src/stores/workbench.ts", + "web/hwlab-cloud-web/src/stores/workbench-colada-queries.ts" ]; -test("Workbench production source has one fixed transactional projection architecture", async () => { +test("Workbench production source fixes Kafka as the realtime and replay authority", async () => { + for (const [relativePath, identifier] of requiredProductionIdentifiers) { + const source = await readFile(`${repositoryRoot}/${relativePath}`, "utf8"); + assert.ok(source.includes(identifier), `${relativePath} must contain ${identifier}`); + } +}); + +test("Workbench production source keeps HTTP repair and sync authority deleted", async () => { const violations: string[] = []; for (const relativePath of productionFiles) { const source = await readFile(`${repositoryRoot}/${relativePath}`, "utf8"); - for (const identifier of forbiddenArchitectureIdentifiers) { - if (source.includes(identifier)) violations.push(`${relativePath}: ${identifier}`); - } + for (const identifier of forbiddenRepairIdentifiers) if (source.includes(identifier)) violations.push(`${relativePath}: ${identifier}`); } assert.deepEqual(violations, []); }); -test("production deploy does not render architecture selection env", async () => { +test("production deploy selects direct publish, live SSE and retention replay without transactional startup authority", async () => { const deploySource = await readFile(`${repositoryRoot}/deploy/deploy.yaml`, "utf8"); - const violations = obsoleteArchitectureEnv.filter((name) => deploySource.includes(name)); - assert.deepEqual(violations, []); + for (const expected of [ + "HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: \"true\"", + "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: \"true\"", + "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED: \"true\"", + "HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED: \"false\"", + "HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED: \"false\"", + "HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: \"false\"" + ]) assert.ok(deploySource.includes(expected), `deploy must contain ${expected}`); }); -test("raw hwlab Kafka fixtures remain isolated to the explicit debug path", async () => { - const fixtureSource = await readFile(`${repositoryRoot}/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts`, "utf8"); - const debugSource = await readFile(`${repositoryRoot}/web/hwlab-cloud-web/src/stores/workbench-isolated-kafka-debug.ts`, "utf8"); - assert.match(fixtureSource, /historical hwlab\.event\.debug\.v1 fixtures/u); - assert.match(fixtureSource, /never imported by the product Workbench SSE path/u); - assert.match(debugSource, /workbench-live-kafka-event/u); +test("feature schema remains product-only and excludes architecture environment keys", async () => { + const schemaSource = await readFile(`${repositoryRoot}/config/feature-config.schema.json`, "utf8"); + for (const name of [ + "HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED", + "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED", + "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_ENABLED", + "HWLAB_KAFKA_TRANSACTIONAL_PROJECTOR_ENABLED", + "HWLAB_KAFKA_PROJECTION_OUTBOX_RELAY_ENABLED", + "HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED" + ]) assert.equal(schemaSource.includes(name), false, `${name} must not enter feature schema`); }); diff --git a/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts b/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts index 32dbd148..c8ac44ca 100644 --- a/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts +++ b/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts @@ -8,7 +8,7 @@ import type { ChatMessage } from "../src/types/index.ts"; import { workbenchSessionDetailPathForTest, workbenchSessionMessagesPathForTest } from "../src/api/workbench.ts"; import { workbenchEventStreamPath } from "../src/api/workbench-events.ts"; import type { WorkbenchStreamTransportRecovery } from "../src/utils/workbench-realtime-runtime.ts"; -import { workbenchRealtimeTraceId } from "../src/utils/workbench-stream-transport.ts"; +import { workbenchRealtimeTraceIdForCapabilities } from "../src/utils/workbench-stream-transport.ts"; import { workbenchRuntimePolicy } from "../src/config/workbench-runtime-policy.ts"; import { AsyncQueue, work } from "../src/utils/scheduler/async-queue.ts"; import { createCoalescedEventQueue } from "../src/utils/scheduler/coalesced-event-queue.ts"; @@ -79,16 +79,18 @@ test("Workbench API uses metadata-only session detail and bounded messages paths assert.equal(workbenchSessionMessagesPathForTest("ses_metadata", { limit: 9 }), "/v1/workbench/sessions/ses_metadata/messages?limit=9"); }); -test("the product EventSource always sends the durable outbox cursor", () => { - assert.equal(workbenchEventStreamPath({ sessionId: "ses_projection", traceId: "trc_projection", afterSeq: 42 }), "/v1/workbench/events?sessionId=ses_projection&traceId=trc_projection&afterSeq=42"); +test("the product EventSource omits projection cursor for live Kafka authority", () => { + assert.equal(workbenchEventStreamPath({ realtimeCapabilities: { liveKafkaSse: true, kafkaRefreshReplay: true, projectionRealtime: false }, sessionId: "ses_projection", traceId: null, afterSeq: 42 }), "/v1/workbench/events?sessionId=ses_projection"); }); test("the active projection trace scope follows the current request", () => { - const beforeSubmit = workbenchRealtimeScopeKey("ses_projection", workbenchRealtimeTraceId(null, null)); - const afterSubmit = workbenchRealtimeScopeKey("ses_projection", workbenchRealtimeTraceId("trc_current_request", "trc_message")); + const projectionCapabilities = { liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true }; + const beforeSubmit = workbenchRealtimeScopeKey("ses_projection", workbenchRealtimeTraceIdForCapabilities(projectionCapabilities, null, null)); + const afterSubmit = workbenchRealtimeScopeKey("ses_projection", workbenchRealtimeTraceIdForCapabilities(projectionCapabilities, "trc_current_request", "trc_message")); assert.notEqual(afterSubmit, beforeSubmit); assert.equal(afterSubmit, workbenchRealtimeScopeKey("ses_projection", "trc_current_request")); - assert.equal(workbenchRealtimeTraceId(null, "trc_message"), "trc_message"); + assert.equal(workbenchRealtimeTraceIdForCapabilities(projectionCapabilities, null, "trc_message"), "trc_message"); + assert.equal(workbenchRealtimeTraceIdForCapabilities({ liveKafkaSse: true, kafkaRefreshReplay: true, projectionRealtime: false }, "trc_current_request"), null); }); test("Error runtime owns Workbench message diagnostic view model", () => { diff --git a/web/hwlab-cloud-web/src/api/workbench-events.test.ts b/web/hwlab-cloud-web/src/api/workbench-events.test.ts index dc773e52..5e4f9a3e 100644 --- a/web/hwlab-cloud-web/src/api/workbench-events.test.ts +++ b/web/hwlab-cloud-web/src/api/workbench-events.test.ts @@ -20,6 +20,7 @@ test("trace events keep sequence-specific coalescing keys", () => { test("product EventSource resumes the transactional projection authority with one durable outbox cursor", () => { assert.equal( workbenchEventStreamPath({ + realtimeCapabilities: { liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true }, sessionId: "ses_projection_authority", traceId: "trc_projection_authority", afterSeq: 42 @@ -36,6 +37,7 @@ test("product EventSource tees projection frames once before reducer delivery", globalThis.EventSource = FakeProductEventSource as unknown as typeof EventSource; try { const stream = connectWorkbenchEvents({ + realtimeCapabilities: { liveKafkaSse: false, kafkaRefreshReplay: false, projectionRealtime: true }, sessionId: "ses_raw_tee", flushYieldMs: 1, onIngress: (frame) => ingress.push(frame), diff --git a/web/hwlab-cloud-web/src/api/workbench-events.ts b/web/hwlab-cloud-web/src/api/workbench-events.ts index f86142a1..1653257d 100644 --- a/web/hwlab-cloud-web/src/api/workbench-events.ts +++ b/web/hwlab-cloud-web/src/api/workbench-events.ts @@ -6,6 +6,7 @@ import type { ApiResult, ChatMessage, ProjectionDiagnostic, TraceEvent } from "@ import { createCoalescedEventQueue } from "@/utils/scheduler/coalesced-event-queue"; import { composeWorkbenchScopedKey, firstScopePart } from "@/utils/workbench-key"; import { recordWorkbenchRuntimeDiagnostic, recordWorkbenchSseLifecycle } from "@/utils/workbench-performance"; +import type { WorkbenchRealtimeCapabilities } from "@/config/runtime"; import { decodeWorkbenchRealtimeEventFrame } from "./workbench-realtime-codec"; export interface WorkbenchRealtimeTraceSnapshot { @@ -28,6 +29,11 @@ export interface WorkbenchRealtimeEvent { eventType?: string | null; eventId?: string | null; sourceEventId?: string | null; + capabilities?: WorkbenchRealtimeCapabilities | null; + liveOnly?: boolean | null; + replay?: boolean | null; + deliverySemantics?: string | null; + lossPossible?: boolean | null; type?: string; contractVersion?: string | null; realtimeAuthority?: string | null; @@ -70,6 +76,7 @@ export interface WorkbenchRealtimeEvent { } export interface WorkbenchEventStreamOptions { + realtimeCapabilities: WorkbenchRealtimeCapabilities; sessionId?: string | null; traceId?: string | null; afterSeq?: number | null; @@ -98,6 +105,7 @@ interface QueuedRealtimeEvent { } const WORKBENCH_EVENT_NAMES = [ + "hwlab.event.v1", "workbench.connected", "workbench.trace.snapshot", "workbench.trace.event", @@ -175,11 +183,11 @@ export function connectWorkbenchEvents(options: WorkbenchEventStreamOptions): Wo }; } -export function workbenchEventStreamPath(options: Pick): string { +export function workbenchEventStreamPath(options: Pick): string { const params = new URLSearchParams(); appendParam(params, "sessionId", options.sessionId); appendParam(params, "traceId", options.traceId); - appendNumberParam(params, "afterSeq", options.afterSeq); + if (!options.realtimeCapabilities.liveKafkaSse && options.realtimeCapabilities.projectionRealtime) appendNumberParam(params, "afterSeq", options.afterSeq); return `/v1/workbench/events?${params.toString()}`; } @@ -212,6 +220,10 @@ function scheduleRealtimeFlushYield(flush: () => void, yieldMs: number | null | } export function realtimeCoalesceKey(event: WorkbenchRealtimeEvent, eventName: string): string | null { + if (eventName === "hwlab.event.v1" || event.schema === "hwlab.event.v1") { + const eventId = firstScopePart(event.eventId, event.sourceEventId, event.event?.sourceEventId); + return eventId ? composeWorkbenchScopedKey("workbench.sse.live-kafka-event", eventId) : null; + } const entity = event.entity; const entityFamily = firstScopePart(entity?.family); const entityId = firstScopePart(entity?.id); diff --git a/web/hwlab-cloud-web/src/config/runtime.ts b/web/hwlab-cloud-web/src/config/runtime.ts index a5be0cec..cf6aa3a9 100644 --- a/web/hwlab-cloud-web/src/config/runtime.ts +++ b/web/hwlab-cloud-web/src/config/runtime.ts @@ -15,6 +15,12 @@ export interface WorkbenchDebugCapabilities { rawHwlabEventWindow: WorkbenchRawHwlabEventWindowCapability; } +export interface WorkbenchRealtimeCapabilities { + liveKafkaSse: boolean; + kafkaRefreshReplay: boolean; + projectionRealtime: boolean; +} + export interface WorkbenchRawHwlabEventWindowCapability { enabled: boolean; maxEntries: number; @@ -112,6 +118,22 @@ export function workbenchDebugCapabilities(): WorkbenchDebugCapabilities { }; } +export function workbenchRealtimeCapabilities(): WorkbenchRealtimeCapabilities { + const features = window.HWLAB_CLOUD_WEB_CONFIG?.workbench?.realtimeFeatures; + if ( + typeof features?.liveKafkaSse !== "boolean" || + typeof features?.kafkaRefreshReplay !== "boolean" || + typeof features?.projectionRealtime !== "boolean" + ) { + throw new Error("HWLAB Cloud Web workbench.realtimeFeatures are required"); + } + return { + liveKafkaSse: features.liveKafkaSse, + kafkaRefreshReplay: features.kafkaRefreshReplay, + projectionRealtime: features.projectionRealtime + }; +} + export function traceExplorerHref(traceId: string | null | undefined): string | null { const safeTraceId = normalizedTraceId(traceId); if (!safeTraceId) return null; diff --git a/web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts b/web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts index dcf4da5e..e9984b5c 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-event-reducer.ts @@ -48,7 +48,15 @@ export function reduceWorkbenchRealtimeEvent(event: WorkbenchRealtimeEvent, even }; } +export function workbenchRealtimeEventIsBusinessActivity(event: WorkbenchRealtimeEvent, eventName: string, liveKafkaSse: boolean): boolean { + if (!liveKafkaSse) return true; + return eventName === "hwlab.event.v1" || event.schema === "hwlab.event.v1"; +} + function reduceRealtimeAction(event: WorkbenchRealtimeEvent, eventName: string): WorkbenchRealtimeAction { + if (event.schema === "hwlab.event.v1" && event.event) { + return { type: "trace.event", traceId: realtimeTraceId(event), event: event.event, snapshot: null, realtimeEvent: event }; + } const authority = primaryAuthority(event); if (authority) return authority; switch (event.type) { diff --git a/web/hwlab-cloud-web/src/stores/workbench-isolated-kafka-debug.ts b/web/hwlab-cloud-web/src/stores/workbench-isolated-kafka-debug.ts index 1eb88acc..1fdce381 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-isolated-kafka-debug.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-isolated-kafka-debug.ts @@ -5,7 +5,9 @@ import type { WorkbenchKafkaSseDebugEvent } from "@/api/workbench-debug"; import type { WorkbenchRealtimeEvent } from "@/api/workbench-events"; import type { ChatMessage } from "@/types"; import { firstNonEmptyString } from "@/utils"; +import { reduceWorkbenchRealtimeEvent } from "./workbench-event-reducer"; import { projectWorkbenchLiveKafkaMessage } from "./workbench-live-kafka-event"; +import { planWorkbenchRealtimeApply } from "./workbench-realtime-plan"; export interface WorkbenchIsolatedKafkaDebugLog { id: string; @@ -46,17 +48,22 @@ export function applyWorkbenchIsolatedKafkaDebugEvent( if (!traceId) return rejected(base, debugEvent, null, "debug-envelope-trace-missing"); if (expected && traceId !== expected) return rejected(base, debugEvent, traceId, "debug-envelope-trace-mismatch"); - const traceEvent = debugEvent.event; - if (!traceEvent) return withLog(base, debugEvent, traceId, "ignore", [], false, "debug-envelope-not-a-trace-event"); + const event = { ...debugEvent, schema: "hwlab.event.v1" } as WorkbenchRealtimeEvent; + const reduced = reduceWorkbenchRealtimeEvent(event, "hwlab.event.v1"); + const plan = planWorkbenchRealtimeApply(reduced.action); + const traceStep = plan.steps.find((step) => step.type === "apply-trace-event"); + if (!traceStep || traceStep.type !== "apply-trace-event" || !traceStep.event) { + return withLog(base, event, traceId, reduced.action.type, plan.steps.map((step) => step.type), false, reduced.action.type === "ignore" ? reduced.action.reason : "debug-envelope-not-a-trace-event"); + } const message = projectWorkbenchLiveKafkaMessage({ previous: state.message, traceId, - sessionId: firstNonEmptyString(debugEvent.hwlabSessionId, debugEvent.sessionId, traceEvent.sessionId, state.message?.sessionId) ?? "ses_workbench_isolated_debug", - event: traceEvent, + sessionId: firstNonEmptyString(event.hwlabSessionId, event.sessionId, traceStep.event.sessionId, state.message?.sessionId) ?? "ses_workbench_isolated_debug", + event: traceStep.event, receivedAt: new Date().toISOString(), title: "Code Agent · 隔离调试" }); - return withLog({ ...base, message, appliedCount: state.appliedCount + 1, error: null }, debugEvent, traceId, "debug.trace.event", ["debug-project-message"], true, null); + return withLog({ ...base, message, appliedCount: state.appliedCount + 1, error: null }, event, traceId, reduced.action.type, plan.steps.map((step) => step.type), true, null); } export function workbenchCurrentDebugTraceId(messages: ChatMessage[], sessionLastTraceId?: string | null): string | null { diff --git a/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts b/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts index 0bbd4b72..abee3c58 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-live-kafka-event.ts @@ -1,5 +1,5 @@ -// SPEC: PJ2026-0104010803 isolated Kafka debug fixture. -// Responsibility: project historical hwlab.event.debug.v1 fixtures for the admin-only isolated debugger; never imported by the product Workbench SSE path. +// SPEC: PJ2026-0104010803 Workbench live Kafka SSE. +// Responsibility: project one transparent hwlab.event.v1 business event into visible turn state without replay/finalizer/polling. import { mergeRunnerTrace } from "../composables/workbench-trace-snapshot"; import type { ChatMessage, TraceEvent, WorkbenchTurnTimingProjection } from "../types"; diff --git a/web/hwlab-cloud-web/src/stores/workbench.ts b/web/hwlab-cloud-web/src/stores/workbench.ts index 1c7f909d..9f31b912 100644 --- a/web/hwlab-cloud-web/src/stores/workbench.ts +++ b/web/hwlab-cloud-web/src/stores/workbench.ts @@ -4,13 +4,13 @@ import { computed, nextTick, ref } from "vue"; import { defineStore } from "pinia"; import { api } from "@/api"; -import { workbenchDebugCapabilities } from "@/config/runtime"; +import { workbenchDebugCapabilities, workbenchRealtimeCapabilities } from "@/config/runtime"; import { workbenchRuntimePolicy } from "@/config/workbench-runtime-policy"; import { createKeyedSingleflight } from "@/utils/scheduler/keyed-singleflight"; import { createWorkbenchHealthProbeCache } from "@/utils/workbench-health"; import { agentErrorFromProjection, normalizeApiErrorRecord, normalizeErrorDiagnostic, normalizeProjectionDiagnostic, projectionDiagnosticFromApiFailure, projectionDiagnosticFromFailure } from "@/utils/workbench-error-runtime"; import { readWorkbenchJson, readWorkbenchNumber, readWorkbenchString, removeWorkbenchStorageKey, writeWorkbenchJson, writeWorkbenchString } from "@/utils/workbench-storage-runtime"; -import { createWorkbenchStreamTransportRuntime, workbenchRealtimeTraceId, type WorkbenchRealtimeEvent, type WorkbenchStreamTransportRecovery } from "@/utils/workbench-realtime-runtime"; +import { createWorkbenchStreamTransportRuntime, workbenchRealtimeTraceIdForCapabilities, type WorkbenchRealtimeEvent, type WorkbenchStreamTransportRecovery } from "@/utils/workbench-realtime-runtime"; import { mergeRunnerTrace, snapshotToRunnerTrace, type TraceSnapshot } from "@/composables/workbench-trace-snapshot"; import type { WorkbenchMessagePageResponse, WorkbenchSessionDetailResponse } from "@/api/workbench"; import type { AgentChatResponse, AgentChatResultResponse, AgentRunProvenance, ApiError, ApiResult, ChatMessage, ErrorDiagnostic, LiveSurface, ProjectionBlocker, ProjectionDiagnostic, ProviderProfile, TraceEvent, WorkbenchSessionRecord, WorkbenchTurnTimingProjection } from "@/types"; @@ -21,7 +21,8 @@ import { RECENT_DRAFTS_STORAGE_KEY, appendSessionPage, defaultProviderProfileOpt import { initialWorkbenchSessionIdFromLocation } from "./workbench-projection"; import { cleanupWorkbenchServerStateSessions, selectActiveMessages, selectActiveSession, selectSessionList, selectSessionStatusAuthority, selectTraceAuthorityById, selectTurnStatusAuthority, type WorkbenchServerAction } from "./workbench-server-state"; import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "./workbench-session-cache"; -import { reduceWorkbenchRealtimeEvent, type WorkbenchRealtimeAction } from "./workbench-event-reducer"; +import { reduceWorkbenchRealtimeEvent, workbenchRealtimeEventIsBusinessActivity, type WorkbenchRealtimeAction } from "./workbench-event-reducer"; +import { projectWorkbenchLiveKafkaMessage, projectWorkbenchLiveKafkaUserMessage, workbenchLiveKafkaEnvelope, workbenchLiveKafkaProjectionTarget } from "./workbench-live-kafka-event"; import { projectRejectedWorkbenchAdmission } from "./workbench-admission-failure"; import { messageHasSealedTerminalResult, messageIsSealedTerminal, traceAuthorityIsSealed } from "./workbench-terminal-authority"; import { boundedProjectionMessageLimit } from "./workbench-message-projection-budget"; @@ -114,6 +115,7 @@ function workbenchMessageIdForTrace(traceId: string, role: "user" | "agent"): st export const useWorkbenchStore = defineStore("workbench", () => { const runtimePolicy = workbenchRuntimePolicy(); + const realtimeCapabilities = workbenchRealtimeCapabilities(); const debugCapabilities = workbenchDebugCapabilities(); const workbenchColadaReducer = useWorkbenchColadaReducer(); const workbenchColadaQueries = useWorkbenchColadaQueries(); @@ -143,6 +145,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { const error = ref(null); const activityRef = ref({ lastActivityAt: Date.now(), lastActivityIso: new Date().toISOString(), waitingFor: "idle", lastEventLabel: null as string | null }); const currentRequest = ref<{ traceId: string; sessionId: string | null; threadId: string | null; status?: string | null } | null>(null); + const liveRealtimeReadySessionId = ref(null); const explicitSessionId = ref(initialWorkbenchSessionIdFromLocation()); const activeSelectionSource = ref(explicitSessionId.value ? "route" : "system"); const selectionEpoch = ref(0); @@ -169,7 +172,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { const sessionListLoadedCount = computed(() => sessionTabs.value.length); const sessionListLoading = computed(() => shouldShowSessionListLoading({ loading: loading.value, sessionsReady: sessionsReady.value })); const sessionDetailLoading = computed(() => Boolean(sessionDetailLoadingId.value || switchingSessionId.value || (loading.value && messages.value.length === 0))); - const composer = computed(() => resolveComposerState({ messages: messages.value, sessions: sessions.value, activeSessionId: activeSessionId.value, chatPending: chatPending.value, realtimeReady: true, currentRequest: currentRequest.value, turnStatusAuthority: turnStatusAuthority.value })); + const composer = computed(() => resolveComposerState({ messages: messages.value, sessions: sessions.value, activeSessionId: activeSessionId.value, chatPending: chatPending.value, realtimeReady: !realtimeCapabilities.liveKafkaSse || liveRealtimeReadySessionId.value === activeSessionId.value, currentRequest: currentRequest.value, turnStatusAuthority: turnStatusAuthority.value })); function recordActivity(label = "user-activity"): void { const now = Date.now(); @@ -957,8 +960,10 @@ export const useWorkbenchStore = defineStore("workbench", () => { void reason; const realtimeScopeKey = workbenchRealtimeScopeKey(sessionId, traceId); const scopeChanged = realtimeTransport.currentKey() !== realtimeScopeKey; + if (realtimeCapabilities.liveKafkaSse && scopeChanged) liveRealtimeReadySessionId.value = null; if (debugCapabilities.rawHwlabEventWindow.enabled && scopeChanged) rawHwlabIngress.value = createRawHwlabIngressState(realtimeScopeKey); realtimeTransport.restart({ + realtimeCapabilities, sessionId, traceId, forceReconnect, @@ -969,8 +974,30 @@ export const useWorkbenchStore = defineStore("workbench", () => { onIngress: debugCapabilities.rawHwlabEventWindow.enabled ? (frame) => { rawHwlabIngress.value = appendRawHwlabIngressFrame(rawHwlabIngress.value, frame, debugCapabilities.rawHwlabEventWindow); } : undefined, + onError: () => { + if (realtimeCapabilities.liveKafkaSse) liveRealtimeReadySessionId.value = null; + }, + onState: (state) => { + if (realtimeCapabilities.liveKafkaSse && ["connecting", "error", "closed", "blocked"].includes(state.phase)) liveRealtimeReadySessionId.value = null; + }, onRecovery: (recovery) => handleRealtimeRecovery(recovery), - onEvent: (event, eventName) => applyRealtimeEvent(event, eventName) + onEvent: (event, eventName) => { + if (realtimeCapabilities.liveKafkaSse && eventName === "workbench.connected") { + const filters = recordValue(event.filters); + const connectedSessionId = normalizeWorkbenchSessionId(filters?.sessionId); + const refreshReplay = realtimeCapabilities.kafkaRefreshReplay; + const deliveryValid = refreshReplay + ? event.deliverySemantics === "kafka-retention-then-live" && event.liveOnly === false && event.replay === true && event.lossPossible === false + : event.deliverySemantics === "live-only" && event.liveOnly === true && event.replay === false; + const valid = deliveryValid + && event.capabilities?.liveKafkaSse === true + && event.capabilities?.kafkaRefreshReplay === refreshReplay + && connectedSessionId === sessionId; + liveRealtimeReadySessionId.value = valid ? connectedSessionId : null; + if (!valid) error.value = "workbench_live_realtime_contract_invalid"; + } + applyRealtimeEvent(event, eventName); + } }); } @@ -997,11 +1024,13 @@ export const useWorkbenchStore = defineStore("workbench", () => { } function stopRealtime(): void { + liveRealtimeReadySessionId.value = null; realtimeTransport.stop(); } function realtimeTraceId(): string | null { - return workbenchRealtimeTraceId( + return workbenchRealtimeTraceIdForCapabilities( + realtimeCapabilities, currentRequest.value?.traceId, activeTraceIdFromMessages(messages.value, turnStatusAuthority.value) ); @@ -1086,7 +1115,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { function applyRealtimeEvent(event: WorkbenchRealtimeEvent, eventName: string): void { const reduced = reduceWorkbenchRealtimeEvent(event, eventName); - recordActivity(reduced.activityLabel); + if (workbenchRealtimeEventIsBusinessActivity(event, eventName, realtimeCapabilities.liveKafkaSse)) recordActivity(reduced.activityLabel); applyWorkbenchRealtimeAction(reduced.action); } @@ -1142,18 +1171,67 @@ export const useWorkbenchStore = defineStore("workbench", () => { function applyRealtimeTraceEvent(traceId: string | null | undefined, event: WorkbenchRealtimeEvent["event"], snapshot: WorkbenchRealtimeEvent["snapshot"], realtimeEvent?: WorkbenchRealtimeEvent | null): void { const id = firstNonEmptyString(traceId, event?.traceId, snapshot?.traceId); if (!id) return; + const liveKafkaEnvelope = workbenchLiveKafkaEnvelope(realtimeEvent?.schema); const sessionId = realtimeEvent ? realtimeEventSessionId(realtimeEvent) : traceResultSessionId(snapshot ?? event ?? null); - if (!shouldApplyActiveTraceAuthority(id, sessionId)) return; - if (traceTerminalBodyIsVisible(id, sessionId)) { + if (!shouldApplyActiveTraceAuthority(id, sessionId, liveKafkaEnvelope)) return; + if (traceTerminalBodyIsVisible(id, sessionId) && !liveKafkaEnvelope) { recordWorkbenchRuntimeDiagnostic({ module: "workbench-terminal-priority", sessionId, traceId: id, outcome: "ok", diagnostic: { code: "terminal_low_priority_sse_trace_skip", source: "realtime-trace-event", valuesRedacted: true } }); return; } const events = event ? [event] : Array.isArray(snapshot?.events) ? snapshot.events : []; markWorkbenchTraceEventsReceived({ traceId: id, events, transport: "sse", serverSentAt: realtimeEvent?.serverSentAt, eventCreatedAt: realtimeEvent?.eventCreatedAt, traceSeq: realtimeEvent?.traceSeq ?? realtimeEvent?.cursor?.traceSeq }); + if (liveKafkaEnvelope && event) { + applyLiveKafkaBusinessEvent(id, sessionId, event, new Date().toISOString()); + return; + } const eventSnapshot = snapshot ?? { traceId: id, status: event?.status, events }; applyTraceSnapshot(id, realtimeSnapshotToTraceSnapshot(id, eventSnapshot, events)); } + function applyLiveKafkaBusinessEvent(traceId: string, authoritySessionId: string | null, event: TraceEvent, receivedAt: string): void { + const eventSessionId = firstNonEmptyString(event.sessionId); + const ownerSessionId = traceOwnerSessionId(traceId, authoritySessionId ?? eventSessionId, true); + if (!ownerSessionId) return; + if (workbenchLiveKafkaProjectionTarget(event) === "user") { + const messageId = firstNonEmptyString(event.userMessageId, event.messageId); + const previousUserMessage = messageId + ? (serverState.value.messagesBySessionId[ownerSessionId] ?? []).find((message) => (message.messageId ?? message.id) === messageId) ?? null + : null; + const userMessage = projectWorkbenchLiveKafkaUserMessage({ previous: previousUserMessage, traceId, sessionId: ownerSessionId, event, receivedAt }); + if (userMessage) reduceServerState({ type: "message.upsert", sessionId: ownerSessionId, message: userMessage }); + else error.value = "workbench_live_user_message_invalid"; + return; + } + const previousMessage = (serverState.value.messagesBySessionId[ownerSessionId] ?? []).find((message) => messageMatchesTraceAuthority(message, traceId, authoritySessionId ?? eventSessionId, ownerSessionId, true)) ?? null; + const message = projectWorkbenchLiveKafkaMessage({ previous: previousMessage, traceId, sessionId: ownerSessionId, event, receivedAt }); + reduceServerState({ type: "message.upsert", sessionId: ownerSessionId, message }); + if (message.runnerTrace) rememberTraceAuthority(message.runnerTrace); + markWorkbenchTraceProjected(traceId); + const terminal = message.traceAutoLifecycle === "terminal"; + const status = firstNonEmptyString(message.status) ?? "running"; + const finalResponse = terminal && message.text ? { text: message.text, status, traceId } : undefined; + rememberTurnStatus(traceId, { + traceId, + sessionId: ownerSessionId, + status, + running: !terminal, + terminal, + finalResponse, + timing: message.timing, + startedAt: message.timing?.startedAt, + lastEventAt: message.timing?.lastEventAt, + finishedAt: message.timing?.finishedAt, + durationMs: message.timing?.durationMs, + updatedAt: receivedAt + } as AgentChatResultResponse); + const existing = sessions.value.find((session) => session.sessionId === ownerSessionId) ?? null; + if (existing) rememberSessionList(mergeSessionIntoList(sessions.value, { ...existing, status, lastTraceId: traceId, updatedAt: receivedAt })); + if (terminal && currentRequest.value?.traceId === traceId) { + chatPending.value = false; + currentRequest.value = null; + } + } + function applyRealtimeTurnSnapshot(turn: Record): void { const traceId = firstNonEmptyString(turn.traceId); if (!traceId) return; diff --git a/web/hwlab-cloud-web/src/types/global.d.ts b/web/hwlab-cloud-web/src/types/global.d.ts index 203fbdcc..a331cbaa 100644 --- a/web/hwlab-cloud-web/src/types/global.d.ts +++ b/web/hwlab-cloud-web/src/types/global.d.ts @@ -8,6 +8,11 @@ declare global { label?: string; }; workbench?: { + realtimeFeatures?: { + liveKafkaSse?: boolean; + kafkaRefreshReplay?: boolean; + projectionRealtime?: boolean; + }; debugCapabilities?: { isolatedKafka?: boolean; rawHwlabEventWindow?: { diff --git a/web/hwlab-cloud-web/src/utils/workbench-realtime-runtime.ts b/web/hwlab-cloud-web/src/utils/workbench-realtime-runtime.ts index 771674b5..b1eddb6a 100644 --- a/web/hwlab-cloud-web/src/utils/workbench-realtime-runtime.ts +++ b/web/hwlab-cloud-web/src/utils/workbench-realtime-runtime.ts @@ -1,4 +1,4 @@ // SPEC: PJ2026-0106050514 Workbench实时运行面 draft-2026-06-30-p0-1297-spec-first. // Responsibility: Production Workbench realtime runtime entry point for SSE transport, event coalescing, and scoped cursor ownership. -export { createWorkbenchStreamTransportRuntime, WorkbenchStreamTransportRuntime, workbenchRealtimeTraceId, type WorkbenchStreamTransportRecovery, type WorkbenchStreamTransportRestartInput, type WorkbenchStreamTransportRestartResult, type WorkbenchStreamTransportState, type WorkbenchRealtimeEvent } from "@/utils/workbench-stream-transport"; +export { createWorkbenchStreamTransportRuntime, WorkbenchStreamTransportRuntime, workbenchRealtimeTraceIdForCapabilities, type WorkbenchStreamTransportRecovery, type WorkbenchStreamTransportRestartInput, type WorkbenchStreamTransportRestartResult, type WorkbenchStreamTransportState, type WorkbenchRealtimeEvent } from "@/utils/workbench-stream-transport"; diff --git a/web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts b/web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts index 6417daa9..cfb29b6a 100644 --- a/web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts +++ b/web/hwlab-cloud-web/src/utils/workbench-stream-transport.ts @@ -8,11 +8,13 @@ // - OpenCode stream.transport.ts:504-530 reduce output/trace. import { connectWorkbenchEvents, type WorkbenchEventStream, type WorkbenchRealtimeEvent, type WorkbenchSseIngressFrame } from "@/api/workbench-events"; +import type { WorkbenchRealtimeCapabilities } from "@/config/runtime"; import { workbenchRealtimeScopeKey } from "@/utils/workbench-key"; export type { WorkbenchRealtimeEvent }; export interface WorkbenchStreamTransportRestartInput { + realtimeCapabilities: WorkbenchRealtimeCapabilities; sessionId?: string | null; traceId?: string | null; afterSeq?: number | null; @@ -76,9 +78,11 @@ interface WorkbenchTransportCursor { traceSeq: number | null; } -export function workbenchRealtimeTraceId( +export function workbenchRealtimeTraceIdForCapabilities( + capabilities: WorkbenchRealtimeCapabilities, ...candidates: Array ): string | null { + if (capabilities.liveKafkaSse) return null; for (const candidate of candidates) { const traceId = typeof candidate === "string" ? candidate.trim() : ""; if (traceId) return traceId; @@ -86,6 +90,10 @@ export function workbenchRealtimeTraceId( return null; } +export function workbenchRealtimeTransportEnabled(capabilities: WorkbenchRealtimeCapabilities): boolean { + return capabilities.liveKafkaSse || capabilities.projectionRealtime; +} + export class WorkbenchStreamTransportRuntime { private stream: WorkbenchEventStream | null = null; private key = ""; @@ -100,12 +108,14 @@ export class WorkbenchStreamTransportRuntime { this.stop(); this.key = key; this.armWait(key); + if (!workbenchRealtimeTransportEnabled(input.realtimeCapabilities)) return { key, changed: true, streamStarted: false }; if (!input.sessionId && !input.traceId) return { key, changed: true, streamStarted: false }; this.emitState(input, "connecting", null); this.stream = connectWorkbenchEvents({ + realtimeCapabilities: input.realtimeCapabilities, sessionId: input.sessionId ?? null, traceId: input.traceId ?? null, - afterSeq: input.afterSeq ?? this.cursorByKey.get(key)?.outboxSeq ?? null, + afterSeq: !input.realtimeCapabilities.liveKafkaSse && input.realtimeCapabilities.projectionRealtime ? input.afterSeq ?? this.cursorByKey.get(key)?.outboxSeq ?? null : null, flushMaxItemsPerChunk: input.flushMaxItemsPerChunk, flushMaxChunkMs: input.flushMaxChunkMs, flushYieldMs: input.flushYieldMs, @@ -122,7 +132,7 @@ export class WorkbenchStreamTransportRuntime { this.requestRecovery(input, event.type || "eventsource-error"); }, onEvent: (event, eventName) => { - this.rememberCursor(key, event); + if (!input.realtimeCapabilities.liveKafkaSse && input.realtimeCapabilities.projectionRealtime) this.rememberCursor(key, event); this.markLive(key); this.emitState(input, "event", eventName); input.onEvent(event, eventName);