From 9cd3a083ebef2bead305ab94f63baa1ea6c73b4d Mon Sep 17 00:00:00 2001 From: root Date: Mon, 20 Jul 2026 04:14:30 +0200 Subject: [PATCH] fix(workbench): replay Kafka retention in native SSE --- internal/workbench/http.ts | 109 +++++++++++++++++++++++++++ internal/workbench/workbench.test.ts | 71 +++++++++++++++++ tools/src/workbench-cli.ts | 8 +- 3 files changed, 186 insertions(+), 2 deletions(-) diff --git a/internal/workbench/http.ts b/internal/workbench/http.ts index cf27bb38..fd1ddaf8 100644 --- a/internal/workbench/http.ts +++ b/internal/workbench/http.ts @@ -1,4 +1,5 @@ import type { WorkbenchCommand, WorkbenchMode } from "./contracts.ts"; +import { createWorkbenchKafkaRefreshHandoff, workbenchKafkaRefreshErrorPayload } from "../cloud/workbench-kafka-refresh-handoff.ts"; export function createWorkbenchHttpApp(options: { dispatch: (command: WorkbenchCommand) => Promise; snapshot?: () => Promise<{ sessions: Record>; turns: Record> }>; authorization?: string; mode?: WorkbenchMode; kafkaEventBridge?: any; close?: () => Promise }) { return { @@ -99,6 +100,9 @@ async function nativeEventStream(options: { mode?: WorkbenchMode; kafkaEventBrid if (mode === "agentrun-native" && (bridge?.capabilities?.liveKafkaSse !== true || typeof bridge?.subscribeLiveHwlabEvents !== "function")) { return json(503, { ok: false, error: { code: "workbench_live_kafka_unconfigured", message: "Native Workbench requires the Kafka SSE bridge" } }); } + if (mode === "agentrun-native" && bridge?.capabilities?.kafkaRefreshReplay === true) { + return nativeKafkaRefreshEventStream(bridge, sessionId, traceId, mode); + } const encoder = new TextEncoder(); let heartbeat: ReturnType | undefined; let unsubscribe: (() => void) | undefined; @@ -124,6 +128,111 @@ async function nativeEventStream(options: { mode?: WorkbenchMode; kafkaEventBrid }), { status: 200, headers: { "content-type": "text/event-stream", "cache-control": "no-store", connection: "keep-alive" } }); } +function nativeKafkaRefreshEventStream(bridge: any, sessionId: string, traceId: string, mode: WorkbenchMode) { + const refreshReplay = bridge?.refreshReplay; + if (!refreshReplay || typeof bridge?.queryHwlabEventRetention !== "function") { + return json(503, { ok: false, error: { code: "workbench_kafka_refresh_unconfigured", message: "Native Workbench Kafka refresh replay requires its retention query runtime" } }); + } + const encoder = new TextEncoder(); + let heartbeat: ReturnType | undefined; + let handoff: ReturnType | undefined; + let closed = false; + return new Response(new ReadableStream({ + start(controller) { + const write = (name: string, payload: unknown) => { + if (closed) return false; + controller.enqueue(encoder.encode(`event: ${name}\ndata: ${JSON.stringify(payload)}\n\n`)); + return true; + }; + handoff = createWorkbenchKafkaRefreshHandoff({ + sessionId, + traceId, + liveBufferLimit: refreshReplay.liveBufferLimit, + identityWindowLimit: refreshReplay.matchedEventLimit + refreshReplay.liveBufferLimit, + subscribeLive: (listener: (envelope: any, transport: any) => void) => bridge.subscribeLiveHwlabEvents(listener), + queryRetention: async ({ signal }: { signal: AbortSignal }) => { + const result = await bridge.queryHwlabEventRetention({ + sessionId, + traceId, + limit: refreshReplay.matchedEventLimit, + scanLimit: refreshReplay.scanLimit, + timeoutMs: refreshReplay.timeoutMs, + groupIdPrefix: refreshReplay.groupIdPrefix, + partitionKey: sessionId, + fromBeginning: true, + signal + }); + console.log(JSON.stringify({ + code: "workbench-kafka-refresh-query", + component: "hwlab-workbench-native", + sessionId: sessionId || null, + traceId: traceId || null, + completionReason: result?.completionReason ?? null, + complete: result?.completion?.complete === true, + scannedCount: result?.scannedCount ?? null, + matchedCount: result?.matchedCount ?? null, + partitionKeyScoped: result?.partitionKeyScoped === true, + targetPartition: result?.targetPartition ?? null, + timing: result?.timing ?? null, + valuesPrinted: false + })); + return result; + }, + deliverEvent: (envelope: any) => write("hwlab.event.v1", envelope), + deliverConnected: (summary: any) => write("workbench.connected", { + type: "connected", + status: "connected", + mode, + capabilities: bridge.capabilities, + realtimeSource: "hwlab.event.v1", + deliverySemantics: "kafka-retention-then-live", + liveOnly: false, + replay: true, + replaySupported: true, + lossPossible: false, + filters: { sessionId: sessionId || null, traceId: traceId || null }, + refreshReplay: summary, + valuesPrinted: false + }), + onFailure: (error: any) => { + if (closed) return; + write("workbench.error", workbenchKafkaRefreshErrorPayload(error, { sessionId, traceId })); + closed = true; + if (heartbeat) clearInterval(heartbeat); + controller.close(); + } + }); + void Promise.resolve(bridge.liveReady ?? bridge.ready) + .then(() => handoff?.start()) + .then(() => { + if (closed) return; + heartbeat = setInterval(() => write("workbench.heartbeat", { + type: "heartbeat", + status: "connected", + realtimeSource: "hwlab.event.v1", + deliverySemantics: "kafka-retention-then-live", + liveOnly: false, + replay: true, + lossPossible: false, + serverSentAt: new Date().toISOString(), + valuesPrinted: false + }), 10_000); + }) + .catch((error: any) => { + if (closed) return; + write("workbench.error", workbenchKafkaRefreshErrorPayload(error, { sessionId, traceId })); + closed = true; + controller.close(); + }); + }, + cancel() { + closed = true; + if (heartbeat) clearInterval(heartbeat); + handoff?.stop("connection-closed"); + } + }), { status: 200, headers: { "content-type": "text/event-stream", "cache-control": "no-store", connection: "keep-alive" } }); +} + function nativeEnvelopeMatches(envelope: any, sessionId: string, traceId: string) { const event = envelope?.event && typeof envelope.event === "object" ? envelope.event : {}; const envelopeSessionId = text(envelope?.sessionId ?? envelope?.hwlabSessionId ?? event.sessionId); diff --git a/internal/workbench/workbench.test.ts b/internal/workbench/workbench.test.ts index ff460d84..bc121f89 100644 --- a/internal/workbench/workbench.test.ts +++ b/internal/workbench/workbench.test.ts @@ -152,6 +152,77 @@ describe("Workbench native HTTP adapter", () => { } }); + test("L0 CLI replays retained Kafka events before the same live SSE stream", async () => { + const sessionId = "ses_l0_refresh"; + const traceId = "trc_l0_refresh"; + const bridge = { + capabilities: { liveKafkaSse: true, kafkaRefreshReplay: true }, + refreshReplay: { groupIdPrefix: "hwlab-l0-refresh", timeoutMs: 1000, scanLimit: 100, matchedEventLimit: 20, liveBufferLimit: 20 }, + ready: Promise.resolve(), + subscribeLiveHwlabEvents() { return () => {}; }, + async queryHwlabEventRetention(options: any) { + expect(options).toMatchObject({ sessionId, traceId, partitionKey: sessionId, fromBeginning: true }); + return { + topic: "hwlab.event.v1", + events: [{ + topic: "hwlab.event.v1", + partition: 0, + offset: "4", + value: { + schema: "hwlab.event.v1", + eventType: "hwlab.trace.event.projected", + eventId: "hwlab:evt_l0_refresh", + sourceEventId: "evt_l0_refresh", + sessionId, + hwlabSessionId: sessionId, + traceId, + event: { type: "backend", eventType: "backend", sessionId, traceId, sourceEventId: "evt_l0_refresh" } + } + }], + completionReason: "end-offset", + completion: { reason: "end-offset", complete: true, barrierReached: true, retentionStartVerified: true }, + reachedEndOffsets: true, + endOffsetsAvailable: true, + endOffsets: [{ partition: 0, startOffset: "0", endOffset: "5" }], + scannedCount: 5, + matchedCount: 1 + }; + } + }; + const app = createWorkbenchHttpApp({ + mode: "agentrun-native", + kafkaEventBridge: bridge, + async dispatch() { return { ok: true }; }, + async snapshot() { return { sessions: {}, turns: {} }; } + }); + const server = Bun.serve({ port: 0, fetch: (request) => app.fetch(request) }); + try { + expect(await runWorkbenchCli([ + "events", "inspect", + "--session-id", sessionId, + "--trace-id", traceId, + "--over-api", + "--api-url", `http://127.0.0.1:${server.port}`, + "--timeout-ms", "1000" + ], {})).toMatchObject({ + ok: true, + connected: true, + connectedContract: { + deliverySemantics: "kafka-retention-then-live", + realtimeSource: "hwlab.event.v1", + replay: true, + liveOnly: false, + refreshReplay: { phase: "live", counts: { replayed: 1 } } + }, + eventCount: 1, + eventTypes: ["backend"] + }); + } finally { + server.stop(true); + await app.close(); + } + }); + test("L0 CLI waits for explicit Kafka event semantics in one SSE window", async () => { const subscribers = new Set<(envelope: any) => void>(); const bridge = { diff --git a/tools/src/workbench-cli.ts b/tools/src/workbench-cli.ts index 1226def5..9325951e 100644 --- a/tools/src/workbench-cli.ts +++ b/tools/src/workbench-cli.ts @@ -72,7 +72,7 @@ async function inspectEvents(parsed: Parsed, env: Record undefined); break readStream; } @@ -81,7 +81,7 @@ async function inspectEvents(parsed: Parsed, env: Record entry.name === "workbench.connected")}`); } @@ -176,6 +176,10 @@ function workbenchConnectedContract(data: Record) { }; } function businessFrames(frames: SseFrame[]) { return frames.filter((entry) => entry.name === "hwlab.event.v1"); } +function framesInspectionSatisfied(frames: SseFrame[], minEvents: number, waitFor: EventSemantic[]) { + return frames.some((entry) => entry.name === "workbench.connected") + && inspectionSatisfied(businessFrames(frames), minEvents, waitFor); +} function inspectionSatisfied(business: SseFrame[], minEvents: number, waitFor: EventSemantic[]) { if (business.length < minEvents) return false; const observed = semanticSummary(business).observedSemantics;