diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index 213921d4..4e2d7069 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -8,6 +8,7 @@ import { Kafka, logLevel } from "kafkajs"; const TRUE_VALUES = new Set(["1", "true", "yes", "on"]); const DEFAULT_AGENTRUN_EVENT_TOPIC = "agentrun.event.v1"; const DEFAULT_HWLAB_EVENT_TOPIC = "hwlab.event.v1"; +const DEFAULT_STDIO_TOPIC = "codex-stdio.raw.v1"; const DEFAULT_CLIENT_ID = "hwlab-v03-cloud-api"; const DEFAULT_GROUP_ID = "hwlab-v03-agentrun-event-bridge"; const DEFAULT_QUERY_TIMEOUT_MS = 5000; @@ -305,17 +306,24 @@ function eventFieldCandidates(value, name) { const context = objectValue(value.context); const run = objectValue(value.run); const command = objectValue(value.command); + const payload = objectValue(value.payload); + const metadata = objectValue(value.metadata ?? value.meta); + const trace = objectValue(value.trace); + const ids = objectValue(value.ids); + const stdio = objectValue(value.stdio ?? value.frame ?? value.codexStdio); const sourceEvent = objectValue(value.sourceEvent); - const nestedEvent = objectValue(value.event?.sourceEvent ?? value.agentRunEvent); - if (name === "traceId") return textCandidates(value.traceId, event.traceId, context.traceId, sourceEvent.traceId, nestedEvent.traceId, nestedEvent.payload?.traceId); - if (name === "sessionId") return textCandidates(value.sessionId, event.sessionId, context.sessionId, run.sessionId, nestedEvent.payload?.sessionId); - if (name === "runId") return textCandidates(value.runId, event.runId, context.runId, run.runId, nestedEvent.runId, nestedEvent.payload?.runId); - if (name === "commandId") return textCandidates(value.commandId, event.commandId, context.commandId, command.commandId, nestedEvent.commandId, nestedEvent.payload?.commandId); + const nestedEvent = objectValue(value.event?.sourceEvent ?? value.agentRunEvent ?? payload.event); + const nestedPayload = objectValue(nestedEvent.payload); + if (name === "traceId") return textCandidates(value.traceId, value.trace_id, event.traceId, event.trace_id, context.traceId, context.trace_id, sourceEvent.traceId, sourceEvent.trace_id, nestedEvent.traceId, nestedEvent.trace_id, nestedPayload.traceId, nestedPayload.trace_id, payload.traceId, payload.trace_id, metadata.traceId, metadata.trace_id, trace.traceId, trace.trace_id, ids.traceId, ids.trace_id, stdio.traceId, stdio.trace_id); + if (name === "sessionId") return textCandidates(value.sessionId, value.session_id, event.sessionId, event.session_id, context.sessionId, context.session_id, run.sessionId, run.session_id, nestedPayload.sessionId, nestedPayload.session_id, payload.sessionId, payload.session_id, metadata.sessionId, metadata.session_id, ids.sessionId, ids.session_id, stdio.sessionId, stdio.session_id); + if (name === "runId") return textCandidates(value.runId, value.run_id, event.runId, event.run_id, context.runId, context.run_id, run.runId, run.run_id, nestedEvent.runId, nestedEvent.run_id, nestedPayload.runId, nestedPayload.run_id, payload.runId, payload.run_id, metadata.runId, metadata.run_id, ids.runId, ids.run_id, stdio.runId, stdio.run_id); + if (name === "commandId") return textCandidates(value.commandId, value.command_id, event.commandId, event.command_id, context.commandId, context.command_id, command.commandId, command.command_id, nestedEvent.commandId, nestedEvent.command_id, nestedPayload.commandId, nestedPayload.command_id, payload.commandId, payload.command_id, metadata.commandId, metadata.command_id, ids.commandId, ids.command_id, stdio.commandId, stdio.command_id); return []; } function kafkaTopicForStream(stream, env) { const normalized = stringValue(stream) || "hwlab"; + if (normalized === "stdio" || normalized === "codex-stdio") return stringValue(env.HWLAB_KAFKA_STDIO_TOPIC ?? env.AGENTRUN_KAFKA_STDIO_TOPIC) || DEFAULT_STDIO_TOPIC; if (normalized === "agentrun") return stringValue(env.HWLAB_KAFKA_AGENTRUN_EVENT_TOPIC) || DEFAULT_AGENTRUN_EVENT_TOPIC; if (normalized === "hwlab") return stringValue(env.HWLAB_KAFKA_EVENT_TOPIC) || DEFAULT_HWLAB_EVENT_TOPIC; return normalized; diff --git a/internal/cloud/workbench-kafka-sse-debug.test.ts b/internal/cloud/workbench-kafka-sse-debug.test.ts index 3d300605..b46f2c27 100644 --- a/internal/cloud/workbench-kafka-sse-debug.test.ts +++ b/internal/cloud/workbench-kafka-sse-debug.test.ts @@ -38,11 +38,46 @@ test("workbench Kafka SSE debug endpoint streams filtered raw HWLAB Kafka events } }); +test("workbench Kafka SSE debug endpoint streams filtered codex stdio Kafka events", async () => { + const fakeKafka = createFakeKafkaFactory(); + const env = { + HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092", + HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api", + HWLAB_KAFKA_STDIO_TOPIC: "codex-stdio.raw.v1" + }; + const server = createServer((request, response) => { + const url = new URL(request.url ?? "/", "http://127.0.0.1"); + void handleWorkbenchKafkaSseDebugHttp(request, response, url, { env, kafkaFactory: fakeKafka.factory, logger: null }); + }); + await listen(server); + const abort = new AbortController(); + try { + const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=stdio&traceId=trc_stdio_sse_debug&fromBeginning=true`, { signal: abort.signal }); + assert.equal(response.status, 200); + const reader = response.body?.getReader(); + assert.ok(reader, "SSE response must expose a readable stream"); + const connected = await readUntil(reader, "hwlab.kafka.connected"); + assert.match(connected, /codex-stdio\.raw\.v1/u); + + await fakeKafka.emit({ traceId: "trc_other", stream: "stdout", line: "ignored" }); + await fakeKafka.emit({ trace_id: "trc_stdio_sse_debug", metadata: { session_id: "ses_stdio_sse_debug", run_id: "run_stdio_sse_debug" }, stream: "stdout", line: "visible stdio frame" }); + const streamed = await readUntil(reader, "visible stdio frame"); + assert.match(streamed, /hwlab\.kafka\.event/u); + assert.match(streamed, /codex-stdio\.raw\.v1/u); + assert.match(streamed, /ses_stdio_sse_debug/u); + assert.doesNotMatch(streamed, /trc_other/u); + } finally { + abort.abort(); + await close(server); + } +}); + function createFakeKafkaFactory() { let eachMessage: ((input: any) => Promise) | null = null; + let subscribedTopic = "hwlab.event.v1"; const consumer = { connect: async () => undefined, - subscribe: async () => undefined, + subscribe: async (input: { topic?: string }) => { subscribedTopic = String(input.topic ?? subscribedTopic); }, run: async (input: any) => { eachMessage = input.eachMessage; }, stop: async () => undefined, disconnect: async () => undefined @@ -52,7 +87,7 @@ function createFakeKafkaFactory() { emit: async (value: Record) => { assert.ok(eachMessage, "consumer.run must be called before emitting fake Kafka events"); await eachMessage({ - topic: "hwlab.event.v1", + topic: subscribedTopic, partition: 0, message: { offset: String(value.sessionId === "ses_other" ? 1 : 2), diff --git a/internal/cloud/workbench-kafka-sse-debug.ts b/internal/cloud/workbench-kafka-sse-debug.ts index 39107994..977e93da 100644 --- a/internal/cloud/workbench-kafka-sse-debug.ts +++ b/internal/cloud/workbench-kafka-sse-debug.ts @@ -7,7 +7,7 @@ import { openKafkaEventStream } from "./kafka-event-bridge.ts"; const CONTRACT_VERSION = "workbench-debug-kafka-sse-v1"; const DEFAULT_STREAM = "hwlab"; -const ALLOWED_STREAMS = new Set(["hwlab", "agentrun"]); +const ALLOWED_STREAMS = new Set(["stdio", "agentrun", "hwlab"]); export async function handleWorkbenchKafkaSseDebugHttp(request, response, url, options = {}) { const route = routeSuffix(url.pathname); @@ -137,6 +137,7 @@ 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"; return textValue(env.HWLAB_KAFKA_EVENT_TOPIC) || "hwlab.event.v1"; } diff --git a/web/hwlab-cloud-web/src/api/workbench-debug.ts b/web/hwlab-cloud-web/src/api/workbench-debug.ts index 00ea5af8..68134dc8 100644 --- a/web/hwlab-cloud-web/src/api/workbench-debug.ts +++ b/web/hwlab-cloud-web/src/api/workbench-debug.ts @@ -57,6 +57,8 @@ export interface WorkbenchKafkaSseDebugFilters { commandId?: string; } +export type WorkbenchKafkaSseDebugStreamName = "stdio" | "agentrun" | "hwlab"; + export interface WorkbenchKafkaSseDebugEvent { ok?: boolean; contractVersion?: string; @@ -74,7 +76,7 @@ export interface WorkbenchKafkaSseDebugEvent { } export interface WorkbenchKafkaSseDebugStreamOptions extends WorkbenchKafkaSseDebugFilters { - stream?: "hwlab" | "agentrun"; + stream?: WorkbenchKafkaSseDebugStreamName; fromBeginning?: boolean; onEvent: (event: WorkbenchKafkaSseDebugEvent, eventName: string) => void; onOpen?: () => void; diff --git a/web/hwlab-cloud-web/src/views/workbench/WorkbenchDebugView.vue b/web/hwlab-cloud-web/src/views/workbench/WorkbenchDebugView.vue index 0da7ecda..11687892 100644 --- a/web/hwlab-cloud-web/src/views/workbench/WorkbenchDebugView.vue +++ b/web/hwlab-cloud-web/src/views/workbench/WorkbenchDebugView.vue @@ -3,7 +3,7 @@ // Responsibility: Workbench debug route for backend-driven fake SSE single-step Trace card inspection. import { computed, onBeforeUnmount, onMounted, ref } from "vue"; -import { connectWorkbenchDebugFakeSse, connectWorkbenchKafkaSseDebug, workbenchDebugAPI, type WorkbenchDebugFakeSseQueue, type WorkbenchDebugFakeSseSequence, type WorkbenchDebugFakeSseStream, type WorkbenchKafkaSseDebugEvent } from "@/api"; +import { connectWorkbenchDebugFakeSse, connectWorkbenchKafkaSseDebug, workbenchDebugAPI, type WorkbenchDebugFakeSseQueue, type WorkbenchDebugFakeSseSequence, type WorkbenchDebugFakeSseStream, type WorkbenchKafkaSseDebugEvent, type WorkbenchKafkaSseDebugStreamName } from "@/api"; import StatusBadge from "@/components/common/StatusBadge.vue"; import WorkbenchMessageCard from "@/components/workbench/WorkbenchMessageCard.vue"; import type { WorkbenchRealtimeEvent } from "@/api/workbench-events"; @@ -17,15 +17,26 @@ 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" | "agentrun">("hwlab"); +const kafkaStreamName = ref("hwlab"); const kafkaFromBeginning = ref(false); const kafkaFilters = ref({ traceId: "", sessionId: "", runId: "", commandId: "" }); const kafkaEvents = ref([]); +const kafkaTraceStatuses = ref>(createKafkaStatusMap()); +const kafkaTraceEvents = ref>(createKafkaEventMap()); const busy = ref(false); const error = ref(null); const appendDraft = ref(defaultAppendDraft()); let stream: WorkbenchDebugFakeSseStream | null = null; let kafkaStream: WorkbenchDebugFakeSseStream | null = null; +const kafkaTraceStreams: Partial> = {}; + +const KAFKA_TRACE_STREAMS: Array<{ name: WorkbenchKafkaSseDebugStreamName; 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" } +]; + +type KafkaStreamStatus = "idle" | "connecting" | "open" | "error" | "closed"; interface KafkaRawRow { id: string; @@ -39,6 +50,8 @@ const logs = computed(() => projection.value.logs); const queueSummary = computed(() => queue.value ? `${queue.value.cursor}/${queue.value.eventCount}` : "-"); const streamBadgeStatus = computed(() => streamStatus.value === "open" ? "completed" : streamStatus.value === "error" ? "failed" : "running"); const kafkaBadgeStatus = computed(() => kafkaStatus.value === "open" ? "completed" : kafkaStatus.value === "error" ? "failed" : "running"); +const kafkaTraceEventCount = computed(() => Object.values(kafkaTraceEvents.value).reduce((total, rows) => total + rows.length, 0)); +const kafkaTraceConnectedCount = computed(() => Object.values(kafkaTraceStatuses.value).filter((status) => status === "open").length); const canStep = computed(() => !busy.value && (queue.value?.remaining ?? 0) > 0); onMounted(async () => { @@ -52,6 +65,7 @@ onBeforeUnmount(() => { stream = null; kafkaStream?.close(); kafkaStream = null; + closeKafkaTraceStreams(); streamStatus.value = "closed"; kafkaStatus.value = "closed"; }); @@ -117,6 +131,62 @@ function clearKafkaEvents(): void { kafkaEvents.value = []; } +function openKafkaTraceStreams(): void { + closeKafkaTraceStreams(); + kafkaTraceEvents.value = createKafkaEventMap(); + kafkaTraceStatuses.value = createKafkaStatusMap("connecting"); + for (const streamName of KAFKA_TRACE_STREAMS.map((entry) => entry.name)) { + const opened = connectWorkbenchKafkaSseDebug({ + stream: streamName, + fromBeginning: true, + traceId: cleanFilter(kafkaFilters.value.traceId), + sessionId: cleanFilter(kafkaFilters.value.sessionId), + runId: cleanFilter(kafkaFilters.value.runId), + commandId: cleanFilter(kafkaFilters.value.commandId), + onOpen: () => setKafkaTraceStatus(streamName, "open"), + onError: () => setKafkaTraceStatus(streamName, "error"), + onEvent: (event, eventName) => { + if (eventName === "hwlab.kafka.connected") return; + if (eventName === "hwlab.kafka.error") { + setKafkaTraceStatus(streamName, "error"); + error.value = event.error?.message || "Kafka SSE 三流订阅错误"; + } + pushKafkaTraceEvent(streamName, { id: kafkaRowId(event, eventName), eventName, receivedAt: new Date().toISOString(), payload: event }); + } + }); + if (opened) kafkaTraceStreams[streamName] = opened; + else setKafkaTraceStatus(streamName, "error"); + } +} + +function closeKafkaTraceStreams(): void { + for (const streamName of KAFKA_TRACE_STREAMS.map((entry) => entry.name)) { + 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; +} + +function clearKafkaTraceEvents(): void { + kafkaTraceEvents.value = createKafkaEventMap(); +} + +function setKafkaTraceStatus(streamName: WorkbenchKafkaSseDebugStreamName, status: KafkaStreamStatus): void { + kafkaTraceStatuses.value = { ...kafkaTraceStatuses.value, [streamName]: status }; +} + +function pushKafkaTraceEvent(streamName: WorkbenchKafkaSseDebugStreamName, row: KafkaRawRow): void { + kafkaTraceEvents.value = { ...kafkaTraceEvents.value, [streamName]: [row, ...kafkaTraceEvents.value[streamName]].slice(0, 80) }; +} + +function createKafkaStatusMap(status: KafkaStreamStatus = "idle"): Record { + return { stdio: status, agentrun: status, hwlab: status }; +} + +function createKafkaEventMap(): Record { + return { stdio: [], agentrun: [], hwlab: [] }; +} + async function resetQueue(): Promise { busy.value = true; error.value = null; @@ -237,7 +307,7 @@ function defaultAppendDraft(): string {
- {{ activeTab === 'fake' ? `Queue ${queueSummary}` : `${kafkaEvents.length} events` }} + {{ activeTab === 'fake' ? `Queue ${queueSummary}` : `${kafkaTraceConnectedCount}/3 streams · ${kafkaTraceEventCount + kafkaEvents.length} events` }}
@@ -262,6 +332,7 @@ function defaultAppendDraft(): string {