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
19 changes: 19 additions & 0 deletions src/runtime/webcore/Blob.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5997,6 +5997,25 @@ fn resolve_file_stat(store: &RefPtr<Store>) {
}
}

/// Whether a second Blob over `store` reads the same bytes from the start.
/// Memory and S3 do. A path does when it names a regular file (each read
/// opens it again); it is stat'd here if that is not yet known. A file
/// descriptor never does: its offset, and for a pipe its bytes, are shared.
pub(crate) fn store_reads_repeatably(store: &RefPtr<Store>) -> bool {
match Store::data_mut(store).tag() {
store::DataTag::Bytes | store::DataTag::S3 => true,
store::DataTag::File => {
if let PathOrFileDescriptor::Fd(_) = Store::data_mut(store).as_file().pathlike {
return false;
}
if Store::data_mut(store).as_file().seekable.is_none() {
resolve_file_stat(store);
}
Store::data_mut(store).as_file().seekable != Some(false)
}
}
}

// ──────────────────────────────────────────────────────────────────────────
// toStringWithBytes / toString / toJSON / toFormData / toArrayBuffer{View}
// ──────────────────────────────────────────────────────────────────────────
Expand Down
91 changes: 72 additions & 19 deletions src/runtime/webcore/Body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -370,6 +370,34 @@ impl PendingValue {
None
}

/// [`Self::to_any_blob`] for `clone()`: the stream is the wrapper's cached
/// `.body` when there is one, and it is detached afterwards so a reference
/// the caller still holds reads as locked, the state a tee leaves it in.
fn take_blob_from_unread_stream(
&mut self,
global: &JSGlobalObject,
cached: Option<ReadableStream>,
) -> Option<AnyBlob> {
if self.promise.is_some() || self.on_receive_value.is_some() {
return None;
}
let mut stream = cached.or_else(|| self.readable.get())?;
// Two Blobs over a pipe, a tty, or any fd would compete for its bytes,
// so a stream over such a file stays teed.
if let Some(reader) = stream.ptr.file() {
let webcore::file_reader::Lazy::Blob(store) = reader.lazy.get() else {
return None;
};
if !blob::store_reads_repeatably(store) {
return None;
}
}
let blob = stream.to_any_blob(global)?;
stream.force_detach(global);
self.readable.deinit();
Some(blob)
}

fn set_promise(
&mut self,
global_this: &JSGlobalObject,
Expand Down Expand Up @@ -1526,11 +1554,16 @@ impl Value {
global_this: &JSGlobalObject,
readable: Option<&mut ReadableStream>,
) -> JsResult<Value> {
// Tee a Locked body before any blob extraction: `to_blob_if_possible()`
// would `.done()` an already-materialized `.body` stream, leaving the
// user-visible cached stream empty instead of a live tee branch.
if matches!(self, Value::Locked(_)) {
return self.tee(global_this, readable);
// A native blob, file, or fully buffered byte stream that nothing has
// read goes back to being a Blob, so both bodies share one store (and
// its type) instead of pumping the bytes through a JS tee. The owner
// must then drop its cached `.body` (`sync_body_stream_caches`).
// Anything else is teed.
if let Value::Locked(locked) = self {
match locked.take_blob_from_unread_stream(global_this, readable.as_deref().copied()) {
Some(blob) => *self = Value::from(blob),
None => return self.tee(global_this, readable),
}
}

self.to_blob_if_possible();
Expand All @@ -1541,6 +1574,14 @@ impl Value {
}

if let Value::Blob(b) = self {
if b.store()
.is_some_and(|store| !blob::store_reads_repeatably(store))
{
// A pipe or other fd yields its bytes once: read it as one
// stream and tee that.
self.to_readable_stream(global_this)?;
return self.tee(global_this, None);
}
return Ok(Value::Blob(b.dupe_with_content_type(false)));
}

Expand Down Expand Up @@ -1651,9 +1692,29 @@ pub(crate) trait BodyMixin: BodyOwnerJs + Sized {
}
}

/// After `clone()` replaced this body's stream: point the wrapper's cached
/// `body` at the tee branch now in `Locked.readable`, or, when the clone
/// moved an unread native stream back into a Blob, forget the detached
/// stream so `.body` is rebuilt from that Blob.
fn sync_body_stream_caches(&self, this_value: JSValue, global_this: &JSGlobalObject) {
match self.get_body_value() {
Value::Locked(locked) => {
if let Some(readable) = locked.readable.get() {
Self::body_set_cached(this_value, global_this, readable.value);
}
}
_ => {
if Self::stream_get_cached(this_value).is_some() {
Self::stream_set_cached(this_value, global_this, JSValue::ZERO);
Self::body_set_cached(this_value, global_this, JSValue::ZERO);
}
}
}
}

/// Shared tail of `do_clone`: after the clone's `to_js` ran
/// `check_body_stream_ref`, sync both wrappers' cached `body` slots to
/// their respective teed streams, then migrate the original's
/// their respective streams, then migrate the original's
/// `Locked.readable` into its own `js.gc.stream`.
fn sync_cloned_body_stream_caches(
&self,
Expand All @@ -1666,17 +1727,13 @@ pub(crate) trait BodyMixin: BodyOwnerJs + Sized {
Self::body_set_cached(js_wrapper, global_this, cloned_stream);
}
}
if let Value::Locked(locked) = self.get_body_value() {
if let Some(readable) = locked.readable.get() {
Self::body_set_cached(this_value, global_this, readable.value);
}
}
self.sync_body_stream_caches(this_value, global_this);
self.check_body_stream_ref(global_this);
}

/// Shared body-clone for `clone_into` / `clone_value`: tee through the
/// JS-side cached stream when present, then repoint this owner's
/// `body`/`stream` cache slots at the fresh branch in `locked.readable`.
/// Shared body-clone for `clone_into` / `clone_value`: clone through the
/// JS-side cached stream when present, then resync this owner's
/// `body`/`stream` cache slots with whatever the body now holds.
fn clone_body_value_via_cached_stream(&self, global_this: &JSGlobalObject) -> JsResult<Value> {
let cloned = 'brk: {
if let Some(js_ref) = self.js_ref() {
Expand All @@ -1692,11 +1749,7 @@ pub(crate) trait BodyMixin: BodyOwnerJs + Sized {
self.get_body_value().clone(global_this)?
};
if let Some(js_ref) = self.js_ref() {
if let Value::Locked(locked) = self.get_body_value() {
if let Some(readable) = locked.readable.get() {
Self::body_set_cached(js_ref, global_this, readable.value);
}
}
self.sync_body_stream_caches(js_ref, global_this);
}
self.check_body_stream_ref(global_this);
Ok(cloned)
Expand Down
132 changes: 131 additions & 1 deletion test/js/web/fetch/body-clone.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { describe, expect, test } from "bun:test";
import { bunEnv, bunExe, isASAN, isDebug } from "harness";
import { bunEnv, bunExe, isASAN, isDebug, isWindows, tempDirWithFiles } from "harness";
import { join } from "node:path";

test("Request with streaming body can be cloned", async () => {
const stream = new ReadableStream({
Expand Down Expand Up @@ -908,9 +909,19 @@ describe.concurrent("clone() after `.body` was observed returns a fresh tee bran
// Blob-backed stream
// - new Response(ReadableStream) → Locked with a user stream already
// rooted in the JS-side stream slot
// - new Response(Bun.file()) → Blob over a file store, then .body
// materializes a file-backed stream
// - new Response(Bun.file().stream()) → Locked with an unread
// file-backed stream
const N = 8192;
const payload = Buffer.alloc(N, "a");
const fileDir = tempDirWithFiles("body-clone-observe", { "payload.bin": payload });
const cases: Array<[string, () => Promise<Request | Response>]> = [
["Response with a Bun.file() body", async () => new Response(Bun.file(join(fileDir, "payload.bin")))],
[
"Response with a Bun.file() stream body",
async () => new Response(Bun.file(join(fileDir, "payload.bin")).stream()),
],
[
"fetch() Response with a buffered body",
async () => {
Expand Down Expand Up @@ -1054,6 +1065,125 @@ test("new Request(src, init) with a user ReadableStream body: both derived and s
});
});

// The readers and Bun.serve move an unread Bun.file()/Blob stream back into
// the Blob it came from, type included. clone() must do the same instead of
// teeing it into two plain JS streams, or the clone (and, once `.body` was
// observed, the original too) answers differently from an un-cloned body.
describe("clone() of a body over an unread native stream keeps the Blob behind it", () => {
const dir = tempDirWithFiles("body-clone-type", { "page.html": "<p>hi</p>" });
const file = () => Bun.file(join(dir, "page.html"));
const typed = () => new Blob(["<p>hi</p>"], { type: "text/html;charset=utf-8" });

async function typesAndText(original: Request | Response) {
const clone = original.clone();
const [a, b] = await Promise.all([original.blob(), clone.blob()]);
return { original: [a.type, await a.text()], clone: [b.type, await b.text()] };
}
const expected = {
original: ["text/html;charset=utf-8", "<p>hi</p>"],
clone: ["text/html;charset=utf-8", "<p>hi</p>"],
};

test("Response over Bun.file().stream()", async () => {
expect(await typesAndText(new Response(file().stream()))).toEqual(expected);
});

test("Request over Bun.file().stream()", async () => {
expect(await typesAndText(new Request("http://example.com/", { method: "POST", body: file().stream() }))).toEqual(
expected,
);
});

test("Response over Bun.file() after .body was observed", async () => {
const response = new Response(file());
expect(response.body).toBeInstanceOf(ReadableStream);
expect(await typesAndText(response)).toEqual(expected);
});

test("Response over a typed Blob after .body was observed", async () => {
const response = new Response(typed());
expect(response.body).toBeInstanceOf(ReadableStream);
expect(await typesAndText(response)).toEqual(expected);
});

test("the stream given to the constructor is locked after clone() and .body is a fresh stream", async () => {
const stream = file().stream();
const response = new Response(stream);
expect(response.body).toBe(stream);
const clone = response.clone();
expect(stream.locked).toBe(true);
expect(response.body).not.toBe(stream);
expect(await Promise.all([response.text(), clone.text()])).toEqual(["<p>hi</p>", "<p>hi</p>"]);
});

test("Bun.serve sends the file's Content-Type for a cloned Response over Bun.file().stream()", async () => {
await using server = Bun.serve({
port: 0,
fetch: req =>
new URL(req.url).pathname === "/clone" ? new Response(file().stream()).clone() : new Response(file().stream()),
});
const results: Record<string, [string | null, string]> = {};
for (const path of ["/direct", "/clone"]) {
const res = await fetch(new URL(path, server.url));
results[path] = [res.headers.get("content-type"), await res.text()];
}
expect(results).toEqual({
"/direct": ["text/html;charset=utf-8", "<p>hi</p>"],
"/clone": ["text/html;charset=utf-8", "<p>hi</p>"],
});
});

// A pipe yields its bytes once, so two Blobs over it would compete for
// them. A body over such a store is read as one stream and teed instead,
// whether it was given as a stream or as the Blob itself, and both bodies
// see the whole input.
async function cloneInChild(bodyExpr: string, args: string[] = []) {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`const r = new Response(${bodyExpr});
const c = r.clone();
const [a, b] = await Promise.all([r.text(), c.text()]);
console.log(JSON.stringify([a, b]));`,
...args,
],
env: bunEnv,
stdin: "pipe",
stdout: "pipe",
stderr: "pipe",
});
proc.stdin.write("hello ");
await proc.stdin.flush();
proc.stdin.write("world");
await proc.stdin.end();
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
return { stdout: stdout.trim(), stderr, exitCode };
}
const bothBodiesReadStdin = { stdout: `["hello world","hello world"]`, stderr: "", exitCode: 0 };

test("a body over Bun.stdin.stream() is still teed", async () => {
expect(await cloneInChild("Bun.stdin.stream()")).toEqual(bothBodiesReadStdin);
});

test("a body over Bun.stdin itself is teed, not duped", async () => {
expect(await cloneInChild("Bun.stdin")).toEqual(bothBodiesReadStdin);
});

// The same store kind reached by path: stat says it is not a regular file.
test.skipIf(isWindows)("a body over a FIFO opened by path is still teed", async () => {
const fifo = join(tempDirWithFiles("body-clone-fifo", {}), "body.fifo");
expect(Bun.spawnSync({ cmd: ["mkfifo", fifo] }).exitCode).toBe(0);
await using writer = Bun.spawn({
cmd: ["sh", "-c", `printf 'hello world' > "$1"`, "sh", fifo],
stdout: "ignore",
stderr: "inherit",
});
expect(await cloneInChild("Bun.file(process.argv.at(-1)).stream()", [fifo])).toEqual(bothBodiesReadStdin);
expect(await writer.exited).toBe(0);
});
});

test("Blob type from a consumed Response keeps the original content-type after clones with different content-types are consumed", async () => {
// The Response and its clones share one underlying body store. Consuming a clone
// with a different Content-Type must not change (or invalidate) the type of a Blob
Expand Down
Loading