Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
f4e37fc
fetch: receive an untouched body up to the mark, abort it once its Re…
robobun Aug 19, 2026
2df2c13
Body: make the stream of a body that already failed the body's stream
robobun Aug 20, 2026
2f83da6
fetch: a Bun.write() consumer keeps a collected Response's body loadi…
robobun Aug 20, 2026
f08c478
fetch: one high-water-mark rule for receive backpressure, one abandon…
Jarred-Sumner Aug 24, 2026
18ccd6e
[autofix.ci] apply automated fixes
autofix-ci[bot] Aug 24, 2026
70c94ad
fetch: don't enter BufferAll for readableStreamTo*(res.body)
Jarred-Sumner Aug 24, 2026
dfdbc95
Bun.write(path, response): stream the body into the file
Jarred-Sumner Aug 24, 2026
6f64826
Bun.write(path, stream): reject on write errors, count string chunks,…
Jarred-Sumner Aug 24, 2026
a3ed3c4
test(bun-write): await the rejects assertions
robobun Aug 24, 2026
656599f
Blob: safety comments on the new unsafe blocks, drop a needless borrow
robobun Aug 24, 2026
f2e36e4
Bun.write(dest, readableStream): pipe the stream instead of stringify…
Jarred-Sumner Aug 24, 2026
3bd8ab2
test(bun-write): await the rejects assertion in the ReadableStream so…
Jarred-Sumner Aug 24, 2026
611a586
[autofix.ci] apply automated fixes
autofix-ci[bot] Aug 24, 2026
fdd1c1d
s3: receive backpressure for downloads, same rule as fetch
Jarred-Sumner Aug 24, 2026
962c7bf
One ProducerHold for fetch and S3 body streams; FileSink counts only …
Jarred-Sumner Aug 24, 2026
e31f379
Bun.write(dest, stream): a held reader owns the body; settle on Windo…
Jarred-Sumner Aug 25, 2026
33722c7
[autofix.ci] apply automated fixes
autofix-ci[bot] Aug 25, 2026
868929a
s3: one documented unsafe deref of the task in after_chunk_delivered
robobun Aug 25, 2026
055ad4c
s3: statement-scoped unsafe derefs of the task, each with its safety …
robobun Aug 25, 2026
06cc920
s3: no reference into the download task or its wrapper spans the chun…
robobun Aug 25, 2026
45d00b9
s3: a streamed upload resolves with the bytes written, not 0
Jarred-Sumner Aug 25, 2026
c58a6db
s3: do not pause the transport behind an error body
robobun Aug 25, 2026
bb751b9
s3: don't pause an error body; free the writer() NetworkSink
Jarred-Sumner Aug 25, 2026
93e1718
tests: await events instead of polling; ByteStream: a cancelled sourc…
Jarred-Sumner Aug 25, 2026
c4bd703
NewSource: root_wrapper/unroot_wrapper take the source pointer, not &…
robobun Aug 25, 2026
13dacd4
Merge branch 'main' into farm/7dcb608f/fetch-abandoned-body-drain-cap
robobun Aug 25, 2026
427f1fc
[autofix.ci] apply automated fixes
autofix-ci[bot] Aug 25, 2026
94cb9ca
s3: writer().end() counts the bytes from the upload, not from the sin…
robobun Aug 25, 2026
596869c
s3, fetch: shorter comments on the callback argument and on_stream_ca…
robobun Aug 25, 2026
b899dda
tests: drain both pipes of the two spawned clients; the bun-write ori…
robobun Aug 25, 2026
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
18 changes: 10 additions & 8 deletions packages/bun-types/bun.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2134,7 +2134,7 @@ declare module "bun" {
*/
function write(
destination: BunFile | S3File | PathLike,
input: Blob | NodeJS.TypedArray | ArrayBufferLike | string | BlobPart[] | Archive,
input: Blob | NodeJS.TypedArray | ArrayBufferLike | string | BlobPart[] | Archive | ReadableStream,
options?: {
/**
* If writing to a PathLike, set the permissions of the file.
Expand All @@ -2152,19 +2152,20 @@ declare module "bun" {
): Promise<number>;

/**
* Persist a {@link Response} body to disk.
* Persist a {@link Response} or {@link Request} body to disk. The body is
* streamed into the file as it arrives.
*
* @param destination The file to write to. If the file doesn't exist, it is
* created; if it does, it is overwritten. If `input` is smaller than
* `destination`, `destination` is truncated.
* @param input The `Response` whose body is written
* @param input The `Response` or `Request` whose body is written
* @param options Options for the write
*
* @returns A promise that resolves with the number of bytes written.
*/
function write(
destination: BunFile,
input: Response,
input: Response | Request,
options?: {
/**
* If `true`, create the parent directory if it doesn't exist.
Expand All @@ -2178,17 +2179,18 @@ declare module "bun" {
): Promise<number>;

/**
* Persist a {@link Response} body to disk.
* Persist a {@link Response} or {@link Request} body to disk. The body is
* streamed into the file as it arrives.
*
* @param destinationPath The file path to write to. If the file doesn't
* exist, it is created; if it does, it is overwritten. If `input` is
* smaller than the existing file, the file is truncated.
* @param input The `Response` whose body is written
* @param input The `Response` or `Request` whose body is written
* @returns A promise that resolves with the number of bytes written.
*/
function write(
destinationPath: PathLike,
input: Response,
input: Response | Request,
options?: {
/**
* If `true`, create the parent directory if it doesn't exist.
Expand Down Expand Up @@ -2740,7 +2742,7 @@ declare module "bun" {
* @param options - The options to use for the write.
*/
write(
data: string | ArrayBufferView | ArrayBuffer | SharedArrayBuffer | Request | Response | BunFile,
data: string | ArrayBufferView | ArrayBuffer | SharedArrayBuffer | Request | Response | BunFile | ReadableStream,
options?: { highWaterMark?: number },
): Promise<number>;

Expand Down
58 changes: 40 additions & 18 deletions src/http/Signals.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,17 +12,22 @@ pub struct Signals {
pub body_receive_mode: Option<NonNull<AtomicU8>>,
}

/// Receive backpressure high-water mark: bytes no consumer has taken, on either side of the
/// HTTP→JS hop. A body shorter than this completes unread, which frees its connection.
Comment thread
robobun marked this conversation as resolved.
pub const BODY_HIGH_WATER_MARK: usize = 256 * 1024;

/// Receive backpressure for a body handed to JS. Whichever side holds bytes no consumer has
/// taken moves `Flowing -> Paused` once they reach the high-water mark; whoever takes them
/// moves `Paused -> Flowing` and schedules a resume. The transport applies `Paused` after the
/// next read. Two terminal states: `BufferAll` (a consumer wants the whole body) and
/// `Abandoned` (nothing will read it; the transport is being shut down, drop what arrives).
Comment thread
robobun marked this conversation as resolved.
#[repr(u8)]
#[derive(Copy, Clone, PartialEq, Eq, Debug)]
pub enum BodyReceiveMode {
/// Pause the transport after each delivered body chunk until JS pulls.
AutoPause = 0,
/// `callback` won the CAS; transport should be paused until JS pulls.
Flowing = 0,
Paused = 1,
/// `.arrayBuffer()`/`.text()`/etc attached — never pause.
BufferAll = 2,
/// Cancelled or abandoned — never pause, callback discards bytes.
Ignore = 3,
Abandoned = 3,
}

impl BodyReceiveMode {
Expand All @@ -31,8 +36,8 @@ impl BodyReceiveMode {
match v {
1 => Self::Paused,
2 => Self::BufferAll,
3 => Self::Ignore,
_ => Self::AutoPause,
3 => Self::Abandoned,
_ => Self::Flowing,
}
}
}
Expand Down Expand Up @@ -96,7 +101,7 @@ impl Default for Store {
response_body_streaming: AtomicBool::new(false),
aborted: AtomicBool::new(false),
cert_errors: AtomicBool::new(false),
body_receive_mode: AtomicU8::new(BodyReceiveMode::AutoPause as u8),
body_receive_mode: AtomicU8::new(BodyReceiveMode::Flowing as u8),
}
}
}
Expand Down Expand Up @@ -125,21 +130,38 @@ impl Store {
}

#[inline]
pub fn try_transition_receive_mode(&self, from: BodyReceiveMode, to: BodyReceiveMode) -> bool {
fn try_transition_receive_mode(&self, from: BodyReceiveMode, to: BodyReceiveMode) -> bool {
self.body_receive_mode
.compare_exchange(from as u8, to as u8, Ordering::AcqRel, Ordering::Relaxed)
.is_ok()
}

/// Unconditionally move to a terminal mode (`BufferAll`/`Ignore`).
/// Returns whether the previous state was `Paused`.
/// `Flowing -> Paused`. No-op in the other states.
#[inline]
pub fn pause_receive(&self) {
let _ = self.try_transition_receive_mode(BodyReceiveMode::Flowing, BodyReceiveMode::Paused);
}

/// `Paused -> Flowing`. Returns whether it was paused, i.e. whether the caller has to
/// schedule the transport's resume.
Comment thread
robobun marked this conversation as resolved.
#[inline]
pub fn unpause_receive(&self) -> bool {
self.try_transition_receive_mode(BodyReceiveMode::Paused, BodyReceiveMode::Flowing)
}

/// Terminal: never pause again. Returns whether it was paused.
#[inline]
pub fn receive_all(&self) -> bool {
self.body_receive_mode
.swap(BodyReceiveMode::BufferAll as u8, Ordering::AcqRel)
== BodyReceiveMode::Paused as u8
}

/// Terminal.
#[inline]
pub fn set_receive_mode_terminal(&self, mode: BodyReceiveMode) -> bool {
debug_assert!(matches!(
mode,
BodyReceiveMode::BufferAll | BodyReceiveMode::Ignore
));
self.body_receive_mode.swap(mode as u8, Ordering::AcqRel) == BodyReceiveMode::Paused as u8
pub fn abandon(&self) {
self.body_receive_mode
.store(BodyReceiveMode::Abandoned as u8, Ordering::Release);
}
}

Expand Down
4 changes: 2 additions & 2 deletions src/http/h2_client/ClientSession.rs
Original file line number Diff line number Diff line change
Expand Up @@ -683,8 +683,8 @@ impl ClientSession {
self.by_http_id.get(&async_http_id).copied()
}

/// JS just enabled `response_body_streaming` on the request, so flush any
/// body bytes that arrived between metadata delivery and `getReader()`.
/// A body consumer attached on the JS side: flush any body bytes that arrived between
/// metadata delivery and `getReader()`.
Comment thread
robobun marked this conversation as resolved.
fn drain_response_body(&mut self, async_http_id: u32) {
let Some(stream) = self.stream_for_http_id(async_http_id) else {
return;
Expand Down
4 changes: 4 additions & 0 deletions src/jsc/bindings/webcore/streams/BunStreamSource.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1017,6 +1017,10 @@ static std::optional<bool> rsisWriteChunk(JSC::VM& vm, JSGlobalObject* globalObj
bool shouldSuspend = wrote.isNumber() && wrote.asNumber() < 0;
if (auto* wrotePromise = dynamicDowncast<JSPromise>(wrote)) {
markPromiseAsHandled(vm, wrotePromise);
if (wrotePromise->status() == JSPromise::Status::Rejected) {
throwException(globalObject, scope, wrotePromise->result());
return std::nullopt;
}
shouldSuspend = wrotePromise->status() == JSPromise::Status::Pending;
}
if (shouldSuspend) {
Expand Down
Loading
Loading