Merge pull request #2218 from pikasTech/issue-1078-otel-read-spans

obs(workbench): add OTel read and auth spans for #1078
This commit is contained in:
Lyon
2026-06-27 03:58:08 +08:00
committed by GitHub
3 changed files with 319 additions and 20 deletions
+86 -8
View File
@@ -204,8 +204,7 @@ class AccessController {
return;
}
if (request.method === "GET" && url.pathname === "/auth/session") {
const payload = await this.sessionPayload(request);
sendJson(response, payload.authenticated ? 200 : 401, payload);
await this.handleSessionRoute(request, response);
return;
}
@@ -386,6 +385,47 @@ class AccessController {
return sessionResponseFromAuth(auth, { setupRequired });
}
async handleSessionRoute(request, response) {
const otelContext = authLoginOtelTraceContext({ traceparent: getHeader(request, "traceparent") });
setAuthTraceHeaders(response, otelContext);
const startedAtMs = Date.now();
let httpStatus = 500;
let outcome = "error";
let errorCode = "";
try {
const payload = await this.sessionPayload(request);
httpStatus = payload.authenticated ? 200 : 401;
outcome = payload.authenticated ? "authenticated" : "unauthenticated";
errorCode = payload.error?.code ? textOr(payload.error.code, "") : "";
sendJson(response, httpStatus, payload);
return;
} catch (error) {
httpStatus = Number(error?.statusCode ?? error?.status ?? httpStatus ?? 500);
errorCode = textOr(error?.code, "auth_session_error");
throw error;
} finally {
void emitAuthOtelSpan("auth.session", otelContext, this.env, {
spanId: otelContext.rootSpanId,
parentSpanId: otelContext.parentSpanId,
kind: 2,
startTimeMs: startedAtMs,
endTimeMs: Date.now(),
status: httpStatus >= 500 ? "error" : "ok",
error: httpStatus >= 500 || (errorCode && httpStatus >= 500) ? new Error(errorCode || "auth_session_error") : null,
attributes: {
"http.method": "GET",
"http.route": "/auth/session",
"http.status_code": httpStatus,
"auth.outcome": outcome,
"auth.authenticated": outcome === "authenticated",
"auth.trace_visible": true,
"auth.values_redacted": true,
...(errorCode ? { "error.code": errorCode } : {})
}
});
}
}
async accessStatusPayload(request) {
const session = await this.sessionPayload(request);
return {
@@ -1142,14 +1182,52 @@ class AccessController {
}
async handleLogout(request, response) {
const otelContext = authLoginOtelTraceContext({ traceparent: getHeader(request, "traceparent") });
setAuthTraceHeaders(response, otelContext);
const startedAtMs = Date.now();
const token = sessionCookieFromRequest(request);
if (token?.startsWith("hws_") && this.userBilling?.configured && typeof this.userBilling.logout === "function") {
await this.userBilling.logout(token);
} else if (token) {
await this.store.revokeSessionByTokenHash(sha256(token), this.now());
let httpStatus = 500;
let outcome = "error";
let errorCode = "";
try {
if (token?.startsWith("hws_") && this.userBilling?.configured && typeof this.userBilling.logout === "function") {
await this.userBilling.logout(token);
outcome = "user_billing_session";
} else if (token) {
await this.store.revokeSessionByTokenHash(sha256(token), this.now());
outcome = "local_session";
} else {
outcome = "no_session_cookie";
}
clearSessionCookie(response, request, this.env);
httpStatus = 200;
return sendJson(response, 200, { authenticated: false, loggedOut: true });
} catch (error) {
httpStatus = Number(error?.statusCode ?? error?.status ?? httpStatus ?? 500);
errorCode = textOr(error?.code, "auth_logout_error");
throw error;
} finally {
void emitAuthOtelSpan("auth.logout", otelContext, this.env, {
spanId: otelContext.rootSpanId,
parentSpanId: otelContext.parentSpanId,
kind: 2,
startTimeMs: startedAtMs,
endTimeMs: Date.now(),
status: httpStatus >= 500 || errorCode ? "error" : "ok",
error: errorCode ? new Error(errorCode) : null,
attributes: {
"http.method": "POST",
"http.route": "/auth/logout",
"http.status_code": httpStatus,
"auth.outcome": outcome,
"auth.cookie_present": Boolean(token),
"auth.user_billing_token": Boolean(token?.startsWith("hws_")),
"auth.trace_visible": true,
"auth.values_redacted": true,
...(errorCode ? { "error.code": errorCode } : {})
}
});
}
clearSessionCookie(response, request, this.env);
return sendJson(response, 200, { authenticated: false, loggedOut: true });
}
async codeAgentToolCapabilitiesForOwner(ownerUserId) {
+203 -5
View File
@@ -16,7 +16,7 @@ import {
import { createWorkbenchReadModel } from "./workbench-read-model.ts";
import { createWorkbenchRuntimeClient } from "./workbench-runtime-client.ts";
import { durableTraceStatus, RUNNING_STATUSES, TERMINAL_STATUSES } from "./workbench-turn-projection.ts";
import { emitHttpServerRequestSpan } from "./otel-trace.ts";
import { emitCodeAgentOtelSpan, emitHttpServerRequestSpan } from "./otel-trace.ts";
const DEFAULT_PAGE_LIMIT = 50;
const DEFAULT_SESSION_LIST_LIMIT = 20;
@@ -1682,6 +1682,13 @@ async function handleWorkbenchTurnSnapshot(request, response, url, options, acto
const result = await readModel.queryFacts({ traceId, limit: MAX_PAGE_LIMIT, families: [...WORKBENCH_TRACE_EVENT_METADATA_FAMILIES, ...WORKBENCH_TRACE_EVENT_PAGE_FAMILIES] });
if (result.error) {
const body = workbenchProjectionStoreError(result.error);
emitWorkbenchTurnStatusReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 503,
phase: "query_failed",
turnId,
errorCode: result.error?.code ?? "workbench_turn_query_failed"
}, result.error);
recordWorkbenchTurnReadMetric(options, url, 503, body, startedAt);
return sendJson(response, 503, body);
}
@@ -1691,6 +1698,14 @@ async function handleWorkbenchTurnSnapshot(request, response, url, options, acto
if (!found) {
const body = workbenchError("workbench_turn_not_found", "Workbench turn is not visible to the current actor.", { turnId, traceId });
attachWorkbenchReadModelDiagnosticOtel(request, { route: "/v1/workbench/turns/:id", code: "workbench_turn_not_found", traceId, turnId, count: result.count });
emitWorkbenchTurnStatusReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 404,
phase: "not_found",
turnId,
readCount: result.count,
errorCode: "workbench_turn_not_found"
});
recordWorkbenchTurnReadMetric(options, url, 404, body, startedAt);
return sendJson(response, 404, body);
}
@@ -1718,6 +1733,22 @@ async function handleWorkbenchTurnSnapshot(request, response, url, options, acto
projection: { ...projection, status, eventCount: snapshot.trace?.eventCount ?? projection.lastProjectedSeq ?? 0 },
diagnostic: projection
});
emitWorkbenchTurnStatusReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 200,
phase: "succeeded",
status,
terminal: snapshot.terminal,
sessionId: snapshot.sessionId,
turnId: snapshot.turnId,
traceEventCount: snapshot.trace?.eventCount ?? null,
latestProjectedSeq: projection.lastProjectedSeq,
projectionStatus: projection.projectionStatus,
projectionHealth: projection.projectionHealth,
sourceRunId: projection.sourceRunId,
sourceCommandId: projection.sourceCommandId,
readCount: result.count
});
recordWorkbenchTurnReadMetric(options, url, 200, body, startedAt);
recordWorkbenchCacheMetric(options, url, body, startedAt);
sendJson(response, 200, body);
@@ -1730,14 +1761,43 @@ async function handleWorkbenchTraceEventPage(request, response, url, options, ac
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));
if (metadata.error) {
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 503,
phase: "metadata_query_failed",
pageOptions,
errorCode: metadata.error?.code ?? "workbench_trace_metadata_query_failed"
}, metadata.error);
return sendJson(response, 503, workbenchProjectionStoreError(metadata.error));
}
const resolved = await visibleFactSessionForTrace(readModel, metadata.facts, actor, traceId);
if (resolved.error) return sendJson(response, 503, workbenchProjectionStoreError(resolved.error));
if (resolved.error) {
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 503,
phase: "session_lookup_failed",
pageOptions,
metadataCount: metadata.count,
errorCode: resolved.error?.code ?? "workbench_trace_session_lookup_failed"
}, resolved.error);
return sendJson(response, 503, workbenchProjectionStoreError(resolved.error));
}
const session = resolved.session;
const metadataFacts = resolved.facts ?? metadata.facts;
if (!session) {
const context = await readModel.queryFacts({ traceId, families: ["sessions", "turns", "checkpoints"], limit: MAX_PAGE_LIMIT });
if (context.error) return sendJson(response, 503, workbenchProjectionStoreError(context.error));
if (context.error) {
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 503,
phase: "context_query_failed",
pageOptions,
metadataCount: metadata.count,
errorCode: context.error?.code ?? "workbench_trace_context_query_failed"
}, context.error);
return sendJson(response, 503, workbenchProjectionStoreError(context.error));
}
const contextSessions = visibleFactSessions(context.facts, actor);
const contextSession = contextSessions.find((item) => item.lastTraceId === traceId) ?? contextSessions[0] ?? null;
const contextTurn = factTurnForTrace(context.facts, traceId, traceId);
@@ -1745,14 +1805,47 @@ async function handleWorkbenchTraceEventPage(request, response, url, options, ac
const projection = factProjectionForTrace(context.facts, traceId);
const blocker = traceEventsReadModelBlocker("workbench_trace_metadata_missing", "Workbench trace metadata is missing from the trace events read model.", { traceId, projection, session: contextSession, turn: contextTurn, route: url.pathname });
attachWorkbenchReadModelDiagnosticOtel(request, { route: "/v1/workbench/traces/:id/events", code: "workbench_trace_metadata_missing", sessionId: factSessionId(contextSession), traceId, count: context.count });
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 404,
phase: "metadata_missing",
pageOptions,
projection,
sessionId: factSessionId(contextSession),
turnId: factTurnId(contextTurn, traceId),
readCount: context.count,
metadataCount: metadata.count,
errorCode: "workbench_trace_metadata_missing"
});
return sendJson(response, 404, workbenchTraceEventsReadModelError(blocker, { traceId, projection, session: contextSession, turn: contextTurn }));
}
attachWorkbenchReadModelDiagnosticOtel(request, { route: "/v1/workbench/traces/:id/events", code: "workbench_trace_not_found", traceId, count: context.count });
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 404,
phase: "not_found",
pageOptions,
readCount: context.count,
metadataCount: metadata.count,
errorCode: "workbench_trace_not_found"
});
return sendJson(response, 404, workbenchError("workbench_trace_not_found", "Workbench trace is not visible to the current actor.", { traceId }));
}
const projection = factProjectionForTrace(metadataFacts, traceId);
const pageResult = await readModel.queryFacts({ traceId, families: WORKBENCH_TRACE_EVENT_PAGE_FAMILIES, afterProjectedSeq: pageOptions.afterProjectedSeq, limit: pageOptions.limit + 1 });
if (pageResult.error) return sendJson(response, 503, workbenchProjectionStoreError(pageResult.error));
if (pageResult.error) {
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 503,
phase: "events_query_failed",
pageOptions,
projection,
sessionId: factSessionId(session),
metadataCount: metadata.count,
errorCode: pageResult.error?.code ?? "workbench_trace_events_query_failed"
}, pageResult.error);
return sendJson(response, 503, workbenchProjectionStoreError(pageResult.error));
}
const traceTurn = factTurnForTrace(metadataFacts, traceId, traceId);
const turnTraceStatus = normalizeStatus(traceTurn?.status);
const traceStatus = turnTraceStatus !== "unknown" ? turnTraceStatus : durableTraceStatus(pageResult.facts.traceEvents);
@@ -1761,6 +1854,20 @@ async function handleWorkbenchTraceEventPage(request, response, url, options, ac
if (missingTraceEvents) {
const blocker = traceEventsReadModelBlocker("workbench_trace_events_missing", "Workbench trace events are missing from the durable read model.", { traceId, projection, session, route: url.pathname });
attachWorkbenchReadModelDiagnosticOtel(request, { route: "/v1/workbench/traces/:id/events", code: "workbench_trace_events_missing", sessionId: factSessionId(session), traceId, count: pageResult.count });
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 404,
phase: "events_missing",
pageOptions,
page,
projection,
sessionId: factSessionId(session),
turnId: factTurnId(traceTurn, traceId),
readCount: pageResult.count,
rawEventCount: factArray(pageResult.facts?.traceEvents).length,
metadataCount: metadata.count,
errorCode: "workbench_trace_events_missing"
});
return sendJson(response, 404, workbenchTraceEventsReadModelError(blocker, { traceId, projection, session, turn: traceTurn }));
}
const responseProjection = page.blocker ? blockedTraceEventProjection(projection, page.blocker) : projection;
@@ -1790,9 +1897,100 @@ async function handleWorkbenchTraceEventPage(request, response, url, options, ac
secretMaterialStored: false
};
recordWorkbenchCacheMetric(options, url, body, startedAt);
emitWorkbenchTraceEventsReadOtel(traceId, options, {
startTimeMs: startedAt,
statusCode: 200,
phase: page.blocker ? "blocked" : "succeeded",
pageOptions,
page,
projection: responseProjection,
sessionId: body.sessionId,
turnId: factTurnId(traceTurn, traceId),
readCount: pageResult.count,
rawEventCount: factArray(pageResult.facts?.traceEvents).length,
metadataCount: metadata.count,
traceStatus
});
sendJson(response, 200, body);
}
function emitWorkbenchTurnStatusReadOtel(traceId, options = {}, fields = {}, error = null) {
const safeId = safeTraceId(traceId);
if (!safeId) return;
const statusCode = Number.isInteger(Number(fields.statusCode)) ? Number(fields.statusCode) : 0;
void emitCodeAgentOtelSpan("turn_status_read", safeId, options.env ?? process.env, {
startTimeMs: fields.startTimeMs,
status: statusCode >= 500 || error ? "error" : "ok",
error,
attributes: {
"http.method": "GET",
"http.route": "/v1/workbench/turns/:id",
"http.status_code": statusCode || null,
phase: fields.phase ?? null,
status: fields.status ?? null,
terminal: typeof fields.terminal === "boolean" ? fields.terminal : null,
sessionId: fields.sessionId ?? null,
turnId: fields.turnId ?? safeId,
readCount: nonNegativeInteger(fields.readCount),
traceEventCount: nonNegativeInteger(fields.traceEventCount),
latestProjectedSeq: nonNegativeInteger(fields.latestProjectedSeq),
projectionStatus: fields.projectionStatus ?? null,
projectionHealth: fields.projectionHealth ?? null,
sourceRunId: fields.sourceRunId ?? null,
sourceCommandId: fields.sourceCommandId ?? null,
errorCode: fields.errorCode ?? null,
valuesRedacted: true
}
});
}
function emitWorkbenchTraceEventsReadOtel(traceId, options = {}, fields = {}, error = null) {
const safeId = safeTraceId(traceId);
if (!safeId) return;
const page = fields.page ?? {};
const range = page.range ?? {};
const projection = fields.projection ?? {};
const pageOptions = fields.pageOptions ?? {};
const statusCode = Number.isInteger(Number(fields.statusCode)) ? Number(fields.statusCode) : 0;
const returnedEvents = Array.isArray(page.events) ? page.events.length : nonNegativeInteger(fields.returnedEvents);
const afterProjectedSeq = nonNegativeInteger(pageOptions.afterProjectedSeq ?? range.afterProjectedSeq);
const limit = nonNegativeInteger(pageOptions.limit ?? range.limit);
const hasMore = typeof page.hasMore === "boolean" ? page.hasMore : null;
void emitCodeAgentOtelSpan("trace_events_read", safeId, options.env ?? process.env, {
startTimeMs: fields.startTimeMs,
status: statusCode >= 500 || error ? "error" : "ok",
error,
attributes: {
"http.method": "GET",
"http.route": "/v1/workbench/traces/:id/events",
"http.status_code": statusCode || null,
phase: fields.phase ?? null,
sessionId: fields.sessionId ?? null,
turnId: fields.turnId ?? safeId,
returnedEvents,
rawEventCount: nonNegativeInteger(fields.rawEventCount),
readCount: nonNegativeInteger(fields.readCount),
metadataCount: nonNegativeInteger(fields.metadataCount),
sinceSeq: afterProjectedSeq,
afterProjectedSeq,
limit,
fromSeq: nonNegativeInteger(range.fromProjectedSeq),
toSeq: nonNegativeInteger(range.toProjectedSeq),
totalEvents: nonNegativeInteger(range.total ?? page.eventCount),
hasMore,
fullTraceLoaded: hasMore === null ? null : !hasMore && !page.blocker,
traceLastSeq: nonNegativeInteger(projection.lastProjectedSeq),
projectionStatus: projection.projectionStatus ?? null,
projectionHealth: projection.projectionHealth ?? null,
sourceRunId: projection.sourceRunId ?? null,
sourceCommandId: projection.sourceCommandId ?? null,
traceStatus: fields.traceStatus ?? page.traceStatus ?? null,
errorCode: fields.errorCode ?? null,
valuesRedacted: true
}
});
}
function traceEventPageFromFacts(sourceEvents, options, metadata = {}) {
const rows = factArray(sourceEvents).filter((event) => factProjectedSeq(event)).sort(compareFactTraceEventsAsc);
const collision = projectedSeqCollision(rows);
+30 -7
View File
@@ -112,6 +112,7 @@ export async function writeWorkbenchProjectionEvent({ runtimeStore = null, event
const sourceSeq = nonNegativeInteger(event.sourceSeq ?? event.seq);
const sessionId = textValue(event.sessionId ?? requestMeta.sessionId) || null;
const turnId = textValue(event.turnId) || traceId;
const eventType = textValue(event.eventType ?? event.type ?? event.label) || "event";
const terminal = event.terminal === true || TERMINAL_STATUSES.has(normalizeWorkbenchStatus(event.status));
const status = terminal ? normalizeWorkbenchStatus(event.status) : "running";
const checkpointHint = previousCheckpoint && typeof previousCheckpoint === "object" ? previousCheckpoint : null;
@@ -127,7 +128,9 @@ export async function writeWorkbenchProjectionEvent({ runtimeStore = null, event
"workbench.projection.terminal": false,
"workbench.projection.suppressed_after_seal": true,
"workbench.projection.session_id": sessionId,
"workbench.projection.turn_id": turnId
"workbench.projection.turn_id": turnId,
"workbench.projection.source_seq": sourceSeq,
"workbench.projection.event_type": eventType
});
return { written: false, suppressedAfterSeal: true, valuesPrinted: false };
}
@@ -161,7 +164,7 @@ export async function writeWorkbenchProjectionEvent({ runtimeStore = null, event
commandId: textValue(event.commandId) || null,
sourceSeq,
sourceEventId,
eventType: textValue(event.eventType ?? event.type ?? event.label) || "event",
eventType,
occurredAt,
updatedAt: projectedAt
}, {
@@ -224,7 +227,7 @@ export async function writeWorkbenchProjectionEvent({ runtimeStore = null, event
sourceSeq,
sourceEventId,
projectedSeq,
eventType: textValue(event.eventType ?? event.type ?? event.label) || "event",
eventType,
terminal,
sealed: terminal,
timing: eventTiming,
@@ -280,7 +283,13 @@ export async function writeWorkbenchProjectionEvent({ runtimeStore = null, event
"workbench.projection.terminal": terminal,
"workbench.projection.suppressed_after_seal": suppressedAfterSeal,
"workbench.projection.session_id": sessionId,
"workbench.projection.turn_id": turnId
"workbench.projection.turn_id": turnId,
"workbench.projection.projected_seq": projectedSeq,
"workbench.projection.source_seq": sourceSeq,
"workbench.projection.source_event_id": sourceEventId,
"workbench.projection.event_type": eventType,
"workbench.projection.event_id": eventId,
"workbench.projection.message_id": messageFact?.messageId ?? null
});
return result;
} catch (error) {
@@ -290,7 +299,13 @@ export async function writeWorkbenchProjectionEvent({ runtimeStore = null, event
"workbench.projection.terminal": terminal,
"workbench.projection.suppressed_after_seal": suppressedAfterSeal,
"workbench.projection.session_id": sessionId,
"workbench.projection.turn_id": turnId
"workbench.projection.turn_id": turnId,
"workbench.projection.projected_seq": projectedSeq,
"workbench.projection.source_seq": sourceSeq,
"workbench.projection.source_event_id": sourceEventId,
"workbench.projection.event_type": eventType,
"workbench.projection.event_id": eventId,
"workbench.projection.message_id": messageFact?.messageId ?? null
}, error);
appendProjectionDiagnostic(defaultCodeAgentTraceStore, traceId, {
type: "facts",
@@ -351,7 +366,11 @@ async function writeWorkbenchProjectionFacts({ runtimeStore = null, traceStore =
"workbench.projection.phase": "session-facts",
"workbench.projection.terminal": facts.turns.some((turn) => turn.terminal === true),
"workbench.projection.session_id": facts.sessions[0]?.sessionId ?? sessionId ?? null,
"workbench.projection.turn_id": facts.turns[0]?.turnId ?? null
"workbench.projection.turn_id": facts.turns[0]?.turnId ?? null,
"workbench.projection.projected_seq": facts.checkpoints[0]?.projectedSeq ?? facts.turns[0]?.projectedSeq ?? facts.messages[0]?.projectedSeq ?? null,
"workbench.projection.source_seq": facts.checkpoints[0]?.sourceSeq ?? facts.turns[0]?.sourceSeq ?? facts.messages[0]?.sourceSeq ?? null,
"workbench.projection.source_event_id": facts.checkpoints[0]?.sourceEventId ?? facts.turns[0]?.sourceEventId ?? facts.messages[0]?.sourceEventId ?? null,
"workbench.projection.message_id": facts.messages[0]?.messageId ?? null
});
return result;
} catch (error) {
@@ -361,7 +380,11 @@ async function writeWorkbenchProjectionFacts({ runtimeStore = null, traceStore =
"workbench.projection.phase": "session-facts",
"workbench.projection.terminal": facts.turns.some((turn) => turn.terminal === true),
"workbench.projection.session_id": facts.sessions[0]?.sessionId ?? sessionId ?? null,
"workbench.projection.turn_id": facts.turns[0]?.turnId ?? null
"workbench.projection.turn_id": facts.turns[0]?.turnId ?? null,
"workbench.projection.projected_seq": facts.checkpoints[0]?.projectedSeq ?? facts.turns[0]?.projectedSeq ?? facts.messages[0]?.projectedSeq ?? null,
"workbench.projection.source_seq": facts.checkpoints[0]?.sourceSeq ?? facts.turns[0]?.sourceSeq ?? facts.messages[0]?.sourceSeq ?? null,
"workbench.projection.source_event_id": facts.checkpoints[0]?.sourceEventId ?? facts.turns[0]?.sourceEventId ?? facts.messages[0]?.sourceEventId ?? null,
"workbench.projection.message_id": facts.messages[0]?.messageId ?? null
}, error);
if (safeId) appendProjectionDiagnostic(traceStore, safeId, {
type: "facts",