fix(workbench): read durable facts from runtime service (#1960)

This commit is contained in:
Lyon
2026-06-23 10:05:39 +08:00
committed by GitHub
parent 56acbf7f36
commit 78a57dc98c
+25 -12
View File
@@ -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));