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
45 changes: 35 additions & 10 deletions src/io/PipeReader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -402,6 +402,18 @@ impl PosixBufferedReader {
/// embedding `self` (the shell `PipeReader` does exactly that), so the
/// caller must not touch `self` again after a `false` return.
pub fn register_poll(&mut self) -> bool {
match self.try_register_poll() {
sys::Result::Ok(()) => true,
sys::Result::Err(err) => {
self.vtable.on_reader_error(err);
false
}
}
}

/// Like [`register_poll`] but returns the registration error instead of
/// dispatching it through `on_reader_error`, so the caller can recover.
fn try_register_poll(&mut self) -> sys::Result<()> {
// Hoist vtable-derived scalars and
// normalize self.handle to Poll before taking the single &mut borrow,
// so no raw-pointer escape is needed.
Expand All @@ -411,7 +423,7 @@ impl PosixBufferedReader {

if let PollOrFd::Fd(fd) = self.handle {
if !self.flags.contains(PosixFlags::POLLABLE) {
return true;
return sys::Result::Ok(());
}
self.handle = PollOrFd::Poll(FilePollRef::init(
ev,
Expand All @@ -420,21 +432,15 @@ impl PosixBufferedReader {
));
}
let Some(poll) = self.handle.get_poll_mut() else {
return true;
return sys::Result::Ok(());
};
poll.set_owner(Owner::new(PollTag::BufferedReader, owner_ptr.cast()));

if !poll.has_flag(FilePollFlag::WasEverRegistered) {
poll.enable_keeping_process_alive(ev);
}

match poll.register_with_fd(lp.cast(), FilePollKind::Readable, poll.fd()) {
sys::Result::Err(err) => {
self.vtable.on_reader_error(err);
false
}
sys::Result::Ok(()) => true,
}
poll.register_with_fd(lp.cast(), FilePollKind::Readable, poll.fd())
}

pub fn start(&mut self, fd: Fd, is_pollable: bool) -> sys::Result<()> {
Expand All @@ -449,7 +455,26 @@ impl PosixBufferedReader {
if self.get_fd() != fd {
self.handle = PollOrFd::Fd(fd);
}
self.register_poll();
if let sys::Result::Err(err) = self.try_register_poll() {
// epoll_ctl/kevent reject fds whose driver has no poll support
// (e.g. /dev/null, /dev/zero): EPERM from epoll on Linux, EINVAL
// from kqueue on macOS. Such fds are always-readable, so fall
// back to the non-pollable path instead of tearing the reader
// down. Mirrors IOWriter::__start.
let fd_not_pollable = matches!(err.get_errno(), sys::E::EINVAL)
|| (cfg!(any(target_os = "linux", target_os = "android"))
&& err.get_errno() == sys::E::EPERM);
if fd_not_pollable {
self.flags
.remove(PosixFlags::POLLABLE | PosixFlags::NONBLOCKING);
if matches!(self.handle, PollOrFd::Poll(_)) {
self.handle.close_impl(None, None::<fn(*mut c_void)>, false);
}
self.handle = PollOrFd::Fd(fd);
return sys::Result::Ok(());
}
self.vtable.on_reader_error(err);
}
Comment thread
robobun marked this conversation as resolved.

sys::Result::Ok(())
}
Expand Down
7 changes: 5 additions & 2 deletions src/runtime/server/FileResponseStream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -496,8 +496,11 @@ impl FileResponseStream {
if !self.state.contains(State::RESPONSE_DONE) {
self.state.insert(State::RESPONSE_DONE);
self.detach_resp();
self.resp
.end_without_body(self.resp.should_close_connection());
// `end` (not `end_without_body`): when no Content-Length has been
// written yet (non-regular files) uWS supplies the framing here,
// so an immediately-EOF fd like /dev/null still produces a valid
// empty response instead of leaving the client waiting.
self.resp.end(b"", self.resp.should_close_connection());
(self.on_complete)(self.ctx, self.resp);
}

Expand Down
9 changes: 8 additions & 1 deletion src/runtime/server/RequestContext.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1835,7 +1835,10 @@ where
});
}

self.flags.set_needs_content_length(true);
// Non-regular files (FIFOs, character devices, sockets) have no
// meaningful stat size, so the body length is unknown up front and
// must be framed with chunked encoding.
self.flags.set_needs_content_length(is_regular);
let blob_offset = match &self.blob {
AnyBlob::Blob(b) => b.offset.get(),
_ => unreachable!(),
Expand Down Expand Up @@ -1972,6 +1975,10 @@ where
offset: self.sendfile.offset as u64,
length: if is_regular {
Some(self.sendfile.remain as u64)
} else if original_size != crate::webcore::blob::MAX_SIZE {
// An explicit .slice() on a non-regular file caps the body at
// that many bytes; without it, read until EOF.
Some(original_size as u64)
} else {
None
},
Expand Down
85 changes: 85 additions & 0 deletions test/js/bun/http/bun-serve-file.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import { afterAll, beforeAll, describe, expect, it, mock, test } from "bun:test"
import { bunEnv, bunExe, isASAN, isWindows, rmScope, tempDir, tempDirWithFiles } from "harness";
import { mkfifo } from "mkfifo";
import { unlinkSync } from "node:fs";
import * as net from "node:net";
import { join } from "node:path";

const LARGE_SIZE = 1024 * 1024 * 8;
Expand Down Expand Up @@ -1138,3 +1139,87 @@ console.log("OK");
},
60_000,
);

// On Linux, epoll_ctl(EPOLL_CTL_ADD) on /dev/null-class character devices
// returns EPERM (their file_operations have no .poll). The file response
// stream used to treat that as a fatal reader error and force-close the
// connection before any status line or header byte was written, without ever
// invoking the error() callback.
describe.skipIf(isWindows)("serving a character-device Bun.file from fetch()", () => {
// Use a raw TCP client so we can observe the zero-byte close that fetch()
// would otherwise report as a generic connection error.
async function rawGet(port: number, path: string): Promise<Buffer> {
const { promise, resolve } = Promise.withResolvers<Buffer>();
let acc = Buffer.alloc(0);
const s = net.connect(port, "127.0.0.1", () => {
s.write(`GET ${path} HTTP/1.1\r\nHost: x\r\nConnection: close\r\n\r\n`);
});
s.on("data", d => {
acc = Buffer.concat([acc, d]);
});
// The pre-fix force_close sends RST, so swallow ECONNRESET and report
// whatever bytes arrived; the assertion on the status line catches the
// zero-byte case.
s.on("error", () => {});
s.on("close", () => resolve(acc));
await promise;
return acc;
}

test("/dev/null serves an empty 200 response", async () => {
let errorArg: unknown = undefined;
await using server = Bun.serve({
port: 0,
hostname: "127.0.0.1",
fetch() {
return new Response(Bun.file("/dev/null"));
},
error(e) {
errorArg = e;
return new Response("ERR", { status: 500 });
},
});

const raw = await rawGet(server.port, "/");
const head = raw.toString("latin1").split("\r\n\r\n")[0];
expect(head.split("\r\n")[0]).toBe("HTTP/1.1 200 OK");
expect(errorArg).toBeUndefined();

const res = await fetch(server.url);
expect({
status: res.status,
body: await res.text(),
error: errorArg,
}).toEqual({ status: 200, body: "", error: undefined });
});

test("/dev/zero with .slice() serves the sliced length", async () => {
const len = 4096;
let errorArg: unknown = undefined;
await using server = Bun.serve({
port: 0,
hostname: "127.0.0.1",
fetch() {
return new Response(Bun.file("/dev/zero").slice(0, len));
},
error(e) {
errorArg = e;
return new Response("ERR", { status: 500 });
},
});

const raw = await rawGet(server.port, "/");
// The pre-fix behavior was a zero-byte close; after the fix we must at
// least receive a status line.
expect(raw.toString("latin1", 0, 15)).toBe("HTTP/1.1 200 OK");

const res = await fetch(server.url);
const body = new Uint8Array(await res.arrayBuffer());
expect({
status: res.status,
bodyLength: body.length,
allZero: body.every(b => b === 0),
error: errorArg,
}).toEqual({ status: 200, bodyLength: len, allZero: true, error: undefined });
});
});
Loading