Merge pull request #2456 from pikasTech/fix/kafka-trace-filter-resolution
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success
fix: resolve kafka debug trace filters
This commit is contained in:
@@ -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<void>) | null = null;
|
||||
let subscribedTopic = "hwlab.event.v1";
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user