Skip to content
Open
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
22 changes: 22 additions & 0 deletions bench/snippets/blob-append.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
import { bench, run } from "../runner.mjs";

// Accumulating into a Blob by re-wrapping it. Each case builds the whole chain,
// so the per-iteration cost should grow linearly with the number of chunks.
const chunk64KiB = new Uint8Array(64 * 1024);
const chunk64B = new Uint8Array(64);

function accumulate(chunk, count) {
let blob = new Blob([chunk]);
for (let i = 1; i < count; i++) blob = new Blob([blob, chunk]);
return blob;
}

bench("b = new Blob([b, 64 KiB chunk]) x 64", () => accumulate(chunk64KiB, 64));
bench("b = new Blob([b, 64 KiB chunk]) x 256", () => accumulate(chunk64KiB, 256));
bench("b = new Blob([b, 64 B chunk]) x 1024", () => accumulate(chunk64B, 1024));
bench("b = new Blob([b, 64 B chunk]) x 4096", () => accumulate(chunk64B, 4096));

const oneMiB = new Blob([new Uint8Array(1024 * 1024)]);
bench("new Blob([1 MiB blob, 64 B chunk]) once", () => new Blob([oneMiB, chunk64B]));

await run();
19 changes: 14 additions & 5 deletions src/jsc/webcore_types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -590,6 +590,11 @@ pub mod store {
/// rather than `Vec<u8>` so the memfd-backed path
/// (`LinuxMemFdAllocator::create` → `mmap`'d region freed via `munmap`)
/// can carry its allocator vtable with the buffer.
///
/// `ptr[..len]` is immutable. The allocation may be shared: the stores
/// `bun_runtime`'s `AppendBuffer` builds are several `Bytes` viewing
/// prefixes of one allocation, so writing through `ptr` needs more than
/// `&mut self` (see `as_array_list_leak`).
Comment thread
robobun marked this conversation as resolved.
pub struct Bytes {
pub ptr: Option<NonNull<u8>>,
pub len: SizeType,
Expand All @@ -600,11 +605,11 @@ pub mod store {
pub stored_name: Box<[u8]>,
}

// SAFETY: `Bytes` is morally `Vec<u8>`-with-custom-free. The raw
// `NonNull<u8>` is uniquely owned (`ptr` is the sole alias) and
// `StdAllocator` is `Send + Sync`.
// SAFETY: `Bytes` is morally `Vec<u8>`-with-custom-free over immutable
// bytes; every allocator's `free` is thread-safe and `StdAllocator` is
// `Send + Sync`.
unsafe impl Send for Bytes {}
// SAFETY: `&Bytes` only reads the uniquely-owned slice via `slice()`; no
// SAFETY: `&Bytes` only reads the immutable `ptr[..len]` via `slice()`; no
// interior mutability, so sharing references across threads is sound.
unsafe impl Sync for Bytes {}

Expand Down Expand Up @@ -737,9 +742,13 @@ pub mod store {
self.as_array_list_leak()
}

/// Only write through the result (or hand it out writable) when no
/// other `Bytes` shares the allocation: a fresh store, or one that
/// `AppendBuffer::shares_allocation` clears.
Comment thread
robobun marked this conversation as resolved.
pub fn as_array_list_leak(&mut self) -> &mut [u8] {
match self.ptr {
// SAFETY: `ptr[..len]` is live and uniquely owned by `*self`.
// SAFETY: `ptr[..len]` is live while `*self` is; exclusivity of
// the memory is the caller's (see above).
Some(p) => unsafe {
core::slice::from_raw_parts_mut(p.as_ptr(), self.len as usize)
},
Expand Down
29 changes: 28 additions & 1 deletion src/runtime/webcore/Blob.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,9 @@ use crate::node::types::{PathLikeExt as _, PathOrFdExt as _};
use store::{BytesExt as _, FileExt as _, S3Ext as _, StoreExt as _};
pub use store::{Store, StoreRef};

#[path = "blob/AppendBuffer.rs"]
pub(crate) mod append_buffer;
use append_buffer::AppendBuffer;
#[path = "blob/copy_file.rs"]
pub mod copy_file;
#[cfg(not(windows))]
Expand Down Expand Up @@ -3031,7 +3034,11 @@ impl BlobExt for Blob {
}
}
Lifetime::Transfer => {
if self.store().is_some_and(|s| !s.has_one_ref()) {
// JS gets the bytes writable, so they must not be shared either.
if self
.store()
.is_some_and(|s| !s.has_one_ref() || AppendBuffer::shares_allocation(s))
{
// SAFETY: same `buf` contract as the caller; the `Clone` arm only reads it.
let copied = unsafe {
self.to_array_buffer_view_with_bytes::<{ Lifetime::Clone }, TYPED_ARRAY_VIEW>(
Expand Down Expand Up @@ -3340,6 +3347,9 @@ impl BlobExt for Blob {
let mut stack: Vec<JSValue> = Vec::new();
let mut joiner = bun_core::string_joiner::StringJoiner::default();
let mut could_have_non_ascii = false;
// A leading Blob part is appended onto (`AppendBuffer`) instead of
// joined; the ref keeps it alive while later parts may run user JS.
Comment thread
robobun marked this conversation as resolved.
let mut append_prefix: Option<StoreRef> = None;

loop {
match current.js_type_loose() {
Expand Down Expand Up @@ -3432,6 +3442,12 @@ impl BlobExt for Blob {
if let Some(blob) = item.as_class_ref::<Blob>() {
could_have_non_ascii = could_have_non_ascii
|| blob.charset.get() != strings::AsciiStatus::AllAscii;
if append_prefix.is_none() && joiner.len == 0 {
if let Some(store) = AppendBuffer::prefix_store(blob) {
append_prefix = Some(store.clone());
continue;
}
}
// A later part may run user JS that drops the
// last ref to this Blob's Store before `done()`.
if parts_can_run_js {
Expand Down Expand Up @@ -3506,6 +3522,17 @@ impl BlobExt for Blob {
};
}

if let Some(prefix) = append_prefix {
// As below, only a positive ASCII answer is recorded.
let is_all_ascii = (!could_have_non_ascii).then_some(true);
// The joiner's borrowed parts are still alive here, as for `done()`.
let store = AppendBuffer::concat(&prefix, &joiner, is_all_ascii);
let blob = Blob::init_with_store(store, global);
blob.charset
.set(strings::AsciiStatus::from_bool(is_all_ascii));
return Ok(blob);
}

let joined: Vec<u8> = joiner.done().expect("oom").into_vec();

if !could_have_non_ascii {
Expand Down
252 changes: 252 additions & 0 deletions src/runtime/webcore/blob/AppendBuffer.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,252 @@
//! Backing storage shared by the Blobs produced by `b = new Blob([b, chunk])`,
//! which otherwise re-copies the whole prefix on every step.
//!
//! One allocation with spare capacity is shared by every store that successive
//! appends produce; each store is an ordinary immutable `Bytes` viewing the
//! prefix `[0, len)`. An append onto the store viewing the longest published
//! prefix claims `[len, len + n)` and publishes a new store; bytes an existing
//! store can see are never written again. A full buffer is replaced by a
//! bigger one with headroom, so the bytes copied stay linear overall.
//!
//! Like `LinuxMemFdAllocator`, the buffer travels in `Bytes::allocator`: every
//! store built on it owns one reference, released by the vtable's `free`.
Comment thread
robobun marked this conversation as resolved.

use core::ffi::c_void;
use core::mem::ManuallyDrop;
use core::ptr::NonNull;
use core::sync::atomic::{AtomicUsize, Ordering};

use bun_alloc::{Alignment, AllocatorVTable, StdAllocator};
use bun_core::UnwrapOrOom as _;
use bun_core::string_joiner::StringJoiner;

use super::store::{Bytes, Data, Store, StoreRef};
use super::{Blob, SizeType};

#[derive(bun_ptr::ThreadSafeRefCounted)]
pub(crate) struct AppendBuffer {
ref_count: bun_ptr::ThreadSafeRefCount<AppendBuffer>,
/// `capacity` bytes from the global allocator, never reallocated: stores
/// point straight into it.
Comment thread
robobun marked this conversation as resolved.
ptr: NonNull<u8>,
capacity: usize,
/// Length of the longest prefix published as a store; advanced by CAS so
/// two appends onto the same prefix cannot both claim the tail.
Comment thread
robobun marked this conversation as resolved.
committed: AtomicUsize,
}

impl Drop for AppendBuffer {
fn drop(&mut self) {
// SAFETY: `ptr`/`capacity` describe the `Vec` from `create` (or an
// empty one after `take_unique_storage`); no store points into it.
drop(unsafe { Vec::from_raw_parts(self.ptr.as_ptr(), 0, self.capacity) });
}
}

/// Releases the dropped store's reference; `_buf` is only its prefix view.
unsafe fn free(buffer: *mut c_void, _buf: &mut [u8], _: Alignment, _: usize) {
// SAFETY: `buffer` is the context `store` put in the `Bytes`' allocator,
// together with the reference released here.
unsafe { bun_ptr::ThreadSafeRefCount::<AppendBuffer>::deref(buffer.cast::<AppendBuffer>()) };
}

/// Its address identifies buffer-backed `Bytes`, like the memfd vtable does.
static VTABLE: &AllocatorVTable = &AllocatorVTable::free_only(free);

impl AppendBuffer {
/// `blob`'s store, if the first part of `new Blob(parts)` being this blob
/// can be appended onto: in memory, non-empty, and viewed in full.
Comment thread
robobun marked this conversation as resolved.
pub(crate) fn prefix_store(blob: &Blob) -> Option<&StoreRef> {
let store = blob.store()?;
let Data::Bytes(bytes) = &store.data else {
return None;
};
(blob.offset.get() == 0 && bytes.len() > 0 && blob.size.get() == bytes.len())
.then_some(store)
}

/// Whether other stores view the same memory, so that even a store with a
/// single reference must not hand its bytes out writable.
Comment thread
robobun marked this conversation as resolved.
pub(crate) fn shares_allocation(store: &Store) -> bool {
let Data::Bytes(bytes) = &store.data else {
return false;
};
let Some(buffer) = Self::from_allocator(bytes.allocator()) else {
return false;
};
// SAFETY: `bytes` holds a reference on the buffer, so it is live.
!unsafe { &(*buffer).ref_count }.has_one_ref()
}

/// For `Bytes::to_internal_blob`: if `bytes` is the only store on its
/// buffer, moves the allocation out as a `Vec` of the store's bytes and
/// leaves `bytes` empty. `None` if not buffer-backed or the buffer is shared.
Comment thread
robobun marked this conversation as resolved.
pub(crate) fn take_unique_storage(bytes: &mut Bytes) -> Option<Vec<u8>> {
let buffer = Self::from_allocator(bytes.allocator())?;
// SAFETY: `bytes` holds a reference on the buffer, so it is live.
if !unsafe { &(*buffer).ref_count }.has_one_ref() {
return None;
}
let ptr = bytes.ptr.take()?;
let len = core::mem::take(&mut bytes.len) as usize;
bytes.cap = 0;
bytes.allocator = bun_alloc::basic::C_ALLOCATOR;
// SAFETY: `bytes` owns the only reference, so nothing else reaches the
// buffer; `ptr`/`capacity`/`len` are the `Vec` parts from `create`.
// Emptying the header first makes the deref (the reference `bytes`
// owned; with `ptr` taken its `free` no longer runs) free only the header.
unsafe {
let capacity = core::mem::replace(&mut (*buffer).capacity, 0);
(*buffer).ptr = NonNull::dangling();
bun_ptr::ThreadSafeRefCount::<AppendBuffer>::deref(buffer);
Some(Vec::from_raw_parts(ptr.as_ptr(), len, capacity))
}
}

/// The store holding `prefix`'s bytes followed by the joiner's contents.
/// `prefix` must come from [`Self::prefix_store`].
Comment thread
robobun marked this conversation as resolved.
pub(crate) fn concat(
prefix: &StoreRef,
suffix: &StringJoiner<'_>,
is_all_ascii: Option<bool>,
) -> StoreRef {
let Data::Bytes(prefix_bytes) = &prefix.data else {
unreachable!("AppendBuffer::concat prefix is not an in-memory store")
};
let mut grow = false;
if let Some(buffer) = Self::from_allocator(prefix_bytes.allocator()) {
// SAFETY: `prefix_bytes` holds a reference on the buffer.
if let Some(store) = unsafe { Self::append(buffer, prefix_bytes, suffix, is_all_ascii) }
{
return store;
}
// Full, or the tail was already claimed. Only this second-or-later
// append adds headroom; a one-off `new Blob([blob, x])` stays exact.
Comment thread
robobun marked this conversation as resolved.
grow = true;
}
Self::create(prefix_bytes.slice(), suffix, grow, is_all_ascii)
}

fn from_allocator(allocator: StdAllocator) -> Option<*mut AppendBuffer> {
core::ptr::eq(allocator.vtable, VTABLE).then(|| allocator.ptr.cast::<AppendBuffer>())
}

/// In-place append: succeeds only when `prefix` views exactly the
/// committed bytes and the suffix fits.
///
/// # Safety
/// `prefix` must be a `Bytes` built on `*this` by [`Self::store`], so the
/// buffer is live for the duration of the call.
Comment thread
robobun marked this conversation as resolved.
unsafe fn append(
this: *mut AppendBuffer,
prefix: &Bytes,
suffix: &StringJoiner<'_>,
is_all_ascii: Option<bool>,
) -> Option<StoreRef> {
// SAFETY: live per the caller contract, and only `take_unique_storage`
// writes the header, which `prefix` being shared here rules out.
let buffer = unsafe { &*this };
debug_assert_eq!(prefix.slice().as_ptr(), buffer.ptr.as_ptr());
let old_len = prefix.len() as usize;
let new_len = old_len + suffix.len;
if new_len > buffer.capacity {
return None;
}
buffer
.committed
.compare_exchange(old_len, new_len, Ordering::AcqRel, Ordering::Relaxed)
.ok()?;
// SAFETY: the CAS made this call the only writer of `[old_len, new_len)`,
// which is within capacity and not viewed by any store yet.
unsafe { write_suffix(buffer.ptr.as_ptr().add(old_len), suffix) };

// SAFETY: `this` is live; this reference becomes the new store's.
unsafe { bun_ptr::ThreadSafeRefCount::<AppendBuffer>::ref_(this) };
// SAFETY: `[0, new_len)` is initialized and the reference was just taken.
Some(unsafe { Self::store(this, new_len, is_all_ascii) })
}

/// A fresh buffer holding `prefix` followed by the joiner's contents, with
/// 50% headroom when `grow` is set.
Comment thread
robobun marked this conversation as resolved.
fn create(
prefix: &[u8],
suffix: &StringJoiner<'_>,
grow: bool,
is_all_ascii: Option<bool>,
) -> StoreRef {
let len = prefix.len() + suffix.len;
let wanted = if grow {
len.saturating_add(len / 2)
} else {
len
};
let mut storage: Vec<u8> = Vec::new();
storage
.try_reserve_exact(wanted)
.or_else(|_| storage.try_reserve_exact(len))
.unwrap_or_oom();
let mut storage = ManuallyDrop::new(storage);
let capacity = storage.capacity();
let ptr = storage.as_mut_ptr();
// SAFETY: `len <= capacity`, and the sources live in other allocations.
unsafe {
core::ptr::copy_nonoverlapping(prefix.as_ptr(), ptr, prefix.len());
write_suffix(ptr.add(prefix.len()), suffix);
}
let buffer = bun_core::heap::into_raw(Box::new(AppendBuffer {
ref_count: bun_ptr::ThreadSafeRefCount::init(),
// SAFETY: `Vec::as_mut_ptr` is never null.
ptr: unsafe { NonNull::new_unchecked(ptr) },
capacity,
committed: AtomicUsize::new(len),
}));
// SAFETY: `[0, len)` was just written; `init()`'s reference goes to
// this first store.
unsafe { Self::store(buffer, len, is_all_ascii) }
}

/// A store viewing `[0, len)` of the buffer.
///
/// # Safety
/// `this` must be live with `[0, len)` initialized, and the caller hands
/// over one reference, which the store's `Bytes` releases through [`free`].
Comment thread
robobun marked this conversation as resolved.
unsafe fn store(this: *mut AppendBuffer, len: usize, is_all_ascii: Option<bool>) -> StoreRef {
// SAFETY: `this` is live (caller contract).
let ptr = unsafe { (*this).ptr }.as_ptr();
let len = len as SizeType;
// SAFETY: `ptr[..len]` is initialized and outlives the reference handed
// over; `cap == len` keeps `allocated_slice()` within this store's view.
let bytes = unsafe {
Bytes::from_raw_parts(
ptr,
len,
len,
StdAllocator {
ptr: this.cast::<c_void>(),
vtable: VTABLE,
},
)
};
StoreRef::from(Store::new(Store {
data: Data::Bytes(bytes),
mime_type: bun_http_types::MimeType::NONE,
ref_count: bun_ptr::ThreadSafeRefCount::init(),
is_all_ascii,
}))
}
}

/// Writes the joiner's nodes back to back at `dst`.
///
/// # Safety
/// `dst` must be valid for `suffix.len` bytes of writes, and no node may
/// overlap it (unpublished buffer space never does).
Comment thread
robobun marked this conversation as resolved.
unsafe fn write_suffix(dst: *mut u8, suffix: &StringJoiner<'_>) {
let mut written = 0usize;
for node in suffix.node_slices() {
// SAFETY: `suffix.len` is the sum of the node lengths; see the contract.
unsafe { core::ptr::copy_nonoverlapping(node.as_ptr(), dst.add(written), node.len()) };
written += node.len();
}
debug_assert_eq!(written, suffix.len);
}
Loading
Loading