Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/bak/s3.ts

Raw
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<AwsClient, "fetch">;

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<string> {
	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<Response> {
		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<StoredObject | undefined> {
		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<string | undefined> {
		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"));
	}
}