fix: keep Kafka session SSE live for new traces
This commit is contained in:
@@ -2397,7 +2397,7 @@ test("workbench realtime stream forwards HWLAB Kafka events after initial connec
|
||||
traceId,
|
||||
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 }
|
||||
}); }, 0);
|
||||
}); }, 25);
|
||||
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.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 () => {
|
||||
const traceStore = createCodeAgentTraceStore();
|
||||
const results = createCodeAgentChatResultStore();
|
||||
|
||||
@@ -247,7 +247,7 @@ export async function handleWorkbenchRealtimeHttp(request, response, url, option
|
||||
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" }));
|
||||
}
|
||||
const kafkaFilters = resolveWorkbenchRealtimeKafkaFilters({ traceId: activeTraceId, sessionId: streamSessionId, traceStore });
|
||||
const kafkaFilters = resolveWorkbenchRealtimeKafkaFilters({ traceId: requestedTraceId, sessionId: streamSessionId, traceStore });
|
||||
try {
|
||||
const kafkaStream = await openKafkaEventStream({
|
||||
env: options.env ?? process.env,
|
||||
@@ -379,7 +379,7 @@ function workbenchRealtimeEventFromKafka(record, context = {}) {
|
||||
if (!value) return null;
|
||||
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 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 threadId = safeOpaqueId(context.threadId ?? sourceContext.threadId) ?? textValue(context.threadId ?? sourceContext.threadId);
|
||||
if (!traceId && !sessionId) return null;
|
||||
|
||||
Reference in New Issue
Block a user