Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions workers/iroh-v2/src/dashboard-control.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down Expand Up @@ -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);
}
}
Expand Down
4 changes: 2 additions & 2 deletions workers/iroh-v2/src/dashboard-routing.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Record<string, string>>;
Expand Down Expand Up @@ -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));
}
}
33 changes: 28 additions & 5 deletions workers/iroh-v2/src/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) };
}
3 changes: 2 additions & 1 deletion workers/iroh-v2/src/index.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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 });
Expand Down
15 changes: 12 additions & 3 deletions workers/iroh-v2/src/routing.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -41,13 +41,17 @@ const aliases: Readonly<Record<string, string>> = {
/** Public requests never control the private headers passed through the DO binding. */
export async function routeControl(request: Request, dependencies: RoutingDependencies): Promise<Response> {
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";
const session = url.pathname === "/v2/control/session";
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);
Expand All @@ -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({
Expand All @@ -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);
}
}
Expand Down
14 changes: 11 additions & 3 deletions workers/iroh-v2/src/team-control.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -50,9 +50,13 @@ export class TeamControl extends DurableObject<Environment> {
async fetch(request: Request): Promise<Response> {
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);
Expand All @@ -66,6 +70,7 @@ export class TeamControl extends DurableObject<Environment> {
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;
Expand All @@ -77,15 +82,18 @@ export class TeamControl extends DurableObject<Environment> {
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);
}
}
Expand Down Expand Up @@ -126,7 +134,7 @@ export class TeamControl extends DurableObject<Environment> {
} 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 {
Expand Down
22 changes: 14 additions & 8 deletions workers/iroh-v2/src/user-usage-object.ts
Original file line number Diff line number Diff line change
@@ -1,19 +1,22 @@
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";
import { UserUsageStore, type UsageOperation } from "./storage/user-usage";
import type { ErrorCode } from "./contracts/responses";

type Result<T> = { ok: true; value: T } | { ok: false; code: ErrorCode; status: number; retryable: boolean; retryAfterMs?: number };
function result<T>(action: () => T): Result<T> {
function result<T>(action: () => T, report?: (cause: string) => void): Result<T> {
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 }) };
}
Expand All @@ -36,19 +39,22 @@ export class UserUsage extends DurableObject<Environment> {
}

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);
Expand Down
13 changes: 12 additions & 1 deletion workers/iroh-v2/test/errors.test.ts
Original file line number Diff line number Diff line change
@@ -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(
Expand All @@ -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" });
});
Loading
Loading