fix: derive Workbench rail titles from Kafka index
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success

This commit is contained in:
root
2026-07-20 11:10:13 +02:00
parent 54e823e140
commit f07669fe40
7 changed files with 98 additions and 3 deletions
+16
View File
@@ -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;
@@ -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`
@@ -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
+15 -3
View File
@@ -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<any>; mode?: WorkbenchMode }, url: URL) {
async function nativeSessionList(options: { snapshot?: () => Promise<any>; 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<any>; mode?: WorkbenchMode }, sessionId: string, messagesOnly: boolean) {
const state = await requiredSnapshot(options);
+29
View File
@@ -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 }; },
@@ -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);
@@ -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
};