diff --git a/internal/cloud/server-workbench-http.test.ts b/internal/cloud/server-workbench-http.test.ts index bf1ab963..a04ac5e9 100644 --- a/internal/cloud/server-workbench-http.test.ts +++ b/internal/cloud/server-workbench-http.test.ts @@ -2393,7 +2393,7 @@ test("workbench realtime stream forwards HWLAB Kafka events after initial connec const { port } = serverWithKafka.address(); setTimeout(() => { void fakeKafka.emit({ eventType: "hwlab.trace.event.projected", - sessionId: session.id, + sessionId: "ses_agentrun_workbench_realtime_after_seq", 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 } diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index cf53c44e..618047be 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -340,14 +340,22 @@ function attachWorkbenchRealtimeOtelContext(request, fields = {}) { function resolveWorkbenchRealtimeKafkaFilters({ traceId = null, sessionId = null, traceStore = null } = {}) { const resolved = traceId && typeof traceStore?.snapshot === "function" ? collectTraceLinkedIds(traceStore.snapshot(traceId)) : {}; const hasAgentRunKey = Boolean(resolved.runId || resolved.commandId); + const fallbackSessionId = resolved.sessionId || agentRunScopedSessionIdFromWorkbenchSession(sessionId) || sessionId; return compactObject({ traceId: hasAgentRunKey ? null : traceId, - sessionId: hasAgentRunKey ? null : (resolved.sessionId || sessionId), + sessionId: hasAgentRunKey ? null : fallbackSessionId, runId: resolved.runId, commandId: resolved.commandId }); } +function agentRunScopedSessionIdFromWorkbenchSession(sessionId) { + const safeId = safeSessionId(sessionId); + if (!safeId || isAgentRunAliasSessionId(safeId)) return safeId; + const base = safeId.replace(/^ses_/u, "").replace(/[^A-Za-z0-9_]+/gu, "_").replace(/^_+|_+$/gu, "") || "session"; + return safeSessionId(`ses_agentrun_${base}`); +} + function collectTraceLinkedIds(snapshot) { const out = { sessionId: null, runId: null, commandId: null }; const events = Array.isArray(snapshot?.events) ? snapshot.events : []; diff --git a/internal/cloud/workbench-realtime-authority-contract.test.ts b/internal/cloud/workbench-realtime-authority-contract.test.ts index ef8c3a37..ce2eed94 100644 --- a/internal/cloud/workbench-realtime-authority-contract.test.ts +++ b/internal/cloud/workbench-realtime-authority-contract.test.ts @@ -61,7 +61,7 @@ test("workbench sync delta stays durable while SSE realtime forwards Kafka event setTimeout(() => { void fakeKafka.emit({ eventType: "hwlab.trace.event.projected", - sessionId, + sessionId: "ses_agentrun_realtime_authority_p1", traceId, context: { sourceSeq: 8, runId: "run_realtime_authority_p1", commandId: "cmd_realtime_authority_p1" }, event: { type: "result", eventType: "terminal", status: "completed", label: "agentrun:terminal:completed", message: "redacted final", terminal: true, sourceSeq: 8 }