diff --git a/src/azure.ts b/src/azure.ts index 3ef5e83..e547878 100644 --- a/src/azure.ts +++ b/src/azure.ts @@ -14,10 +14,10 @@ import type { } from "./driver/object.ts"; import { split } from "./chunk.ts"; import { + type FetchType, RequestMetrics, type RequestMetricsType, type RequestPolicyType, - RequestTransportError, sendRequest, } from "./request.ts"; import { type AdapterLimitsType, MetricsModeSchema, type MetricsModeType } from "./schema.ts"; @@ -95,7 +95,7 @@ export interface AzureClientOptionsType { /** Blob REST version. Defaults to {@link AZURE_STORAGE_VERSION}. */ readonly version?: AzureStorageVersionType; /** Fetch implementation. */ - readonly fetch?: typeof fetch; + readonly fetch?: FetchType; /** Clock used by `x-ms-date` and deterministic Shared Key tests. */ readonly now?: () => Date; /** Streaming block size. Defaults to 8 MiB. */ @@ -237,6 +237,39 @@ function compareText(left: string, right: string): number { return left < right ? -1 : left > right ? 1 : 0; } +/** + * Adds Azure metadata only after validating the provider's header contract. + * + * Azure accepts metadata names that start with a letter or underscore and then + * contain only ASCII letters, digits, or underscores. Metadata values must also + * be ASCII. Validate before request construction so a caller gets a local, + * deterministic failure instead of a provider HTTP 400 after bytes may already + * have been staged for a multipart write. `Headers.set()` remains responsible + * for the ordinary HTTP header-value syntax, such as rejecting embedded CR/LF. + */ +function setMetadata(headers: Headers, metadata: Readonly> | undefined): void { + const names = new Set(); + for (const [name, value] of Object.entries(metadata ?? {})) { + if (!/^[A-Za-z_][A-Za-z0-9_]*$/.test(name)) { + throw new TypeError( + `Azure metadata key ${JSON.stringify(name)} must start with a letter or underscore and contain only ASCII ` + + "letters, numbers, or underscores.", + ); + } + const normalized = name.toLowerCase(); + if (names.has(normalized)) { + throw new TypeError(`Azure metadata contains the case-insensitive duplicate key ${JSON.stringify(name)}.`); + } + names.add(normalized); + for (const character of value) { + if (character.codePointAt(0)! > 0x7f) { + throw new TypeError(`Azure metadata value for ${JSON.stringify(name)} must contain only ASCII characters.`); + } + } + headers.set(`x-ms-meta-${name}`, value); + } +} + /** * Normalizes HTTP linear whitespace for Azure Shared Key canonicalization. * @@ -517,7 +550,7 @@ class AzureClient implements AzureClientType { /** REST service version sent on every authorized request. */ readonly #version: AzureStorageVersionType; /** Fetch implementation used for all provider traffic. */ - readonly #fetch: typeof fetch; + readonly #fetch: FetchType; /** Clock used for request authorization. */ readonly #now: () => Date; /** Block size used by streamed uploads. */ @@ -697,13 +730,9 @@ class AzureClient implements AzureClientType { ...(signal === undefined ? {} : { signal }), }; if (options.body instanceof ReadableStream) init.duplex = "half"; - try { - return await this.#fetch(url, init); - } catch (error) { - if (options.signal?.aborted) throw error; - throw new RequestTransportError(error); - } + return { input: url, init }; }, { + fetch: this.#fetch, ...(this.#requestPolicy === undefined ? {} : { policy: this.#requestPolicy }), ...(options.signal === undefined ? {} : { signal: options.signal }), replayable, @@ -771,9 +800,7 @@ class AzureClient implements AzureClientType { } if (options.ifMatch !== undefined) headers.set("if-match", options.ifMatch); if (options.ifNoneMatch !== undefined) headers.set("if-none-match", options.ifNoneMatch); - if ("metadata" in options) { - for (const [name, value] of Object.entries(options.metadata ?? {})) headers.set(`x-ms-meta-${name}`, value); - } + if ("metadata" in options) setMetadata(headers, options.metadata); return headers; } @@ -993,7 +1020,7 @@ class AzureClient implements AzureClientType { const headers = this.#getWriteHeaders(options); headers.set("content-type", "application/xml"); if (source.mediaType !== undefined) headers.set("x-ms-blob-content-type", source.mediaType); - for (const [name, value] of Object.entries(source.metadata ?? {})) headers.set(`x-ms-meta-${name}`, value); + setMetadata(headers, source.metadata); await assertResponse( await this.request({ method: "PUT", diff --git a/src/error.ts b/src/error.ts index 5b2d87f..8ed9b35 100644 --- a/src/error.ts +++ b/src/error.ts @@ -1,4 +1,4 @@ -import type { ErrorCodeType } from "./schema.ts"; +import { ErrorCodeSchema, type ErrorCodeType } from "./schema.ts"; /** * Error returned by the high-level filesystem and first-party adapters. @@ -37,6 +37,30 @@ export function getErrorName(error: unknown): string { return "Error"; } +/** + * Reconstructs a package error thrown by another worker or iframe realm. + * + * `instanceof` is realm-specific. The public error fields are deliberately + * serializable, so the normalizer recognizes that stable shape and creates a + * local {@link FileSystemError} without degrading its code to `unknown`. + */ +function fromForeignFileSystemError(error: unknown): FileSystemError | undefined { + if (typeof error !== "object" || error === null || getErrorName(error) !== "FileSystemError") return undefined; + const code = ErrorCodeSchema.safeParse(Reflect.get(error, "code")); + const operation = Reflect.get(error, "operation"); + const path = Reflect.get(error, "path"); + if (!code.success || typeof operation !== "string") return undefined; + if (path !== undefined && typeof path !== "string") return undefined; + const cause = Reflect.get(error, "cause"); + return new FileSystemError( + code.data, + operation, + path, + getErrorMessage(error), + cause === undefined ? error : cause, + ); +} + /** Returns a runtime error code such as `ENOENT` when one is exposed. */ function getRuntimeErrorCode(error: unknown): string | undefined { if (typeof error !== "object" || error === null || !("code" in error)) return undefined; @@ -61,6 +85,8 @@ export function getErrorMessage(error: unknown): string { */ export function toFileSystemError(error: unknown, operation: string, path?: string): FileSystemError { if (error instanceof FileSystemError) return error; + const foreign = fromForeignFileSystemError(error); + if (foreign !== undefined) return foreign; const name = getErrorName(error); const runtimeCode = getRuntimeErrorCode(error); diff --git a/src/lock.ts b/src/lock.ts index d3d555e..10fed09 100644 --- a/src/lock.ts +++ b/src/lock.ts @@ -1,4 +1,4 @@ -import { FileSystemError, throwIfAborted } from "./error.ts"; +import { FileSystemError, throwIfAborted, toFileSystemError } from "./error.ts"; import type { CoordinationModeType } from "./schema.ts"; /** Lock access required for one operation. */ @@ -260,8 +260,15 @@ class WebLockCoordinator implements LockCoordinatorType { acquired.resolve(); await hold.promise; }); - await Promise.race([acquired.promise, request]); - return new WebHeldLock(hold.resolve, request); + try { + await Promise.race([acquired.promise, request]); + return new WebHeldLock(hold.resolve, request); + } catch (error) { + // Web Locks rejects a queued request with a realm-native AbortError. + // Normalize at the coordinator seam because several facade operations + // acquire their lock before entering the adapter error-normalization path. + throw toFileSystemError(error, "lock"); + } } } diff --git a/src/path.ts b/src/path.ts index 3203fbf..658174b 100644 --- a/src/path.ts +++ b/src/path.ts @@ -19,6 +19,11 @@ function failPath(path: string, message: string): never { * interpreted as separators. This keeps the same virtual path on Windows, * Unix, OPFS, databases, and key-value stores. * + * Unicode code points are preserved exactly. The function does not apply NFC, + * NFD, or another Unicode normalization form. Two canonically equivalent names + * therefore remain distinct virtual paths unless the selected backend itself + * aliases them. + * * @example * ```ts * normalizePath("reports/../cache/data.bin"); // "/cache/data.bin" diff --git a/src/request.ts b/src/request.ts index 5e116c7..1367c09 100644 --- a/src/request.ts +++ b/src/request.ts @@ -1,4 +1,4 @@ -import { retry } from "@std/async/retry"; +import { RetryError, retry } from "@std/async/retry"; import { z } from "zod"; /** @@ -25,15 +25,26 @@ export const RequestPolicySchema = z.object({ /** A validated direct-client request policy. */ export type RequestPolicyType = z.output; +/** + * Callable Web Fetch contract used by storage clients. + * + * This intentionally models only the standard call signature. Runtime-specific + * globals can attach unrelated properties to `fetch`. Bun, for example, adds + * `fetch.preconnect()`. Using `typeof fetch` here would make that Bun extension + * part of every injected Fetch implementation while Deno type-checks the same + * source. A normal test double only needs to be callable. + */ +export type FetchType = (input: RequestInfo | URL, init?: RequestInit) => Promise; + /** Detached counters for one direct protocol client. */ export interface RequestMetricsType { /** Total HTTP requests actually sent, including retries. */ readonly requests: number; - /** Additional HTTP attempts after an initial failure/status. */ + /** Additional HTTP attempts after an earlier concrete Fetch attempt. */ readonly retries: number; - /** Terminal request failures after retry policy is exhausted. */ + /** Terminal logical request failures after retry policy is exhausted or canceled. */ readonly failures: number; - /** Responses returned to the protocol layer, including non-2xx service responses. */ + /** HTTP responses received, including non-2xx service responses. */ readonly responses: number; /** Total wall-clock milliseconds spent inside Fetch when timing is enabled. */ readonly durationMs: number; @@ -45,7 +56,7 @@ export class RequestMetrics { readonly #timing: boolean; /** Concrete Fetch attempts, including retries. */ #requests = 0; - /** Fetch attempts made after the first attempt for one logical request. */ + /** Fetch attempts made after an earlier concrete Fetch call for one logical request. */ #retries = 0; /** Logical requests that exhausted retry policy or were canceled. */ #failures = 0; @@ -60,24 +71,24 @@ export class RequestMetrics { } /** Records one concrete Fetch call and returns a start timestamp when needed. */ - request(retryAttempt: boolean): number { + request(retryAttempt: boolean): number | undefined { this.#requests += 1; if (retryAttempt) this.#retries += 1; - return this.#timing ? performance.now() : 0; + return this.#timing ? performance.now() : undefined; } /** Records one Fetch response. */ - response(started: number): void { + response(started: number | undefined): void { this.#responses += 1; - if (started !== 0) this.#durationMs += Math.max(0, performance.now() - started); + if (started !== undefined) this.#durationMs += Math.max(0, performance.now() - started); } /** Records elapsed Fetch time for an attempt that rejected before a response arrived. */ - rejected(started: number): void { - if (started !== 0) this.#durationMs += Math.max(0, performance.now() - started); + rejected(started: number | undefined): void { + if (started !== undefined) this.#durationMs += Math.max(0, performance.now() - started); } - /** Records one terminal request failure after retry policy is exhausted or canceled. */ + /** Records one terminal logical request failure. */ failure(): void { this.#failures += 1; } @@ -95,7 +106,7 @@ export class RequestMetrics { } /** Marker for a failure thrown by the concrete Fetch transport after request construction succeeded. */ -export class RequestTransportError extends Error { +class RequestTransportError extends Error { constructor(cause: unknown) { super("Storage request transport failed."); this.name = "RequestTransportError"; @@ -105,7 +116,7 @@ export class RequestTransportError extends Error { /** Internal marker used to make retryable HTTP responses flow through `retry()`. */ class RetryResponseError extends Error { - /** Response retained so the final retry can return it to the protocol parser. */ + /** Response retained for diagnostics while a later attempt is scheduled. */ readonly response: Response; constructor(response: Response) { @@ -115,6 +126,14 @@ class RetryResponseError extends Error { } } +/** Request values prepared before the shared layer owns the concrete Fetch call. */ +interface RequestAttemptType { + /** Fully prepared URL or RequestInfo for this attempt. */ + readonly input: RequestInfo | URL; + /** Fully prepared request initialization for this attempt. */ + readonly init?: RequestInit; +} + /** Validates integer policy values once before a request loop starts. */ function integer(value: number | undefined, fallback: number, name: string, minimum: number): number { const resolved = value ?? fallback; @@ -188,81 +207,129 @@ function getSignal(signal: AbortSignal | undefined, timeoutMs: number | false | }; } -/** Extracts the original error from `@std/async/retry` without coupling to its error class. */ -function cause(error: unknown): unknown { - if (typeof error === "object" && error !== null && "cause" in error) { - return (error as { cause?: unknown }).cause ?? error; +/** + * Waits for request preparation while making the scoped attempt signal authoritative. + * + * Signing and credential callbacks are normally fast, but they are still part + * of one request attempt. Racing preparation with the scoped signal means a + * configured per-attempt timeout also limits a slow credential source. The + * preparation promise can continue internally if that source has no cancellation + * API, but its eventual result can no longer publish an HTTP request. + */ +async function prepare( + create: (signal?: AbortSignal) => Promise, + signal: AbortSignal | undefined, +): Promise { + if (signal === undefined) return await create(); + signal.throwIfAborted(); + + let onAbort!: () => void; + const aborted = new Promise((_resolve, reject) => { + onAbort = () => reject(signal.reason ?? new DOMException("The request attempt was aborted.", "AbortError")); + signal.addEventListener("abort", onAbort, { once: true }); + }); + + try { + return await Promise.race([create(signal), aborted]); + } finally { + signal.removeEventListener("abort", onAbort); } - return error; +} + +/** Returns the error callers should observe after one internal retry marker escapes. */ +function unwrap(error: unknown): unknown { + const original = error instanceof RetryError ? error.cause : error; + return original instanceof RequestTransportError ? original.cause : original; } /** - * Sends a replayable request through `@std/async/retry` while preserving the final HTTP response. + * Sends one storage request through a shared retry and timeout policy. * - * `create` runs for every attempt. This is essential for signed storage - * protocols because credentials and timestamps can change between attempts. - * A non-replayable stream must pass `replayable: false`; it receives exactly - * one attempt rather than risking a second request with an already-consumed body. + * `create` prepares a new URL and `RequestInit` for every attempt. This is + * required for signed protocols because credentials and timestamps can change + * between attempts. The shared layer owns the actual Fetch call so metrics count + * concrete network attempts instead of deterministic signing failures. * - * Request construction, credential, and signing failures are deterministic at - * this layer and are not retried. A client that reaches Fetch and gets a - * transport failure wraps that failure in {@link RequestTransportError}. This - * distinction prevents a malformed signature or invalid request option from - * consuming the retry budget as if it were a transient network failure. + * A non-replayable body or `retry: false` path passes `replayable: false`. That + * path bypasses `@std/async/retry` completely and therefore cannot fail because + * retry-only delay options are invalid for a request that will never retry. + * + * A zero-delay retry policy is supported. `@std/async/retry` requires a positive + * `maxTimeout`, so the shared layer passes `1` as the validation ceiling when the + * project policy requests `0`. With `minTimeout: 0`, the actual retry delay stays + * zero because exponential backoff starts from zero. */ export async function sendRequest( - create: (signal?: AbortSignal) => Promise, + create: (signal?: AbortSignal) => Promise, options: { + /** Concrete Fetch implementation. Runtime globals and ordinary test doubles both satisfy this callable contract. */ + readonly fetch: FetchType; + /** Retry, delay, jitter, and optional attempt-timeout policy. */ readonly policy?: RequestPolicyType; + /** Caller cancellation authority for the complete logical request. */ readonly signal?: AbortSignal; + /** Whether the request can be rebuilt and sent again after a transient failure. */ readonly replayable?: boolean; + /** Optional concrete HTTP counters owned by the protocol client. */ readonly metrics?: RequestMetrics; - } = {}, + }, ): Promise { const policy = getRequestPolicy(options.policy); const attempts = options.replayable === false ? 1 : policy.retries! + 1; let attempt = 0; - let lastStarted = 0; + let fetches = 0; - try { - return await retry(async () => { - attempt += 1; - const scoped = getSignal(options.signal, policy.timeoutMs); - const started = options.metrics?.request(attempt > 1) ?? 0; - lastStarted = started; + const run = async (): Promise => { + attempt += 1; + const scoped = getSignal(options.signal, policy.timeoutMs); + try { + const request = await prepare(create, scoped.signal); + if (options.signal?.aborted) { + throw options.signal.reason ?? new DOMException("The request was aborted.", "AbortError"); + } + scoped.signal?.throwIfAborted(); + + const started = options.metrics?.request(fetches > 0); + fetches += 1; try { - const response = await create(scoped.signal); + const response = await options.fetch(request.input, request.init); options.metrics?.response(started); - lastStarted = 0; if (attempt < attempts && isRetryStatus(response.status)) { await response.body?.cancel().catch(() => undefined); throw new RetryResponseError(response); } return response; } catch (error) { - // RetryResponseError already has a concrete response and its duration was - // recorded above. Network/timeout failures have no Response, so record - // the failed Fetch attempt here without counting it as a terminal failure. if (!(error instanceof RetryResponseError)) options.metrics?.rejected(started); - lastStarted = 0; - throw error; - } finally { - scoped.cleanup(); + if (options.signal?.aborted) throw error; + if (error instanceof RetryResponseError) throw error; + throw new RequestTransportError(scoped.signal?.aborted ? scoped.signal.reason ?? error : error); } - }, { + } catch (error) { + if (options.signal?.aborted) throw error; + if (error instanceof RetryResponseError || error instanceof RequestTransportError) throw error; + if (scoped.signal?.aborted) throw new RequestTransportError(scoped.signal.reason ?? error); + // Request preparation, credentials, canonicalization, and signing failures + // are deterministic at this layer. Do not spend the network retry budget. + throw error; + } finally { + scoped.cleanup(); + } + }; + + try { + if (attempts === 1) return await run(); + return await retry(run, { maxAttempts: attempts, minTimeout: policy.minDelayMs!, - maxTimeout: policy.maxDelayMs!, + maxTimeout: Math.max(1, policy.maxDelayMs!), multiplier: policy.multiplier!, jitter: policy.jitter!, ...(options.signal === undefined ? {} : { signal: options.signal }), isRetriable: (error: unknown) => error instanceof RetryResponseError || error instanceof RequestTransportError, }); } catch (error) { - const original = cause(error); - if (original instanceof RetryResponseError) return original.response; - if (lastStarted !== 0) options.metrics?.rejected(lastStarted); options.metrics?.failure(); - throw original instanceof RequestTransportError ? original.cause : original; + throw unwrap(error); } } diff --git a/src/s3.ts b/src/s3.ts index 1ff357f..6fa9753 100644 --- a/src/s3.ts +++ b/src/s3.ts @@ -4,10 +4,10 @@ import { z } from "zod"; import { split } from "./chunk.ts"; import { + type FetchType, RequestMetrics, type RequestMetricsType, type RequestPolicyType, - RequestTransportError, sendRequest, } from "./request.ts"; import { type AdapterLimitsType, MetricsModeSchema, type MetricsModeType } from "./schema.ts"; @@ -84,7 +84,7 @@ export interface S3ClientOptionsType { /** URL addressing style. Path style is the compatibility-oriented default. */ readonly addressing?: S3AddressingType; /** Fetch implementation. The global Web Fetch API is used by default. */ - readonly fetch?: typeof fetch; + readonly fetch?: FetchType; /** Clock used for Signature Version 4 timestamps. */ readonly now?: () => Date; /** Multipart part size. Defaults to 8 MiB and must be between 5 MiB and 5 GiB. */ @@ -503,7 +503,7 @@ class S3Client implements S3ClientType { /** Path-style or virtual-hosted-style URL strategy. */ readonly #addressing: S3AddressingType; /** Fetch implementation used for every request. */ - readonly #fetch: typeof fetch; + readonly #fetch: FetchType; /** Clock injected for deterministic signing and tests. */ readonly #now: () => Date; /** Configured minimum multipart upload size. */ @@ -788,13 +788,9 @@ class S3Client implements S3ClientType { ...(signal === undefined ? {} : { signal }), }; if (options.body instanceof ReadableStream) init.duplex = "half"; - try { - return await this.#fetch(url, init); - } catch (error) { - if (options.signal?.aborted) throw error; - throw new RequestTransportError(error); - } + return { input: url, init }; }, { + fetch: this.#fetch, ...(this.#requestPolicy === undefined ? {} : { policy: this.#requestPolicy }), ...(options.signal === undefined ? {} : { signal: options.signal }), replayable, diff --git a/src/stream.ts b/src/stream.ts index abd206d..dff02da 100644 --- a/src/stream.ts +++ b/src/stream.ts @@ -170,7 +170,10 @@ class AbortByteSource implements UnderlyingDefaultSource { } } catch (error) { controller.error(error); - await this.close(error); + // Reader cancellation is cleanup after the producer has already failed. + // Do not let a second cleanup rejection replace the producer's original + // failure, especially on runtimes whose native streams reject cancel(). + await this.close(error).catch(() => undefined); } } @@ -206,7 +209,10 @@ class AbortByteSource implements UnderlyingDefaultSource { this.#signal.reason, ); this.#controller?.error(error); - void this.close(error); + // The abort failure above is authoritative. Cancellation only releases the + // borrowed reader, so a runtime-specific cancel() rejection must not become + // a second unhandled failure after callers have already observed `aborted`. + void this.close(error).catch(() => undefined); }; }