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
67 changes: 14 additions & 53 deletions src/runtime/server/FileResponseStream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,22 +60,22 @@ pub enum Mode {
}

struct Sendfile {
#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))]
#[cfg(any(target_os = "linux", target_os = "android"))]
socket_fd: Fd,
remain: u64,
offset: u64,
#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))]
#[cfg(any(target_os = "linux", target_os = "android"))]
has_set_on_writable: bool,
}

impl Default for Sendfile {
fn default() -> Self {
Self {
#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))]
#[cfg(any(target_os = "linux", target_os = "android"))]
socket_fd: Fd::INVALID,
remain: 0,
offset: 0,
#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))]
#[cfg(any(target_os = "linux", target_os = "android"))]
has_set_on_writable: false,
}
}
Expand Down Expand Up @@ -165,11 +165,11 @@ impl FileResponseStream {

if use_sendfile {
this.sendfile = Sendfile {
#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))]
#[cfg(any(target_os = "linux", target_os = "android"))]
socket_fd: opts.resp.get_native_handle(),
offset: opts.offset,
remain: opts.length.expect("can_sendfile gates None"),
#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))]
#[cfg(any(target_os = "linux", target_os = "android"))]
has_set_on_writable: false,
};
this.resp.prepare_for_sendfile();
Expand Down Expand Up @@ -402,55 +402,13 @@ impl FileResponseStream {
}
}
}
#[cfg(target_os = "macos")]
loop {
let mut sbytes: libc::off_t =
i64::try_from(self.sendfile.remain.min(i32::MAX as u64)).expect("int cast");
// SAFETY: both fds are valid open file descriptors owned by `self`;
// `sbytes` is a stack local; hdtr is null per spec.
let errno = sys::get_errno(unsafe {
sys::c::sendfile(
self.fd.native(),
self.sendfile.socket_fd.native(),
i64::try_from(self.sendfile.offset).expect("int cast"),
&raw mut sbytes,
core::ptr::null_mut(),
0,
)
});
let sent: u64 = u64::try_from(sbytes).expect("int cast");
self.sendfile.offset += sent;
self.sendfile.remain = self.sendfile.remain.saturating_sub(sent);

match errno {
sys::E::SUCCESS => {
if self.sendfile.remain == 0 || sent == 0 {
self.end_sendfile();
return false;
}
return self.arm_sendfile_writable();
}
sys::E::EINTR => continue,
sys::E::EAGAIN => return self.arm_sendfile_writable(),
sys::E::EPIPE | sys::E::ENOTCONN => {
self.end_sendfile();
return false;
}
_ => {
self.fail_with(
sys::Error::from_code(errno, sys::Tag::sendfile).with_fd(self.fd),
);
return false;
}
}
}
#[cfg(not(any(target_os = "linux", target_os = "android", target_os = "macos")))]
#[cfg(not(any(target_os = "linux", target_os = "android")))]
{
unreachable!() // can_sendfile gates this
}
Comment thread
robobun marked this conversation as resolved.
}

#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))]
#[cfg(any(target_os = "linux", target_os = "android"))]
fn arm_sendfile_writable(&mut self) -> bool {
bun_output::scoped_log!(FileResponseStream, "armSendfileWritable");
if !self.sendfile.has_set_on_writable {
Expand All @@ -467,7 +425,7 @@ impl FileResponseStream {
true
}

#[cfg(any(target_os = "linux", target_os = "android", target_os = "macos"))]
#[cfg(any(target_os = "linux", target_os = "android"))]
fn end_sendfile(&mut self) {
bun_output::scoped_log!(FileResponseStream, "endSendfile");
if self.state.contains(State::RESPONSE_DONE) {
Expand Down Expand Up @@ -594,12 +552,15 @@ impl Drop for FileResponseStream {
}

fn can_sendfile(resp: AnyResponse, file_type: FileType, length: Option<u64>) -> bool {
#[cfg(windows)]
// Matches the cfg on `on_sendfile`. macOS is excluded: XNU's sendfile can
// sleep uninterruptibly under mbuf pressure, leaving the process unkillable;
// the BufferedReader path stays non-blocking.
#[cfg(not(any(target_os = "linux", target_os = "android")))]
{
let _ = (resp, file_type, length);
return false;
}
#[cfg(not(windows))]
#[cfg(any(target_os = "linux", target_os = "android"))]
{
// sendfile() needs a real socket fd; SSL writes go through BIO and H3
// through lsquic stream frames — neither has one.
Expand Down
41 changes: 19 additions & 22 deletions test/js/bun/http/serve.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2139,39 +2139,36 @@ it.concurrent("should work with dispose keyword", async () => {
expect(fetch(url)).rejects.toThrow();
});

// prettier-ignore
// Fixture serves a >1 MB file (the sendfile threshold). Each iteration reads
// one chunk from several streams so the server is provably mid-send when
// killed; on macOS the old sendfile(2) path could hang uninterruptibly here.
it("should be able to stop in the middle of a file response", async () => {
async function doRequest(url: string) {
try {
const response = await fetch(url, { signal: AbortSignal.timeout(10) });
const read = (response.body as ReadableStream<any>).getReader();
while (true) {
const { value, done } = await read.read();
if (done) break;
}
expect(response.status).toBe(200);
} catch {}
}
const fixture = join(import.meta.dir, "server-bigfile-send.fixture.js");
for (let i = 0; i < 3; i++) {
const process = Bun.spawn([bunExe(), fixture], {
await using proc = Bun.spawn({
cmd: [bunExe(), fixture],
env: bunEnv,
stderr: "inherit",
stdout: "pipe",
stdin: "ignore",
});
const { value } = await process.stdout.getReader().read();
const { value } = await proc.stdout.getReader().read();
const url = new TextDecoder().decode(value).trim();
Comment thread
coderabbitai[bot] marked this conversation as resolved.
const requests = [];
for (let j = 0; j < 5_000; j++) {
requests.push(doRequest(url));
// Deliberately small so macOS CI runners never approach mbuf exhaustion.
const readers: ReadableStreamDefaultReader[] = [];
for (let j = 0; j < 16; j++) {
const res = await fetch(url);
expect(res.status).toBe(200);
readers.push((res.body as ReadableStream).getReader());
}
// only await for 1k requests (and kill the process)
await Promise.all(requests.slice(0, 1_000));
expect(process.exitCode || 0).toBe(0);
process.kill();
await Promise.all(readers.map(r => r.read()));
expect(proc.exitCode).toBe(null);
proc.kill();
await proc.exited;
expect(proc.signalCode).toBe("SIGTERM");
for (const r of readers) await r.cancel().catch(() => {});
}
}, 60_000);
});

it("should be able to abrupt stop the server", async () => {
for (let i = 0; i < 10; i++) {
Expand Down
Loading