Merge pull request #2461 from pikasTech/fix/workbench-kafka-session-filter
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success

修复 Workbench Kafka session SSE 新 trace 实时订阅
This commit is contained in:
Lyon
2026-07-10 04:38:26 +08:00
committed by GitHub
2 changed files with 75 additions and 3 deletions
+73 -1
View File
@@ -2397,7 +2397,7 @@ test("workbench realtime stream forwards HWLAB Kafka events after initial connec
traceId, traceId,
context: { sourceSeq: 7, runId: "run_workbench_realtime_after_seq", commandId: "cmd_workbench_realtime_after_seq" }, context: { sourceSeq: 7, runId: "run_workbench_realtime_after_seq", commandId: "cmd_workbench_realtime_after_seq" },
event: { type: "backend", eventType: "backend", status: "running", label: "agentrun:event:test", message: "Kafka realtime event", sourceSeq: 7 } event: { type: "backend", eventType: "backend", status: "running", label: "agentrun:event:test", message: "Kafka realtime event", sourceSeq: 7 }
}); }, 0); }); }, 25);
const events = await getSseEvents(port, `/v1/workbench/events?sessionId=${encodeURIComponent(session.id)}&traceId=${encodeURIComponent(traceId)}&afterSeq=10`, 2); const events = await getSseEvents(port, `/v1/workbench/events?sessionId=${encodeURIComponent(session.id)}&traceId=${encodeURIComponent(traceId)}&afterSeq=10`, 2);
assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.trace.event"]); assert.deepEqual(events.map((event) => event.event), ["workbench.connected", "workbench.trace.event"]);
assert.equal(events[0].id, "10"); assert.equal(events[0].id, "10");
@@ -2416,6 +2416,78 @@ test("workbench realtime stream forwards HWLAB Kafka events after initial connec
} }
}); });
test("workbench session realtime Kafka stream does not pin subscription to stale lastTraceId", async () => {
const fakeKafka = createFakeKafkaFactory();
const staleTraceId = "trc_workbench_realtime_stale_trace";
const liveTraceId = "trc_workbench_realtime_live_trace";
const session = {
id: "ses_workbench_realtime_session_wide",
projectId: "prj_hwpod_workbench",
agentId: "hwlab-code-agent",
status: "running",
ownerUserId: ACTOR.id,
conversationId: "cnv_workbench_realtime_session_wide",
threadId: "thread-workbench-realtime-session-wide",
lastTraceId: staleTraceId,
updatedAt: "2026-06-24T14:00:00.000Z",
session: { sessionStatus: "running", lastTraceId: staleTraceId }
};
const accessController = {
store: {
async getAgentSession(sessionId) { return sessionId === session.id ? session : null; },
async getAgentSessionByTraceId() { return null; }
},
async ensureBootstrap() {},
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
};
const serverWithKafka = createCloudApiServer({
accessController,
workbenchRuntime: {
async queryWorkbenchFacts(params = {}) {
return {
facts: {
sessions: [{ sessionId: session.id, ownerUserId: ACTOR.id, threadId: session.threadId, lastTraceId: staleTraceId, status: "running", valuesRedacted: true }],
messages: [],
parts: [],
turns: [],
checkpoints: []
},
count: 1,
persistence: { adapter: "test-session-wide-kafka", durable: true },
params
};
}
},
kafkaFactory: fakeKafka.factory,
env: {
HWLAB_WORKBENCH_SSE_HEARTBEAT_MS: "60000",
HWLAB_KAFKA_BOOTSTRAP_SERVERS: "127.0.0.1:9092",
HWLAB_KAFKA_CLIENT_ID: "test-hwlab-cloud-api",
HWLAB_KAFKA_EVENT_TOPIC: "hwlab.event.v1"
}
});
await new Promise((resolve) => serverWithKafka.listen(0, "127.0.0.1", resolve));
try {
const { port } = serverWithKafka.address();
setTimeout(() => { void fakeKafka.emit({
eventType: "hwlab.trace.event.projected",
sessionId: "ses_agentrun_workbench_realtime_session_wide",
traceId: liveTraceId,
context: { runId: "run_workbench_realtime_session_wide", commandId: "cmd_workbench_realtime_session_wide" },
event: { type: "backend", eventType: "backend", status: "running", label: "agentrun:event:live", message: "Kafka live trace event" }
}); }, 100);
const events = await getSseEvents(port, `/v1/workbench/events?sessionId=${encodeURIComponent(session.id)}`, 4);
assert.equal(events[0].event, "workbench.connected");
const liveEvent = events.find((event) => event.event === "workbench.trace.event" && event.data?.traceId === liveTraceId);
assert.ok(liveEvent, JSON.stringify(events.map((event) => ({ event: event.event, traceId: event.data?.traceId, message: event.data?.event?.message }))));
assert.equal(liveEvent.data.event.message, "Kafka live trace event");
assert.equal(liveEvent.data.realtimeAuthority, "workbench-realtime-authority-v2");
} finally {
await new Promise((resolve, reject) => serverWithKafka.close((error) => error ? reject(error) : resolve()));
}
});
test("workbench read model exposes runtime trace projection query failures as projection blockers", async () => { test("workbench read model exposes runtime trace projection query failures as projection blockers", async () => {
const traceStore = createCodeAgentTraceStore(); const traceStore = createCodeAgentTraceStore();
const results = createCodeAgentChatResultStore(); const results = createCodeAgentChatResultStore();
+2 -2
View File
@@ -247,7 +247,7 @@ export async function handleWorkbenchRealtimeHttp(request, response, url, option
if (activeTraceId && requestedAfterSeq <= 0) { if (activeTraceId && requestedAfterSeq <= 0) {
await (perf ? perf.measure("workbench_initial_trace", () => writeTraceRealtimeSnapshotSafe({ writeEvent, options, actor: auth.actor, traceId: activeTraceId, reason: "initial" })) : writeTraceRealtimeSnapshotSafe({ writeEvent, options, actor: auth.actor, traceId: activeTraceId, reason: "initial" })); await (perf ? perf.measure("workbench_initial_trace", () => writeTraceRealtimeSnapshotSafe({ writeEvent, options, actor: auth.actor, traceId: activeTraceId, reason: "initial" })) : writeTraceRealtimeSnapshotSafe({ writeEvent, options, actor: auth.actor, traceId: activeTraceId, reason: "initial" }));
} }
const kafkaFilters = resolveWorkbenchRealtimeKafkaFilters({ traceId: activeTraceId, sessionId: streamSessionId, traceStore }); const kafkaFilters = resolveWorkbenchRealtimeKafkaFilters({ traceId: requestedTraceId, sessionId: streamSessionId, traceStore });
try { try {
const kafkaStream = await openKafkaEventStream({ const kafkaStream = await openKafkaEventStream({
env: options.env ?? process.env, env: options.env ?? process.env,
@@ -379,7 +379,7 @@ function workbenchRealtimeEventFromKafka(record, context = {}) {
if (!value) return null; if (!value) return null;
const hwlabEvent = value.event && typeof value.event === "object" && !Array.isArray(value.event) ? value.event : {}; const hwlabEvent = value.event && typeof value.event === "object" && !Array.isArray(value.event) ? value.event : {};
const sourceContext = value.context && typeof value.context === "object" && !Array.isArray(value.context) ? value.context : {}; const sourceContext = value.context && typeof value.context === "object" && !Array.isArray(value.context) ? value.context : {};
const traceId = safeTraceId(context.traceId ?? value.traceId ?? hwlabEvent.traceId) ?? textValue(context.traceId ?? value.traceId ?? hwlabEvent.traceId); const traceId = safeTraceId(value.traceId ?? hwlabEvent.traceId ?? context.traceId) ?? textValue(value.traceId ?? hwlabEvent.traceId ?? context.traceId);
const sessionId = safeSessionId(context.sessionId ?? value.sessionId ?? hwlabEvent.sessionId) ?? textValue(context.sessionId ?? value.sessionId ?? hwlabEvent.sessionId); const sessionId = safeSessionId(context.sessionId ?? value.sessionId ?? hwlabEvent.sessionId) ?? textValue(context.sessionId ?? value.sessionId ?? hwlabEvent.sessionId);
const threadId = safeOpaqueId(context.threadId ?? sourceContext.threadId) ?? textValue(context.threadId ?? sourceContext.threadId); const threadId = safeOpaqueId(context.threadId ?? sourceContext.threadId) ?? textValue(context.threadId ?? sourceContext.threadId);
if (!traceId && !sessionId) return null; if (!traceId && !sessionId) return null;