diff --git a/internal/cloud/server-workbench-http.test.ts b/internal/cloud/server-workbench-http.test.ts index 1351e790..ecc79f6d 100644 --- a/internal/cloud/server-workbench-http.test.ts +++ b/internal/cloud/server-workbench-http.test.ts @@ -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(); diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index 7ba4a0d4..be23deaa 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -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;