From 7ab7761f1cd239bec5dcc3e1d09c7afbe1a44271 Mon Sep 17 00:00:00 2001 From: Anton Sauchyk Date: Thu, 1 Oct 2026 23:15:09 +0200 Subject: [PATCH] fix(publish): upload large files in parts, retrying each part A single PutObject of a large parquet fails the whole file on one dropped TLS record, and publish with it. Bodies over 16 MiB now go up as a multipart upload, one part at a time, each retried up to four times; a part that keeps failing aborts the upload and reports a retryable error. Conditional writes stay a single PutObject. Co-Authored-By: Claude Opus 5.5 --- src/publish/s3.ts | 122 ++++++++++++++++++++++++++-- tests/publish/s3multipart.test.ts | 130 ++++++++++++++++++++++++++++++ 2 files changed, 245 insertions(+), 7 deletions(-) create mode 100644 tests/publish/s3multipart.test.ts diff --git a/src/publish/s3.ts b/src/publish/s3.ts index b3da6f9..b25da1a 100644 --- a/src/publish/s3.ts +++ b/src/publish/s3.ts @@ -1,6 +1,10 @@ import { S3Client, + AbortMultipartUploadCommand, + CompleteMultipartUploadCommand, + CreateMultipartUploadCommand, PutObjectCommand, + UploadPartCommand, GetObjectCommand, HeadObjectCommand, DeleteObjectCommand, @@ -91,8 +95,43 @@ export interface S3Ops { delete(key: string): Promise; } -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; +interface S3Sender { + send(command: object): Promise; +} + +export function makeS3Ops( + env: S3TargetEnv, + client: S3Sender = new S3Client({ endpoint: env.endpoint, region: env.region ?? "auto", credentials: { @@ -100,7 +139,10 @@ export function makeS3Ops(env: S3TargetEnv): S3Ops { 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: { @@ -126,8 +168,14 @@ 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); @@ -135,7 +183,7 @@ export function makeS3Ops(env: S3TargetEnv): S3Ops { }, 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 ?? ""), @@ -147,7 +195,7 @@ 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; @@ -155,11 +203,71 @@ export function makeS3Ops(env: S3TargetEnv): S3Ops { } }, 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(tuning: UploadTuning, run: () => Promise): Promise { + 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 diff --git a/tests/publish/s3multipart.test.ts b/tests/publish/s3multipart.test.ts new file mode 100644 index 0000000..ce0e09f --- /dev/null +++ b/tests/publish/s3multipart.test.ts @@ -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 }; + +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; + 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"]); + }); +}); +