diff --git a/src/runtime/server/AnyRequestContext.rs b/src/runtime/server/AnyRequestContext.rs index 0024f14c3d14..eae501bef63a 100644 --- a/src/runtime/server/AnyRequestContext.rs +++ b/src/runtime/server/AnyRequestContext.rs @@ -191,6 +191,10 @@ impl AnyRequestContext { }) } + pub(crate) fn set_pathname(self, url: &bun_core::String) { + dispatch!(self, (), |_T, ctx| ctx.set_pathname(url)) + } + /// Wont actually set anything if `self` is `.none` pub(crate) fn set_request(self, req: *mut uws::Request) { dispatch!(self, (), |T, ctx| { diff --git a/src/runtime/server/RequestContext.rs b/src/runtime/server/RequestContext.rs index d33de89ed240..3655e97ffc74 100644 --- a/src/runtime/server/RequestContext.rs +++ b/src/runtime/server/RequestContext.rs @@ -2350,7 +2350,7 @@ where // For HTTP/3, prepareJsRequestContextFor() already eagerly // populated url+headers (the lazy getRequest() path is H1-only), - // so the guards below short-circuit and `req` is never read. + // so the copy below short-circuits and `req` is never read. if !MUX { // `Req` is erased to `c_void`; for !MUX the concrete // type is `uws::Request`, so the cast is nominal. @@ -2359,29 +2359,21 @@ where .set_request(req.cast::()); } - if request_object.ensure_url().is_err() { - request_object.url.set(BunString::EMPTY); - } + // `req` dies with this stack frame, so what is still read lazily from it is copied out now. + request_object.detach_uws_request_head(); + } - // we have to clone the request headers here since they will soon belong to a different request - if !request_object.has_fetch_headers() { - if !MUX { - // `HeadersRef::create_from_uws` adopts the freshly-allocated +1 ref. - request_object.set_fetch_headers(Some(response::HeadersRef::create_from_uws(req))); - } + /// The path the development-mode error page prints once the uWS request is detached. + pub(crate) fn set_pathname(&self, url: &BunString) { + if DEBUG_MODE { + self.pathname.set(url.clone()); } - - // This object dies after the stack frame is popped - // so we have to clear it in here too - request_object.request_context.detach_request(); } pub(crate) fn to_async(&self, req: *mut Req, request_object: &mut Request) { ctx_log!("toAsync"); self.to_async_without_abort_handler(req, request_object); - if DEBUG_MODE { - self.pathname.set(request_object.url.get().clone()); - } + self.set_pathname(request_object.url.get()); self.set_abort_handler(); } diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index b96201e5dce9..2a944a32fb8c 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -419,6 +419,21 @@ impl Drop for DetachRequestOnDrop { } } +/// The url and headers that `request_object` still reads lazily from the uWS receive buffer when its handler runs. +#[inline] +fn borrow_request_head( + request_object: *mut crate::webcore::Request, +) -> bun_uws_sys::loop_::RecvBufferBorrow { + unsafe fn release(owner: *mut c_void) { + // SAFETY: whoever registered the borrow keeps `owner`, the heap `webcore::Request`, live until it is unregistered. + let request = unsafe { &*owner.cast::() }; + request.detach_uws_request_head(); + request.request_context.set_pathname(request.url.get()); + } + + bun_uws_sys::loop_::RecvBufferBorrow::new(release, request_object.cast::()) +} + impl NewServer { pub(crate) const HAS_H3: bool = SSL; @@ -973,6 +988,11 @@ impl NewServer { args.push(prepared.js_request); args.extend_from_slice(&extra_args); + let mut head_borrow = borrow_request_head(prepared.request_object); + // SAFETY: `request_object` stays allocated for this whole frame (its JS wrapper, `prepared.js_request`, is on this + // stack), and the guard is a local of the frame. + let _head_borrow = unsafe { head_borrow.register() }; + // SAFETY: `this` is the live server backref for this request. let server = unsafe { &*this }; let _entered = server.vm().enter_event_loop_scope_without_checkpoint(); @@ -1120,6 +1140,11 @@ impl NewServer { return; }; + let mut head_borrow = borrow_request_head(prepared.request_object); + // SAFETY: `request_object` stays allocated for this whole frame (its JS wrapper, `prepared.js_request`, is on this + // stack), and the guard is a local of the frame. + let _head_borrow = unsafe { head_borrow.register() }; + // SAFETY: `this` is the live server backref for this request. let server = unsafe { &*this }; let _entered = server.vm().enter_event_loop_scope_without_checkpoint(); @@ -1172,6 +1197,11 @@ impl NewServer { return; }; + let mut head_borrow = borrow_request_head(prepared.request_object); + // SAFETY: `request_object` stays allocated for this whole frame (its JS wrapper, `prepared.js_request`, is on this + // stack), and the guard is a local of the frame. + let _head_borrow = unsafe { head_borrow.register() }; + // SAFETY: `server` is the live backref stored in `user_route`. let server_ref = unsafe { &*server }; let _entered = server_ref.vm().enter_event_loop_scope_without_checkpoint(); diff --git a/src/runtime/server/server_body.rs b/src/runtime/server/server_body.rs index 5a3f6ff0a1e0..79ef24dc1ec9 100644 --- a/src/runtime/server/server_body.rs +++ b/src/runtime/server/server_body.rs @@ -3183,6 +3183,10 @@ where .upgrade_context .set(UpgradeState::Pending(NonNull::from(upgrade_ctx))) }; + let mut head_borrow = super::borrow_request_head(prepared.request_object); + // SAFETY: `request_object` stays allocated for this whole frame (its JS wrapper, `prepared.js_request`, is on this + // stack), and the guard is a local of the frame. + let _head_borrow = unsafe { head_borrow.register() }; let _entered = server_ref.vm().enter_event_loop_scope_without_checkpoint(); let server_request_list = Self::js_route_list_get_cached(server_js).unwrap(); // S008: `JSGlobalObject` is an `opaque_ffi!` ZST — safe deref. @@ -3323,6 +3327,11 @@ where let args = [unsafe { (*request_object_ptr).to_js(&global) }, server_js]; args[0].ensure_still_alive(); + let mut head_borrow = super::borrow_request_head(request_object_ptr); + // SAFETY: `request_object_ptr` stays allocated for this whole frame (its JS wrapper, `args[0]`, is on this stack), + // and the guard is a local of the frame. + let _head_borrow = unsafe { head_borrow.register() }; + let response_value = match this.config.on_request.call(&global, server_js, &args) { Ok(v) => v, Err(err) => global.take_exception(err), diff --git a/src/runtime/webcore/Request.rs b/src/runtime/webcore/Request.rs index 8ceaadd119dd..151fc9cafa2d 100644 --- a/src/runtime/webcore/Request.rs +++ b/src/runtime/webcore/Request.rs @@ -236,6 +236,25 @@ impl Request { self.headers.set(headers); } + /// Copies the url and headers still read lazily from the `uWS::HttpRequest` into `self`, then forgets that request. + pub(crate) fn detach_uws_request_head(&self) { + let Some(req) = self.request_context.get_request() else { + return; + }; + + if self.ensure_url().is_err() { + self.url.set(BunString::EMPTY); + } + + if !self.has_fetch_headers() { + self.set_fetch_headers(Some(HeadersRef::create_from_uws( + req.cast::(), + ))); + } + + self.request_context.detach_request(); + } + /// Returns the headers of the request. If the headers are not already cached, it will create a new FetchHeaders object. /// If the headers are empty, it will look at request_context to get the headers. /// If the headers are empty and request_context is null, it will create an empty FetchHeaders object. diff --git a/src/uws_sys/Loop.rs b/src/uws_sys/Loop.rs index 843061721f4b..291024116cb6 100644 --- a/src/uws_sys/Loop.rs +++ b/src/uws_sys/Loop.rs @@ -1,3 +1,4 @@ +use core::cell::Cell; use core::ffi::{c_int, c_uint, c_void}; use core::ptr::NonNull; @@ -13,6 +14,73 @@ bun_core::declare_scope!(Loop, visible); /// it reaches the idle sweep, so passing this costs nothing on the paths that never park. pub const NOW_NS_UNKNOWN: u64 = 0; +// ─────────────────── borrows of the receive buffer ─────────────────── + +/// Receive-buffer bytes that a callback still reads after it called into JS. `release` copies them out. +pub struct RecvBufferBorrow { + next: Cell<*const RecvBufferBorrow>, + release: unsafe fn(*mut c_void), + owner: *mut c_void, +} + +thread_local! { + /// The borrows registered on this thread, innermost first. + static RECV_BUFFER_BORROWS: Cell<*const RecvBufferBorrow> = + const { Cell::new(core::ptr::null()) }; +} + +impl RecvBufferBorrow { + pub fn new(release: unsafe fn(*mut c_void), owner: *mut c_void) -> Self { + Self { + next: Cell::new(core::ptr::null()), + release, + owner, + } + } + + /// Safety: `release(owner)` stays sound to call, repeatedly, until the guard drops. Guards drop in reverse order. + pub unsafe fn register(&mut self) -> RecvBufferBorrowGuard<'_> { + self.next.set(RECV_BUFFER_BORROWS.get()); + RECV_BUFFER_BORROWS.set(core::ptr::from_ref(self)); + RecvBufferBorrowGuard(self) + } +} + +/// Unregisters its `RecvBufferBorrow` on drop. +pub struct RecvBufferBorrowGuard<'a>(&'a RecvBufferBorrow); + +impl Drop for RecvBufferBorrowGuard<'_> { + fn drop(&mut self) { + debug_assert!( + core::ptr::eq(RECV_BUFFER_BORROWS.get(), core::ptr::from_ref(self.0)), + "RecvBufferBorrowGuard dropped out of order" + ); + RECV_BUFFER_BORROWS.set(self.0.next.get()); + } +} + +/// Runs before every run of this thread's loop, which is what reads a socket into the buffer again. +#[inline] +fn release_recv_buffer_borrows() { + if RECV_BUFFER_BORROWS.get().is_null() { + return; + } + release_registered_borrows(); +} + +#[cold] +fn release_registered_borrows() { + let mut node = RECV_BUFFER_BORROWS.get(); + while !node.is_null() { + // SAFETY: `register`'s caller keeps each listed node live and its `release(owner)` sound to call. + unsafe { + let borrow = &*node; + node = borrow.next.get(); + (borrow.release)(borrow.owner); + } + } +} + // ───────────────────────────── PosixLoop ───────────────────────────── // Mirrors C `struct us_loop_t` (packages/bun-usockets/src/internal/eventing/ @@ -245,11 +313,13 @@ impl PosixLoop { } pub fn tick(&mut self) { + release_recv_buffer_borrows(); // SAFETY: self is a valid loop pointer unsafe { c::us_loop_run_bun_tick(self, core::ptr::null(), NOW_NS_UNKNOWN) }; } pub fn tick_without_idle(&mut self) { + release_recv_buffer_borrows(); let timespec = Timespec { sec: 0, nsec: 0 }; // SAFETY: self is a valid loop pointer; ×pec lives for the call unsafe { c::us_loop_run_bun_tick(self, &raw const timespec, NOW_NS_UNKNOWN) }; @@ -259,6 +329,7 @@ impl PosixLoop { /// `timer::All::get_timeout`), reused by the JS park hook's idle-sweep rate limit rather /// than read again. `NOW_NS_UNKNOWN` if the caller has none to share. pub fn tick_with_timeout(&mut self, timespec: Option<&Timespec>, now_ns: u64) { + release_recv_buffer_borrows(); // SAFETY: self is a valid loop pointer unsafe { c::us_loop_run_bun_tick( @@ -386,11 +457,13 @@ impl WindowsLoop { /// Windows the park hook is driven from `us_loop_run` (libuv.c), which reads libuv's /// already-refreshed clock via `uv_now` rather than taking one of its own. pub fn tick_with_timeout(&mut self, _: Option<&Timespec>, _now_ns: u64) { + release_recv_buffer_borrows(); // SAFETY: self is a valid loop pointer unsafe { c::us_loop_run(self) }; } pub fn tick_without_idle(&mut self) { + release_recv_buffer_borrows(); // SAFETY: self is a valid loop pointer unsafe { c::us_loop_pump(self) }; } @@ -418,6 +491,7 @@ impl WindowsLoop { } pub fn run(&mut self) { + release_recv_buffer_borrows(); // SAFETY: self is a valid loop pointer unsafe { c::us_loop_run(self) }; } diff --git a/test/js/bun/http/bun-serve-nested-event-loop.test.ts b/test/js/bun/http/bun-serve-nested-event-loop.test.ts new file mode 100644 index 000000000000..1624d9b8071d --- /dev/null +++ b/test/js/bun/http/bun-serve-nested-event-loop.test.ts @@ -0,0 +1,232 @@ +import { expect, test } from "bun:test"; +import { bunEnv, bunExe, tempDir } from "harness"; +import { connect } from "node:net"; + +// `req.url` and `req.headers` are read lazily from the `uWS::HttpRequest`, +// whose views point into the per-loop receive buffer. A `fetch` handler that +// runs the event loop inside itself lets the next socket read overwrite those +// bytes in place, so the handler must not end up with another request's head. +// +// The nested run here is `Bun.build` with a plugin whose `setup()` returns a +// pending promise: the call waits for that promise with the event loop. + +const FIRST = "1111111111111111"; +const SECOND = "2222222222222222"; + +function head(token: string) { + return ( + `GET /u-${token} HTTP/1.1\r\n` + + `Host: h-${token}.example\r\n` + + `Authorization: Bearer au-${token}\r\n` + + `Cookie: c=ck-${token}\r\n` + + `\r\n` + ); +} + +async function connectTo(port: number) { + const socket = connect(port, "127.0.0.1"); + await new Promise((resolve, reject) => { + socket.once("connect", resolve); + socket.once("error", reject); + }); + socket.on("error", () => {}); + socket.on("data", () => {}); + return socket; +} + +type Head = { url: string; authorization: string | null; cookie: string | null }; + +// Runs the event loop from inside the handler. The build does not return until +// the plugin's `setup()` promise settles, and the second request settles it, so +// the tests need no timing. +function eventLoopRunner(entrypoint: string) { + let resolveSetup: (() => void) | undefined; + return { + end: () => resolveSetup!(), + run: () => + Bun.build({ + entrypoints: [entrypoint], + plugins: [ + { + name: "pending-setup", + setup: () => + new Promise(resolve => { + resolveSetup = resolve; + }), + }, + ], + }), + }; +} + +async function handlerReadsItsOwnHead(readAfterAwait: boolean) { + using dir = tempDir("serve-nested-event-loop", { "entry.js": "export default 1;\n" }); + + const recorded = Promise.withResolvers(); + const nested = eventLoopRunner(`${dir}/entry.js`); + let second: ReturnType | undefined; + let requests = 0; + + await using server = Bun.serve({ + port: 0, + hostname: "127.0.0.1", + async fetch(req) { + if (++requests > 1) { + // The second request was parsed inside the first handler's nested + // event loop run, so its head sits in the receive buffer now. + nested.end(); + return new Response("second"); + } + + second!.write(head(SECOND)); + + if (readAfterAwait) { + // The handler reads nothing before it suspends, so the head is the + // copy the server makes when it hands the request to the async path. + await nested.run(); + } else { + nested.run().then( + () => {}, + () => {}, + ); + } + + recorded.resolve({ + url: req.url, + authorization: req.headers.get("authorization"), + cookie: req.headers.get("cookie"), + }); + return new Response("first"); + }, + }); + + const first = await connectTo(server.port); + second = await connectTo(server.port); + try { + first.write(head(FIRST)); + expect(await recorded.promise).toEqual({ + url: `http://h-${FIRST}.example/u-${FIRST}`, + authorization: `Bearer au-${FIRST}`, + cookie: `c=ck-${FIRST}`, + }); + } finally { + first.destroy(); + second.destroy(); + } +} + +test("a handler that runs the event loop reads its own url and headers", async () => { + await handlerReadsItsOwnHead(false); +}); + +test("a handler that runs the event loop and then suspends reads its own url and headers", async () => { + await handlerReadsItsOwnHead(true); +}); + +test("server.upgrade after the handler runs the event loop still reads its own handshake", async () => { + using dir = tempDir("serve-nested-event-loop-ws", { "entry.js": "export default 1;\n" }); + + const nested = eventLoopRunner(`${dir}/entry.js`); + let second: ReturnType | undefined; + let requests = 0; + + await using server = Bun.serve({ + port: 0, + hostname: "127.0.0.1", + fetch(req, server) { + if (++requests > 1) { + nested.end(); + return new Response("second"); + } + + second!.write(head(SECOND)); + nested.run().then( + () => {}, + () => {}, + ); + + // `Sec-WebSocket-Key` and `Upgrade` come from the same head. Reading + // the second request's values here fails the handshake, and the client + // that did send a handshake gets an error instead of a socket. + if (server.upgrade(req)) return undefined; + return new Response("upgrade failed", { status: 400 }); + }, + websocket: { + open(ws) { + ws.send("open"); + }, + message() {}, + }, + }); + + second = await connectTo(server.port); + const opened = Promise.withResolvers(); + const ws = new WebSocket(`ws://127.0.0.1:${server.port}/`); + ws.onmessage = event => opened.resolve(String(event.data)); + ws.onerror = () => opened.reject(new Error("the upgrade failed")); + ws.onclose = () => opened.reject(new Error("the socket closed before it opened")); + try { + expect(await opened.promise).toBe("open"); + } finally { + ws.close(); + second.destroy(); + } +}); + +test("the development error page of a handler that runs the event loop names its own request", async () => { + using dir = tempDir("serve-nested-event-loop-dev", { "entry.js": "export default 1;\n" }); + + // A subprocess, because the development error log goes to stderr. + const script = ` + import { connect } from "node:net"; + const head = token => + "GET /u-" + token + " HTTP/1.1\\r\\nHost: h-" + token + ".example\\r\\nConnection: close\\r\\n\\r\\n"; + let resolveSetup, second, requests = 0; + const server = Bun.serve({ + port: 0, + hostname: "127.0.0.1", + development: true, + fetch() { + if (++requests > 1) { + resolveSetup(); + return new Response("second"); + } + second.write(head("${SECOND}")); + Bun.build({ + entrypoints: [process.env.ENTRY], + plugins: [{ name: "pending-setup", setup: () => new Promise(resolve => (resolveSetup = resolve)) }], + }).then(() => {}, () => {}); + throw new Error("boom"); + }, + }); + const open = () => + new Promise((resolve, reject) => { + const socket = connect(server.port, "127.0.0.1", () => resolve(socket)); + socket.on("error", reject); + }); + const first = await open(); + second = await open(); + second.on("data", () => {}); + let page = ""; + first.on("data", chunk => (page += chunk.toString("latin1"))); + first.on("close", () => { + console.log(page.match(/GET - \\S+ failed/)?.[0]); + server.stop(true); + process.exit(0); + }); + first.write(head("${FIRST}")); + `; + + await using proc = Bun.spawn({ + cmd: [bunExe(), "-e", script], + env: { ...bunEnv, ENTRY: `${dir}/entry.js` }, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + + const failed = `GET - http://h-${FIRST}.example/u-${FIRST} failed`; + expect(stdout.trim()).toBe(failed); + expect(stderr).toContain(failed); + expect(exitCode).toBe(0); +});