import { AwsClient } from "aws4fetch"; import type { S3Config } from "./config.js"; export { type S3Config, s3Config } from "./config.js"; const ERROR_TEXT_LIMIT = 500; const MAX_ATTEMPTS = 3; const REQUEST_TIMEOUT_MS = 30_000; export interface StoredObject { body: Uint8Array; etag?: string; } export class S3Error extends Error { constructor( readonly status: number, message: string, ) { super(message); this.name = "S3Error"; } } type SignedClient = Pick; function objectUrl(config: S3Config, key: string): string { const endpoint = config.endpoint.replace(/\/+$/, ""); const path = [config.bucket, ...key.split("/")] .map((part) => encodeURIComponent(part)) .join("/"); return `${endpoint}/${path}`; } function shouldRetry(status: number): boolean { return status === 408 || status === 429 || status >= 500; } function strongEtag(value: string | null): string | undefined { const etag = value?.trim().replace(/^W\//, ""); return etag || undefined; } async function boundedError(response: Response): Promise { const text = (await response.text()).replace(/\s+/g, " ").trim(); return ( text.slice(0, ERROR_TEXT_LIMIT) || response.statusText || "request failed" ); } export class S3Store { private readonly client: SignedClient; constructor( readonly config: S3Config, client?: SignedClient, ) { this.client = client ?? new AwsClient({ accessKeyId: config.accessKeyId, secretAccessKey: config.secretAccessKey, service: "s3", region: config.region, retries: 0, }); } private async request( key: string, init: RequestInit, signal?: AbortSignal, ): Promise { let lastError: unknown; for (let attempt = 0; attempt < MAX_ATTEMPTS; attempt++) { signal?.throwIfAborted(); try { const timeout = AbortSignal.timeout(REQUEST_TIMEOUT_MS); const response = await this.client.fetch(objectUrl(this.config, key), { ...init, signal: signal ? AbortSignal.any([signal, timeout]) : timeout, }); if (!shouldRetry(response.status) || attempt === MAX_ATTEMPTS - 1) return response; await response.body?.cancel(); } catch (error) { lastError = error; signal?.throwIfAborted(); if (attempt === MAX_ATTEMPTS - 1) throw error; } await new Promise((resolve) => setTimeout(resolve, 250 * 2 ** attempt)); } throw lastError instanceof Error ? lastError : new Error("S3 request failed"); } async get( key: string, options: { signal?: AbortSignal } = {}, ): Promise { const response = await this.request( key, { method: "GET", headers: { "accept-encoding": "identity" } }, options.signal, ); if (response.status === 404) return undefined; if (!response.ok) throw new S3Error(response.status, await boundedError(response)); return { body: new Uint8Array(await response.arrayBuffer()), etag: strongEtag(response.headers.get("etag")), }; } async put( key: string, body: string | Uint8Array, options: { contentType?: string; ifMatch?: string; ifNoneMatch?: string; signal?: AbortSignal; } = {}, ): Promise { const headers = new Headers({ "content-type": options.contentType ?? "application/octet-stream", }); const ifMatch = strongEtag(options.ifMatch ?? null); if (ifMatch) headers.set("if-match", ifMatch); if (options.ifNoneMatch) headers.set("if-none-match", options.ifNoneMatch); const requestBody = typeof body === "string" ? body : Uint8Array.from(body).buffer; const response = await this.request( key, { method: "PUT", headers, body: requestBody, }, options.signal, ); if (!response.ok) throw new S3Error(response.status, await boundedError(response)); return strongEtag(response.headers.get("etag")); } }