diff --git a/workers/iroh-v2/src/dashboard-control.ts b/workers/iroh-v2/src/dashboard-control.ts index 77cf7d16725f..d607135f4af9 100644 --- a/workers/iroh-v2/src/dashboard-control.ts +++ b/workers/iroh-v2/src/dashboard-control.ts @@ -7,7 +7,7 @@ import { canonicalJSON, hash } from "./crypto"; import { DASHBOARD_AUTHORITY_HEADER, DashboardClaimsSchema, type DashboardClaims } from "./dashboard-auth"; import { acknowledgeDelivery, DeliveryStateSchema, deliveryUsage, emptyDeliveryState, prepareDelivery } from "./delivery"; import type { Environment } from "./environment"; -import { OperationError } from "./errors"; +import { failureDiagnostics, OperationError } from "./errors"; import { observe } from "./observability"; import { unwrap, type UserUsage } from "./user-usage-object"; @@ -66,7 +66,7 @@ export class DashboardControl { } finally { this.services.opening.delete(sessionId); } } catch (error) { const failure = errorResponse(error, requestId).failure; - observe(this.ctx, this.env, { event: "iroh.dashboard.failure", requestId, code: failure.code, status: failure.status }); + observe(this.ctx, this.env, { event: "iroh.dashboard.failure", requestId, code: failure.code, status: failure.status, ...failureDiagnostics(error) }); return httpFailure(error, requestId); } } diff --git a/workers/iroh-v2/src/dashboard-routing.ts b/workers/iroh-v2/src/dashboard-routing.ts index 11631563e480..1ff19aaac906 100644 --- a/workers/iroh-v2/src/dashboard-routing.ts +++ b/workers/iroh-v2/src/dashboard-routing.ts @@ -2,7 +2,7 @@ import type { StackAuthority } from "./auth"; import { encodeResponse, errorResponse, httpFailure, inputRequestId, parseInput, readBoundedBody } from "./boundary"; import { DashboardOpenSchema } from "./contracts/common"; import { DASHBOARD_AUTHORITY_HEADER, issueDashboardTicket, verifyDashboardTicket } from "./dashboard-auth"; -import { OperationError } from "./errors"; +import { failureDiagnostics, OperationError } from "./errors"; interface Dependencies { environment: string; projectId: string; keys: Readonly>; @@ -69,7 +69,7 @@ export async function routeDashboard(request: Request, services: Dependencies): return cors(await services.dispatchTeam(claims.authority.teamId, forwarded)); } catch (error) { const failure = errorResponse(error, requestId).failure; - services.observe?.({ event: "iroh.dashboard.failure", requestId, code: failure.code, status: failure.status }); + services.observe?.({ event: "iroh.dashboard.failure", requestId, code: failure.code, status: failure.status, ...failureDiagnostics(error) }); return cors(httpFailure(error, requestId)); } } diff --git a/workers/iroh-v2/src/errors.ts b/workers/iroh-v2/src/errors.ts index 6b654cce754a..83ee4e15ae05 100644 --- a/workers/iroh-v2/src/errors.ts +++ b/workers/iroh-v2/src/errors.ts @@ -20,16 +20,39 @@ export function publicError(error: unknown): OperationError { * records only error names and known capacity markers from the cause chain. */ export function errorSummary(error: unknown): string { - const markers = ["socket_output_capacity", "socket_capacity"] as const; const parts: string[] = []; let current: unknown = error; for (let depth = 0; current !== undefined && current !== null && depth < 4; depth += 1) { const rawName = current instanceof Error ? current.name : typeof current; - const name = ["Error", "TypeError", "RangeError", "SyntaxError", "DrizzleError", "DrizzleQueryError"].includes(rawName) ? rawName : "Error"; + const name = ["Error", "TypeError", "RangeError", "SyntaxError", "DrizzleError", "DrizzleQueryError", "ZodError"].includes(rawName) ? rawName : "Error"; const text = current instanceof Error ? current.message : String(current); - const marker = markers.find(candidate => text.includes(candidate)); - parts.push(marker ? `${name}:${marker}` : name); + const marker = FAILURE_MARKERS.find(([needle]) => text.includes(needle))?.[1]; + // Workers runtime errors carry these booleans; they name the platform + // failure class without exposing the message text. + const flags = ["retryable", "overloaded", "remote"].filter(flag => current !== null && typeof current === "object" && Reflect.get(current, flag) === true); + parts.push([marker ? `${name}:${marker}` : name, ...flags].join("+")); current = current instanceof Error ? current.cause : undefined; } - return parts.join(" <- ").slice(0, 160); + return parts.join(" <- ").slice(0, 200); +} + +/** + * Known failure texts mapped to fixed tags. Storage guards raise the first + * group; the rest are Workers and Durable Object runtime failures. Only the + * tag is recorded, never the matched message. + */ +const FAILURE_MARKERS: readonly (readonly [string, string])[] = [ + ["socket_output_capacity", "socket_output_capacity"], ["socket_capacity", "socket_capacity"], + ["audit_limit", "audit_limit"], ["device_limit", "device_limit"], ["storage_limit", "storage_limit"], + ["SQLITE_BUSY", "sqlite_busy"], ["SQLITE_FULL", "sqlite_full"], ["CONSTRAINT", "sqlite_constraint"], + ["code was updated", "do_code_updated"], ["Durable Object reset", "do_reset"], + ["Network connection lost", "network_lost"], ["overloaded", "overloaded"], + ["exceeded timeout", "storage_timeout"], ["memory limit", "memory_limit"], ["CPU time", "cpu_limit"], + ["transient issue", "do_transient"], ["too many subrequests", "subrequest_limit"], + ["The operation was aborted", "aborted"], ["timed out", "timed_out"], ["fetch failed", "fetch_failed"], +]; + +/** Telemetry fields for a failure: a cause only when the error was not classified. */ +export function failureDiagnostics(error: unknown): { cause?: string } { + return error instanceof OperationError ? {} : { cause: errorSummary(error) }; } diff --git a/workers/iroh-v2/src/index.ts b/workers/iroh-v2/src/index.ts index 997b89e8e913..6eb80ee36b07 100644 --- a/workers/iroh-v2/src/index.ts +++ b/workers/iroh-v2/src/index.ts @@ -1,4 +1,5 @@ import { errorResponse, httpFailure } from "./boundary"; +import { failureDiagnostics } from "./errors"; import { runtime, type Environment } from "./environment"; import { routeControl, objectName } from "./routing"; import { unwrap } from "./user-usage-object"; @@ -30,7 +31,7 @@ export default { }); } catch (error) { const failure = errorResponse(error, "unidentified").failure; - observe(ctx, env, { event: "iroh.http.failure", environment: env.ENVIRONMENT, path: new URL(request.url).pathname, code: failure.code, status: failure.status, retryable: failure.retryable }); + observe(ctx, env, { event: "iroh.http.failure", environment: env.ENVIRONMENT, path: new URL(request.url).pathname, code: failure.code, status: failure.status, retryable: failure.retryable, ...failureDiagnostics(error) }); response = httpFailure(error); } observe(ctx, env, { event: "iroh.http.response", environment: env.ENVIRONMENT, path: new URL(request.url).pathname, status: response.status, durationMs: Date.now() - started }); diff --git a/workers/iroh-v2/src/routing.ts b/workers/iroh-v2/src/routing.ts index 3272ba5bafe2..0bff72c56553 100644 --- a/workers/iroh-v2/src/routing.ts +++ b/workers/iroh-v2/src/routing.ts @@ -1,13 +1,13 @@ import { z } from "zod"; import type { StackAuthority, VerifiedAuthority } from "./auth"; -import { decodeJSON, errorResponse, httpFailure, inputRequestId, parseInput, parseSocketSetup, readBoundedBody, INPUT_BYTES } from "./boundary"; +import { decodeJSON, errorResponse, httpFailure, inputOperation, inputRequestId, parseInput, parseSocketSetup, readBoundedBody, INPUT_BYTES } from "./boundary"; import { identifier, timestamp } from "./contracts/common"; import { SocketSetupSchema, type SocketSetup } from "./contracts/requests"; import { STORAGE_SCHEMA_VERSION, STORAGE_WRITE_SCHEMA_VERSION } from "./storage/migrations"; import { HealthSchema } from "./health"; import { CONTROL_PLANE_RULES, sourceRevision } from "./rules"; import { API_TICKET_SECONDS, canonicalJSON, decodeBase64URL, encodeBase64URL, verifyTicket } from "./crypto"; -import { OperationError } from "./errors"; +import { failureDiagnostics, OperationError } from "./errors"; export const SETUP_HEADER = "x-cmux-v2-setup"; const INTERNAL_HEADER = "x-cmux-v2-verified-authority"; @@ -41,6 +41,9 @@ const aliases: Readonly> = { /** Public requests never control the private headers passed through the DO binding. */ export async function routeControl(request: Request, dependencies: RoutingDependencies): Promise { let requestId = "unidentified"; + // Where a failure happened, for telemetry only: route kind, the stage + // reached and the requested operation, never request contents. + let route = "unknown", stage = "parse", operationName = "none"; try { const url = new URL(request.url); const socket = url.pathname === "/v2/control/socket"; @@ -48,6 +51,7 @@ export async function routeControl(request: Request, dependencies: RoutingDepend const operation = url.pathname === "/v2/requests" || Object.hasOwn(aliases, url.pathname); const health = url.pathname === "/v2/health"; if ((!socket && !session && !operation && !health) || url.search) throw new OperationError("unsupported_method", 404); + route = socket ? "socket" : session ? "session" : operation ? "request" : health ? "health" : "unknown"; if (health) return healthResponse(request, dependencies); if (request.method !== (socket ? "GET" : "POST")) throw new OperationError("unsupported_method", 405); if (socket && request.headers.get("upgrade")?.toLowerCase() !== "websocket") throw new OperationError("invalid_request", 400); @@ -65,7 +69,10 @@ export async function routeControl(request: Request, dependencies: RoutingDepend throw new OperationError("invalid_request", 400); } } + if (operation) operationName = inputOperation(input); + stage = "authenticate"; const authorization = await authenticate(request.headers.get("authorization"), setup, dependencies); + stage = "charge"; if (!operation) await dependencies.chargeOpen(authorization.authority.userId); // A fresh Request deliberately copies no caller headers, cookies or credentials. const headers = new Headers({ @@ -79,10 +86,12 @@ export async function routeControl(request: Request, dependencies: RoutingDepend method: socket ? "GET" : "POST", headers, ...(socket ? {} : { body: JSON.stringify({ setup, ...(operation ? { input } : {}) }) }), }); + stage = "dispatch"; return await dependencies.dispatchTeam(setup.device.identity.teamId, forwarded); } catch (error) { const failure = errorResponse(error, requestId).failure; - dependencies.observe?.({ event: "iroh.control.failure", requestId, code: failure.code, status: failure.status, retryable: failure.retryable }); + dependencies.observe?.({ event: "iroh.control.failure", requestId, code: failure.code, status: failure.status, retryable: failure.retryable, + route, stage, operation: operationName, ...failureDiagnostics(error) }); return httpFailure(error, requestId); } } diff --git a/workers/iroh-v2/src/team-control.ts b/workers/iroh-v2/src/team-control.ts index 15a439bc95e8..72265960e067 100644 --- a/workers/iroh-v2/src/team-control.ts +++ b/workers/iroh-v2/src/team-control.ts @@ -7,7 +7,7 @@ import type { ControlResponse } from "./contracts/responses"; import { identityKey, issueTicket } from "./crypto"; import { acknowledgeDelivery, DeliveryStateSchema, deliveryUsage, emptyDeliveryState, prepareDelivery } from "./delivery"; import { environmentScope, runtime, type Environment } from "./environment"; -import { errorSummary, OperationError, publicError } from "./errors"; +import { failureDiagnostics, OperationError, publicError } from "./errors"; import { AuthoritySchema, objectName, readInternalRequest } from "./routing"; import { applyStorageMigrations } from "./storage/migrations"; import { TeamStore } from "./storage/team-store"; @@ -50,9 +50,13 @@ export class TeamControl extends DurableObject { async fetch(request: Request): Promise { if (new URL(request.url).pathname === "/dashboard/socket") return this.dashboard.fetch(request); let requestId = "unidentified"; + const pathname = new URL(request.url).pathname; + const route = ["/request", "/session", "/socket"].includes(pathname) ? pathname.slice(1) : "unknown"; + let stage = "parse"; try { const incoming = await readInternalRequest(request); requestId = incoming.setup.requestId; + stage = incoming.path === "/request" ? "execute" : "open"; const broker = this.broker(incoming.authority.teamId); if (incoming.path === "/request") { const session = await broker.authorizeHTTP(incoming.setup, incoming.input, incoming.authority, incoming.expiresAt, incoming.issueTicket); @@ -66,6 +70,7 @@ export class TeamControl extends DurableObject { this.scheduleChanges(result, incoming.authority.teamId); observe(this.ctx, this.env, { event: "iroh.team.operation", environment: this.env.ENVIRONMENT, operation: result.response.schemaId, requestId, status: 200 }); if (incoming.path === "/session") return this.json(result.response); + stage = "accept"; if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") throw new OperationError("invalid_request", 400); if (this.ctx.getWebSockets().length >= TEAM_SOCKET_LIMIT) throw new OperationError("rate_limited", 429, true, 5000); const session = result.session; @@ -77,15 +82,18 @@ export class TeamControl extends DurableObject { const client = pair[0], server = pair[1]; this.ctx.acceptWebSocket(server, ["user:" + session.identity.userId, "device:" + deviceKey]); this.save(server, { version: 1, session, deviceKey, delivery: emptyDeliveryState(), outputRevision: 0, closed: false }); + stage = "send"; try { await this.enqueue(server, 0, () => this.send(server, result.response)); } catch (error) { this.close(server, "slow_consumer"); throw error; } + stage = "ready"; // The replacement is accepted and ready before any previous socket closes. for (const old of this.ctx.getWebSockets("device:" + deviceKey)) if (old !== server) this.close(old, "session_replaced"); return new Response(null, { status: 101, webSocket: client }); } finally { this.opening.delete(session.sessionId); } } catch (error) { const failure = publicError(error); - observe(this.ctx, this.env, { event: "iroh.team.failure", environment: this.env.ENVIRONMENT, requestId, code: failure.code, status: failure.status, retryable: failure.retryable }); + observe(this.ctx, this.env, { event: "iroh.team.failure", environment: this.env.ENVIRONMENT, requestId, code: failure.code, status: failure.status, retryable: failure.retryable, + route, stage, ...failureDiagnostics(error) }); return httpFailure(error, requestId); } } @@ -126,7 +134,7 @@ export class TeamControl extends DurableObject { } catch (error) { const failure = errorResponse(error, inputRequestId(input)); status = failure.failure.status; code = failure.failure.code; - if (!(error instanceof OperationError)) cause = errorSummary(error); + cause = failureDiagnostics(error).cause; try { await this.send(ws, failure.body); } catch { this.close(ws, "slow_consumer"); } if (["device_revoked", "team_access_revoked", "identity_mismatch", "key_replacement_required"].includes(code)) this.close(ws, code); } finally { diff --git a/workers/iroh-v2/src/user-usage-object.ts b/workers/iroh-v2/src/user-usage-object.ts index e6bb85a40e1c..5a6c1c0f9e3d 100644 --- a/workers/iroh-v2/src/user-usage-object.ts +++ b/workers/iroh-v2/src/user-usage-object.ts @@ -1,7 +1,8 @@ import { DurableObject } from "cloudflare:workers"; import { identifier } from "./contracts/common"; import { environmentScope, type Environment } from "./environment"; -import { errorSummary, OperationError, publicError } from "./errors"; +import { failureDiagnostics, OperationError, publicError } from "./errors"; +import { observe } from "./observability"; import { objectName } from "./routing"; import { applyStorageMigrations } from "./storage/migrations"; import { UserSocketStore } from "./storage/socket-store"; @@ -9,11 +10,13 @@ import { UserUsageStore, type UsageOperation } from "./storage/user-usage"; import type { ErrorCode } from "./contracts/responses"; type Result = { ok: true; value: T } | { ok: false; code: ErrorCode; status: number; retryable: boolean; retryAfterMs?: number }; -function result(action: () => T): Result { +function result(action: () => T, report?: (cause: string) => void): Result { try { return { ok: true, value: action() }; } catch (error) { const failure = publicError(error); - if (!(error instanceof OperationError)) console.error(JSON.stringify({ event: "iroh.user_usage.unclassified", cause: errorSummary(error) })); + // The caller only receives the public code across RPC; report the cause here. + const { cause } = failureDiagnostics(error); + if (cause) report?.(cause); return { ok: false, code: failure.code, status: failure.status, retryable: failure.retryable, ...(failure.retryAfterMs === undefined ? {} : { retryAfterMs: failure.retryAfterMs }) }; } @@ -36,19 +39,22 @@ export class UserUsage extends DurableObject { } consume(userId: string, operation: UsageOperation) { - return result(() => { this.assertUser(userId); return this.usage.consume(userId, operation); }); + return result(() => { this.assertUser(userId); return this.usage.consume(userId, operation); }, this.report("consume")); } reserveSocket(input: { userId: string; teamId: string; sessionId: string; deviceKey: string }) { - return result(() => { this.assertUser(input.userId); this.sockets.reserveSocket(input); }); + return result(() => { this.assertUser(input.userId); this.sockets.reserveSocket(input); }, this.report("reserveSocket")); } setOutput(userId: string, sessionId: string, revision: number, bytes: number, messages: number) { - return result(() => { this.assertUser(userId); this.sockets.setOutput(userId, sessionId, revision, bytes, messages); }); + return result(() => { this.assertUser(userId); this.sockets.setOutput(userId, sessionId, revision, bytes, messages); }, this.report("setOutput")); } releaseSocket(userId: string, sessionId: string) { - return result(() => { this.assertUser(userId); this.sockets.releaseSocket(userId, sessionId); }); + return result(() => { this.assertUser(userId); this.sockets.releaseSocket(userId, sessionId); }, this.report("releaseSocket")); } listSocketReservations(userId: string) { - return result(() => { this.assertUser(userId); return this.sockets.listSocketReservations(userId); }); + return result(() => { this.assertUser(userId); return this.sockets.listSocketReservations(userId); }, this.report("listSocketReservations")); + } + private report(operation: string): (cause: string) => void { + return cause => observe(this.ctx, this.env, { event: "iroh.user_usage.unclassified", status: 500, environment: this.env.ENVIRONMENT, operation, cause }); } private assertUser(userId: string): void { identifier.parse(userId); diff --git a/workers/iroh-v2/test/errors.test.ts b/workers/iroh-v2/test/errors.test.ts index 5b9901eb2013..4ef5f6b9d44c 100644 --- a/workers/iroh-v2/test/errors.test.ts +++ b/workers/iroh-v2/test/errors.test.ts @@ -1,5 +1,5 @@ import { expect, test } from "bun:test"; -import { errorSummary } from "../src/errors"; +import { errorSummary, failureDiagnostics, OperationError } from "../src/errors"; test("error summaries keep safe cause categories without SQL or bound parameters", () => { const error = new Error( @@ -17,3 +17,14 @@ test("diagnostics never include a custom error name or message", () => { error.name = "credential-value"; expect(errorSummary(error)).toBe("Error"); }); + +test("runtime failures map to fixed tags and platform flags", () => { + const overloaded = Object.assign(new Error("Durable Object is overloaded. Requests queued for too long."), { retryable: true, overloaded: true }); + expect(errorSummary(overloaded)).toBe("Error:overloaded+retryable+overloaded"); + expect(errorSummary(new TypeError("Network connection lost."))).toBe("TypeError:network_lost"); +}); + +test("classified operation errors carry no cause", () => { + expect(failureDiagnostics(new OperationError("ticket_expired", 401, true))).toEqual({}); + expect(failureDiagnostics(new Error("x"))).toEqual({ cause: "Error" }); +}); diff --git a/workers/iroh-v2/test/routing.test.ts b/workers/iroh-v2/test/routing.test.ts index 012aa829401a..1553c2113cf6 100644 --- a/workers/iroh-v2/test/routing.test.ts +++ b/workers/iroh-v2/test/routing.test.ts @@ -94,3 +94,26 @@ test("cross-environment setup never reaches Stack or a team object", async () => expect(result.status).toBe(403); expect(calls).toEqual({ stack: 0, open: 0, team: 0 }); }); + +test("control failures record the route, stage, operation and an unclassified cause", async () => { + const { dependencies } = fixture(); + const events: Record[] = []; + const ticket = await issueTicket(device, "test", key, 99); + const reset = Object.assign(new Error("Durable Object reset because its code was updated."), { retryable: true }); + const result = await routeControl(new Request("https://api.example/v2/requests", { + method: "POST", headers: { "content-type": "application/json", authorization: "IrohTicket " + ticket.token, [SETUP_HEADER]: encodedSetup }, + body: JSON.stringify({ schemaId: "directory.request.v1", requestId: "request" }), + }), { ...dependencies, observe: event => events.push(event), dispatchTeam: async () => { throw reset; } }); + expect(result.status).toBe(500); + expect(events).toEqual([expect.objectContaining({ + event: "iroh.control.failure", code: "internal_error", route: "request", stage: "dispatch", + operation: "directory.request", cause: "Error:do_code_updated+retryable", + })]); + expect(JSON.stringify(events)).not.toContain("code was updated"); + + events.length = 0; + await routeControl(new Request("https://api.example/v2/control/session", { + method: "POST", headers: { "content-type": "application/json", authorization: "Bearer stack-token" }, body: JSON.stringify(setup), + }), { ...dependencies, observe: event => events.push(event), stack: { verify: async () => { throw new Error("private upstream detail"); } } }); + expect(events).toEqual([expect.objectContaining({ route: "session", stage: "authenticate", operation: "none", cause: "Error" })]); +});