diff --git a/deploy/deploy.yaml b/deploy/deploy.yaml index 492fe82e..59f8c88b 100644 --- a/deploy/deploy.yaml +++ b/deploy/deploy.yaml @@ -619,6 +619,7 @@ lanes: HWLAB_CLOUD_WEB_DISPLAY_TIME_LABEL: 北京时间 HWLAB_WORKBENCH_LIVE_KAFKA_SSE_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 HWLAB_CLOUD_WEB_OPENCODE_PASSWORD: secretRef:hwlab-opencode-server-auth/password @@ -849,6 +850,7 @@ lanes: 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 HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC: agentrun.event.v1 @@ -857,6 +859,10 @@ lanes: HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID: hwlab-v03-agentrun-event-direct-publish HWLAB_KAFKA_PROJECTOR_GROUP_ID: hwlab-v03-agentrun-event-projector HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID: hwlab-v03-workbench-live-sse + HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: hwlab.event.debug.v1 + HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: hwlab-v03-workbench-isolated-debug + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "200" + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "15000" HWLAB_KAFKA_PROJECTOR_HEARTBEAT_INTERVAL_MS: "2000" HWLAB_KAFKA_OUTBOX_RELAY_INTERVAL_MS: "250" HWLAB_KAFKA_OUTBOX_RELAY_BATCH_SIZE: "100" diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index e8be13e6..6c2ebe06 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -977,7 +977,7 @@ export async function queryKafkaEventStream({ env = process.env, stream = "hwlab }; } -export async function openKafkaEventStream({ env = process.env, stream = "hwlab", topic = null, traceId = null, sessionId = null, runId = null, commandId = null, fromBeginning = false, onEvent, onError = null, kafkaFactory = defaultKafkaFactory } = {}) { +export async function openKafkaEventStream({ env = process.env, stream = "hwlab", topic = null, groupIdPrefix = null, traceId = null, sessionId = null, runId = null, commandId = null, fromBeginning = false, onEvent, onError = null, kafkaFactory = defaultKafkaFactory } = {}) { if (typeof onEvent !== "function") throw new Error("onEvent callback is required for Kafka event streaming."); const brokers = csv(env.HWLAB_KAFKA_BOOTSTRAP_SERVERS); if (brokers.length === 0) throw new Error("HWLAB_KAFKA_BOOTSTRAP_SERVERS is required for Kafka event streams."); @@ -985,7 +985,8 @@ export async function openKafkaEventStream({ env = process.env, stream = "hwlab" const resolvedTopic = stringValue(topic) || kafkaTopicForStream(stream, env); const filters = compactObject({ traceId, sessionId, runId, commandId }); const kafka = kafkaFactory({ brokers, clientId: `${clientId}-debug-sse` }); - const groupId = `${clientId}-debug-sse-${Date.now()}-${randomUUID().slice(0, 8)}`; + const resolvedGroupIdPrefix = stringValue(groupIdPrefix) || `${clientId}-debug-sse`; + const groupId = `${resolvedGroupIdPrefix}-${Date.now()}-${randomUUID().slice(0, 8)}`; const consumer = kafka.consumer({ groupId, allowAutoTopicCreation: false }); let running = false; let stopped = false; diff --git a/internal/cloud/workbench-kafka-debug-capability.ts b/internal/cloud/workbench-kafka-debug-capability.ts new file mode 100644 index 00000000..e8a311a0 --- /dev/null +++ b/internal/cloud/workbench-kafka-debug-capability.ts @@ -0,0 +1,49 @@ +// Responsibility: validate the independently YAML-owned Workbench Kafka debug stream. +// This capability is intentionally separate from the five product realtime capabilities. + +export const WORKBENCH_KAFKA_DEBUG_ENVS = Object.freeze({ + enabled: "HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED", + topic: "HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC", + groupIdPrefix: "HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX", + replayLimit: "HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT", + replayTimeoutMs: "HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS" +}); + +export function workbenchKafkaDebugCapability(env = process.env) { + const enabled = requiredBoolean(env?.[WORKBENCH_KAFKA_DEBUG_ENVS.enabled], WORKBENCH_KAFKA_DEBUG_ENVS.enabled); + const topic = requiredKafkaName(env?.[WORKBENCH_KAFKA_DEBUG_ENVS.topic], WORKBENCH_KAFKA_DEBUG_ENVS.topic, 249); + const groupIdPrefix = requiredKafkaName(env?.[WORKBENCH_KAFKA_DEBUG_ENVS.groupIdPrefix], WORKBENCH_KAFKA_DEBUG_ENVS.groupIdPrefix, 180); + const replayLimit = requiredPositiveInteger(env?.[WORKBENCH_KAFKA_DEBUG_ENVS.replayLimit], WORKBENCH_KAFKA_DEBUG_ENVS.replayLimit); + const replayTimeoutMs = requiredPositiveInteger(env?.[WORKBENCH_KAFKA_DEBUG_ENVS.replayTimeoutMs], WORKBENCH_KAFKA_DEBUG_ENVS.replayTimeoutMs); + return { enabled, topic, groupIdPrefix, replayLimit, replayTimeoutMs, valuesPrinted: false }; +} + +function requiredBoolean(value, name) { + const normalized = String(value ?? "").trim().toLowerCase(); + if (normalized === "true" || normalized === "1") return true; + if (normalized === "false" || normalized === "0") return false; + throw capabilityError(`${name} is required and must be explicitly true or false.`, name); +} + +function requiredKafkaName(value, name, maxLength) { + const text = String(value ?? "").trim(); + if (!text || text.length > maxLength || !/^[A-Za-z0-9._-]+$/u.test(text)) { + throw capabilityError(`${name} is required and must be a valid Kafka identifier.`, name); + } + return text; +} + +function requiredPositiveInteger(value, name) { + const number = Number(value); + if (!Number.isInteger(number) || number <= 0) throw capabilityError(`${name} is required and must be a positive integer.`, name); + return number; +} + +function capabilityError(message, envName) { + return Object.assign(new Error(message), { + code: "workbench_kafka_debug_capability_invalid", + envName, + statusCode: 503, + valuesRedacted: true + }); +} diff --git a/internal/cloud/workbench-kafka-sse-debug.test.ts b/internal/cloud/workbench-kafka-sse-debug.test.ts index f9e200a5..6d7ba0d1 100644 --- a/internal/cloud/workbench-kafka-sse-debug.test.ts +++ b/internal/cloud/workbench-kafka-sse-debug.test.ts @@ -40,6 +40,114 @@ test("workbench Kafka SSE debug endpoint streams filtered raw HWLAB Kafka events } }); +test("isolated Workbench Kafka replay uses its YAML-owned topic, unique group, and bound", async () => { + const fakeKafka = createFakeKafkaFactory(); + const env = { + HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092", + HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true", + HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.debug.v1", + HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-test-workbench-isolated-debug", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "10", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "3000" + }; + const server = createServer((request, response) => { + const url = new URL(request.url ?? "/", "http://127.0.0.1"); + void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, accessController: fakeAccessController(), logger: null }); + }); + await listen(server); + const abort = new AbortController(); + try { + const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_isolated_debug&fromBeginning=true`, { signal: abort.signal }); + assert.equal(response.status, 200); + const reader = response.body?.getReader(); + assert.ok(reader); + const connected = await readUntil(reader, "hwlab.kafka.connected"); + assert.match(connected, /hwlab\.event\.debug\.v1/u); + assert.match(connected, /"debugIsolation":true/u); + assert.match(connected, /"deliverySemantics":"debug-replay"/u); + assert.match(connected, /"liveOnly":false/u); + assert.match(connected, /"replay":true/u); + assert.match(fakeKafka.groupId ?? "", /^hwlab-test-workbench-isolated-debug-/u); + assert.equal(fakeKafka.subscribedTopic, "hwlab.event.debug.v1"); + assert.equal(fakeKafka.fromBeginning, true); + + await fakeKafka.emit({ schema: "hwlab.event.debug.v1", eventType: "assistant", traceId: "trc_other", event: { traceId: "trc_other", type: "assistant", text: "ignored" } }); + await fakeKafka.emit({ schema: "hwlab.event.debug.v1", eventType: "assistant", traceId: "trc_isolated_debug", event: { traceId: "trc_isolated_debug", type: "assistant", text: "isolated-visible" } }); + await fakeKafka.emit({ schema: "hwlab.event.debug.v1", eventType: "terminal", traceId: "trc_isolated_debug", event: { traceId: "trc_isolated_debug", eventType: "terminal", type: "result", terminal: true, status: "completed" } }); + const streamed = await readUntil(reader, "hwlab.kafka.replay-complete"); + assert.match(streamed, /isolated-visible/u); + assert.match(streamed, /"reason":"terminal"/u); + assert.match(streamed, /"count":2/u); + assert.match(streamed, /"terminalObserved":true/u); + assert.doesNotMatch(streamed, /trc_other/u); + } finally { + abort.abort(); + await close(server); + } +}); + +test("isolated Workbench Kafka debug fails closed when disabled and rejects replay", async () => { + const fakeKafka = createFakeKafkaFactory(); + const env = { + HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "false", + HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.debug.v1", + HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-test-workbench-isolated-debug", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "10", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "3000" + }; + const server = createServer((request, response) => { + const url = new URL(request.url ?? "/", "http://127.0.0.1"); + void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, accessController: fakeAccessController(), logger: null }); + }); + await listen(server); + try { + const disabled = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_disabled&fromBeginning=true`); + assert.equal(disabled.status, 503); + assert.match(await disabled.text(), /workbench_kafka_debug_disabled/u); + assert.equal(fakeKafka.consumerCreated, false); + + env.HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED = "true"; + const replay = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_replay`); + assert.equal(replay.status, 400); + assert.match(await replay.text(), /workbench_kafka_debug_replay_required/u); + assert.equal(fakeKafka.consumerCreated, false); + } finally { + await close(server); + } +}); + +test("isolated Workbench Kafka replay closes as incomplete when terminal is absent", async () => { + const fakeKafka = createFakeKafkaFactory(); + const env = { + HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true", + HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.debug.v1", + HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-test-workbench-isolated-debug", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "10", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "30" + }; + const server = createServer((request, response) => { + const url = new URL(request.url ?? "/", "http://127.0.0.1"); + void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, accessController: fakeAccessController(), logger: null }); + }); + await listen(server); + try { + const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&traceId=trc_incomplete&fromBeginning=true`); + assert.equal(response.status, 200); + const reader = response.body?.getReader(); + assert.ok(reader); + const replay = await readUntil(reader, "hwlab.kafka.replay-complete"); + assert.match(replay, /"reason":"timeout"/u); + assert.match(replay, /"count":0/u); + assert.match(replay, /"terminalObserved":false/u); + await waitUntil(() => fakeKafka.stopCalls > 0); + } finally { + await close(server); + } +}); + test("workbench Kafka SSE debug endpoint streams filtered codex stdio Kafka events", async () => { const fakeKafka = createFakeKafkaFactory(); const env = { @@ -213,6 +321,7 @@ test("workbench Kafka SSE debug endpoint stops a consumer that becomes ready aft function createFakeKafkaFactory(options: { deferRun?: boolean } = {}) { let eachMessage: ((input: any) => Promise) | null = null; let subscribedTopic = "hwlab.event.v1"; + let fromBeginning = false; let releaseRun: (() => void) | null = null; const runGate = options.deferRun ? new Promise((resolve) => { releaseRun = resolve; }) : null; let groupId: string | null = null; @@ -221,7 +330,10 @@ function createFakeKafkaFactory(options: { deferRun?: boolean } = {}) { let consumerCreated = false; const consumer = { connect: async () => undefined, - subscribe: async (input: { topic?: string }) => { subscribedTopic = String(input.topic ?? subscribedTopic); }, + subscribe: async (input: { topic?: string; fromBeginning?: boolean }) => { + subscribedTopic = String(input.topic ?? subscribedTopic); + fromBeginning = input.fromBeginning === true; + }, run: async (input: any) => { if (runGate) await runGate; eachMessage = input.eachMessage; @@ -238,6 +350,8 @@ function createFakeKafkaFactory(options: { deferRun?: boolean } = {}) { } }), get groupId() { return groupId; }, + get subscribedTopic() { return subscribedTopic; }, + get fromBeginning() { return fromBeginning; }, get stopCalls() { return stopCalls; }, get disconnectCalls() { return disconnectCalls; }, get consumerCreated() { return consumerCreated; }, diff --git a/internal/cloud/workbench-kafka-sse-debug.ts b/internal/cloud/workbench-kafka-sse-debug.ts index 10d4b13a..c858abfd 100644 --- a/internal/cloud/workbench-kafka-sse-debug.ts +++ b/internal/cloud/workbench-kafka-sse-debug.ts @@ -6,10 +6,11 @@ import { sendJson } from "./server-http-utils.ts"; import { openKafkaEventStream } from "./kafka-event-bridge.ts"; import { defaultCodeAgentTraceStore } from "./code-agent-trace-store.ts"; import { authenticateWorkbenchRead } from "./server-workbench-read-http.ts"; +import { workbenchKafkaDebugCapability } from "./workbench-kafka-debug-capability.ts"; const CONTRACT_VERSION = "workbench-debug-kafka-sse-v1"; const DEFAULT_STREAM = "hwlab"; -const ALLOWED_STREAMS = new Set(["stdio", "agentrun", "hwlab"]); +const ALLOWED_STREAMS = new Set(["stdio", "agentrun", "hwlab", "hwlab-debug"]); export async function handleWorkbenchKafkaSseDebugHttp(request, response, url, options = {}) { const route = routeSuffix(url.pathname); @@ -19,13 +20,16 @@ export async function handleWorkbenchKafkaSseDebugHttp(request, response, url, o if (auth.actor?.role !== "admin") { return sendJson(response, 403, debugError("admin_required", "Only admin users can access raw Workbench Kafka debug streams.")); } + const stream = streamFromUrl(url); + const isolatedDebug = stream === "hwlab-debug" ? isolatedDebugConfig(response, url, options.env ?? process.env, { requireReplay: route === "/events" }) : null; + if (isolatedDebug === false) return; if (route === "") { if (request.method !== "GET") return methodNotAllowed(response, "GET"); - return sendJson(response, 200, describeKafkaSseDebug(url, options.env ?? process.env)); + return sendJson(response, 200, describeKafkaSseDebug(url, options.env ?? process.env, isolatedDebug)); } if (route === "/events") { if (request.method !== "GET") return methodNotAllowed(response, "GET"); - return openKafkaDebugSse(request, response, url, options); + return openKafkaDebugSse(request, response, url, options, isolatedDebug); } return sendJson(response, 404, debugError("workbench_debug_kafka_sse_route_not_found", "Workbench debug Kafka SSE route is not implemented.", { route })); } catch (error) { @@ -41,15 +45,21 @@ export async function handleWorkbenchKafkaSseDebugHttp(request, response, url, o } } -function describeKafkaSseDebug(url, env) { +function describeKafkaSseDebug(url, env, isolatedDebug = null) { const stream = streamFromUrl(url); const requestedFilters = filtersFromUrl(url); - const resolvedFilters = resolveKafkaDebugFilters(requestedFilters, { traceStore: defaultCodeAgentTraceStore }); + const resolvedFilters = stream === "hwlab-debug" ? compactObject({ traceId: requestedFilters.traceId }) : resolveKafkaDebugFilters(requestedFilters, { traceStore: defaultCodeAgentTraceStore }); return { ok: true, contractVersion: CONTRACT_VERSION, stream, - topic: topicForStream(stream, env), + topic: isolatedDebug?.topic ?? topicForStream(stream, env), + debugIsolation: stream === "hwlab-debug", + liveOnly: false, + replay: stream === "hwlab-debug", + replayLimit: isolatedDebug?.replayLimit ?? null, + replayTimeoutMs: isolatedDebug?.replayTimeoutMs ?? null, + groupIdPrefix: isolatedDebug?.groupIdPrefix ?? null, filters: requestedFilters, resolvedFilters, eventsRoute: `/v1/workbench/debug/kafka-sse/events${url.search || ""}`, @@ -57,11 +67,12 @@ function describeKafkaSseDebug(url, env) { }; } -async function openKafkaDebugSse(request, response, url, options) { +async function openKafkaDebugSse(request, response, url, options, isolatedDebug = null) { const env = options.env ?? process.env; const stream = streamFromUrl(url); const filters = filtersFromUrl(url); - const resolvedFilters = resolveKafkaDebugFilters(filters, options); + const isolatedReplay = stream === "hwlab-debug"; + const resolvedFilters = isolatedReplay ? compactObject({ traceId: filters.traceId }) : resolveKafkaDebugFilters(filters, options); response.writeHead(200, { "content-type": "text/event-stream; charset=utf-8", "cache-control": "no-store, no-transform", @@ -72,35 +83,81 @@ async function openKafkaDebugSse(request, response, url, options) { if (typeof response.flushHeaders === "function") response.flushHeaders(); let kafkaStream = null; let closed = false; + let consumerReady = false; + let replayFinished = false; + let replayTimer = null; + let replayCount = 0; + let replayTerminalObserved = false; + const bufferedRecords = []; const close = () => { if (closed) return; closed = true; + if (replayTimer) clearTimeout(replayTimer); void kafkaStream?.stop?.(); }; response.once("close", close); request.once?.("aborted", close); request.socket?.once?.("close", close); + const finishReplay = (reason) => { + if (!isolatedReplay || replayFinished || closed) return; + replayFinished = true; + if (replayTimer) clearTimeout(replayTimer); + writeSse(response, "hwlab.kafka.replay-complete", { + ok: true, + contractVersion: CONTRACT_VERSION, + stream, + topic: isolatedDebug.topic, + traceId: resolvedFilters.traceId, + replayComplete: true, + reason, + count: replayCount, + limit: isolatedDebug.replayLimit, + timeoutMs: isolatedDebug.replayTimeoutMs, + terminalObserved: replayTerminalObserved, + valuesPrinted: false + }); + if (!response.writableEnded) response.end(); + }; + const writeRecord = (record) => { + if (closed || replayFinished) return; + writeSse(response, "hwlab.kafka.event", { + ok: true, + contractVersion: CONTRACT_VERSION, + stream, + topic: record.topic, + partition: record.partition, + offset: record.offset, + key: record.key, + timestamp: record.timestamp, + value: record.value, + serverSentAt: new Date().toISOString(), + valuesPrinted: false + }); + if (!isolatedReplay) return; + replayCount += 1; + replayTerminalObserved = debugRecordIsTerminal(record); + if (replayTerminalObserved) finishReplay("terminal"); + else if (replayCount >= isolatedDebug.replayLimit) finishReplay("limit"); + }; + if (isolatedReplay) { + replayTimer = setTimeout(() => finishReplay("timeout"), isolatedDebug.replayTimeoutMs); + replayTimer.unref?.(); + } try { kafkaStream = await openKafkaEventStream({ env, stream, - fromBeginning: url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true", + topic: isolatedDebug?.topic ?? null, + groupIdPrefix: isolatedDebug?.groupIdPrefix ?? null, + fromBeginning: isolatedReplay || url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true", ...resolvedFilters, kafkaFactory: options.kafkaFactory, onEvent: async (record) => { - writeSse(response, "hwlab.kafka.event", { - ok: true, - contractVersion: CONTRACT_VERSION, - stream, - topic: record.topic, - partition: record.partition, - offset: record.offset, - key: record.key, - timestamp: record.timestamp, - value: record.value, - serverSentAt: new Date().toISOString(), - valuesPrinted: false - }); + if (isolatedReplay && !consumerReady) { + if (bufferedRecords.length < isolatedDebug.replayLimit) bufferedRecords.push(record); + return; + } + writeRecord(record); } }); if (closed) { @@ -113,14 +170,27 @@ async function openKafkaDebugSse(request, response, url, options) { consumerReady: true, groupId: kafkaStream.groupId, stream, - topic: kafkaStream.topic ?? topicForStream(stream, env), + topic: kafkaStream.topic ?? isolatedDebug?.topic ?? topicForStream(stream, env), + debugIsolation: isolatedReplay, + deliverySemantics: isolatedReplay ? "debug-replay" : "diagnostic", + liveOnly: false, + replay: isolatedReplay, + replayLimit: isolatedDebug?.replayLimit ?? null, + replayTimeoutMs: isolatedDebug?.replayTimeoutMs ?? null, filters, resolvedFilters, serverSentAt: new Date().toISOString(), valuesPrinted: false }); + consumerReady = true; + for (const record of bufferedRecords.splice(0)) { + if (replayFinished || closed) break; + writeRecord(record); + } } catch (error) { + if (replayTimer) clearTimeout(replayTimer); writeSse(response, "hwlab.kafka.error", { ok: false, error: errorMessagePayload(error), valuesPrinted: false }); + if (isolatedReplay && !response.writableEnded) response.end(); } } @@ -191,9 +261,33 @@ function filtersFromUrl(url) { function topicForStream(stream, env) { if (stream === "stdio") return textValue(env.HWLAB_KAFKA_STDIO_TOPIC ?? env.AGENTRUN_KAFKA_STDIO_TOPIC) || "codex-stdio.raw.v1"; if (stream === "agentrun") return textValue(env.HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC) || "agentrun.event.v1"; + if (stream === "hwlab-debug") return workbenchKafkaDebugCapability(env).topic; return textValue(env.HWLAB_KAFKA_EVENT_TOPIC) || "hwlab.event.v1"; } +function isolatedDebugConfig(response, url, env, options = {}) { + let config; + try { + config = workbenchKafkaDebugCapability(env); + } catch (error) { + sendJson(response, Number(error?.statusCode) || 503, debugError(error?.code || "workbench_kafka_debug_capability_invalid", error instanceof Error ? error.message : "Workbench isolated Kafka debug capability is invalid.")); + return false; + } + if (!config.enabled) { + sendJson(response, 503, debugError("workbench_kafka_debug_disabled", "Workbench isolated Kafka debug capability is disabled.")); + return false; + } + if (options.requireReplay && !safeId(url.searchParams.get("traceId") || url.searchParams.get("trace-id"))) { + sendJson(response, 400, debugError("workbench_kafka_debug_trace_required", "Isolated Workbench Kafka replay requires traceId.")); + return false; + } + if (options.requireReplay && !["1", "true"].includes(String(url.searchParams.get("fromBeginning") || "").toLowerCase())) { + sendJson(response, 400, debugError("workbench_kafka_debug_replay_required", "Isolated Workbench Kafka debug requires explicit fromBeginning=true.")); + return false; + } + return config; +} + function methodNotAllowed(response, allow) { response.setHeader("allow", allow); return sendJson(response, 405, debugError("method_not_allowed", `Method not allowed; expected ${allow}.`)); @@ -207,6 +301,13 @@ function errorMessagePayload(error) { return { code: "kafka_sse_error", message: error instanceof Error ? error.message : String(error ?? "unknown") }; } +function debugRecordIsTerminal(record) { + const value = record?.value && typeof record.value === "object" && !Array.isArray(record.value) ? record.value : {}; + const event = value.event && typeof value.event === "object" && !Array.isArray(value.event) ? value.event : {}; + const eventType = textValue(value.eventType ?? event.eventType ?? event.type); + return event.terminal === true || eventType === "terminal" || eventType === "result"; +} + function safeId(value) { const text = textValue(value); return text && /^[A-Za-z0-9_.:-]{3,220}$/u.test(text) ? text : null; diff --git a/internal/dev-entrypoint/cloud-web-runtime.mjs b/internal/dev-entrypoint/cloud-web-runtime.mjs index 8fbb4a73..e63230a0 100644 --- a/internal/dev-entrypoint/cloud-web-runtime.mjs +++ b/internal/dev-entrypoint/cloud-web-runtime.mjs @@ -394,6 +394,9 @@ function workbenchRuntimeConfigFromEnv() { realtimeFeatures: { liveKafkaSse: requiredWorkbenchRealtimeFeature(process.env.HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED, "HWLAB_WORKBENCH_LIVE_KAFKA_SSE_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") } }; const autoExpandRunning = parseEnvBoolean(process.env.HWLAB_WORKBENCH_TRACE_AUTO_EXPAND_RUNNING); diff --git a/internal/dev-entrypoint/cloud-web-runtime.test.mjs b/internal/dev-entrypoint/cloud-web-runtime.test.mjs index d16bc98c..5a78aa2b 100644 --- a/internal/dev-entrypoint/cloud-web-runtime.test.mjs +++ b/internal/dev-entrypoint/cloud-web-runtime.test.mjs @@ -431,7 +431,8 @@ test("cloud web serves client deep links through the Vue shell", async () => { HWLAB_CLOUD_WEB_DISPLAY_TIME_LOCALE: "zh-CN", HWLAB_CLOUD_WEB_DISPLAY_TIME_LABEL: "北京时间", HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true" }); const root = await mkdtemp(path.join(os.tmpdir(), "hwlab-cloud-web-runtime-")); await writeFile(path.join(root, "index.html"), "
\n", "utf8"); @@ -462,6 +463,7 @@ test("cloud web serves client deep links through the Vue shell", async () => { const html = await response.text(); assert.match(html, /
<\/div>/u, route); assert.match(html, /HWLAB_CLOUD_WEB_CONFIG/u, route); + assert.match(html, /"debugCapabilities":\{"isolatedKafka":true\}/u, route); } const assetResponse = await fetch(`${serverUrl(cloudWeb)}/asset.txt`, { @@ -488,7 +490,8 @@ test("cloud web injects trace explorer runtime config from env", async () => { HWLAB_CLOUD_WEB_DISPLAY_TIME_LABEL: "北京时间", HWLAB_WORKBENCH_TRACE_EXPLORER_URL_TEMPLATE: "/v1/workbench/traces/{trace_id}/events", HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true" }); const root = await mkdtemp(path.join(os.tmpdir(), "hwlab-cloud-web-runtime-")); await writeFile(path.join(root, "index.html"), "
\n", "utf8"); @@ -594,7 +597,8 @@ test("cloud web OpenCode frame-url endpoint mints traced tickets", async () => { HWLAB_CLOUD_WEB_DISPLAY_TIME_LABEL: "北京时间", HWLAB_CLOUD_WEB_OPENCODE_URL: "https://opencode.example.test", HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true" }); const cloudApiRequests = []; const opencodeRequests = []; @@ -1036,7 +1040,8 @@ test("cloud web OpenCode proxy accepts short-lived tickets minted by the shell", HWLAB_CLOUD_WEB_DISPLAY_TIME_LABEL: "北京时间", HWLAB_CLOUD_WEB_OPENCODE_URL: "https://opencode.example.test", HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED: "true", - HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false" + HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED: "false", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true" }); const root = await mkdtemp(path.join(os.tmpdir(), "hwlab-cloud-web-runtime-")); await writeFile(path.join(root, "index.html"), "
\n", "utf8"); diff --git a/scripts/gitops-render.test.ts b/scripts/gitops-render.test.ts index b1248245..e9ffcce3 100644 --- a/scripts/gitops-render.test.ts +++ b/scripts/gitops-render.test.ts @@ -640,6 +640,7 @@ test("v03 render keeps node identity as data instead of generated structure", as 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"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC"), "agentrun.event.v1"); @@ -648,12 +649,17 @@ test("v03 render keeps node identity as data instead of generated structure", as assert.equal(cloudApiEnv.get("HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID"), "hwlab-v03-agentrun-event-direct-publish"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_PROJECTOR_GROUP_ID"), "hwlab-v03-agentrun-event-projector"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID"), "hwlab-v03-workbench-live-sse"); + assert.equal(cloudApiEnv.get("HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC"), "hwlab.event.debug.v1"); + assert.equal(cloudApiEnv.get("HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX"), "hwlab-v03-workbench-isolated-debug"); + assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT"), "200"); + assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS"), "15000"); assert.equal(cloudApiEnv.has("HWLAB_SESSION_COOKIE_DOMAIN"), false); 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])); assert.equal(cloudWebEnv.get("HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED"), "true"); assert.equal(cloudWebEnv.get("HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED"), "false"); + assert.equal(cloudWebEnv.get("HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED"), "true"); assert.match(cloudApiEnv.get("HWLAB_ENVIRONMENT_IMAGE") ?? "", /hwlab-cloud-api-env:env-stale/u); assert.deepEqual(cloudApiEnvEntries.get("HWLAB_SECRET_PLANE_SMOKE")?.valueFrom?.secretKeyRef, { name: "hwlab-secret-plane-smoke", key: "password", optional: true }); assert.ok((cloudApi?.spec?.template?.spec?.volumes ?? []).some((volume) => volume.name === "hwpod-preinstalled-specs" && volume.configMap?.name === "hwlab-v03-hwpod-preinstalled-specs")); diff --git a/scripts/runtime-realtime-capabilities-config.test.mjs b/scripts/runtime-realtime-capabilities-config.test.mjs index 2a336105..77c0ed62 100644 --- a/scripts/runtime-realtime-capabilities-config.test.mjs +++ b/scripts/runtime-realtime-capabilities-config.test.mjs @@ -20,6 +20,14 @@ const LIVE_KAFKA_CONTRACT = Object.freeze({ HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID: "hwlab-v03-workbench-live-sse" }); +const ISOLATED_DEBUG_CONTRACT = Object.freeze({ + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED: "true", + HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC: "hwlab.event.debug.v1", + HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX: "hwlab-v03-workbench-isolated-debug", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT: "200", + HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_TIMEOUT_MS: "15000" +}); + test("v03 YAML keeps pure Kafka realtime capabilities independently composable", async () => { const deploy = await readStructuredFile(process.cwd(), "deploy/deploy.yaml"); const cloudApi = deploy.lanes?.v03?.services?.find((service) => service.serviceId === "hwlab-cloud-api"); @@ -33,9 +41,13 @@ test("v03 YAML keeps pure Kafka realtime capabilities independently composable", for (const [name, expected] of Object.entries(LIVE_KAFKA_CONTRACT)) { assert.equal(cloudApi.env?.[name], expected, `${name} must preserve the pure Kafka live contract`); } + for (const [name, expected] of Object.entries(ISOLATED_DEBUG_CONTRACT)) { + assert.equal(cloudApi.env?.[name], expected, `${name} must be independently YAML-owned for debug only`); + } assert.equal(cloudApi.env?.HWLAB_SESSION_COOKIE_DOMAIN, undefined, "internal and public origins require host-scoped session cookies"); assert.notEqual(cloudApi.env?.HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID, "hwlab-v03-agentrun-event-bridge"); assert.equal(cloudWeb.env?.HWLAB_WORKBENCH_LIVE_KAFKA_SSE_ENABLED, "true"); assert.equal(cloudWeb.env?.HWLAB_WORKBENCH_PROJECTION_REALTIME_ENABLED, "false"); + assert.equal(cloudWeb.env?.HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_ENABLED, "true"); }); diff --git a/web/hwlab-cloud-web/scripts/workbench-e2e-server.ts b/web/hwlab-cloud-web/scripts/workbench-e2e-server.ts index 4ea06a6e..924f0a75 100644 --- a/web/hwlab-cloud-web/scripts/workbench-e2e-server.ts +++ b/web/hwlab-cloud-web/scripts/workbench-e2e-server.ts @@ -102,7 +102,7 @@ const terminalAssistantFinalText = [ ].join("\n"); const canonicalTraceFinalText = "messageProjection 的 sealed final response 才是主消息正文。"; const traceDetailConflictText = "TRACE_DETAIL_CONFLICT_SHOULD_NOT_REPLACE_MESSAGE"; -const runtimeConfigScript = ``; +const runtimeConfigScript = ``; let state = createScenarioState("baseline"); const sseClients = new Set(); diff --git a/web/hwlab-cloud-web/src/api/workbench-debug.test.ts b/web/hwlab-cloud-web/src/api/workbench-debug.test.ts index 20991c39..4e8cf354 100644 --- a/web/hwlab-cloud-web/src/api/workbench-debug.test.ts +++ b/web/hwlab-cloud-web/src/api/workbench-debug.test.ts @@ -17,6 +17,17 @@ test("Kafka SSE debug path preserves stream and explicit correlation filters", ( ); }); +test("isolated Kafka debug replay path is trace-scoped and explicit", () => { + assert.equal( + workbenchKafkaSseDebugPath({ + stream: "hwlab-debug", + fromBeginning: true, + traceId: "trc_isolated_debug" + }), + "/v1/workbench/debug/kafka-sse/events?stream=hwlab-debug&fromBeginning=true&traceId=trc_isolated_debug" + ); +}); + test("Kafka SSE debug correlation ids normalize stdio snake_case metadata", () => { assert.deepEqual( workbenchKafkaSseDebugCorrelationIds({ diff --git a/web/hwlab-cloud-web/src/api/workbench-debug.ts b/web/hwlab-cloud-web/src/api/workbench-debug.ts index a3b61806..35e91b5d 100644 --- a/web/hwlab-cloud-web/src/api/workbench-debug.ts +++ b/web/hwlab-cloud-web/src/api/workbench-debug.ts @@ -50,12 +50,23 @@ export interface WorkbenchDebugFakeSseStream { close: () => void; } -export type WorkbenchKafkaSseDebugStreamName = "stdio" | "agentrun" | "hwlab"; +export type WorkbenchKafkaSseDebugStreamName = "stdio" | "agentrun" | "hwlab" | "hwlab-debug"; export interface WorkbenchKafkaSseDebugEvent { ok?: boolean; contractVersion?: string; consumerReady?: boolean; + debugIsolation?: boolean; + deliverySemantics?: string; + liveOnly?: boolean; + replay?: boolean; + replayComplete?: boolean; + reason?: string; + count?: number; + limit?: number; + timeoutMs?: number; + terminalObserved?: boolean; + traceId?: string | null; groupId?: string; stream?: WorkbenchKafkaSseDebugStreamName; topic?: string; @@ -97,7 +108,7 @@ const DEBUG_FAKE_SSE_EVENTS = [ "workbench.error" ]; -const DEBUG_KAFKA_SSE_EVENTS = ["hwlab.kafka.connected", "hwlab.kafka.event", "hwlab.kafka.error"]; +const DEBUG_KAFKA_SSE_EVENTS = ["hwlab.kafka.connected", "hwlab.kafka.event", "hwlab.kafka.replay-complete", "hwlab.kafka.error"]; export const workbenchDebugAPI = { describe: (queueId = "trace-card", options: ApiRequestOptions = {}): Promise> => fetchJson(debugPath("", { queueId }), { ...options, timeoutName: "workbench debug fake sse" }), diff --git a/web/hwlab-cloud-web/src/components/workbench/WorkbenchKafkaDebugPanel.test.ts b/web/hwlab-cloud-web/src/components/workbench/WorkbenchKafkaDebugPanel.test.ts new file mode 100644 index 00000000..2d897d7e --- /dev/null +++ b/web/hwlab-cloud-web/src/components/workbench/WorkbenchKafkaDebugPanel.test.ts @@ -0,0 +1,102 @@ +import { afterEach, beforeEach, describe, expect, it } from "vitest"; +import { mount } from "@vue/test-utils"; +import { nextTick } from "vue"; + +import WorkbenchKafkaDebugPanel from "./WorkbenchKafkaDebugPanel.vue"; + +describe("WorkbenchKafkaDebugPanel", () => { + const originalEventSource = globalThis.EventSource; + + beforeEach(() => { + FakeEventSource.instances = []; + globalThis.EventSource = FakeEventSource as unknown as typeof EventSource; + }); + + afterEach(() => { + globalThis.EventSource = originalEventSource; + }); + + it("replays only after an explicit click, completes on terminal, and clears locally", async () => { + const wrapper = mount(WorkbenchKafkaDebugPanel, { props: { traceId: "trc_component_debug" } }); + expect(FakeEventSource.instances).toHaveLength(0); + + await wrapper.get('[data-testid="workbench-kafka-debug-replay"]').trigger("click"); + const source = FakeEventSource.instances[0]; + expect(source.url).toContain("stream=hwlab-debug"); + expect(source.url).toContain("fromBeginning=true"); + expect(source.url).toContain("traceId=trc_component_debug"); + + source.emit("hwlab.kafka.connected", { consumerReady: true, debugIsolation: true, deliverySemantics: "debug-replay", liveOnly: false, replay: true, topic: "hwlab.event.debug.v1", groupId: "debug-group-unique" }); + source.emit("hwlab.kafka.event", debugRecord("assistant", { type: "assistant", assistantText: "debug answer", status: "running" })); + source.emit("hwlab.kafka.event", debugRecord("terminal", { type: "result", eventType: "terminal", terminal: true, status: "completed" })); + source.emit("hwlab.kafka.replay-complete", { reason: "terminal", terminalObserved: true, count: 2 }); + await nextTick(); + + expect(wrapper.get('[data-testid="workbench-kafka-debug-panel"]').attributes("data-replay-status")).toBe("completed"); + expect(wrapper.get('[data-testid="workbench-kafka-debug-counts"]').text()).toContain("2/2"); + expect(wrapper.get(".message-card").attributes("data-status")).toBe("completed"); + expect(wrapper.get(".message-card").text()).toContain("debug answer"); + expect(source.closed).toBe(true); + + await wrapper.get('[data-testid="workbench-kafka-debug-clear"]').trigger("click"); + expect(wrapper.find(".message-card").exists()).toBe(false); + expect(wrapper.get('[data-testid="workbench-kafka-debug-counts"]').text()).toContain("0/0"); + + await wrapper.get('[data-testid="workbench-kafka-debug-replay"]').trigger("click"); + const incompleteSource = FakeEventSource.instances[1]; + incompleteSource.emit("hwlab.kafka.connected", { consumerReady: true, debugIsolation: true, deliverySemantics: "debug-replay", liveOnly: false, replay: true }); + incompleteSource.emit("hwlab.kafka.replay-complete", { reason: "timeout", terminalObserved: false, count: 0 }); + await nextTick(); + expect(wrapper.get('[data-testid="workbench-kafka-debug-panel"]').attributes("data-replay-status")).toBe("incomplete"); + expect(wrapper.text()).toContain("重放未观测到 terminal(timeout)"); + expect(incompleteSource.closed).toBe(true); + }); +}); + +class FakeEventSource { + static instances: FakeEventSource[] = []; + readonly url: string; + closed = false; + onerror: ((event: Event) => void) | null = null; + onmessage: ((event: MessageEvent) => void) | null = null; + private listeners = new Map void>>(); + + constructor(url: string | URL) { + this.url = String(url); + FakeEventSource.instances.push(this); + } + + addEventListener(name: string, listener: (event: MessageEvent) => void) { + const listeners = this.listeners.get(name) ?? new Set(); + listeners.add(listener); + this.listeners.set(name, listeners); + } + + removeEventListener(name: string, listener: (event: MessageEvent) => void) { + this.listeners.get(name)?.delete(listener); + } + + close() { + this.closed = true; + } + + emit(name: string, payload: Record) { + const event = new MessageEvent(name, { data: JSON.stringify(payload) }); + for (const listener of this.listeners.get(name) ?? []) listener(event); + } +} + +function debugRecord(eventType: string, event: Record) { + return { + stream: "hwlab-debug", + topic: "hwlab.event.debug.v1", + offset: eventType === "terminal" ? "2" : "1", + value: { + schema: "hwlab.event.debug.v1", + eventType, + traceId: "trc_component_debug", + hwlabSessionId: "ses_component_debug", + event: { traceId: "trc_component_debug", sessionId: "ses_component_debug", sourceEventId: `evt_${eventType}`, ...event } + } + }; +} diff --git a/web/hwlab-cloud-web/src/components/workbench/WorkbenchKafkaDebugPanel.vue b/web/hwlab-cloud-web/src/components/workbench/WorkbenchKafkaDebugPanel.vue new file mode 100644 index 00000000..3dd7ba20 --- /dev/null +++ b/web/hwlab-cloud-web/src/components/workbench/WorkbenchKafkaDebugPanel.vue @@ -0,0 +1,314 @@ + + + + + diff --git a/web/hwlab-cloud-web/src/config/runtime.ts b/web/hwlab-cloud-web/src/config/runtime.ts index 673cb10c..95825f41 100644 --- a/web/hwlab-cloud-web/src/config/runtime.ts +++ b/web/hwlab-cloud-web/src/config/runtime.ts @@ -15,6 +15,10 @@ export interface WorkbenchRealtimeCapabilities { projectionRealtime: boolean; } +export interface WorkbenchDebugCapabilities { + isolatedKafka: boolean; +} + const DEFAULT_DISPLAY_DATE_TIME_OPTIONS: Intl.DateTimeFormatOptions = { year: "numeric", month: "2-digit", @@ -93,6 +97,14 @@ export function workbenchRealtimeCapabilities(): WorkbenchRealtimeCapabilities { return { liveKafkaSse: features.liveKafkaSse, projectionRealtime: features.projectionRealtime }; } +export function workbenchDebugCapabilities(): WorkbenchDebugCapabilities { + const features = window.HWLAB_CLOUD_WEB_CONFIG?.workbench?.debugCapabilities; + if (typeof features?.isolatedKafka !== "boolean") { + throw new Error("HWLAB Cloud Web workbench.debugCapabilities.isolatedKafka is required"); + } + return { isolatedKafka: features.isolatedKafka }; +} + 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-isolated-kafka-debug.test.ts b/web/hwlab-cloud-web/src/stores/workbench-isolated-kafka-debug.test.ts new file mode 100644 index 00000000..efc91299 --- /dev/null +++ b/web/hwlab-cloud-web/src/stores/workbench-isolated-kafka-debug.test.ts @@ -0,0 +1,109 @@ +import assert from "node:assert/strict"; +import { test } from "bun:test"; + +import type { ChatMessage } from "@/types"; +import { applyWorkbenchIsolatedKafkaDebugEvent, createWorkbenchIsolatedKafkaDebugState, workbenchCurrentDebugTraceId } from "./workbench-isolated-kafka-debug"; + +test("isolated Kafka debug projects assistant and terminal events through production reducers", () => { + const traceId = "trc_isolated_reducer"; + const assistant = debugEnvelope("41", { + schema: "hwlab.event.debug.v1", + eventId: "hwevt_assistant", + eventType: "assistant", + traceId, + hwlabSessionId: "ses_isolated_reducer", + event: { + sourceEventId: "agevt_assistant", + projectedSeq: 1, + traceId, + sessionId: "ses_isolated_reducer", + type: "assistant", + eventType: "assistant", + status: "running", + assistantText: "isolated partial response", + createdAt: "2026-07-10T10:00:01.000Z" + } + }); + const terminal = debugEnvelope("42", { + schema: "hwlab.event.debug.v1", + eventId: "hwevt_terminal", + eventType: "terminal", + traceId, + hwlabSessionId: "ses_isolated_reducer", + event: { + sourceEventId: "agevt_terminal", + projectedSeq: 2, + traceId, + sessionId: "ses_isolated_reducer", + type: "result", + eventType: "terminal", + terminal: true, + status: "completed", + createdAt: "2026-07-10T10:00:02.000Z", + elapsedMs: 1000 + } + }); + + const running = applyWorkbenchIsolatedKafkaDebugEvent(createWorkbenchIsolatedKafkaDebugState(), assistant, traceId); + assert.equal(running.appliedCount, 1); + assert.equal(running.message?.status, "running"); + assert.equal(running.message?.text, "isolated partial response"); + assert.equal(running.message?.runnerTrace?.eventSource, "hwlab-kafka-sse"); + assert.equal(running.message?.runnerTrace?.eventCount, 1); + + const completed = applyWorkbenchIsolatedKafkaDebugEvent(running, terminal, traceId); + assert.equal(completed.receivedCount, 2); + assert.equal(completed.appliedCount, 2); + assert.equal(completed.lastOffset, "42"); + assert.equal(completed.message?.status, "completed"); + assert.equal(completed.message?.text, "isolated partial response"); + assert.equal(completed.message?.traceAutoLifecycle, "terminal"); + assert.equal(completed.message?.finalResponse?.text, "isolated partial response"); + assert.equal(completed.message?.runnerTrace?.eventCount, 2); + assert.deepEqual(completed.message?.runnerTrace?.events?.map((event) => event.sourceEventId), ["agevt_assistant", "agevt_terminal"]); + assert.equal(completed.logs.every((entry) => entry.accepted), true); +}); + +test("isolated Kafka debug rejects another trace without touching a canonical message", () => { + const canonical: ChatMessage = Object.freeze({ + id: "msg_product", + role: "agent", + title: "Code Agent", + text: "product state", + status: "running", + traceId: "trc_product", + createdAt: "2026-07-10T10:00:00.000Z" + }) as ChatMessage; + const state = applyWorkbenchIsolatedKafkaDebugEvent(createWorkbenchIsolatedKafkaDebugState(), debugEnvelope("51", { + schema: "hwlab.event.debug.v1", + eventType: "assistant", + traceId: "trc_other", + event: { traceId: "trc_other", type: "assistant", text: "must stay isolated" } + }), "trc_expected"); + + assert.equal(state.message, null); + assert.equal(state.receivedCount, 1); + assert.equal(state.appliedCount, 0); + assert.equal(state.error, "debug-envelope-trace-mismatch"); + assert.equal(canonical.text, "product state"); +}); + +test("current Workbench debug trace prefers the active agent turn", () => { + const messages = [ + { id: "msg_old", role: "agent", title: "Code Agent", text: "done", status: "completed", traceId: "trc_old", createdAt: "2026-07-10T10:00:00.000Z" }, + { id: "msg_current", role: "agent", title: "Code Agent", text: "", status: "running", traceId: "trc_current", createdAt: "2026-07-10T10:01:00.000Z" } + ] as ChatMessage[]; + assert.equal(workbenchCurrentDebugTraceId(messages), "trc_current"); +}); + +function debugEnvelope(offset: string, value: Record) { + return { + ok: true, + stream: "hwlab-debug" as const, + topic: "hwlab.event.debug.v1", + partition: 0, + offset, + value, + valuesPrinted: false + }; +} 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 new file mode 100644 index 00000000..288a8870 --- /dev/null +++ b/web/hwlab-cloud-web/src/stores/workbench-isolated-kafka-debug.ts @@ -0,0 +1,156 @@ +// Responsibility: pure, isolated projection of hwlab.event.debug.v1 records through the production Workbench reducers. +// This module never reads or writes the canonical Workbench store. + +import type { WorkbenchKafkaSseDebugEvent } from "@/api/workbench-debug"; +import type { WorkbenchRealtimeEvent } from "@/api/workbench-events"; +import { mergeRunnerTrace, snapshotToRunnerTrace } from "@/composables/workbench-trace-snapshot"; +import type { ChatMessage, TraceEvent } from "@/types"; +import { firstNonEmptyString } from "@/utils"; +import { reduceWorkbenchRealtimeEvent } from "./workbench-event-reducer"; +import { reduceWorkbenchLiveKafkaMessageState, workbenchLiveKafkaMessageId } from "./workbench-live-kafka-event"; +import { planWorkbenchRealtimeApply } from "./workbench-realtime-plan"; +import { realtimeSnapshotToTraceSnapshot } from "./workbench-trace-detail"; + +export interface WorkbenchIsolatedKafkaDebugLog { + id: string; + eventType: string; + actionType: string; + stepTypes: string[]; + traceId: string | null; + accepted: boolean; + reason: string | null; + observedAt: string; +} + +export interface WorkbenchIsolatedKafkaDebugState { + message: ChatMessage | null; + receivedCount: number; + appliedCount: number; + lastOffset: string | null; + logs: WorkbenchIsolatedKafkaDebugLog[]; + error: string | null; +} + +export function createWorkbenchIsolatedKafkaDebugState(): WorkbenchIsolatedKafkaDebugState { + return { message: null, receivedCount: 0, appliedCount: 0, lastOffset: null, logs: [], error: null }; +} + +export function applyWorkbenchIsolatedKafkaDebugEvent( + state: WorkbenchIsolatedKafkaDebugState, + envelope: WorkbenchKafkaSseDebugEvent, + expectedTraceId: string | null | undefined +): WorkbenchIsolatedKafkaDebugState { + const value = recordValue(envelope.value); + const debugEvent = value as WorkbenchRealtimeEvent; + const traceId = realtimeTraceId(debugEvent); + const expected = firstNonEmptyString(expectedTraceId); + const receivedCount = state.receivedCount + 1; + const base = { ...state, receivedCount, lastOffset: firstNonEmptyString(envelope.offset, state.lastOffset) ?? null }; + if (debugEvent.schema !== "hwlab.event.debug.v1") return rejected(base, debugEvent, traceId, "debug-envelope-schema-mismatch"); + 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 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 = projectDebugMessage(state.message, traceId, event, traceStep.event); + 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[]): string | null { + const agentMessages = messages.filter((message) => message.role === "agent"); + const running = [...agentMessages].reverse().find((message) => ["pending", "running"].includes(String(message.status ?? "").trim().toLowerCase()) && firstNonEmptyString(message.traceId, message.runnerTrace?.traceId)); + return firstNonEmptyString(running?.traceId, running?.runnerTrace?.traceId, ...[...agentMessages].reverse().flatMap((message) => [message.traceId, message.runnerTrace?.traceId])) ?? null; +} + +function projectDebugMessage(previous: ChatMessage | null, traceId: string, realtimeEvent: WorkbenchRealtimeEvent, event: TraceEvent): ChatMessage { + const liveState = reduceWorkbenchLiveKafkaMessageState({ + text: previous?.text ?? "", + status: previous?.status ?? "running", + terminal: false + }, event); + const sessionId = firstNonEmptyString(realtimeEvent.hwlabSessionId, realtimeEvent.sessionId, event.sessionId, previous?.sessionId) ?? "ses_workbench_isolated_debug"; + const updatedAt = firstNonEmptyString(event.createdAt, event.updatedAt, realtimeEvent.eventCreatedAt, realtimeEvent.serverSentAt) ?? new Date().toISOString(); + const startedAt = firstNonEmptyString(previous?.timing?.startedAt, updatedAt) ?? updatedAt; + const snapshot = realtimeSnapshotToTraceSnapshot(traceId, { + traceId, + sessionId, + status: liveState.status, + events: [event], + eventCount: 1, + startedAt, + lastEventAt: updatedAt, + finishedAt: liveState.terminal ? updatedAt : null, + durationMs: finiteNumber(event.elapsedMs) + }, [event], { eventSource: "hwlab-kafka-sse" }); + const runnerTrace = mergeRunnerTrace(previous?.runnerTrace ?? null, snapshotToRunnerTrace(snapshot)); + const messageId = workbenchLiveKafkaMessageId(traceId); + return { + ...(previous ?? {}), + id: previous?.id ?? messageId, + messageId: previous?.messageId ?? messageId, + role: "agent", + title: "Code Agent · 隔离调试", + text: liveState.text, + status: liveState.status as ChatMessage["status"], + traceId, + turnId: traceId, + sessionId, + createdAt: previous?.createdAt ?? startedAt, + updatedAt, + runnerTrace, + traceAutoLifecycle: liveState.terminal ? "terminal" : "running", + timing: { + startedAt, + lastEventAt: updatedAt, + finishedAt: liveState.terminal ? updatedAt : null, + durationMs: liveState.terminal ? finiteNumber(event.elapsedMs) : null, + valuesRedacted: true + }, + ...(liveState.terminal && liveState.text ? { finalResponse: { text: liveState.text, status: liveState.status, traceId } } : {}) + } as ChatMessage; +} + +function rejected(state: WorkbenchIsolatedKafkaDebugState, event: WorkbenchRealtimeEvent, traceId: string | null, reason: string): WorkbenchIsolatedKafkaDebugState { + return withLog(state, event, traceId, "ignore", [], false, reason); +} + +function withLog( + state: WorkbenchIsolatedKafkaDebugState, + event: WorkbenchRealtimeEvent, + traceId: string | null, + actionType: string, + stepTypes: string[], + accepted: boolean, + reason: string | null +): WorkbenchIsolatedKafkaDebugState { + const observedAt = new Date().toISOString(); + const log: WorkbenchIsolatedKafkaDebugLog = { + id: `${observedAt}:${state.receivedCount}:${traceId ?? "none"}`, + eventType: firstNonEmptyString(event.eventType, event.event?.type, event.type, "unknown") ?? "unknown", + actionType, + stepTypes, + traceId, + accepted, + reason, + observedAt + }; + return { ...state, error: accepted ? null : reason, logs: [log, ...state.logs].slice(0, 40) }; +} + +function realtimeTraceId(event: WorkbenchRealtimeEvent): string | null { + return firstNonEmptyString(event.traceId, event.event?.traceId) ?? null; +} + +function recordValue(value: unknown): Record { + return value && typeof value === "object" && !Array.isArray(value) ? value as Record : {}; +} + +function finiteNumber(value: unknown): number | null { + const number = Number(value); + return Number.isFinite(number) && number >= 0 ? Math.trunc(number) : null; +} diff --git a/web/hwlab-cloud-web/src/types/global.d.ts b/web/hwlab-cloud-web/src/types/global.d.ts index 8bbded04..7336bbb3 100644 --- a/web/hwlab-cloud-web/src/types/global.d.ts +++ b/web/hwlab-cloud-web/src/types/global.d.ts @@ -12,6 +12,9 @@ declare global { liveKafkaSse?: boolean; projectionRealtime?: boolean; }; + debugCapabilities?: { + isolatedKafka?: boolean; + }; traceTimeline?: { autoExpandRunning?: boolean; autoCollapseTerminal?: boolean; diff --git a/web/hwlab-cloud-web/src/views/workbench/CodeWorkbenchView.vue b/web/hwlab-cloud-web/src/views/workbench/CodeWorkbenchView.vue index 614e39b2..328ccd83 100644 --- a/web/hwlab-cloud-web/src/views/workbench/CodeWorkbenchView.vue +++ b/web/hwlab-cloud-web/src/views/workbench/CodeWorkbenchView.vue @@ -7,15 +7,21 @@ import { useRoute, useRouter } from "vue-router"; import CommandComposer from "@/components/workbench/CommandComposer.vue"; import ConversationPanel from "@/components/workbench/ConversationPanel.vue"; import SessionRail from "@/components/workbench/SessionRail.vue"; +import WorkbenchKafkaDebugPanel from "@/components/workbench/WorkbenchKafkaDebugPanel.vue"; import CaseRunPanel from "@/components/caserun/CaseRunPanel.vue"; import HwpodNodeOpsPanel from "@/components/hwpod/HwpodNodeOpsPanel.vue"; import { useAutoRefresh } from "@/composables/useAutoRefresh"; import { shouldReflectWorkbenchSessionUrl } from "@/router/workbench-navigation"; +import { workbenchDebugCapabilities } from "@/config/runtime"; +import { useAuthStore } from "@/stores/auth"; import { useWorkbenchStore } from "@/stores/workbench"; +import { workbenchCurrentDebugTraceId } from "@/stores/workbench-isolated-kafka-debug"; import { normalizeWorkbenchSessionId, normalizeWorkbenchSessionRouteId } from "@/utils"; import { finishWorkbenchOpenFullLoad, startWorkbenchOpenJourney } from "@/utils/workbench-performance"; const workbench = useWorkbenchStore(); +const auth = useAuthStore(); +const debugCapabilities = workbenchDebugCapabilities(); const route = useRoute(); const router = useRouter(); const routeRequestId = computed(() => normalizeWorkbenchSessionRouteId(route.params.sessionId)); @@ -23,6 +29,9 @@ const rawRouteSessionId = computed(() => typeof route.params.sessionId === "stri const routeReflectionTarget = computed(() => routeRequestId.value ?? (rawRouteSessionId.value || null)); const applyingRouteSession = ref(Boolean(rawRouteSessionId.value)); const componentActive = ref(false); +const isolatedDebugOpen = ref(false); +const isolatedDebugAvailable = computed(() => debugCapabilities.isolatedKafka && auth.isAdmin); +const isolatedDebugTraceId = computed(() => workbenchCurrentDebugTraceId(workbench.activeMessages)); type LaunchContextRecord = { sourceId?: string | null; hwpodId?: string | null; @@ -64,6 +73,7 @@ onBeforeUnmount(() => { }); watch(() => route.params.sessionId, () => void applyRouteSession()); watch(() => workbench.activeSessionId, (sessionId) => void reflectActiveSessionInUrl(sessionId)); +watch(isolatedDebugAvailable, (available) => { if (!available) isolatedDebugOpen.value = false; }); useAutoRefresh(refreshLiveWhenIdle, 120_000, { immediate: false, initialDelayMs: 90_000 }); async function refreshLiveWhenIdle(): Promise { @@ -117,7 +127,21 @@ async function reflectActiveSessionInUrl(value: string | null): Promise { MDTODO: {{ workbenchLaunchContext.mdtodoRootRef }} Workspace: {{ workbenchLaunchContext.workspaceRootLabel }}
- +
+
+ + {{ isolatedDebugTraceId ? `当前 trace ${isolatedDebugTraceId}` : "等待当前 traceId" }} +
+ + +
@@ -127,3 +151,57 @@ async function reflectActiveSessionInUrl(value: string | null): Promise {
+ + diff --git a/web/hwlab-cloud-web/src/views/workbench/WorkbenchDebugView.vue b/web/hwlab-cloud-web/src/views/workbench/WorkbenchDebugView.vue index 68fa1a13..03ea8ca9 100644 --- a/web/hwlab-cloud-web/src/views/workbench/WorkbenchDebugView.vue +++ b/web/hwlab-cloud-web/src/views/workbench/WorkbenchDebugView.vue @@ -10,6 +10,7 @@ import type { WorkbenchRealtimeEvent } from "@/api/workbench-events"; import { applyWorkbenchDebugFakeSseEvent, createWorkbenchDebugFakeSseState } from "@/stores/workbench-debug-fake-sse"; const queueId = "trace-card"; +type CanonicalKafkaStreamName = Exclude; const activeTab = ref<"fake" | "kafka">("fake"); const selectedSequenceId = ref("trace-card-basic"); const sequences = ref([]); @@ -17,22 +18,22 @@ const queue = ref(null); const projection = ref(createWorkbenchDebugFakeSseState()); const streamStatus = ref<"connecting" | "open" | "error" | "closed">("connecting"); const kafkaStatus = ref<"idle" | "connecting" | "open" | "error" | "closed">("idle"); -const kafkaStreamName = ref("hwlab"); +const kafkaStreamName = ref("hwlab"); const kafkaFromBeginning = ref(false); const kafkaFilters = ref({ traceId: "", sessionId: "", runId: "", commandId: "" }); const kafkaEvents = ref([]); const kafkaConnection = ref(null); -const kafkaTraceStatuses = ref>(createKafkaStatusMap()); -const kafkaTraceEvents = ref>(createKafkaEventMap()); -const kafkaTraceConnections = ref>(createKafkaConnectionMap()); +const kafkaTraceStatuses = ref>(createKafkaStatusMap()); +const kafkaTraceEvents = ref>(createKafkaEventMap()); +const kafkaTraceConnections = ref>(createKafkaConnectionMap()); const busy = ref(false); const error = ref(null); const appendDraft = ref(defaultAppendDraft()); let stream: WorkbenchDebugFakeSseStream | null = null; let kafkaStream: WorkbenchKafkaSseDebugStream | null = null; -const kafkaTraceStreams: Partial> = {}; +const kafkaTraceStreams: Partial> = {}; -const KAFKA_TRACE_STREAMS: Array<{ name: WorkbenchKafkaSseDebugStreamName; label: string; topic: string }> = [ +const KAFKA_TRACE_STREAMS: Array<{ name: CanonicalKafkaStreamName; label: string; topic: string }> = [ { name: "stdio", label: "codex stdio", topic: "codex-stdio.raw.v1" }, { name: "agentrun", label: "AgentRun event", topic: "agentrun.event.v1" }, { name: "hwlab", label: "HWLAB event", topic: "hwlab.event.v1" } @@ -174,30 +175,30 @@ function closeKafkaTraceStreams(): void { kafkaTraceStreams[streamName]?.close(); delete kafkaTraceStreams[streamName]; } - kafkaTraceStatuses.value = Object.fromEntries(KAFKA_TRACE_STREAMS.map((entry) => [entry.name, kafkaTraceStatuses.value[entry.name] === "idle" ? "idle" : "closed"])) as Record; + kafkaTraceStatuses.value = Object.fromEntries(KAFKA_TRACE_STREAMS.map((entry) => [entry.name, kafkaTraceStatuses.value[entry.name] === "idle" ? "idle" : "closed"])) as Record; } function clearKafkaTraceEvents(): void { kafkaTraceEvents.value = createKafkaEventMap(); } -function setKafkaTraceStatus(streamName: WorkbenchKafkaSseDebugStreamName, status: KafkaStreamStatus): void { +function setKafkaTraceStatus(streamName: CanonicalKafkaStreamName, status: KafkaStreamStatus): void { kafkaTraceStatuses.value = { ...kafkaTraceStatuses.value, [streamName]: status }; } -function pushKafkaTraceEvent(streamName: WorkbenchKafkaSseDebugStreamName, row: KafkaRawRow): void { +function pushKafkaTraceEvent(streamName: CanonicalKafkaStreamName, row: KafkaRawRow): void { kafkaTraceEvents.value = { ...kafkaTraceEvents.value, [streamName]: [row, ...kafkaTraceEvents.value[streamName]].slice(0, 80) }; } -function createKafkaStatusMap(status: KafkaStreamStatus = "idle"): Record { +function createKafkaStatusMap(status: KafkaStreamStatus = "idle"): Record { return { stdio: status, agentrun: status, hwlab: status }; } -function createKafkaEventMap(): Record { +function createKafkaEventMap(): Record { return { stdio: [], agentrun: [], hwlab: [] }; } -function createKafkaConnectionMap(): Record { +function createKafkaConnectionMap(): Record { return { stdio: null, agentrun: null, hwlab: null }; }