Merge pull request #2456 from pikasTech/fix/kafka-trace-filter-resolution
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success

fix: resolve kafka debug trace filters
This commit is contained in:
Lyon
2026-07-10 03:08:00 +08:00
committed by GitHub
2 changed files with 84 additions and 2 deletions
@@ -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() { function createFakeKafkaFactory() {
let eachMessage: ((input: any) => Promise<void>) | null = null; let eachMessage: ((input: any) => Promise<void>) | null = null;
let subscribedTopic = "hwlab.event.v1"; let subscribedTopic = "hwlab.event.v1";
+41 -2
View File
@@ -4,6 +4,7 @@
*/ */
import { sendJson } from "./server-http-utils.ts"; import { sendJson } from "./server-http-utils.ts";
import { openKafkaEventStream } from "./kafka-event-bridge.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 CONTRACT_VERSION = "workbench-debug-kafka-sse-v1";
const DEFAULT_STREAM = "hwlab"; const DEFAULT_STREAM = "hwlab";
@@ -36,12 +37,15 @@ export async function handleWorkbenchKafkaSseDebugHttp(request, response, url, o
function describeKafkaSseDebug(url, env) { function describeKafkaSseDebug(url, env) {
const stream = streamFromUrl(url); const stream = streamFromUrl(url);
const requestedFilters = filtersFromUrl(url);
const resolvedFilters = resolveKafkaDebugFilters(requestedFilters, { traceStore: defaultCodeAgentTraceStore });
return { return {
ok: true, ok: true,
contractVersion: CONTRACT_VERSION, contractVersion: CONTRACT_VERSION,
stream, stream,
topic: topicForStream(stream, env), topic: topicForStream(stream, env),
filters: filtersFromUrl(url), filters: requestedFilters,
resolvedFilters,
eventsRoute: `/v1/workbench/debug/kafka-sse/events${url.search || ""}`, eventsRoute: `/v1/workbench/debug/kafka-sse/events${url.search || ""}`,
valuesPrinted: false valuesPrinted: false
}; };
@@ -51,6 +55,7 @@ async function openKafkaDebugSse(request, response, url, options) {
const env = options.env ?? process.env; const env = options.env ?? process.env;
const stream = streamFromUrl(url); const stream = streamFromUrl(url);
const filters = filtersFromUrl(url); const filters = filtersFromUrl(url);
const resolvedFilters = resolveKafkaDebugFilters(filters, options);
response.writeHead(200, { response.writeHead(200, {
"content-type": "text/event-stream; charset=utf-8", "content-type": "text/event-stream; charset=utf-8",
"cache-control": "no-store, no-transform", "cache-control": "no-store, no-transform",
@@ -65,6 +70,7 @@ async function openKafkaDebugSse(request, response, url, options) {
stream, stream,
topic: topicForStream(stream, env), topic: topicForStream(stream, env),
filters, filters,
resolvedFilters,
serverSentAt: new Date().toISOString(), serverSentAt: new Date().toISOString(),
valuesPrinted: false valuesPrinted: false
}); });
@@ -81,7 +87,7 @@ async function openKafkaDebugSse(request, response, url, options) {
env, env,
stream, stream,
fromBeginning: url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true", fromBeginning: url.searchParams.get("fromBeginning") === "1" || url.searchParams.get("fromBeginning") === "true",
...filters, ...resolvedFilters,
kafkaFactory: options.kafkaFactory, kafkaFactory: options.kafkaFactory,
onEvent: async (record) => { onEvent: async (record) => {
writeSse(response, "hwlab.kafka.event", { 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) { function writeSse(response, eventName, payload) {
if (response.destroyed || response.writableEnded) return false; if (response.destroyed || response.writableEnded) return false;
try { try {