diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index 087bd100..0f3eae12 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -42,28 +42,28 @@ export async function handleWorkbenchReadModelHttp(request, response, url, optio const sessionMessagesMatch = url.pathname.match(/^\/v1\/workbench\/sessions\/([^/]+)\/messages$/u); if (sessionMessagesMatch) { if (request.method !== "GET") return methodNotAllowed(response, "GET"); - await (perf ? perf.measure("workbench_message_page", () => handleWorkbenchMessagePage(response, url, options, auth.actor, decodeURIComponent(sessionMessagesMatch[1]))) : handleWorkbenchMessagePage(response, url, options, auth.actor, decodeURIComponent(sessionMessagesMatch[1]))); + await (perf ? perf.measure("workbench_message_page", () => handleWorkbenchMessagePage(request, response, url, options, auth.actor, decodeURIComponent(sessionMessagesMatch[1]))) : handleWorkbenchMessagePage(request, response, url, options, auth.actor, decodeURIComponent(sessionMessagesMatch[1]))); return; } const sessionMatch = url.pathname.match(/^\/v1\/workbench\/sessions\/([^/]+)$/u); if (sessionMatch) { if (request.method !== "GET") return methodNotAllowed(response, "GET"); - await (perf ? perf.measure("workbench_session_detail", () => handleWorkbenchSessionDetail(response, options, auth.actor, decodeURIComponent(sessionMatch[1]))) : handleWorkbenchSessionDetail(response, options, auth.actor, decodeURIComponent(sessionMatch[1]))); + await (perf ? perf.measure("workbench_session_detail", () => handleWorkbenchSessionDetail(request, response, options, auth.actor, decodeURIComponent(sessionMatch[1]))) : handleWorkbenchSessionDetail(request, response, options, auth.actor, decodeURIComponent(sessionMatch[1]))); return; } const turnMatch = url.pathname.match(/^\/v1\/workbench\/turns\/([^/]+)$/u); if (turnMatch) { if (request.method !== "GET") return methodNotAllowed(response, "GET"); - await (perf ? perf.measure("workbench_turn_snapshot", () => handleWorkbenchTurnSnapshot(response, url, options, auth.actor, decodeURIComponent(turnMatch[1]))) : handleWorkbenchTurnSnapshot(response, url, options, auth.actor, decodeURIComponent(turnMatch[1]))); + await (perf ? perf.measure("workbench_turn_snapshot", () => handleWorkbenchTurnSnapshot(request, response, url, options, auth.actor, decodeURIComponent(turnMatch[1]))) : handleWorkbenchTurnSnapshot(request, response, url, options, auth.actor, decodeURIComponent(turnMatch[1]))); return; } const traceEventsMatch = url.pathname.match(/^\/v1\/workbench\/traces\/([^/]+)\/events$/u); if (traceEventsMatch) { if (request.method !== "GET") return methodNotAllowed(response, "GET"); - await (perf ? perf.measure("workbench_trace_events", () => handleWorkbenchTraceEventPage(response, url, options, auth.actor, decodeURIComponent(traceEventsMatch[1]))) : handleWorkbenchTraceEventPage(response, url, options, auth.actor, decodeURIComponent(traceEventsMatch[1]))); + await (perf ? perf.measure("workbench_trace_events", () => handleWorkbenchTraceEventPage(request, response, url, options, auth.actor, decodeURIComponent(traceEventsMatch[1]))) : handleWorkbenchTraceEventPage(request, response, url, options, auth.actor, decodeURIComponent(traceEventsMatch[1]))); return; } @@ -499,6 +499,19 @@ function runtimeDependencyError(code, message, retryable = true) { return error; } +function workbenchRuntimeFactsReadModel(request, options, actor) { + const runtime = options.workbenchRuntime ?? createWorkbenchRuntimeClient({ env: options.env ?? process.env, fetch: options.fetch, logger: options.logger, traceparent: request?.hwlabHttpRequestContext?.traceparent }); + if (typeof runtime?.queryWorkbenchFacts !== "function") { + throw runtimeDependencyError("workbench_runtime_unconfigured", "Workbench runtime service is required for Workbench facts reads.", false); + } + return { + queryFacts: (query = {}) => runtime.queryWorkbenchFacts({ + ...query, + actor: { id: actor.id, role: actor.role ?? "user", valuesRedacted: true } + }) + }; +} + async function handleWorkbenchSessionListBunLegacy(response, url, options, actor) { if (url.searchParams.has("projectId") || url.searchParams.has("workspaceId")) return sendJson(response, 400, workbenchError("workbench_authority_removed", "Workbench session list is keyed by sessionId only.")); const perf = options.backendPerformance; @@ -1138,8 +1151,8 @@ function uniqueText(values = []) { return [...new Set(values.map((value) => textValue(value)).filter(Boolean))]; } -async function handleWorkbenchSessionDetail(response, options, actor, sessionId) { - const readModel = createWorkbenchReadModel(options, actor); +async function handleWorkbenchSessionDetail(request, response, options, actor, sessionId) { + const readModel = workbenchRuntimeFactsReadModel(request, options, actor); const result = await queryFactsByRouteId(readModel, sessionId, { families: WORKBENCH_SESSION_DETAIL_FAMILIES }); if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error)); const session = visibleFactSessions(result.facts, actor)[0] ?? null; @@ -1154,8 +1167,8 @@ async function handleWorkbenchSessionDetail(response, options, actor, sessionId) }); } -async function handleWorkbenchMessagePage(response, url, options, actor, sessionId) { - const readModel = createWorkbenchReadModel(options, actor); +async function handleWorkbenchMessagePage(request, response, url, options, actor, sessionId) { + const readModel = workbenchRuntimeFactsReadModel(request, options, actor); const result = await queryFactsByRouteId(readModel, sessionId, { families: WORKBENCH_SESSION_MESSAGE_PAGE_FAMILIES }); if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error)); const session = visibleFactSessions(result.facts, actor)[0] ?? null; @@ -1182,7 +1195,7 @@ async function handleWorkbenchMessagePage(response, url, options, actor, session }); } -async function handleWorkbenchTurnSnapshot(response, url, options, actor, rawTurnId) { +async function handleWorkbenchTurnSnapshot(request, response, url, options, actor, rawTurnId) { const startedAt = nowMs(); const queryTraceId = safeTraceId(url.searchParams.get("traceId")); const traceId = safeTraceId(rawTurnId) ?? queryTraceId; @@ -1192,7 +1205,7 @@ async function handleWorkbenchTurnSnapshot(response, url, options, actor, rawTur recordWorkbenchTurnReadMetric(options, url, 400, body, startedAt); return sendJson(response, 400, body); } - const readModel = createWorkbenchReadModel(options, actor); + const readModel = workbenchRuntimeFactsReadModel(request, options, actor); const result = await readModel.queryFacts({ traceId, limit: MAX_PAGE_LIMIT, families: WORKBENCH_TRACE_EVENT_METADATA_FAMILIES }); if (result.error) { const body = workbenchProjectionStoreError(result.error); @@ -1236,11 +1249,11 @@ async function handleWorkbenchTurnSnapshot(response, url, options, actor, rawTur sendJson(response, 200, body); } -async function handleWorkbenchTraceEventPage(response, url, options, actor, rawTraceId) { +async function handleWorkbenchTraceEventPage(request, response, url, options, actor, rawTraceId) { const startedAt = nowMs(); 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 readModel = workbenchRuntimeFactsReadModel(request, options, actor); const pageOptions = tracePageOptions(url); const metadata = await readModel.queryFacts({ traceId, families: WORKBENCH_TRACE_EVENT_METADATA_FAMILIES, limit: 1 }); if (metadata.error) return sendJson(response, 503, workbenchProjectionStoreError(metadata.error));