diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index f78df3f6..e6ab6ff5 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -416,6 +416,14 @@ function startLiveHwlabKafkaEventBridge({ config, env = process.env, logger = co } }; }, + workbenchSessionSummaries() { + return sessionIndex?.sessionSummaries() ?? { + ok: false, + sessions: [], + warning: { code: "workbench_kafka_session_index_disabled", message: "Session summary index is disabled.", blocking: false, valuesRedacted: true }, + valuesRedacted: true + }; + }, liveSubscriberCount() { return subscribers.size; }, valuesPrinted: false }; @@ -487,6 +495,14 @@ function combineKafkaEventBridgeComponents(config, components) { if (typeof liveOwner?.queryHwlabEventRetention !== "function") throw contractError("hwlab_kafka_refresh_query_unconfigured", "Kafka refresh retention query is not configured."); return liveOwner.queryHwlabEventRetention(params); }, + workbenchSessionSummaries() { + return liveOwner?.workbenchSessionSummaries?.() ?? { + ok: false, + sessions: [], + warning: { code: "workbench_kafka_session_index_disabled", message: "Session summary index is disabled.", blocking: false, valuesRedacted: true }, + valuesRedacted: true + }; + }, liveSubscriberCount() { const owner = components.find((component) => typeof component.liveSubscriberCount === "function"); return owner?.liveSubscriberCount() ?? 0; diff --git a/internal/cloud/workbench-kafka-session-index.test.ts b/internal/cloud/workbench-kafka-session-index.test.ts index 5d0018f7..faa352b0 100644 --- a/internal/cloud/workbench-kafka-session-index.test.ts +++ b/internal/cloud/workbench-kafka-session-index.test.ts @@ -36,6 +36,9 @@ test("process session index bootstraps once and serves repeated scoped queries f assert.equal(first.result.index.hit, true); assert.equal(first.result.index.indexedEventCount, 3); assert.equal(queryCount, 1); + assert.deepEqual(index.sessionSummaries().sessions, [ + { sessionId: "ses_a", firstUserMessagePreview: "prompt ses_a", lastUserMessageAt: "2026-07-20T00:00:01.000Z", valuesRedacted: true } + ]); index.observeLive(envelope("ses_a", "trc_a", 4), { topic: "hwlab.event.v1", partition: 0, offset: "3" }); const live = index.query({ sessionId: "ses_a", limit: 10 }); @@ -106,6 +109,7 @@ function envelope(sessionId: string, traceId: string, seq: number) { traceId, event: { type: seq === 1 ? "user" : "assistant", + text: seq === 1 ? `prompt ${sessionId}` : "reply", sessionId, traceId, createdAt: `2026-07-20T00:00:0${seq}.000Z` diff --git a/internal/cloud/workbench-kafka-session-index.ts b/internal/cloud/workbench-kafka-session-index.ts index f34935f0..9037a7b3 100644 --- a/internal/cloud/workbench-kafka-session-index.ts +++ b/internal/cloud/workbench-kafka-session-index.ts @@ -227,6 +227,34 @@ export function createWorkbenchKafkaSessionIndex(options = {}) { }; } + function sessionSummaries() { + if (!ready || !state.globalComplete) { + const reason = !ready ? "index-not-ready" : "global-capacity-exceeded"; + return { + ok: false, + sessions: [], + warning: warning("workbench_kafka_session_summary_unavailable", `Session summary index is unavailable (${reason}).`), + status: status(), + valuesRedacted: true + }; + } + const sessions = []; + for (const [id, records] of state.bySession.entries()) { + if (state.incompleteSessions.has(id)) continue; + const userEvents = records.map((record) => record.value?.event).filter((event) => firstText(event?.type, event?.eventType) === "user"); + const firstPreview = firstText(userEvents[0]?.text, userEvents[0]?.message); + if (!firstPreview) continue; + const lastUserEvent = userEvents.at(-1); + sessions.push({ + sessionId: id, + firstUserMessagePreview: firstPreview.slice(0, 240), + lastUserMessageAt: firstText(lastUserEvent?.createdAt, lastUserEvent?.submittedAt), + valuesRedacted: true + }); + } + return { ok: true, sessions, warning: null, status: status(), valuesRedacted: true }; + } + function emptyState() { return { records: [], @@ -295,6 +323,7 @@ export function createWorkbenchKafkaSessionIndex(options = {}) { stop, observeLive, query, + sessionSummaries, status, rebuild, valuesPrinted: false diff --git a/internal/workbench/http.ts b/internal/workbench/http.ts index b3f23624..aeef38ea 100644 --- a/internal/workbench/http.ts +++ b/internal/workbench/http.ts @@ -73,14 +73,26 @@ function requireAuthorization(request: Request, options: { snapshot?: unknown; a if (!expected || actual !== expected) throw Object.assign(new Error("Workbench API internal authentication failed"), { code: "auth_required" }); } -async function nativeSessionList(options: { snapshot?: () => Promise; mode?: WorkbenchMode }, url: URL) { +async function nativeSessionList(options: { snapshot?: () => Promise; mode?: WorkbenchMode; kafkaEventBridge?: any }, url: URL) { const state = await requiredSnapshot(options); const owner = text(url.searchParams.get("ownerUserId")); + const kafkaSummary = options.mode === "agentrun-native" ? options.kafkaEventBridge?.workbenchSessionSummaries?.() : null; + const kafkaSummaryBySession = new Map((Array.isArray(kafkaSummary?.sessions) ? kafkaSummary.sessions : []).map((session: any) => [text(session.sessionId), session])); const sessions = Object.values(state.sessions) .filter((session: any) => !owner || session.ownerUserId === owner) - .map((session: any) => ({ ...session, firstUserMessagePreview: firstUserMessagePreview(session.messages) ?? (text(session.firstUserMessagePreview) || null) })) + .map((session: any) => { + const summary = kafkaSummaryBySession.get(text(session.sessionId)) as any; + const nativeTestPreview = options.mode === "agentrun-native" ? null : firstUserMessagePreview(session.messages) ?? (text(session.firstUserMessagePreview) || null); + return { + ...session, + firstUserMessagePreview: text(summary?.firstUserMessagePreview) || nativeTestPreview, + lastUserMessageAt: text(summary?.lastUserMessageAt) || null, + updatedAt: text(summary?.lastUserMessageAt) || session.updatedAt + }; + }) .sort((left: any, right: any) => String(right.updatedAt).localeCompare(String(left.updatedAt))); - return json(200, { ok: true, status: "ready", sessions, count: sessions.length, mode: options.mode ?? "native-test" }); + const warnings = kafkaSummary?.warning ? [kafkaSummary.warning] : []; + return json(200, { ok: true, status: "ready", sessions, count: sessions.length, warnings, mode: options.mode ?? "native-test" }); } async function nativeSessionDetail(options: { snapshot?: () => Promise; mode?: WorkbenchMode }, sessionId: string, messagesOnly: boolean) { const state = await requiredSnapshot(options); diff --git a/internal/workbench/workbench.test.ts b/internal/workbench/workbench.test.ts index 0a81b02c..2ba0d097 100644 --- a/internal/workbench/workbench.test.ts +++ b/internal/workbench/workbench.test.ts @@ -103,6 +103,35 @@ describe("Workbench native HTTP adapter", () => { }); }); + test("AgentRun native session list uses Kafka index summaries instead of admission messages", async () => { + const app = createWorkbenchHttpApp({ + mode: "agentrun-native", + kafkaEventBridge: { + workbenchSessionSummaries() { + return { ok: true, sessions: [{ sessionId: "ses_title", firstUserMessagePreview: "Kafka title", lastUserMessageAt: "2026-07-20T01:00:00.000Z" }] }; + } + }, + async dispatch() { return { ok: true }; }, + async snapshot() { + return { + sessions: { + ses_title: { + sessionId: "ses_title", + updatedAt: "2026-07-20T00:00:00.000Z", + messages: [{ role: "user", content: "admission snapshot title" }] + } + }, + turns: {} + }; + } + }); + + const response = await app.fetch(new Request("http://native.test/v1/workbench/sessions")); + expect(await response.json()).toMatchObject({ + sessions: [{ sessionId: "ses_title", firstUserMessagePreview: "Kafka title", lastUserMessageAt: "2026-07-20T01:00:00.000Z" }] + }); + }); + test("authorizes only the Workbench navigation entry", async () => { const app = createWorkbenchHttpApp({ async dispatch() { return { ok: true }; }, diff --git a/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts b/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts index 80ed1aa0..aed036ef 100644 --- a/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts +++ b/web/hwlab-cloud-web/src/stores/workbench-message-projection-runtime.test.ts @@ -151,6 +151,8 @@ test("workbench active state uses Kafka SSE without HTTP session or terminal pro assert.doesNotMatch(source, /fetchSessionDetailPage|fetchSessionMessagesPage/u); assert.doesNotMatch(loadBlock, /fetchSession\(|fetchSessionMessages\(/u); assert.match(loadBlock, /workbenchSessionNavigationSeed\(id/u); + assert.match(source, /firstUserMessagePreview: firstNonEmptyString\(source\?\.firstUserMessagePreview\)/u); + assert.match(source, /lastUserMessageAt: firstNonEmptyString\(source\?\.lastUserMessageAt\)/u); assert.doesNotMatch(submitBlock, /applyTurnStatusSnapshot|completeTrace/u); assert.doesNotMatch(traceDetailApplyBlock, /rememberTurnStatus|projectTurnAuthorityToMessages|chatPending\.value\s*=\s*false/u); assert.doesNotMatch(source, /sealRestoredActiveTurnMessages|messageNeedsRestoredTurnSeal|refreshSessionStatusAuthority|readTerminalTraceDetailGaps/u); diff --git a/web/hwlab-cloud-web/src/stores/workbench.ts b/web/hwlab-cloud-web/src/stores/workbench.ts index a7badf7c..8b3c0b35 100644 --- a/web/hwlab-cloud-web/src/stores/workbench.ts +++ b/web/hwlab-cloud-web/src/stores/workbench.ts @@ -1424,6 +1424,9 @@ function workbenchSessionNavigationSeed(sessionId: string, source: WorkbenchSess return { sessionId, threadId: firstNonEmptyString(source?.threadId), + firstUserMessagePreview: firstNonEmptyString(source?.firstUserMessagePreview), + lastUserMessageAt: firstNonEmptyString(source?.lastUserMessageAt), + updatedAt: firstNonEmptyString(source?.lastUserMessageAt, source?.updatedAt), messages: [], messageCount: 0 };