From ff0aa2b1d1fabcae3ad368dc75b1b08231172a34 Mon Sep 17 00:00:00 2001 From: root Date: Thu, 9 Jul 2026 21:07:31 +0200 Subject: [PATCH] fix: resolve kafka debug trace filters --- .../cloud/workbench-kafka-sse-debug.test.ts | 43 +++++++++++++++++++ internal/cloud/workbench-kafka-sse-debug.ts | 43 ++++++++++++++++++- 2 files changed, 84 insertions(+), 2 deletions(-) diff --git a/internal/cloud/workbench-kafka-sse-debug.test.ts b/internal/cloud/workbench-kafka-sse-debug.test.ts index b46f2c27..62f54c46 100644 --- a/internal/cloud/workbench-kafka-sse-debug.test.ts +++ b/internal/cloud/workbench-kafka-sse-debug.test.ts @@ -72,6 +72,49 @@ test("workbench Kafka SSE debug endpoint streams filtered codex stdio Kafka even } }); +test("workbench Kafka SSE debug endpoint resolves traceId to AgentRun Kafka ids", async () => { + const fakeKafka = createFakeKafkaFactory(); + const traceStore = { + snapshot: () => ({ + traceId: "trc_resolved_kafka_debug", + events: [ + { type: "backend", runId: "run_resolved_kafka_debug", commandId: "cmd_resolved_kafka_debug", sessionId: "ses_agentrun_resolved_kafka_debug" } + ], + lastEvent: null + }) + }; + const env = { + HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092", + HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api", + HWLAB_KAFKA_EVENT_TOPIC: "hwlab.event.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, traceStore, logger: null }); + }); + await listen(server); + const abort = new AbortController(); + try { + const response = await fetch(`${serverUrl(server)}/v1/workbench/debug/kafka-sse/events?stream=hwlab&traceId=trc_resolved_kafka_debug`, { 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, "resolvedFilters"); + assert.match(connected, /trc_resolved_kafka_debug/u); + assert.match(connected, /run_resolved_kafka_debug/u); + + await fakeKafka.emit({ eventType: "hwlab.trace.event.projected", traceId: null, sessionId: null, context: { runId: "run_other", commandId: "cmd_other" }, event: { label: "ignored" } }); + await fakeKafka.emit({ eventType: "hwlab.trace.event.projected", traceId: null, sessionId: null, context: { runId: "run_resolved_kafka_debug", commandId: "cmd_resolved_kafka_debug" }, event: { label: "resolved", status: "running" } }); + const streamed = await readUntil(reader, "run_resolved_kafka_debug"); + assert.match(streamed, /hwlab\.kafka\.event/u); + assert.match(streamed, /resolved/u); + assert.doesNotMatch(streamed, /run_other/u); + } finally { + abort.abort(); + await close(server); + } +}); + function createFakeKafkaFactory() { let eachMessage: ((input: any) => Promise) | null = null; let subscribedTopic = "hwlab.event.v1"; diff --git a/internal/cloud/workbench-kafka-sse-debug.ts b/internal/cloud/workbench-kafka-sse-debug.ts index 977e93da..ca7df2b7 100644 --- a/internal/cloud/workbench-kafka-sse-debug.ts +++ b/internal/cloud/workbench-kafka-sse-debug.ts @@ -4,6 +4,7 @@ */ import { sendJson } from "./server-http-utils.ts"; import { openKafkaEventStream } from "./kafka-event-bridge.ts"; +import { defaultCodeAgentTraceStore } from "./code-agent-trace-store.ts"; const CONTRACT_VERSION = "workbench-debug-kafka-sse-v1"; const DEFAULT_STREAM = "hwlab"; @@ -36,12 +37,15 @@ export async function handleWorkbenchKafkaSseDebugHttp(request, response, url, o function describeKafkaSseDebug(url, env) { const stream = streamFromUrl(url); + const requestedFilters = filtersFromUrl(url); + const resolvedFilters = resolveKafkaDebugFilters(requestedFilters, { traceStore: defaultCodeAgentTraceStore }); return { ok: true, contractVersion: CONTRACT_VERSION, stream, topic: topicForStream(stream, env), - filters: filtersFromUrl(url), + filters: requestedFilters, + resolvedFilters, eventsRoute: `/v1/workbench/debug/kafka-sse/events${url.search || ""}`, valuesPrinted: false }; @@ -51,6 +55,7 @@ async function openKafkaDebugSse(request, response, url, options) { const env = options.env ?? process.env; const stream = streamFromUrl(url); const filters = filtersFromUrl(url); + const resolvedFilters = resolveKafkaDebugFilters(filters, options); response.writeHead(200, { "content-type": "text/event-stream; charset=utf-8", "cache-control": "no-store, no-transform", @@ -65,6 +70,7 @@ async function openKafkaDebugSse(request, response, url, options) { stream, topic: topicForStream(stream, env), filters, + resolvedFilters, serverSentAt: new Date().toISOString(), valuesPrinted: false }); @@ -81,7 +87,7 @@ async function openKafkaDebugSse(request, response, url, options) { env, stream, fromBeginning: url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true", - ...filters, + ...resolvedFilters, kafkaFactory: options.kafkaFactory, onEvent: async (record) => { writeSse(response, "hwlab.kafka.event", { @@ -105,6 +111,39 @@ async function openKafkaDebugSse(request, response, url, options) { } } +function resolveKafkaDebugFilters(filters = {}, options = {}) { + const traceId = textValue(filters.traceId); + if (!traceId) return filters; + const traceStore = options.traceStore ?? defaultCodeAgentTraceStore; + const snapshot = typeof traceStore?.snapshot === "function" ? traceStore.snapshot(traceId) : null; + const resolved = collectTraceLinkedIds(snapshot); + const hasResolvedKafkaKey = Boolean(resolved.runId || resolved.commandId || resolved.sessionId); + if (!hasResolvedKafkaKey) return filters; + const resolvedSessionId = resolved.runId || resolved.commandId ? null : resolved.sessionId; + return compactObject({ + sessionId: filters.sessionId || resolvedSessionId, + runId: filters.runId || resolved.runId, + commandId: filters.commandId || resolved.commandId + }); +} + +function collectTraceLinkedIds(snapshot) { + const out = { sessionId: null, runId: null, commandId: null }; + const events = Array.isArray(snapshot?.events) ? snapshot.events : []; + for (const event of events) collectIdsFromRecord(event, out); + collectIdsFromRecord(snapshot?.lastEvent, out); + return out; +} + +function collectIdsFromRecord(record, out) { + const value = record && typeof record === "object" && !Array.isArray(record) ? record : {}; + const agentRun = value.agentRun && typeof value.agentRun === "object" && !Array.isArray(value.agentRun) ? value.agentRun : {}; + const payload = value.payload && typeof value.payload === "object" && !Array.isArray(value.payload) ? value.payload : {}; + out.sessionId ||= textValue(value.sessionId ?? value.sourceSessionId ?? agentRun.sessionId ?? payload.sessionId); + out.runId ||= textValue(value.runId ?? value.sourceRunId ?? agentRun.runId ?? payload.runId); + out.commandId ||= textValue(value.commandId ?? value.sourceCommandId ?? agentRun.commandId ?? payload.commandId); +} + function writeSse(response, eventName, payload) { if (response.destroyed || response.writableEnded) return false; try {