diff --git a/internal/cloud/server-workbench-http.test.ts b/internal/cloud/server-workbench-http.test.ts index 0690a482..5017ab0b 100644 --- a/internal/cloud/server-workbench-http.test.ts +++ b/internal/cloud/server-workbench-http.test.ts @@ -138,12 +138,7 @@ test("workbench trace event page exposes monotonic cursor range for restored mix async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - const runtimeStore = { - async queryAgentTraceEvents(input = {}) { - assert.equal(input.traceId, traceId); - return { events: mixedEvents, count: mixedEvents.length }; - } - }; + const runtimeStore = createDurableFactsRuntimeStore({ sessions: [{ session, events: mixedEvents, status: "completed" }] }); const server = createCloudApiServer({ accessController, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); @@ -207,7 +202,7 @@ test("workbench read model exposes session, messages, turn, and trace without wr secretMaterialStored: false } }; - const listInputs = []; + const factQueries = []; traceStore.append(traceId, { type: "request", status: "accepted", label: "request:accepted" }); traceStore.append(traceId, { type: "result", status: "completed", label: "result:completed", terminal: true }); results.set(traceId, { @@ -224,7 +219,6 @@ test("workbench read model exposes session, messages, turn, and trace without wr const accessController = { store: { async listAgentSessionsForUser(input = {}) { - listInputs.push(input); return [session]; }, async getAgentSession(sessionId) { return sessionId === session.id ? session : null; }, @@ -235,7 +229,22 @@ test("workbench read model exposes session, messages, turn, and trace without wr async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: results }); + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ + session, + events: [ + { seq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-17T00:00:00.000Z" }, + { seq: 2, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-17T00:00:03.000Z" } + ], + status: "completed", + finalText: "pong", + runId: "run_workbench_read_model", + commandId: "cmd_workbench_read_model", + lastProjectedSeq: 2 + }], + queries: factQueries + }); + const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); try { @@ -245,13 +254,13 @@ test("workbench read model exposes session, messages, turn, and trace without wr assert.equal(sessions.body.contractVersion, "workbench-sessions-v1"); assert.equal(sessions.body.sessions[0].sessionId, session.id); assert.equal(sessions.body.sessions[0].turnSummary.traceId, traceId); - assert.equal(listInputs[0].projectId, undefined); - assert.equal(listInputs[0].limit, 21); + assert.equal(factQueries[0].projectId, undefined); + assert.equal(factQueries[0].limit, 21); const staleProject = await getJson(port, `/v1/workbench/sessions?projectId=prj_stale_filter&includeSessionId=${encodeURIComponent(session.id)}`); assert.equal(staleProject.status, 400); assert.equal(staleProject.body.error.code, "workbench_authority_removed"); - assert.equal(listInputs.length, 1); + assert.equal(factQueries.length, 1); const detail = await getJson(port, `/v1/workbench/sessions/${encodeURIComponent(session.id)}`); assert.equal(detail.status, 200); @@ -333,14 +342,11 @@ test("workbench read model recovers trace events from durable projection without async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - let traceQueryCount = 0; - const runtimeStore = { - async queryAgentTraceEvents(params = {}) { - traceQueryCount += 1; - assert.equal(params.traceId, traceId); - return { events: durableEvents, count: durableEvents.length }; - } - }; + const factQueries = []; + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ session, events: durableEvents, status: "completed", finalText: "durable pong", lastProjectedSeq: 2 }], + queries: factQueries + }); const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); @@ -362,7 +368,7 @@ test("workbench read model recovers trace events from durable projection without assert.equal(trace.body.events.length, 2); assert.equal(trace.body.traceStatus, "completed"); assert.equal(trace.body.hasMore, false); - assert.equal(traceQueryCount, 1); + assert.ok(factQueries.length >= 3); } finally { await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); } @@ -415,7 +421,21 @@ test("workbench read model projects terminal result atomically across session, m async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: results }); + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ + session, + events: [ + { seq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-17T02:09:58.000Z" }, + { seq: 2, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-17T02:10:00.000Z" } + ], + status: "completed", + finalText, + runId: "run_workbench_terminal_atomic", + commandId: "cmd_workbench_terminal_atomic", + lastProjectedSeq: 2 + }] + }); + const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); try { @@ -498,7 +518,15 @@ test("workbench read model projects current turn running state from trace projec async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: createCodeAgentChatResultStore() }); + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ + session, + events: [{ seq: 1, type: "backend", status: "running", label: "runner:created", terminal: false, createdAt: session.updatedAt }], + status: "running", + lastProjectedSeq: 1 + }] + }); + const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); try { @@ -536,7 +564,15 @@ test("workbench read model projects current turn running state for idle session async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: createCodeAgentChatResultStore() }); + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ + session, + events: [{ seq: 1, type: "backend", status: "running", label: "runner:created", terminal: false, createdAt: session.updatedAt }], + status: "running", + lastProjectedSeq: 1 + }] + }); + const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); try { @@ -598,7 +634,21 @@ test("workbench read model projects completed current turn for idle session summ async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; const backendPerformanceStore = createBackendPerformanceStore({ env: { POD_NAMESPACE: "hwlab-v03", HWLAB_GITOPS_TARGET: "v03", HWLAB_RUNTIME_LANE: "v03", HWLAB_CODE_AGENT_AGENTRUN_PROVIDER_ID: "D601" } }); - const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: results, backendPerformanceStore }); + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ + session, + events: [ + { seq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-17T00:02:58.000Z" }, + { seq: 2, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-17T00:03:00.000Z" } + ], + status: "completed", + finalText, + runId: "run_workbench_idle_completed", + commandId: "cmd_workbench_idle_completed", + lastProjectedSeq: 2 + }] + }); + const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results, backendPerformanceStore }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); try { @@ -685,7 +735,21 @@ test("workbench read model lets terminal result override stale running session s async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: results }); + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ + session, + events: [ + { seq: 1, type: "backend", status: "running", label: "runner:created", terminal: false, createdAt: "2026-06-18T16:30:00.000Z" }, + { seq: 2, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-18T16:31:04.000Z" } + ], + status: "completed", + finalText, + runId: "run_workbench_stale_running_completed", + commandId: "cmd_workbench_stale_running_completed", + lastProjectedSeq: 2 + }] + }); + const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); try { @@ -720,7 +784,7 @@ test("workbench read model lets terminal result override stale running session s } }); -test("workbench session list uses compact trace result without durable trace hydration", async () => { +test("workbench session list uses durable projection checkpoint without trace/result hydration", async () => { const traceStore = createCodeAgentTraceStore(); const traceId = "trc_workbench_compact_list_summary"; const finalText = "compact list summary OK"; @@ -771,13 +835,22 @@ test("workbench session list uses compact trace result without durable trace hyd async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - let durableQueryCount = 0; - const runtimeStore = { - async queryAgentTraceEvents() { - durableQueryCount += 1; - throw new Error("session list must not hydrate durable trace events"); - } - }; + const factQueries = []; + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ + session, + events: [ + { seq: 1, sourceSeq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-19T15:39:10.000Z" }, + { seq: 42, sourceSeq: 42, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-19T15:40:00.000Z" } + ], + status: "completed", + finalText, + runId: "run_workbench_compact_list_summary", + commandId: "cmd_workbench_compact_list_summary", + lastProjectedSeq: 42 + }], + queries: factQueries + }); const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); @@ -791,13 +864,13 @@ test("workbench session list uses compact trace result without durable trace hyd assert.equal(sessions.body.sessions[0].turnSummary.status, "completed"); assert.equal(sessions.body.sessions[0].turnSummary.eventCount, 42); assert.equal(sessions.body.sessions[0].projectionStatus, "caught-up"); - assert.equal(durableQueryCount, 0); + assert.ok(factQueries.length >= 1); } finally { await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); } }); -test("workbench read model recovers terminal session status from durable trace", async () => { +test("workbench read model reads terminal session status from durable facts", async () => { const traceStore = createCodeAgentTraceStore(); const traceId = "trc_workbench_durable_terminal_after_memory_running"; const session = { @@ -833,14 +906,9 @@ test("workbench read model recovers terminal session status from durable trace", async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - let durableQueryCount = 0; - const runtimeStore = { - async queryAgentTraceEvents(params = {}) { - durableQueryCount += 1; - assert.equal(params.traceId, traceId); - return { events: durableEvents, count: durableEvents.length }; - } - }; + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ session, events: durableEvents, status: "completed", finalText: "OK", lastProjectedSeq: 2 }] + }); const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); @@ -848,17 +916,15 @@ test("workbench read model recovers terminal session status from durable trace", const { port } = server.address(); const sessions = await getJson(port, `/v1/workbench/sessions?includeSessionId=${encodeURIComponent(session.id)}`); assert.equal(sessions.status, 200); - assert.equal(sessions.body.sessions[0].status, "running"); - assert.equal(sessions.body.sessions[0].running, true); - assert.equal(sessions.body.sessions[0].terminal, false); - assert.equal(durableQueryCount, 0); + assert.equal(sessions.body.sessions[0].status, "completed"); + assert.equal(sessions.body.sessions[0].running, false); + assert.equal(sessions.body.sessions[0].terminal, true); const detail = await getJson(port, `/v1/workbench/sessions/${encodeURIComponent(session.id)}`); assert.equal(detail.status, 200); assert.equal(detail.body.session.status, "completed"); assert.equal(detail.body.session.running, false); assert.equal(detail.body.session.terminal, true); - assert.equal(durableQueryCount, 1); const messages = await getJson(port, `/v1/workbench/sessions/${encodeURIComponent(session.id)}/messages?limit=10`); assert.equal(messages.status, 200); @@ -927,12 +993,17 @@ test("workbench read model keeps non-terminal durable tool completed events out async ensureBootstrap() {}, async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; - const runtimeStore = { - async queryAgentTraceEvents(params = {}) { - assert.equal(params.traceId, traceId); - return { events: durableEvents, count: durableEvents.length }; - } - }; + const runtimeStore = createDurableFactsRuntimeStore({ + sessions: [{ + session, + events: durableEvents, + status: "running", + runId: "run_workbench_nonterminal_tool_completed", + commandId: "cmd_workbench_nonterminal_tool_completed", + projectionStatus: "projecting", + lastProjectedSeq: 2 + }] + }); const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results }); await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); @@ -1053,7 +1124,7 @@ test("workbench read model exposes runtime trace projection query failures as pr async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; } }; const runtimeStore = { - async queryAgentTraceEvents() { + async queryWorkbenchFacts() { const error = new Error("runtime query failed"); error.code = "TEST_RUNTIME_QUERY_FAILED"; throw error; @@ -1065,18 +1136,12 @@ test("workbench read model exposes runtime trace projection query failures as pr try { const { port } = server.address(); const turn = await getJson(port, `/v1/workbench/turns/${encodeURIComponent(traceId)}`); - assert.equal(turn.status, 200); - assert.equal(turn.body.turn.status, "running"); - assert.equal(turn.body.projectionHealth, "unavailable"); - assert.equal(turn.body.projection.blocker.code, "projection_store_unavailable"); - assert.equal(turn.body.blocker.code, "projection_store_unavailable"); + assert.equal(turn.status, 503); + assert.equal(turn.body.error.code, "projection_store_unavailable"); const trace = await getJson(port, `/v1/workbench/traces/${encodeURIComponent(traceId)}/events?limit=10`); - assert.equal(trace.status, 200); - assert.equal(trace.body.events.length, 1); - assert.equal(trace.body.traceStatus, "running"); - assert.equal(trace.body.projectionHealth, "unavailable"); - assert.equal(trace.body.projection.blocker.code, "projection_store_unavailable"); + assert.equal(trace.status, 503); + assert.equal(trace.body.error.code, "projection_store_unavailable"); } finally { await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve())); } @@ -1090,6 +1155,198 @@ async function getJson(port, path) { }; } +function createDurableFactsRuntimeStore({ sessions = [], queryError = null, queries = [] } = {}) { + const facts = emptyFacts(); + for (const item of sessions) mergeFacts(facts, buildDurableFactsForSession(item)); + return { + queries, + async queryWorkbenchFacts(params = {}) { + queries.push(params); + if (queryError) throw queryError; + const filtered = filterFacts(facts, params); + return { + facts: filtered, + count: Object.values(filtered).reduce((sum, rows) => sum + rows.length, 0), + persistence: { adapter: "test-durable-workbench-facts", durable: true } + }; + } + }; +} + +function buildDurableFactsForSession({ session, events = [], status = null, finalText = null, runId = null, commandId = null, projectionStatus = null, projectionHealth = "healthy", lastProjectedSeq = null } = {}) { + const traceId = session.lastTraceId ?? session.session?.lastTraceId; + const projectedStatus = normalizeTestStatus(status ?? session.status); + const terminal = testTerminalStatuses.has(projectedStatus); + const normalizedMessages = normalizeTestMessages(session, { traceId, status: projectedStatus, finalText, terminal }); + const normalizedEvents = normalizeTestEvents(events, session, { traceId }); + const projectedSeq = Number.isFinite(Number(lastProjectedSeq)) ? Number(lastProjectedSeq) : Math.max(normalizedEvents.length, normalizedMessages.length, 0); + const checkpointStatus = projectionStatus ?? (terminal ? "caught_up" : "projecting"); + return { + sessions: [{ + sessionId: session.id, + ownerUserId: session.ownerUserId, + ownerRole: session.ownerRole ?? null, + projectId: session.projectId ?? null, + conversationId: session.conversationId ?? null, + threadId: session.threadId ?? null, + agentId: session.agentId ?? "hwlab-code-agent", + status: projectedStatus, + lastTraceId: traceId, + providerProfile: session.session?.providerProfile ?? null, + projectedSeq, + sourceSeq: projectedSeq, + sourceEventId: traceId, + terminal, + sealed: terminal, + createdAt: session.startedAt ?? session.createdAt ?? session.updatedAt ?? null, + updatedAt: session.updatedAt ?? session.endedAt ?? session.startedAt ?? null, + valuesRedacted: true + }], + messages: normalizedMessages, + parts: normalizedMessages.filter((message) => message.text).map((message, index) => ({ + partId: `prt_${session.id}_${index}`, + messageId: message.messageId, + sessionId: session.id, + turnId: message.turnId, + traceId: message.traceId, + partIndex: 0, + partType: "text", + status: message.status, + text: message.text, + projectedSeq: message.projectedSeq, + sourceSeq: message.sourceSeq, + sourceEventId: message.sourceEventId, + terminal: message.terminal, + sealed: message.sealed, + updatedAt: message.updatedAt, + valuesRedacted: true + })), + turns: traceId ? [{ + turnId: traceId, + sessionId: session.id, + traceId, + messageId: normalizedMessages.find((message) => message.role !== "user")?.messageId ?? null, + status: projectedStatus, + projectedSeq, + sourceSeq: projectedSeq, + sourceEventId: traceId, + terminal, + sealed: terminal, + finalResponse: finalText ? { text: finalText, status: projectedStatus, traceId, valuesPrinted: false } : null, + diagnostic: { projectionStatus: checkpointStatus, projectionHealth, valuesRedacted: true }, + createdAt: session.startedAt ?? session.updatedAt ?? null, + updatedAt: session.updatedAt ?? null, + valuesRedacted: true + }] : [], + traceEvents: normalizedEvents, + checkpoints: traceId ? [{ + traceId, + sessionId: session.id, + turnId: traceId, + runId, + commandId, + projectedSeq, + sourceSeq: projectedSeq, + sourceEventId: traceId, + projectionStatus: checkpointStatus, + projectionHealth, + terminal, + sealed: terminal, + diagnostic: { projectionStatus: checkpointStatus, projectionHealth, valuesRedacted: true }, + createdAt: session.startedAt ?? session.updatedAt ?? null, + updatedAt: session.updatedAt ?? null, + valuesRedacted: true + }] : [] + }; +} + +function normalizeTestMessages(session, { traceId, status, finalText, terminal }) { + const raw = Array.isArray(session.session?.messages) ? session.session.messages : []; + return raw.map((message, index) => { + const role = message.role ?? (index % 2 === 0 ? "user" : "agent"); + const isAssistant = role === "agent" || role === "assistant"; + const text = isAssistant && finalText ? finalText : String(message.text ?? message.content ?? ""); + const messageStatus = isAssistant && finalText ? status : normalizeTestStatus(message.status ?? (role === "user" ? "sent" : status)); + return { + messageId: message.messageId ?? `msg_${session.id}_${index}`, + sessionId: session.id, + turnId: message.turnId ?? traceId, + traceId: message.traceId ?? traceId, + role, + status: messageStatus, + projectedSeq: index + 1, + sourceSeq: index + 1, + sourceEventId: `${session.id}:message:${index}`, + terminal: isAssistant && terminal, + sealed: isAssistant && terminal, + text, + createdAt: message.createdAt ?? session.startedAt ?? session.updatedAt ?? null, + updatedAt: message.updatedAt ?? session.updatedAt ?? null, + valuesRedacted: true + }; + }); +} + +function normalizeTestEvents(events, session, { traceId }) { + return events.map((event, index) => { + const seq = Number.isFinite(Number(event.seq)) ? Number(event.seq) : index + 1; + const sourceSeq = Number.isFinite(Number(event.sourceSeq)) ? Number(event.sourceSeq) : seq; + return { + ...event, + id: event.id ?? `wte_${session.id}_${index}`, + traceId: event.traceId ?? traceId, + sessionId: event.sessionId ?? session.id, + turnId: event.turnId ?? traceId, + messageId: event.messageId ?? null, + sourceSeq, + sourceEventId: event.sourceEventId ?? `${traceId}:${sourceSeq}:${index}`, + projectedSeq: seq, + eventType: event.eventType ?? event.type ?? event.label ?? "event", + terminal: event.terminal === true, + sealed: event.terminal === true, + occurredAt: event.occurredAt ?? event.createdAt ?? session.updatedAt ?? null, + updatedAt: event.updatedAt ?? event.createdAt ?? session.updatedAt ?? null, + valuesRedacted: true + }; + }); +} + +function filterFacts(facts, params = {}) { + const filtered = { + sessions: facts.sessions.filter((record) => matchesFact(record, params, ["sessionId", "ownerUserId", "projectId", "conversationId", "threadId", "status"]) && (!params.traceId || record.lastTraceId === params.traceId)), + messages: facts.messages.filter((record) => matchesFact(record, params, ["messageId", "sessionId", "turnId", "traceId", "role", "status"])), + parts: facts.parts.filter((record) => matchesFact(record, params, ["partId", "messageId", "sessionId", "turnId", "traceId", "partType", "status"])), + turns: facts.turns.filter((record) => matchesFact(record, params, ["turnId", "sessionId", "traceId", "messageId", "status"])), + traceEvents: facts.traceEvents.filter((record) => matchesFact(record, params, ["id", "traceId", "sessionId", "turnId", "messageId", "eventType"])), + checkpoints: facts.checkpoints.filter((record) => matchesFact(record, params, ["traceId", "sessionId", "turnId", "runId", "commandId", "projectionStatus", "projectionHealth"])) + }; + const limit = Number.parseInt(params.limit ?? "", 10); + if (!Number.isInteger(limit) || limit <= 0) return filtered; + return Object.fromEntries(Object.entries(filtered).map(([key, value]) => [key, value.slice(0, limit)])); +} + +function matchesFact(record, params, fields) { + return fields.every((field) => params[field] === undefined || params[field] === null || record[field] === params[field]); +} + +function mergeFacts(target, source) { + for (const key of Object.keys(target)) target[key].push(...source[key]); + return target; +} + +function emptyFacts() { + return { sessions: [], messages: [], parts: [], turns: [], traceEvents: [], checkpoints: [] }; +} + +function normalizeTestStatus(value) { + const status = String(value ?? "").trim().toLowerCase().replace(/_/gu, "-"); + if (status === "active" || status === "busy" || status === "processing") return "running"; + if (status === "cancelled") return "canceled"; + return status || "unknown"; +} + +const testTerminalStatuses = new Set(["completed", "failed", "blocked", "timeout", "canceled", "idle"]); + async function getSseEvents(port, path, count) { const controller = new AbortController(); const response = await fetch(`http://127.0.0.1:${port}${path}`, { signal: controller.signal }); diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index f2caddad..97caea78 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -14,7 +14,7 @@ import { sendJson } from "./server-http-utils.ts"; import { createWorkbenchReadModel } from "./workbench-read-model.ts"; -import { createWorkbenchTurnProjection, RUNNING_STATUSES, TERMINAL_STATUSES } from "./workbench-turn-projection.ts"; +import { createWorkbenchTurnProjection, durableTraceStatus, RUNNING_STATUSES, TERMINAL_STATUSES } from "./workbench-turn-projection.ts"; const DEFAULT_PAGE_LIMIT = 50; const DEFAULT_SESSION_LIST_LIMIT = 20; @@ -318,14 +318,21 @@ async function handleWorkbenchSessionList(response, url, options, actor) { const includeSessionId = safeSessionId(url.searchParams.get("includeSessionId")); const includeRouteId = includeSessionId ?? safeConversationId(url.searchParams.get("includeSessionId")); const readModel = createWorkbenchReadModel(options, actor); - const naturalPage = await (perf ? perf.measure("workbench_session_page_query", () => readModel.listSessions({ limit: limit + 1, offset })) : readModel.listSessions({ limit: limit + 1, offset })); - const pageSessions = naturalPage.slice(0, limit); - let responseSessions = pageSessions; - if (includeRouteId && !pageSessions.some((session) => sessionMatchesRouteId(session, includeRouteId))) { - const included = await (perf ? perf.measure("workbench_session_include_query", () => readModel.getSessionByRouteId(includeRouteId)) : readModel.getSessionByRouteId(includeRouteId)); - if (included) responseSessions = [included, ...pageSessions]; + const pageResult = await (perf ? perf.measure("workbench_session_page_query", () => readModel.queryFacts({ ownerUserId: actor.role === "admin" ? undefined : actor.id, limit: limit + offset + 1 })) : readModel.queryFacts({ ownerUserId: actor.role === "admin" ? undefined : actor.id, limit: limit + offset + 1 })); + if (pageResult.error) return sendJson(response, 503, workbenchProjectionStoreError(pageResult.error)); + let pageFacts = pageResult.facts; + let naturalPage = visibleFactSessions(pageFacts, actor).slice(offset, offset + limit + 1); + if (includeRouteId && !naturalPage.some((session) => factSessionMatchesRouteId(session, includeRouteId))) { + const includeResult = await (perf ? perf.measure("workbench_session_include_query", () => queryFactsByRouteId(readModel, includeRouteId)) : queryFactsByRouteId(readModel, includeRouteId)); + if (includeResult.error) return sendJson(response, 503, workbenchProjectionStoreError(includeResult.error)); + const included = visibleFactSessions(includeResult.facts, actor)[0] ?? null; + if (included) { + pageFacts = mergeFactSets(pageFacts, includeResult.facts); + naturalPage = [included, ...naturalPage]; + } } - const summaries = await (perf ? perf.measure("workbench_session_summary_projection", () => sessionListSummaries(readModel, responseSessions, options)) : sessionListSummaries(readModel, responseSessions, options)); + const pageSessions = naturalPage.slice(0, limit); + const summaries = pageSessions.map((session) => factSessionSummary(session, pageFacts)).filter(Boolean); const hasMore = naturalPage.length > limit; sendJson(response, 200, { ok: true, @@ -348,16 +355,428 @@ function sessionMatchesRouteId(session, routeId) { return Boolean(conversationId && session?.conversationId === conversationId); } +async function queryFactsByRouteId(readModel, routeId) { + const sessionId = safeSessionId(routeId); + if (sessionId) return readModel.queryFacts({ sessionId, limit: MAX_PAGE_LIMIT }); + const conversationId = safeConversationId(routeId); + if (conversationId) return readModel.queryFacts({ conversationId, limit: MAX_PAGE_LIMIT }); + return { facts: emptyFactSet(), count: 0, durable: false, persistence: null }; +} + +function visibleFactSessions(facts, actor) { + return factArray(facts?.sessions) + .filter((session) => canActorReadFactSession(session, actor)) + .sort(compareFactSessionsDesc); +} + +function canActorReadFactSession(session, actor) { + if (!session || !actor) return false; + if (normalizeStatus(session.status) === "archived") return false; + if (actor.role === "admin") return true; + return session.ownerUserId === actor.id; +} + +function factSessionMatchesRouteId(session, routeId) { + const sessionId = safeSessionId(routeId); + if (sessionId && factSessionId(session) === sessionId) return true; + const conversationId = safeConversationId(routeId); + return Boolean(conversationId && session?.conversationId === conversationId); +} + +function factSessionSummary(session, facts = {}) { + const sessionId = factSessionId(session); + if (!sessionId) return null; + const traceId = factLastTraceId(session) ?? factLatestTraceIdForSession(facts, sessionId); + const projection = traceId ? factProjectionForTrace(facts, traceId) : null; + const turn = traceId ? factTurnForTrace(facts, traceId) : null; + const messages = factMessagesForSession(session, facts); + const status = normalizeStatus(turn?.status ?? session?.status); + return { + sessionId, + threadId: safeOpaqueId(session?.threadId) ?? (textValue(session?.threadId) || null), + agentId: session?.agentId ?? "hwlab-code-agent", + status, + running: RUNNING_STATUSES.has(status), + terminal: TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status), + lastTraceId: traceId, + projection, + projectionStatus: projection?.projectionStatus ?? null, + projectionHealth: projection?.projectionHealth ?? null, + staleMs: projection?.staleMs ?? null, + blocker: projection?.blocker ?? null, + providerProfile: textValue(session?.providerProfile ?? session?.sessionJson?.providerProfile) || null, + messageCount: messages.length, + firstUserMessagePreview: firstUserPreview(messages), + updatedAt: factUpdatedAt(session), + turnSummary: turn ? { + turnId: factTurnId(turn, traceId), + traceId, + status: normalizeStatus(turn.status ?? status), + running: RUNNING_STATUSES.has(normalizeStatus(turn.status ?? status)), + terminal: TERMINAL_STATUSES.has(normalizeStatus(turn.status ?? status)) && !RUNNING_STATUSES.has(normalizeStatus(turn.status ?? status)), + eventCount: projection?.lastProjectedSeq ?? factTraceSnapshot(facts, traceId).eventCount, + projection, + projectionStatus: projection?.projectionStatus ?? null, + projectionHealth: projection?.projectionHealth ?? null, + staleMs: projection?.staleMs ?? null, + blocker: projection?.blocker ?? null, + updatedAt: factUpdatedAt(turn) + } : null, + valuesRedacted: true + }; +} + +function factSessionDetail(session, facts = {}) { + const sessionId = factSessionId(session); + return { + ...factSessionSummary(session, facts), + metadata: { + startedAt: session?.startedAt ?? session?.createdAt ?? null, + endedAt: session?.endedAt ?? null, + ownerUserId: session?.ownerUserId ?? null + }, + messagePageUrl: sessionId ? `/v1/workbench/sessions/${encodeURIComponent(sessionId)}/messages` : null, + valuesRedacted: true, + secretMaterialStored: false + }; +} + +function factMessagesForSession(session, facts = {}) { + const sessionId = factSessionId(session); + if (!sessionId) return []; + const partsByMessageId = new Map(); + for (const part of factArray(facts.parts)) { + if (part.sessionId !== sessionId) continue; + const messageId = textValue(part.messageId); + if (!messageId) continue; + const existing = partsByMessageId.get(messageId) ?? []; + existing.push(part); + partsByMessageId.set(messageId, existing); + } + return factArray(facts.messages) + .filter((message) => message.sessionId === sessionId) + .sort(compareFactMessagesAsc) + .map((message) => factMessageDto(message, partsByMessageId.get(message.messageId) ?? [])) + .filter(Boolean); +} + +function factMessageDto(message, parts = []) { + const messageId = safeMessageId(message?.messageId) || textValue(message?.messageId); + if (!messageId) return null; + const traceId = safeTraceId(message?.traceId) ?? null; + const text = projectionText(message?.text, message?.content, message?.message, message?.finalResponse); + const normalizedParts = parts.length > 0 + ? [...parts].sort(compareFactPartsAsc).map((part) => factPartDto(part, messageId, traceId)).filter(Boolean) + : text ? [partFact({ type: "text", text, status: message?.status }, 0, messageId, traceId)] : []; + return { + messageId, + role: textValue(message?.role) || "agent", + sessionId: (safeSessionId(message?.sessionId) ?? textValue(message?.sessionId)) || null, + traceId, + turnId: safeTurnId(message?.turnId) || traceId, + status: normalizeStatus(message?.status), + parts: normalizedParts, + text: text || "", + textPreview: text ? text.slice(0, 240) : null, + createdAt: textValue(message?.createdAt) || null, + updatedAt: factUpdatedAt(message), + projectionStatus: null, + projectionHealth: null, + valuesRedacted: message?.valuesRedacted !== false + }; +} + +function factPartDto(part, messageId, traceId) { + const partId = safePartId(part?.partId) || textValue(part?.partId) || `prt_${hash(`${messageId}:${part?.partIndex ?? 0}:${part?.partType ?? part?.type ?? "text"}`).slice(0, 24)}`; + const text = projectionText(part?.text, part?.content, part?.message); + return { + partId, + messageId, + traceId: safeTraceId(part?.traceId) ?? traceId, + type: textValue(part?.partType ?? part?.type) || "text", + text: text || null, + status: normalizeStatus(part?.status), + toolName: textValue(part?.toolName ?? part?.name) || null, + createdAt: textValue(part?.createdAt ?? part?.occurredAt) || null, + valuesRedacted: part?.valuesRedacted !== false + }; +} + +function factTurnForTrace(facts = {}, traceId, turnId = null) { + const safeTrace = safeTraceId(traceId); + if (!safeTrace) return null; + const requestedTurn = safeTurnId(turnId); + const matches = factArray(facts.turns) + .filter((turn) => turn.traceId === safeTrace && (!requestedTurn || turn.turnId === requestedTurn || turn.turnId === safeTrace)) + .sort(compareFactRecordsDesc); + return matches[0] ?? null; +} + +function factTurnSnapshot({ turn = null, session = null, facts = {}, traceId, turnId = null } = {}) { + const safeTrace = safeTraceId(traceId); + const resolvedTurnId = factTurnId(turn, safeTrace) ?? safeTurnId(turnId) ?? safeTrace; + const status = normalizeStatus(turn?.status ?? session?.status); + const messages = session ? factMessagesForSession(session, facts) : []; + const userMessage = messages.find((message) => message.role === "user") ?? null; + const assistantMessage = [...messages].reverse().find((message) => isAssistantLikeRole(message.role)) ?? null; + const checkpoint = safeTrace ? factCheckpointForTrace(facts, safeTrace) : null; + const finalText = projectionText(turn?.finalResponse, turn?.assistantText, turn?.text, assistantMessage?.text); + const trace = factTraceSnapshot(facts, safeTrace); + return { + turnId: resolvedTurnId, + traceId: safeTrace, + status, + running: RUNNING_STATUSES.has(status), + terminal: TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status), + sessionId: factSessionId(session) ?? turn?.sessionId ?? checkpoint?.sessionId ?? null, + threadId: safeOpaqueId(session?.threadId) ?? (textValue(session?.threadId) || null), + userMessageId: userMessage?.messageId ?? null, + assistantMessageId: assistantMessage?.messageId ?? turn?.messageId ?? null, + assistantText: finalText ?? null, + finalResponse: turn?.finalResponse ?? (finalText ? { text: finalText, status, traceId: safeTrace, valuesPrinted: false } : null), + agentRun: checkpoint ? { + runId: textValue(checkpoint.runId) || null, + commandId: textValue(checkpoint.commandId) || null, + status, + lastSeq: factSeq(checkpoint), + valuesRedacted: true + } : null, + trace: { + traceId: safeTrace, + status: trace.status, + eventCount: trace.eventCount, + updatedAt: trace.updatedAt + }, + urls: { + self: resolvedTurnId ? `/v1/workbench/turns/${encodeURIComponent(resolvedTurnId)}` : null, + traceEvents: safeTrace ? `/v1/workbench/traces/${encodeURIComponent(safeTrace)}/events` : null + } + }; +} + +function factTraceSnapshot(facts = {}, traceId) { + const safeTrace = safeTraceId(traceId); + const events = factArray(facts.traceEvents) + .filter((event) => !safeTrace || event.traceId === safeTrace) + .sort(compareFactTraceEventsAsc) + .map(factTraceEventDto); + const status = durableTraceStatus(events); + const lastEvent = events.at(-1) ?? null; + return { + traceId: safeTrace, + status, + eventCount: events.length, + events, + lastEvent, + updatedAt: lastEvent?.updatedAt ?? lastEvent?.createdAt ?? factCheckpointForTrace(facts, safeTrace)?.updatedAt ?? null, + valuesRedacted: true, + secretMaterialStored: false + }; +} + +function factTraceEventDto(event, index) { + const seq = Number.isFinite(Number(event?.seq)) && Number(event.seq) > 0 ? Math.trunc(Number(event.seq)) : factSeq(event) || index + 1; + return { + ...event, + id: textValue(event?.id) || textValue(event?.sourceEventId) || null, + seq, + sourceSeq: factSeq(event) || seq, + type: textValue(event?.type ?? event?.eventType) || "event", + label: textValue(event?.label ?? event?.eventType ?? event?.type) || null, + status: normalizeStatus(event?.status), + createdAt: textValue(event?.createdAt ?? event?.occurredAt) || null, + updatedAt: factUpdatedAt(event), + valuesRedacted: event?.valuesRedacted !== false + }; +} + +function factProjectionForTrace(facts = {}, traceId) { + const checkpoint = factCheckpointForTrace(facts, traceId); + if (!checkpoint) { + return { + projectionStatus: "unknown", + projectionHealth: "unknown", + lastProjectedSeq: null, + sourceRunId: null, + sourceCommandId: null, + staleMs: null, + blocker: null, + updatedAt: null, + valuesRedacted: true + }; + } + const projectionStatus = normalizeProjectionStatus(checkpoint.projectionStatus); + const projectionHealth = normalizeProjectionHealth(checkpoint.projectionHealth, projectionStatus); + const diagnostic = objectValue(checkpoint.diagnostic); + return { + projectionStatus, + projectionHealth, + lastProjectedSeq: factSeq(checkpoint), + sourceRunId: textValue(checkpoint.runId ?? checkpoint.sourceRunId) || null, + sourceCommandId: textValue(checkpoint.commandId ?? checkpoint.sourceCommandId) || null, + staleMs: null, + blocker: diagnostic.blocker ?? checkpoint.blocker ?? null, + updatedAt: factUpdatedAt(checkpoint), + valuesRedacted: true + }; +} + +function factCheckpointForTrace(facts = {}, traceId) { + const safeTrace = safeTraceId(traceId); + if (!safeTrace) return null; + return factArray(facts.checkpoints) + .filter((checkpoint) => checkpoint.traceId === safeTrace) + .sort(compareFactRecordsDesc)[0] ?? null; +} + +function workbenchProjectionStoreError(error) { + return workbenchError("projection_store_unavailable", "Workbench durable projection store is unavailable.", { causeCode: error?.code ?? null }); +} + +function factSessionId(session) { + return safeSessionId(session?.sessionId ?? session?.id) ?? null; +} + +function factLastTraceId(session) { + return safeTraceId(session?.lastTraceId ?? session?.traceId) ?? null; +} + +function factLatestTraceIdForSession(facts = {}, sessionId) { + const turn = factArray(facts.turns).filter((item) => item.sessionId === sessionId).sort(compareFactRecordsDesc)[0] ?? null; + return safeTraceId(turn?.traceId) ?? null; +} + +function factTurnId(turn, fallbackTraceId = null) { + return safeTurnId(turn?.turnId) ?? safeTraceId(turn?.turnId) ?? safeTraceId(fallbackTraceId) ?? null; +} + +function factUpdatedAt(record) { + return textValue(record?.updatedAt ?? record?.occurredAt ?? record?.createdAt) || null; +} + +function factSeq(record) { + for (const value of [record?.projectedSeq, record?.sourceSeq, record?.seq]) { + const parsed = Number(value); + if (Number.isFinite(parsed) && parsed >= 0) return Math.trunc(parsed); + } + return null; +} + +function normalizeProjectionStatus(value) { + const status = normalizeStatus(value); + if (status === "caught-up") return "caught-up"; + if (status === "caughtup") return "caught-up"; + return ["projecting", "blocked", "stalled", "unknown"].includes(status) ? status : "unknown"; +} + +function normalizeProjectionHealth(value, projectionStatus = "unknown") { + const health = normalizeStatus(value); + if (health === "healthy" && projectionStatus === "caught-up") return "caught-up"; + if (health === "healthy" && projectionStatus === "projecting") return "projecting"; + if (["caught-up", "projecting", "degraded", "stalled", "unavailable", "unknown"].includes(health)) return health; + return projectionStatus === "caught-up" || projectionStatus === "projecting" ? projectionStatus : "unknown"; +} + +function compareFactSessionsDesc(left, right) { + return compareTimestampDesc(factUpdatedAt(left), factUpdatedAt(right)) || compareText(factSessionId(left), factSessionId(right)); +} + +function compareFactMessagesAsc(left, right) { + return compareNumberAsc(factSeq(left), factSeq(right)) || compareTimestampAsc(left?.createdAt ?? left?.updatedAt, right?.createdAt ?? right?.updatedAt) || compareText(left?.messageId, right?.messageId); +} + +function compareFactPartsAsc(left, right) { + return compareNumberAsc(left?.partIndex, right?.partIndex) || compareNumberAsc(factSeq(left), factSeq(right)) || compareText(left?.partId, right?.partId); +} + +function compareFactTraceEventsAsc(left, right) { + return compareNumberAsc(factSeq(left), factSeq(right)) || compareTimestampAsc(left?.occurredAt ?? left?.createdAt, right?.occurredAt ?? right?.createdAt) || compareText(left?.id, right?.id); +} + +function compareFactRecordsDesc(left, right) { + return compareNumberDesc(factSeq(left), factSeq(right)) || compareTimestampDesc(factUpdatedAt(left), factUpdatedAt(right)); +} + +function compareTimestampAsc(left, right) { + return timestampMs(left) - timestampMs(right); +} + +function compareTimestampDesc(left, right) { + return timestampMs(right) - timestampMs(left); +} + +function compareNumberAsc(left, right) { + return numericValue(left, Number.MAX_SAFE_INTEGER) - numericValue(right, Number.MAX_SAFE_INTEGER); +} + +function compareNumberDesc(left, right) { + return numericValue(right, -1) - numericValue(left, -1); +} + +function compareText(left, right) { + return textValue(left).localeCompare(textValue(right)); +} + +function timestampMs(value) { + const parsed = Date.parse(String(value ?? "")); + return Number.isFinite(parsed) ? parsed : 0; +} + +function numericValue(value, fallback) { + const parsed = Number(value); + return Number.isFinite(parsed) ? parsed : fallback; +} + +function mergeFactSets(...sets) { + const merged = emptyFactSet(); + const keys = { + sessions: "sessionId", + messages: "messageId", + parts: "partId", + turns: "turnId", + traceEvents: "id", + checkpoints: "traceId" + }; + for (const set of sets) { + for (const [key, idKey] of Object.entries(keys)) { + const seen = new Set(merged[key].map((item) => textValue(item?.[idKey]))); + for (const item of factArray(set?.[key])) { + const id = textValue(item?.[idKey]); + if (id && seen.has(id)) continue; + merged[key].push(item); + if (id) seen.add(id); + } + } + } + return merged; +} + +function emptyFactSet() { + return { + sessions: [], + messages: [], + parts: [], + turns: [], + traceEvents: [], + checkpoints: [] + }; +} + +function factArray(value) { + return Array.isArray(value) ? value : []; +} + async function handleWorkbenchSessionDetail(response, options, actor, sessionId) { const readModel = createWorkbenchReadModel(options, actor); - const session = await readModel.getSessionByRouteId(sessionId); + const result = await queryFactsByRouteId(readModel, sessionId); + if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error)); + const session = visibleFactSessions(result.facts, actor)[0] ?? null; if (!session) return sendJson(response, 404, workbenchError("workbench_session_not_found", "Workbench session is not visible to the current actor.", { sessionId })); - const projectionOptions = await sessionProjectionOptions(readModel, session, options); sendJson(response, 200, { ok: true, status: "found", contractVersion: "workbench-session-detail-v1", - session: sessionDetail(session, projectionOptions), + session: factSessionDetail(session, result.facts), valuesRedacted: true, secretMaterialStored: false }); @@ -365,19 +784,21 @@ async function handleWorkbenchSessionDetail(response, options, actor, sessionId) async function handleWorkbenchMessagePage(response, url, options, actor, sessionId) { const readModel = createWorkbenchReadModel(options, actor); - const session = await readModel.getSessionByRouteId(sessionId); + const result = await queryFactsByRouteId(readModel, sessionId); + if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error)); + const session = visibleFactSessions(result.facts, actor)[0] ?? null; if (!session) return sendJson(response, 404, workbenchError("workbench_session_not_found", "Workbench session is not visible to the current actor.", { sessionId })); - const projectionOptions = await sessionProjectionOptions(readModel, session, options); - const messages = sessionMessages(session, projectionOptions); + const messages = factMessagesForSession(session, result.facts); const limit = boundedLimit(url.searchParams.get("limit")); const offset = cursorOffset(url.searchParams.get("cursor") ?? url.searchParams.get("after")); const page = messages.slice(offset, offset + limit); const nextOffset = offset + page.length; + const resolvedSessionId = factSessionId(session); sendJson(response, 200, { ok: true, status: "succeeded", contractVersion: "workbench-message-page-v1", - sessionId: session.id, + sessionId: resolvedSessionId, messages: page, count: page.length, total: messages.length, @@ -400,30 +821,24 @@ async function handleWorkbenchTurnSnapshot(response, url, options, actor, rawTur return sendJson(response, 400, body); } const readModel = createWorkbenchReadModel(options, actor); - const result = readModel.resultForTrace(traceId); - const session = await readModel.getSessionByTraceId(traceId); - if (result?.ownerUserId && !readModel.canReadOwner(result.ownerUserId)) { - const body = workbenchError("agent_session_owner_required", "Only the session owner or admin can read this Workbench turn.", { traceId }); - recordWorkbenchTurnReadMetric(options, url, 403, body, startedAt); - return sendJson(response, 403, body); - } - const resultVisible = result && (session || actor.role === "admin" || Boolean(result.ownerUserId)); - const trace = await readModel.traceSnapshot(traceId); - const turnProjection = createWorkbenchTurnProjection({ turnId, traceId, result, session, trace }); - const status = turnProjection.status; - const found = Boolean(resultVisible || session); + const result = await readModel.queryFacts({ traceId, limit: MAX_PAGE_LIMIT }); + if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error)); + const session = visibleFactSessions(result.facts, actor).find((item) => item.lastTraceId === traceId) ?? visibleFactSessions(result.facts, actor)[0] ?? null; + const turn = factTurnForTrace(result.facts, traceId, turnId); + const found = Boolean(session || turn); if (!found) { const body = workbenchError("workbench_turn_not_found", "Workbench turn is not visible to the current actor.", { turnId, traceId }); recordWorkbenchTurnReadMetric(options, url, 404, body, startedAt); return sendJson(response, 404, body); } - const projection = readModel.projectionDiagnostics({ traceId, result, trace, projection: turnProjection }); - recordWorkbenchProjectionMetric(options, { result, trace, projection: turnProjection, diagnostic: projection }); + const projection = factProjectionForTrace(result.facts, traceId); + const snapshot = factTurnSnapshot({ turn, session, facts: result.facts, traceId, turnId }); + const status = snapshot.status; const body = { ok: true, status, contractVersion: "workbench-turn-snapshot-v1", - turn: turnSnapshot({ projection: turnProjection, result, session, trace }), + turn: snapshot, projection, projectionStatus: projection.projectionStatus, projectionHealth: projection.projectionHealth, @@ -435,6 +850,10 @@ async function handleWorkbenchTurnSnapshot(response, url, options, actor, rawTur valuesRedacted: true, secretMaterialStored: false }; + recordWorkbenchProjectionMetric(options, { + projection: { ...projection, status, eventCount: snapshot.trace?.eventCount ?? projection.lastProjectedSeq ?? 0 }, + diagnostic: projection + }); recordWorkbenchTurnReadMetric(options, url, 200, body, startedAt); sendJson(response, 200, body); } @@ -443,27 +862,21 @@ async function handleWorkbenchTraceEventPage(response, url, options, actor, rawT const traceId = safeTraceId(rawTraceId); if (!traceId) return sendJson(response, 400, workbenchError("invalid_trace_id", "traceId must start with trc_.", { traceId: rawTraceId })); const readModel = createWorkbenchReadModel(options, actor); - const result = readModel.resultForTrace(traceId); - const session = await readModel.getSessionByTraceId(traceId); - if (result?.ownerUserId && !readModel.canReadOwner(result.ownerUserId)) { - return sendJson(response, 403, workbenchError("agent_session_owner_required", "Only the session owner or admin can read this Workbench trace.", { traceId })); - } - const resultVisible = result && (session || actor.role === "admin" || Boolean(result.ownerUserId)); - if (!resultVisible && !session) { + const result = await readModel.queryFacts({ traceId, limit: MAX_PAGE_LIMIT }); + if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error)); + const session = visibleFactSessions(result.facts, actor).find((item) => item.lastTraceId === traceId) ?? visibleFactSessions(result.facts, actor)[0] ?? null; + if (!session && result.facts.traceEvents.length === 0) { return sendJson(response, 404, workbenchError("workbench_trace_not_found", "Workbench trace is not visible to the current actor.", { traceId })); } - const trace = await readModel.traceSnapshot(traceId); - const page = traceEventPage(trace, tracePageOptions(url)); - const turnProjection = createWorkbenchTurnProjection({ traceId, result, session, trace }); - const projection = readModel.projectionDiagnostics({ traceId, result, trace, projection: turnProjection }); - recordWorkbenchProjectionMetric(options, { result, trace, projection: turnProjection, diagnostic: projection }); + const page = traceEventPage(factTraceSnapshot(result.facts, traceId), tracePageOptions(url)); + const projection = factProjectionForTrace(result.facts, traceId); sendJson(response, 200, { ok: true, status: "succeeded", contractVersion: "workbench-trace-events-v1", traceId, - sessionId: safeSessionId(result?.sessionId ?? result?.session?.sessionId ?? session?.id) ?? null, - threadId: safeOpaqueId(result?.threadId ?? result?.session?.threadId ?? session?.threadId) ?? (textValue(session?.threadId) || null), + sessionId: factSessionId(session), + threadId: safeOpaqueId(session?.threadId) ?? (textValue(session?.threadId) || null), ...page, projection, projectionStatus: projection.projectionStatus, diff --git a/internal/cloud/workbench-facts-store.ts b/internal/cloud/workbench-facts-store.ts index 9f029d6d..7e26aa0e 100644 --- a/internal/cloud/workbench-facts-store.ts +++ b/internal/cloud/workbench-facts-store.ts @@ -16,6 +16,34 @@ export function createWorkbenchFactsStore(options = {}, actor = null) { const runtimeStore = options.runtimeStore ?? null; const results = options.codeAgentChatResults ?? null; + async function queryFacts(params = {}) { + if (typeof runtimeStore?.queryWorkbenchFacts === "function") { + try { + const result = await runtimeStore.queryWorkbenchFacts(params); + return { + facts: normalizeWorkbenchFactsResult(result?.facts), + count: Number.isFinite(Number(result?.count)) ? Number(result.count) : 0, + durable: true, + persistence: result?.persistence ?? null + }; + } catch (error) { + return { + facts: emptyWorkbenchFacts(), + count: 0, + durable: true, + error, + projection: projectionStoreUnavailableTrace(safeTraceId(params.traceId) ?? "trc_unassigned", error) + }; + } + } + return { + facts: emptyWorkbenchFacts(), + count: 0, + durable: false, + persistence: null + }; + } + async function getSessionById(sessionId) { const safeId = safeSessionId(sessionId); if (!safeId) return null; @@ -96,6 +124,7 @@ export function createWorkbenchFactsStore(options = {}, actor = null) { } return { + queryFacts, getSessionById, getSessionByRouteId, getSessionByTraceId, @@ -250,3 +279,22 @@ function normalizeStatus(value) { if (text === "cancelled") return "canceled"; return text || "unknown"; } + +function normalizeWorkbenchFactsResult(facts = {}) { + const normalized = emptyWorkbenchFacts(); + for (const key of Object.keys(normalized)) { + normalized[key] = Array.isArray(facts?.[key]) ? facts[key].map((item) => ({ ...item })) : []; + } + return normalized; +} + +function emptyWorkbenchFacts() { + return { + sessions: [], + messages: [], + parts: [], + turns: [], + traceEvents: [], + checkpoints: [] + }; +} diff --git a/internal/cloud/workbench-read-model.ts b/internal/cloud/workbench-read-model.ts index 7219bdc3..9fc1db8f 100644 --- a/internal/cloud/workbench-read-model.ts +++ b/internal/cloud/workbench-read-model.ts @@ -9,6 +9,7 @@ export function createWorkbenchReadModel(options = {}, actor = null) { const facts = createWorkbenchFactsStore(options, actor); return { facts, + queryFacts: (input = {}) => facts.queryFacts(input), listSessions: (input = {}) => facts.listSessions(input), getSessionById: (sessionId) => facts.getSessionById(sessionId), getSessionByRouteId: (routeId) => facts.getSessionByRouteId(routeId), diff --git a/internal/db/runtime-store.ts b/internal/db/runtime-store.ts index 01bc92a0..26d66788 100644 --- a/internal/db/runtime-store.ts +++ b/internal/db/runtime-store.ts @@ -604,7 +604,7 @@ export class CloudRuntimeStore { } queryWorkbenchFacts(params = {}) { - const sessions = [...this.workbenchSessions.values()].filter((record) => matchesQuery(record, params, ["sessionId", "ownerUserId", "projectId", "conversationId", "threadId", "status"])); + const sessions = [...this.workbenchSessions.values()].filter((record) => matchesWorkbenchSessionFactQuery(record, params)); const messages = [...this.workbenchMessages.values()].filter((record) => matchesQuery(record, params, ["messageId", "sessionId", "turnId", "traceId", "role", "status"])); const parts = [...this.workbenchParts.values()].filter((record) => matchesQuery(record, params, ["partId", "messageId", "sessionId", "turnId", "traceId", "partType", "status"])); const turns = [...this.workbenchTurns.values()].filter((record) => matchesQuery(record, params, ["turnId", "sessionId", "traceId", "messageId", "status"])); @@ -1213,7 +1213,7 @@ export class PostgresCloudRuntimeStore { async queryWorkbenchFacts(params = {}) { await this.assertReadyForDurableReads("workbench.facts.query"); const facts = { - sessions: await this.queryWorkbenchFactRows("workbench.sessions.query", "workbench_sessions", "session_json", params, { sessionId: "session_id", ownerUserId: "owner_user_id", projectId: "project_id", conversationId: "conversation_id", threadId: "thread_id", status: "status" }), + sessions: await this.queryWorkbenchFactRows("workbench.sessions.query", "workbench_sessions", "session_json", params, { sessionId: "session_id", ownerUserId: "owner_user_id", projectId: "project_id", conversationId: "conversation_id", threadId: "thread_id", traceId: "last_trace_id", lastTraceId: "last_trace_id", status: "status" }), messages: await this.queryWorkbenchFactRows("workbench.messages.query", "workbench_messages", "message_json", params, { messageId: "message_id", sessionId: "session_id", turnId: "turn_id", traceId: "trace_id", role: "role", status: "status" }), parts: await this.queryWorkbenchFactRows("workbench.parts.query", "workbench_parts", "part_json", params, { partId: "part_id", messageId: "message_id", sessionId: "session_id", turnId: "turn_id", traceId: "trace_id", partType: "part_type", status: "status" }), turns: await this.queryWorkbenchFactRows("workbench.turns.query", "workbench_turns", "turn_json", params, { turnId: "turn_id", sessionId: "session_id", traceId: "trace_id", messageId: "message_id", status: "status" }), @@ -2166,6 +2166,13 @@ function matchesQuery(record, params, fields) { return true; } +function matchesWorkbenchSessionFactQuery(record, params) { + if (!matchesQuery(record, params, ["sessionId", "ownerUserId", "projectId", "conversationId", "threadId", "status"])) return false; + const traceId = textOr(params.traceId ?? params.lastTraceId, ""); + if (traceId && record.lastTraceId !== traceId && record.traceId !== traceId) return false; + return true; +} + function limitResults(records, limit) { const parsed = Number.parseInt(limit ?? "", 10); if (!Number.isInteger(parsed) || parsed <= 0) {