From f1fb53ab1b47d79e92a66dae58afbfb0c0278742 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Mon, 21 Sep 2026 10:47:17 +0000 Subject: [PATCH 1/5] s3: a multipart upload part owns the bytes it sends An upload's `process_multi_part` took the buffered bytes out of the upload and handed the new `UploadPart` a raw pointer to them, then transferred the allocation with `ManuallyDrop` only after `enqueue_part` returned `Ok`. An `Err` skipped that transfer, so the part and the local `StreamBuffer` both owned the bytes: the upload's failure freed them through the part, and the buffer freed them again as the error unwound. The part now holds a `Vec` it takes when it is created, so the buffer gives the bytes up at the same point the part receives them. The queue's `Drop` frees what a part still owns, which the raw pointer never did. --- src/runtime/webcore/s3/multipart.rs | 107 ++++++--------------- test/js/bun/s3/s3-upload-terminate.test.ts | 102 ++++++++++++++++++++ 2 files changed, 134 insertions(+), 75 deletions(-) create mode 100644 test/js/bun/s3/s3-upload-terminate.test.ts diff --git a/src/runtime/webcore/s3/multipart.rs b/src/runtime/webcore/s3/multipart.rs index 581c8ddec911..ba0db7509fbe 100644 --- a/src/runtime/webcore/s3/multipart.rs +++ b/src/runtime/webcore/s3/multipart.rs @@ -245,11 +245,9 @@ pub(crate) enum PartState { } pub(crate) struct UploadPart { - /// Raw owned slice; backing allocation length is `allocated_size` (may exceed `data.len()`). - /// Freed via `free_allocated_slice`. Default is a static empty slice. - pub(crate) data: Cell<*const [u8]>, + /// The part's bytes, owned from `get_create_part` until `free_data`. + pub(crate) data: JsCell>, pub ctx: bun_ptr::BackRef, // BACKREF (LIFETIMES.tsv) - pub(crate) allocated_size: Cell, pub(crate) state: Cell, pub(crate) part_number: Cell, // max is 10,000 pub(crate) retry: Cell, // auto retry, decrement until 0 and fail after this @@ -262,25 +260,16 @@ pub(crate) struct UploadPartResult { } impl UploadPart { - fn free_allocated_slice(&self) { - let allocated_size = self.allocated_size.get(); - if allocated_size > 0 { - // SAFETY: `data.ptr` was allocated by the global allocator with capacity == allocated_size - // (either via `to_vec().into_boxed_slice()` where len==cap, or by taking ownership of - // StreamBuffer's backing allocation). Reconstruct and drop. - unsafe { - let ptr = (*self.data.get()).as_ptr().cast_mut(); - drop(Vec::from_raw_parts(ptr, allocated_size, allocated_size)); - } - } - self.data.set(std::ptr::from_ref::<[u8]>(b"" as &[u8])); - self.allocated_size.set(0); + /// Release the bytes. A part that has sent them does not need them again. + fn free_data(&self) { + self.data.set(Vec::new()); } + /// The request `perform` hands these bytes to copies them, and can fail the upload before + /// it returns, which releases them. So this borrow ends with the statement that takes it. #[inline] fn data(&self) -> &[u8] { - // SAFETY: data is either a static empty slice or a live heap slice owned by this part - unsafe { &*self.data.get() } + self.data.get() } fn on_part_response(result: S3PartResult, this: *mut c_void) -> bun_jsc::JsResult<()> { @@ -293,7 +282,7 @@ impl UploadPart { if this.state.get() == PartState::Canceled || ctx.state.get() == State::Finished { scoped_log!(S3MultiPartUpload, "onPartResponse {} canceled", part_number); - this.free_allocated_slice(); + this.free_data(); MultiPartUpload::deref_(ctx_ptr); return Ok(()); } @@ -311,7 +300,7 @@ impl UploadPart { } else { scoped_log!(S3MultiPartUpload, "onPartResponse {} failed", part_number); this.state.set(PartState::NotAssigned); - this.free_allocated_slice(); + this.free_data(); // The ctx deref must run after fail(): let r = ctx.fail(err); MultiPartUpload::deref_(ctx_ptr); @@ -321,7 +310,7 @@ impl UploadPart { S3PartResult::Etag(etag) => { scoped_log!(S3MultiPartUpload, "onPartResponse {} success", part_number); let sent = this.data().len(); - this.free_allocated_slice(); + this.free_data(); // we will need to order this ctx.multipart_etags.with_mut(|etags| { etags.push(UploadPartResult { @@ -390,7 +379,7 @@ impl UploadPart { match state { PartState::Pending => { - self.free_allocated_slice(); + self.free_data(); } // if is not pending we will free later or is already freed _ => {} @@ -401,7 +390,7 @@ impl UploadPart { impl Drop for MultiPartUpload { fn drop(&mut self) { scoped_log!(S3MultiPartUpload, "deinit"); - // queue: Box<[UploadPart]> — dropped automatically (parts' raw `data` already freed during lifecycle) + // queue: Box<[UploadPart]> — dropped automatically, with any `data` a part still owns // KeepAlive::unref takes an `EventLoopCtx` (aio cycle-break vtable), // not `&VirtualMachine`. Route through the global hook like simple_request does. let _ = self.vm; @@ -479,12 +468,9 @@ impl MultiPartUpload { } /// This is the only place we allocate the queue or the parts, this is responsible for the flow of parts and the max allowed concurrency - fn get_create_part( - &self, - chunk: &[u8], - allocated_size: usize, - needs_clone: bool, - ) -> Option<&UploadPart> { + /// + /// `take_data` runs only when the queue has a slot: the part owns the bytes it returns. + fn get_create_part(&self, take_data: impl FnOnce() -> Vec) -> Option<&UploadPart> { let mut available = self.available.get(); let Some(index) = available.find_first_set() else { // this means that the queue is full and we cannot flush it @@ -507,8 +493,7 @@ impl MultiPartUpload { // zero set just in case for _ in 0..queue_size { queue.push(UploadPart { - data: Cell::new(std::ptr::from_ref::<[u8]>(b"" as &[u8])), - allocated_size: Cell::new(0), + data: JsCell::new(Vec::new()), part_number: Cell::new(0), ctx: self_ref, index: Cell::new(0), @@ -518,21 +503,12 @@ impl MultiPartUpload { } self.queue.set(Some(queue.into_boxed_slice())); } - let (data, allocated_len): (*const [u8], usize) = if needs_clone { - let owned = Box::<[u8]>::from(chunk); - let len = owned.len(); - (bun_core::heap::into_raw(owned).cast_const(), len) - } else { - (std::ptr::from_ref::<[u8]>(chunk), allocated_size) - }; - let part_number = self.current_part_number.get(); self.current_part_number.set(part_number + 1); let queue = self.queue.get().as_deref().expect("queue allocated above"); let queue_item = &queue[index]; - queue_item.data.set(data); - queue_item.allocated_size.set(allocated_len); + queue_item.data.set(take_data()); queue_item.part_number.set(part_number); queue_item.index.set(index as u8); // @truncate queue_item.retry.set(self.options.get().retry); @@ -888,13 +864,10 @@ impl MultiPartUpload { ) } - fn enqueue_part( - &self, - chunk: &[u8], - allocated_size: usize, - needs_clone: bool, - ) -> bun_jsc::JsResult { - let Some(part) = self.get_create_part(chunk, allocated_size, needs_clone) else { + /// `Ok(false)`: the queue is full and `take_data` did not run. Otherwise a part owns the + /// bytes, also when starting its request failed. + fn enqueue_part(&self, take_data: impl FnOnce() -> Vec) -> bun_jsc::JsResult { + let Some(part) = self.get_create_part(take_data) else { return Ok(false); }; @@ -958,42 +931,26 @@ impl MultiPartUpload { } // if is one big chunk we can pass ownership and avoid dupe if self.buffered.get().cursor == 0 && self.buffered.get().size() == len { - let owned = self.buffered.replace(StreamBuffer::default()); - // we need to know the allocated size to free the memory later - let allocated_size = owned.memory_cost(); - let slice_len = owned.slice().len(); - // we dont care about the result because we are sending everything - if self.enqueue_part(owned.slice(), allocated_size, false)? { + if self.enqueue_part(|| self.buffered.take().list)? { scoped_log!( S3MultiPartUpload, "processMultiPart {} {} full buffer enqueued", BStr::new(&self.path), - slice_len + len + ); + } else { + scoped_log!( + S3MultiPartUpload, + "processMultiPart {} {} queue full", + BStr::new(&self.path), + len ); - let _ = core::mem::ManuallyDrop::new(owned); - return Ok(()); - } - let appended = self.buffered.replace(owned); - if appended.is_not_empty() { - self.buffered - .with_mut(|b| b.write(appended.slice()).map(|_| ())) - .unwrap_or(()); } - scoped_log!( - S3MultiPartUpload, - "processMultiPart {} {} queue full", - BStr::new(&self.path), - slice_len - ); - return Ok(()); } - let slice_ptr = std::ptr::from_ref::<[u8]>(&self.buffered.get().slice()[..len]); - // allocated size is the slice len because we dupe the buffer - // SAFETY: slice_ptr points at self.buffered's storage which is not mutated until after enqueue_part dupes it - if self.enqueue_part(unsafe { &*slice_ptr }, len, true)? { + if self.enqueue_part(|| self.buffered.get().slice()[..len].to_vec())? { scoped_log!( S3MultiPartUpload, "processMultiPart {} {} slice enqueued", diff --git a/test/js/bun/s3/s3-upload-terminate.test.ts b/test/js/bun/s3/s3-upload-terminate.test.ts new file mode 100644 index 000000000000..659935ff536c --- /dev/null +++ b/test/js/bun/s3/s3-upload-terminate.test.ts @@ -0,0 +1,102 @@ +// A streaming multipart upload buffers bytes until it has a part, then hands +// that part the bytes to send. When the buffer held exactly one part it handed +// the part the buffer's own allocation and gave the buffer up only after the +// part's request had been started, so a request that failed inside the enqueue +// left two owners of those bytes: the part, which the failure frees, and the +// local buffer, freed again when the error unwinds. +// +// terminate() is what makes the enqueue fail that way. The S3 client sends +// nothing new for a VM that is stopping, so it fails the request in place, and +// the upload reports that failure to its promise. Each worker below starts its +// uploads and is terminated as its first parts are being enqueued. A run that +// frees a part's bytes twice aborts the process under ASAN, so a clean exit is +// the check. +import { expect, test } from "bun:test"; +import { bunEnv, bunExe } from "harness"; + +// The terminate has to land while a worker enqueues a part, so each worker +// starts several uploads and several workers run at once. +const WORKERS = 4; +const UPLOADS_PER_WORKER = 8; +// The minimum S3 part size. One chunk is one part, so the first chunk of each +// upload fills the buffer with exactly one part. +const PART_SIZE = 5 * 1024 * 1024; + +const worker = /* js */ ` + const { parentPort, workerData } = require("node:worker_threads"); + const client = new Bun.S3Client({ + endpoint: workerData.endpoint, + bucket: "bucket", + accessKeyId: "key", + secretAccessKey: "secret", + region: "us-east-1", + }); + const chunk = new Uint8Array(${PART_SIZE}).fill(97); + for (let i = 0; i < ${UPLOADS_PER_WORKER}; i++) { + const stream = new ReadableStream({ pull(controller) { controller.enqueue(chunk); } }); + client + .write(workerData.key + "-" + i, new Response(stream), { partSize: ${PART_SIZE}, queueSize: 1, retry: 0 }) + .catch(() => {}); + } + parentPort.postMessage("started"); +`; + +const host = /* js */ ` + const { Worker } = require("node:worker_threads"); + const server = Bun.serve({ + port: 0, + async fetch(req) { + const url = new URL(req.url); + if (req.method === "POST" && url.searchParams.has("uploads")) { + return new Response( + "upload-1", + { headers: { "content-type": "application/xml" } }, + ); + } + if (req.method === "PUT") { + try { await req.arrayBuffer(); } catch {} + return new Response("", { headers: { ETag: '"etag"' } }); + } + try { await req.text(); } catch {} + return new Response( + '"etag-1"', + { headers: { "content-type": "application/xml" } }, + ); + }, + }); + await Promise.all( + Array.from({ length: ${WORKERS} }, (_, i) => { + const w = new Worker(${JSON.stringify(worker)}, { + eval: true, + workerData: { endpoint: server.url.origin, key: "key-" + i }, + }); + // Terminate as soon as this worker reports that its uploads are started. + return new Promise(resolve => w.once("message", resolve)).then(() => w.terminate()); + }), + ); + server.stop(true); + console.log("ok"); +`; + +test("terminate() while a multipart upload enqueues a part frees the part's bytes once", async () => { + await using proc = Bun.spawn({ + cmd: [bunExe(), "-e", host], + env: { + ...bunEnv, + // The S3 client does not honor NO_PROXY, so an inherited proxy would + // hijack the loopback stand-in. + HTTP_PROXY: undefined, + HTTPS_PROXY: undefined, + http_proxy: undefined, + https_proxy: undefined, + // An upload still open when its worker goes is not freed (#39692). That + // leak is not what this checks, and LeakSanitizer would report it. + ASAN_OPTIONS: [bunEnv.ASAN_OPTIONS, "detect_leaks=0"].filter(Boolean).join(":"), + }, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stdout: stdout.trim(), stderr: stderr.trim() }).toEqual({ stdout: "ok", stderr: "" }); + expect(exitCode).toBe(0); +}); From 77c5c0a4925bdbf9e42a4237ca752641c441e7e4 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Mon, 21 Sep 2026 20:23:20 +0000 Subject: [PATCH 2/5] s3: keep a part's bytes alive for the request perform lends them to `UploadPart::perform` lent the part's bytes to `execute_simple_s3_request`, which reports a request that fails before it leaves inside that call. The report frees the part, so the lent slice pointed at freed memory until the call returned. Nothing read it, and only a comment said so. The bytes are now an `Rc>`. `perform` holds a count for the length of the call, so the report drops the part's count and the bytes go when the call returns. The test clears `ALL_PROXY` as well, which the S3 client falls back to, and its comment no longer says that the client ignores `NO_PROXY`. --- src/runtime/webcore/s3/multipart.rs | 26 ++++++++++------------ test/js/bun/s3/s3-upload-terminate.test.ts | 6 +++-- 2 files changed, 16 insertions(+), 16 deletions(-) diff --git a/src/runtime/webcore/s3/multipart.rs b/src/runtime/webcore/s3/multipart.rs index ba0db7509fbe..b992c2fdc379 100644 --- a/src/runtime/webcore/s3/multipart.rs +++ b/src/runtime/webcore/s3/multipart.rs @@ -92,6 +92,7 @@ use core::cell::Cell; use core::ffi::c_void; use std::io::Write as _; +use std::rc::Rc; use bstr::BStr; @@ -245,8 +246,9 @@ pub(crate) enum PartState { } pub(crate) struct UploadPart { - /// The part's bytes, owned from `get_create_part` until `free_data`. - pub(crate) data: JsCell>, + /// The part's bytes, from `get_create_part` until `free_data`. Counted, so that `perform` + /// keeps them for a request that frees the part before it returns. + pub(crate) data: JsCell>>>, pub ctx: bun_ptr::BackRef, // BACKREF (LIFETIMES.tsv) pub(crate) state: Cell, pub(crate) part_number: Cell, // max is 10,000 @@ -262,14 +264,7 @@ pub(crate) struct UploadPartResult { impl UploadPart { /// Release the bytes. A part that has sent them does not need them again. fn free_data(&self) { - self.data.set(Vec::new()); - } - - /// The request `perform` hands these bytes to copies them, and can fail the upload before - /// it returns, which releases them. So this borrow ends with the statement that takes it. - #[inline] - fn data(&self) -> &[u8] { - self.data.get() + self.data.set(None); } fn on_part_response(result: S3PartResult, this: *mut c_void) -> bun_jsc::JsResult<()> { @@ -309,7 +304,7 @@ impl UploadPart { } S3PartResult::Etag(etag) => { scoped_log!(S3MultiPartUpload, "onPartResponse {} success", part_number); - let sent = this.data().len(); + let sent = this.data.get().as_deref().map_or(0, Vec::len); this.free_data(); // we will need to order this ctx.multipart_etags.with_mut(|etags| { @@ -334,6 +329,9 @@ impl UploadPart { fn perform(&self) -> bun_jsc::JsResult<()> { let ctx = self.ctx.get(); + // A request that fails before it leaves reports that inside the call below, which frees + // this part. This count keeps the bytes the call borrows until it returns. + let data = self.data.get().clone(); let mut params_buffer = [0u8; 2048]; let written = { let mut w: &mut [u8] = &mut params_buffer[..]; @@ -354,7 +352,7 @@ impl UploadPart { path: &ctx.path, method: bun_http::Method::PUT, proxy_url: ctx.proxy_url(), - body: self.data(), + body: data.as_deref().map_or(&[][..], Vec::as_slice), search_params: Some(search_params), request_payer: ctx.request_payer, ..Default::default() @@ -493,7 +491,7 @@ impl MultiPartUpload { // zero set just in case for _ in 0..queue_size { queue.push(UploadPart { - data: JsCell::new(Vec::new()), + data: JsCell::new(None), part_number: Cell::new(0), ctx: self_ref, index: Cell::new(0), @@ -508,7 +506,7 @@ impl MultiPartUpload { let queue = self.queue.get().as_deref().expect("queue allocated above"); let queue_item = &queue[index]; - queue_item.data.set(take_data()); + queue_item.data.set(Some(Rc::new(take_data()))); queue_item.part_number.set(part_number); queue_item.index.set(index as u8); // @truncate queue_item.retry.set(self.options.get().retry); diff --git a/test/js/bun/s3/s3-upload-terminate.test.ts b/test/js/bun/s3/s3-upload-terminate.test.ts index 659935ff536c..98d5f4377e7e 100644 --- a/test/js/bun/s3/s3-upload-terminate.test.ts +++ b/test/js/bun/s3/s3-upload-terminate.test.ts @@ -83,12 +83,14 @@ test("terminate() while a multipart upload enqueues a part frees the part's byte cmd: [bunExe(), "-e", host], env: { ...bunEnv, - // The S3 client does not honor NO_PROXY, so an inherited proxy would - // hijack the loopback stand-in. + // An inherited proxy would take the requests away from the loopback + // stand-in. HTTP_PROXY: undefined, HTTPS_PROXY: undefined, + ALL_PROXY: undefined, http_proxy: undefined, https_proxy: undefined, + all_proxy: undefined, // An upload still open when its worker goes is not freed (#39692). That // leak is not what this checks, and LeakSanitizer would report it. ASAN_OPTIONS: [bunEnv.ASAN_OPTIONS, "detect_leaks=0"].filter(Boolean).join(":"), From 0c3bdf31630ca39a0dea105750f00d032d6af577 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Mon, 21 Sep 2026 20:31:44 +0000 Subject: [PATCH 3/5] s3: shorten the comments on the multipart part's bytes One line each. The field's comment says why the bytes are counted, so `perform` does not repeat it. --- src/runtime/webcore/s3/multipart.rs | 11 +++-------- 1 file changed, 3 insertions(+), 8 deletions(-) diff --git a/src/runtime/webcore/s3/multipart.rs b/src/runtime/webcore/s3/multipart.rs index b992c2fdc379..7dbade5d82f9 100644 --- a/src/runtime/webcore/s3/multipart.rs +++ b/src/runtime/webcore/s3/multipart.rs @@ -246,8 +246,7 @@ pub(crate) enum PartState { } pub(crate) struct UploadPart { - /// The part's bytes, from `get_create_part` until `free_data`. Counted, so that `perform` - /// keeps them for a request that frees the part before it returns. + /// Counted so that `perform` can keep the bytes across a request that frees the part. pub(crate) data: JsCell>>>, pub ctx: bun_ptr::BackRef, // BACKREF (LIFETIMES.tsv) pub(crate) state: Cell, @@ -329,8 +328,6 @@ impl UploadPart { fn perform(&self) -> bun_jsc::JsResult<()> { let ctx = self.ctx.get(); - // A request that fails before it leaves reports that inside the call below, which frees - // this part. This count keeps the bytes the call borrows until it returns. let data = self.data.get().clone(); let mut params_buffer = [0u8; 2048]; let written = { @@ -466,8 +463,7 @@ impl MultiPartUpload { } /// This is the only place we allocate the queue or the parts, this is responsible for the flow of parts and the max allowed concurrency - /// - /// `take_data` runs only when the queue has a slot: the part owns the bytes it returns. + /// `take_data` runs only when the queue has a slot, and the part owns the bytes it returns. fn get_create_part(&self, take_data: impl FnOnce() -> Vec) -> Option<&UploadPart> { let mut available = self.available.get(); let Some(index) = available.find_first_set() else { @@ -862,8 +858,7 @@ impl MultiPartUpload { ) } - /// `Ok(false)`: the queue is full and `take_data` did not run. Otherwise a part owns the - /// bytes, also when starting its request failed. + /// `Ok(false)`: the queue is full and `take_data` did not run. On `Err` a part has the bytes. fn enqueue_part(&self, take_data: impl FnOnce() -> Vec) -> bun_jsc::JsResult { let Some(part) = self.get_create_part(take_data) else { return Ok(false); From e0bca9ce3e72d7ba699b14c32728c6e2f19b55d8 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Mon, 21 Sep 2026 20:51:54 +0000 Subject: [PATCH 4/5] s3: name a new part's bytes with an enum instead of a closure `get_create_part` and `enqueue_part` were generic over the closure that produced the part's bytes, so each was compiled once per call site although almost none of their code used the closure. The bytes always come from the upload's buffer, either all of them or a copy of the first part of them, so `PartBytes` says which and `get_create_part` takes them itself, still only once the queue has a slot. --- src/runtime/webcore/s3/multipart.rs | 29 +++++++++++++++++++++-------- 1 file changed, 21 insertions(+), 8 deletions(-) diff --git a/src/runtime/webcore/s3/multipart.rs b/src/runtime/webcore/s3/multipart.rs index 7dbade5d82f9..676268069c32 100644 --- a/src/runtime/webcore/s3/multipart.rs +++ b/src/runtime/webcore/s3/multipart.rs @@ -260,6 +260,15 @@ pub(crate) struct UploadPartResult { pub(crate) etag: Box<[u8]>, } +/// Which bytes of `buffered` a new part takes. +#[derive(Clone, Copy)] +enum PartBytes { + /// All of them: the part takes the buffer's allocation. + Whole, + /// A copy of the first `usize` bytes. + First(usize), +} + impl UploadPart { /// Release the bytes. A part that has sent them does not need them again. fn free_data(&self) { @@ -463,8 +472,8 @@ impl MultiPartUpload { } /// This is the only place we allocate the queue or the parts, this is responsible for the flow of parts and the max allowed concurrency - /// `take_data` runs only when the queue has a slot, and the part owns the bytes it returns. - fn get_create_part(&self, take_data: impl FnOnce() -> Vec) -> Option<&UploadPart> { + /// The bytes leave `buffered` only when the queue has a slot, and the part owns them. + fn get_create_part(&self, bytes: PartBytes) -> Option<&UploadPart> { let mut available = self.available.get(); let Some(index) = available.find_first_set() else { // this means that the queue is full and we cannot flush it @@ -502,7 +511,11 @@ impl MultiPartUpload { let queue = self.queue.get().as_deref().expect("queue allocated above"); let queue_item = &queue[index]; - queue_item.data.set(Some(Rc::new(take_data()))); + let data = match bytes { + PartBytes::Whole => self.buffered.take().list, + PartBytes::First(len) => self.buffered.get().slice()[..len].to_vec(), + }; + queue_item.data.set(Some(Rc::new(data))); queue_item.part_number.set(part_number); queue_item.index.set(index as u8); // @truncate queue_item.retry.set(self.options.get().retry); @@ -858,9 +871,9 @@ impl MultiPartUpload { ) } - /// `Ok(false)`: the queue is full and `take_data` did not run. On `Err` a part has the bytes. - fn enqueue_part(&self, take_data: impl FnOnce() -> Vec) -> bun_jsc::JsResult { - let Some(part) = self.get_create_part(take_data) else { + /// `Ok(false)`: the queue is full and `buffered` is untouched. On `Err` a part has the bytes. + fn enqueue_part(&self, bytes: PartBytes) -> bun_jsc::JsResult { + let Some(part) = self.get_create_part(bytes) else { return Ok(false); }; @@ -925,7 +938,7 @@ impl MultiPartUpload { // if is one big chunk we can pass ownership and avoid dupe if self.buffered.get().cursor == 0 && self.buffered.get().size() == len { // we dont care about the result because we are sending everything - if self.enqueue_part(|| self.buffered.take().list)? { + if self.enqueue_part(PartBytes::Whole)? { scoped_log!( S3MultiPartUpload, "processMultiPart {} {} full buffer enqueued", @@ -943,7 +956,7 @@ impl MultiPartUpload { return Ok(()); } - if self.enqueue_part(|| self.buffered.get().slice()[..len].to_vec())? { + if self.enqueue_part(PartBytes::First(len))? { scoped_log!( S3MultiPartUpload, "processMultiPart {} {} slice enqueued", From 0765a97af0635b5911037ac761dfae8dee830bb5 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Mon, 21 Sep 2026 20:59:46 +0000 Subject: [PATCH 5/5] test: fail the s3 upload terminate test when a worker fails before it starts The host waited only for each worker's message, so a worker that threw or exited first left the run to the runner's timeout. Its `error` and `exit` events now reject the wait, and the host prints the error and exits 1. --- test/js/bun/s3/s3-upload-terminate.test.ts | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/test/js/bun/s3/s3-upload-terminate.test.ts b/test/js/bun/s3/s3-upload-terminate.test.ts index 98d5f4377e7e..9fcb9d2bdac4 100644 --- a/test/js/bun/s3/s3-upload-terminate.test.ts +++ b/test/js/bun/s3/s3-upload-terminate.test.ts @@ -71,9 +71,17 @@ const host = /* js */ ` workerData: { endpoint: server.url.origin, key: "key-" + i }, }); // Terminate as soon as this worker reports that its uploads are started. - return new Promise(resolve => w.once("message", resolve)).then(() => w.terminate()); + // A worker that fails before that fails the run instead of hanging it. + return new Promise((resolve, reject) => { + w.once("message", resolve); + w.once("error", reject); + w.once("exit", code => reject(new Error("worker exited with code " + code + " before it started its uploads"))); + }).then(() => w.terminate()); }), - ); + ).catch(error => { + console.error(error); + process.exit(1); + }); server.stop(true); console.log("ok"); `;