diff --git a/internal/cloud/access-control.ts b/internal/cloud/access-control.ts index 39465d7b..a00e172b 100644 --- a/internal/cloud/access-control.ts +++ b/internal/cloud/access-control.ts @@ -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) { diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index 8795bc01..3ce62a16 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -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); diff --git a/internal/cloud/workbench-projection-writer.ts b/internal/cloud/workbench-projection-writer.ts index 020af106..d37038b8 100644 --- a/internal/cloud/workbench-projection-writer.ts +++ b/internal/cloud/workbench-projection-writer.ts @@ -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",