repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/bak/s3.ts
Rawimport { 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"));
}
}