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
122 changes: 115 additions & 7 deletions src/publish/s3.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,10 @@
import {
S3Client,
AbortMultipartUploadCommand,
CompleteMultipartUploadCommand,
CreateMultipartUploadCommand,
PutObjectCommand,
UploadPartCommand,
GetObjectCommand,
HeadObjectCommand,
DeleteObjectCommand,
Expand Down Expand Up @@ -91,16 +95,54 @@ export interface S3Ops {
delete(key: string): Promise<void>;
}

export function makeS3Ops(env: S3TargetEnv): S3Ops {
const client = new S3Client({
/**
* How large bodies go up.
*
* A dataset_referenced release uploads parquet in the hundreds of megabytes.
* As one PutObject, a single dropped TLS record anywhere in the stream fails
* the whole file. Above `partSize` a body goes up as a multipart upload
* instead, and each part is retried on its own. 16 MiB keeps a retry cheap and
* stays well inside S3's 10,000-part limit; R2 wants every part but the last
* the same size, which a fixed `partSize` gives.
*/
export interface UploadTuning {
partSize: number;
attempts: number;
backoffMs: number;
}

const DEFAULT_TUNING: UploadTuning = {
partSize: 16 * 1024 * 1024,
attempts: 4,
backoffMs: 1000,
};

// The one SDK method used, and the response fields read; tests hand in a fake.
interface S3Response {
ETag?: string;
UploadId?: string;
ContentLength?: number;
Body?: unknown;
}
type Send = (command: object) => Promise<S3Response>;
interface S3Sender {
send(command: object): Promise<unknown>;
}

export function makeS3Ops(
env: S3TargetEnv,
client: S3Sender = new S3Client({
endpoint: env.endpoint,
region: env.region ?? "auto",
credentials: {
accessKeyId: env.accessKeyId,
secretAccessKey: env.secretAccessKey,
},
forcePathStyle: true,
});
}),
tuning: UploadTuning = DEFAULT_TUNING,
): S3Ops {
const send = client.send.bind(client) as Send;
return {
async put(key, body, conditions) {
const input: {
Expand All @@ -126,16 +168,22 @@ export function makeS3Ops(env: S3TargetEnv): S3Ops {
if (conditions?.cacheControl !== undefined) {
input.CacheControl = conditions.cacheControl;
}
// A conditional write is a pointer, never a dataset, and a multipart
// upload cannot carry the condition.
const conditional = input.IfMatch !== undefined || input.IfNoneMatch !== undefined;
if (typeof body !== "string" && body.byteLength > tuning.partSize && !conditional) {
return putMultipart(send, input, body, tuning);
}
try {
const out = await client.send(new PutObjectCommand(input));
const out = await send(new PutObjectCommand(input));
return { etag: String(out.ETag ?? "") };
} catch (err) {
throw mapS3Error(err);
}
},
async get(key) {
try {
const out = await client.send(new GetObjectCommand({ Bucket: env.bucket, Key: key }));
const out = await send(new GetObjectCommand({ Bucket: env.bucket, Key: key }));
return {
body: await streamToBuffer(out.Body),
etag: String(out.ETag ?? ""),
Expand All @@ -147,19 +195,79 @@ export function makeS3Ops(env: S3TargetEnv): S3Ops {
},
async head(key) {
try {
const out = await client.send(new HeadObjectCommand({ Bucket: env.bucket, Key: key }));
const out = await send(new HeadObjectCommand({ Bucket: env.bucket, Key: key }));
return { size: Number(out.ContentLength ?? 0), etag: String(out.ETag ?? "") };
} catch (err) {
if (isNotFound(err)) return null;
throw mapS3Error(err);
}
},
async delete(key) {
await client.send(new DeleteObjectCommand({ Bucket: env.bucket, Key: key }));
await send(new DeleteObjectCommand({ Bucket: env.bucket, Key: key }));
},
};
}

async function putMultipart(
send: Send,
input: { Bucket: string; Key: string; ContentType?: string; CacheControl?: string },
body: Uint8Array,
tuning: UploadTuning,
): Promise<{ etag: string }> {
const { Bucket, Key } = input;
const created = await send(
new CreateMultipartUploadCommand({
Bucket,
Key,
ContentType: input.ContentType,
CacheControl: input.CacheControl,
}),
);
const UploadId = String(created.UploadId);
const count = Math.ceil(body.byteLength / tuning.partSize);
try {
const parts: { PartNumber: number; ETag: string }[] = [];
// One part at a time: the body is already in memory, and parallel parts
// would only put more of it on a connection that is dropping records.
for (let i = 0; i < count; i++) {
const PartNumber = i + 1;
const Body = body.subarray(i * tuning.partSize, (i + 1) * tuning.partSize);
const out = await withRetries(tuning, () =>
send(new UploadPartCommand({ Bucket, Key, UploadId, PartNumber, Body })),
);
parts.push({ PartNumber, ETag: String(out.ETag ?? "") });
}
const done = await send(
new CompleteMultipartUploadCommand({ Bucket, Key, UploadId, MultipartUpload: { Parts: parts } }),
);
return { etag: String(done.ETag ?? "") };
} catch (err) {
// An upload left open keeps its parts, and its storage bill, until the
// bucket's lifecycle rule (if any) clears it.
await send(new AbortMultipartUploadCommand({ Bucket, Key, UploadId })).catch(() => {});
throw commandError(
"transient_dependency",
`upload of ${Key} failed after ${tuning.attempts} attempts at one part: ${errorText(err)}`,
{ retryable: true, suggested_next: "publish again; parts already sent are not reused" },
);
}
}

async function withRetries<T>(tuning: UploadTuning, run: () => Promise<T>): Promise<T> {
for (let attempt = 1; ; attempt++) {
try {
return await run();
} catch (err) {
if (attempt >= tuning.attempts) throw err;
await new Promise((resolve) => setTimeout(resolve, tuning.backoffMs * attempt));
}
}
}

function errorText(err: unknown): string {
return err instanceof Error ? err.message.trim() : String(err);
}

function isNotFound(err: unknown): boolean {
const name = (err as { name?: string })?.name;
const status = (err as { $metadata?: { httpStatusCode?: number } })?.$metadata
Expand Down
130 changes: 130 additions & 0 deletions tests/publish/s3multipart.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
import { describe, expect, it } from "vitest";
import {
CompleteMultipartUploadCommand,
CreateMultipartUploadCommand,
PutObjectCommand,
UploadPartCommand,
} from "@aws-sdk/client-s3";
import { makeS3Ops } from "../../src/publish/s3.js";

// A dataset_referenced release uploads parquet in the hundreds of megabytes.
// As one PutObject, a single dropped TLS record anywhere in the stream fails
// the whole file, and publish with it. Large bodies go up in parts instead,
// each retried on its own.

const ENV = {
endpoint: "http://localhost",
bucket: "b",
accessKeyId: "a",
secretAccessKey: "s",
};

const MiB = 1024 * 1024;
const TUNING = { partSize: 5 * MiB, attempts: 3, backoffMs: 0 };

type Sent = { name: string; input: Record<string, unknown> };

function fakeClient(failPart?: { number: number; times: number }) {
const sent: Sent[] = [];
let failures = 0;
return {
sent,
async send(command: { constructor: { name: string }; input: object }) {
const name = command.constructor.name;
const input = command.input as Record<string, unknown>;
sent.push({ name, input });
if (command instanceof CreateMultipartUploadCommand) return { UploadId: "u1" };
if (command instanceof UploadPartCommand) {
if (failPart && input.PartNumber === failPart.number && failures < failPart.times) {
failures++;
throw new Error("ssl3_read_bytes:ssl/tls alert bad record mac");
}
return { ETag: `"p${String(input.PartNumber)}"` };
}
if (command instanceof CompleteMultipartUploadCommand) return { ETag: '"whole"' };
if (command instanceof PutObjectCommand) return { ETag: '"single"' };
return {};
},
};
}

const names = (sent: Sent[]) => sent.map((s) => s.name);

describe("S3 uploads of large bodies", () => {
it("sends a body under the threshold as one PutObject", async () => {
const client = fakeClient();
const ops = makeS3Ops(ENV, client, TUNING);
const out = await ops.put("k", new Uint8Array(1024), { contentType: "application/json" });
expect(out.etag).toBe('"single"');
expect(names(client.sent)).toEqual(["PutObjectCommand"]);
});

it("splits a large body into ordered parts and completes with their ETags", async () => {
const client = fakeClient();
const ops = makeS3Ops(ENV, client, TUNING);
const body = new Uint8Array(12 * MiB);
body[0] = 1;
body[12 * MiB - 1] = 2;
const out = await ops.put("data.parquet", body, { contentType: "application/vnd.apache.parquet" });

expect(out.etag).toBe('"whole"');
expect(names(client.sent)).toEqual([
"CreateMultipartUploadCommand",
"UploadPartCommand",
"UploadPartCommand",
"UploadPartCommand",
"CompleteMultipartUploadCommand",
]);
const create = client.sent[0]!.input;
expect(create).toMatchObject({ Bucket: "b", Key: "data.parquet", ContentType: "application/vnd.apache.parquet" });
const parts = client.sent.filter((s) => s.name === "UploadPartCommand").map((s) => s.input);
expect(parts.map((p) => p.PartNumber)).toEqual([1, 2, 3]);
expect(parts.map((p) => (p.Body as Uint8Array).byteLength)).toEqual([5 * MiB, 5 * MiB, 2 * MiB]);
expect((parts[0]!.Body as Uint8Array)[0]).toBe(1);
expect((parts[2]!.Body as Uint8Array)[2 * MiB - 1]).toBe(2);
expect(client.sent[4]!.input).toMatchObject({
UploadId: "u1",
MultipartUpload: {
Parts: [
{ PartNumber: 1, ETag: '"p1"' },
{ PartNumber: 2, ETag: '"p2"' },
{ PartNumber: 3, ETag: '"p3"' },
],
},
});
});

it("retries a failed part without restarting the upload", async () => {
const client = fakeClient({ number: 2, times: 2 });
const ops = makeS3Ops(ENV, client, TUNING);
await ops.put("data.parquet", new Uint8Array(12 * MiB));
const partNumbers = client.sent
.filter((s) => s.name === "UploadPartCommand")
.map((s) => s.input.PartNumber);
expect(partNumbers).toEqual([1, 2, 2, 2, 3]);
expect(names(client.sent).filter((n) => n === "CreateMultipartUploadCommand")).toHaveLength(1);
expect(names(client.sent).at(-1)).toBe("CompleteMultipartUploadCommand");
});

it("aborts the upload when a part keeps failing, and reports it as retryable", async () => {
const client = fakeClient({ number: 2, times: 99 });
const ops = makeS3Ops(ENV, client, TUNING);
await expect(ops.put("data.parquet", new Uint8Array(12 * MiB))).rejects.toMatchObject({
code: "transient_dependency",
retryable: true,
});
expect(names(client.sent).at(-1)).toBe("AbortMultipartUploadCommand");
expect(names(client.sent)).not.toContain("CompleteMultipartUploadCommand");
const abort = client.sent.at(-1)!.input;
expect(abort).toMatchObject({ Bucket: "b", Key: "data.parquet", UploadId: "u1" });
expect(client.sent.filter((s) => s.input.PartNumber === 2)).toHaveLength(TUNING.attempts);
});

it("keeps a conditional write a single PutObject whatever its size", async () => {
const client = fakeClient();
const ops = makeS3Ops(ENV, client, TUNING);
await ops.put("latest.json", new Uint8Array(12 * MiB), { ifNoneMatch: "*" });
expect(names(client.sent)).toEqual(["PutObjectCommand"]);
});
});

Loading