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
11 changes: 11 additions & 0 deletions src/jsc/VirtualMachine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -974,6 +974,17 @@ impl VirtualMachine {
unsafe { EventLoop::enter_scope(self.event_loop) }
}

/// `event_loop().enter()` now, `.exit_without_checkpoint()` on drop, for a
/// dispatcher that drains microtasks itself: see
/// [`EventLoop::enter_scope_without_checkpoint`].
#[inline]
pub fn enter_event_loop_scope_without_checkpoint(
&self,
) -> crate::event_loop::EventLoopEnterNoCheckpointGuard {
// SAFETY: as `enter_event_loop_scope`.
unsafe { EventLoop::enter_scope_without_checkpoint(self.event_loop) }
}

/// Safe shared-reference accessor for the process-lifetime dotenv loader
/// (`vm.transpiler.env`). The loader is allocated once during VM init and
/// never freed; callers previously open-coded `unsafe { &*vm.transpiler.env }`.
Expand Down
61 changes: 61 additions & 0 deletions src/jsc/event_loop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -271,6 +271,25 @@ impl Drop for EventLoopEnterGuard {
}
}

/// RAII pairing for [`EventLoop::enter`] / [`EventLoop::exit_without_checkpoint`].
///
/// Holds the raw pointer for the same reason as [`EventLoopEnterGuard`].
/// Construct via [`EventLoop::enter_scope_without_checkpoint`].
#[must_use = "dropping immediately exits the event loop scope"]
pub struct EventLoopEnterNoCheckpointGuard {
loop_: *mut EventLoop,
}

impl Drop for EventLoopEnterNoCheckpointGuard {
#[inline]
fn drop(&mut self) {
// SAFETY: as `EventLoopEnterGuard`: `loop_` was live at
// `enter_scope_without_checkpoint` and the VM owns it for the process
// lifetime; short-lived `&mut` only.
unsafe { (*self.loop_).exit_without_checkpoint() };
}
}

impl EventLoop {
/// Before your code enters JavaScript at the top of the event loop, call
/// `loop.enter()`. If running a single callback, prefer `runCallback` instead.
Expand Down Expand Up @@ -315,6 +334,48 @@ impl EventLoop {
EventLoopEnterGuard { loop_ }
}

/// Balance an [`enter`](Self::enter) without the checkpoint [`exit`](Self::exit)
/// runs at the outermost level. See [`Self::enter_scope_without_checkpoint`].
#[inline]
pub fn exit_without_checkpoint(&mut self) {
bun_core::scoped_log!(
EventLoop,
"exit_without_checkpoint() = {}",
self.entered_event_loop_count - 1
);
self.entered_event_loop_count -= 1;
}

/// `enter()` now, [`exit_without_checkpoint`](Self::exit_without_checkpoint)
/// on drop.
///
/// For a dispatcher that runs the checkpoint itself once the callback has
/// returned, at points of its own choosing: the HTTP request paths drain
/// explicitly so that they can look at a returned promise that the drain
/// settled (`RequestContext::on_response`, the node:http dispatch), and a
/// checkpoint on exit would add an empty one per request.
///
/// What the scope is for is the count. Only while it is above zero is the
/// callback's frame safe from a checkpoint in the middle of it: a native
/// call made from inside the callback that dispatches another callback
/// through `enter()`/`exit()` (`server.upgrade()` running `open()`,
/// `ws.close()` running `close()`) is then a nested pair, not the outermost
/// one, so its exit does not run the nextTicks and promise reactions the
/// callback queued before its next statement. The dispatcher's explicit
/// drains are unconditional, so the held count does not skip them, and the
/// continuations they run are covered by it as well.
///
/// # Safety
/// As [`Self::enter_scope`].
#[inline]
pub unsafe fn enter_scope_without_checkpoint(
loop_: *mut EventLoop,
) -> EventLoopEnterNoCheckpointGuard {
// SAFETY: caller contract — `loop_` is live; short-lived `&mut` only.
unsafe { (*loop_).enter() };
EventLoopEnterNoCheckpointGuard { loop_ }
}

pub fn exit_maybe_drain_microtasks(
&mut self,
allow_drain_microtask: bool,
Expand Down
9 changes: 9 additions & 0 deletions src/runtime/server/RequestContext.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1352,6 +1352,9 @@ where
let server = this.server();
let vm = server.vm();
let global_this = server.global_this();
// Entered for the abort listeners below, and (dropped last) for the
// drains below and in the release of `_ref`.
let _entered = vm.enter_event_loop_scope_without_checkpoint();
let _ref = RequestContextRef::adopt(this.as_ctx_ptr());
// This is a task in the event loop.
// If we called into JavaScript, we must drain the microtask queue.
Expand Down Expand Up @@ -2576,6 +2579,12 @@ where
//
// - If you return a Promise, we drain the microtask queue once
// - If you return a streaming Response, we drain the microtask queue (possibly the 2nd time this task!)
//
// Like a task, the handler and these drains run with the event loop entered
// (the dispatchers hold `enter_event_loop_scope_without_checkpoint`), so a
// callback the handler dispatches synchronously through `enter()`/`exit()`
// (`server.upgrade()` -> `open()`, `ws.close()` -> `close()`) does not
// drain in the middle of the handler.
pub(crate) fn on_response(
&self,
this: &ThisServer,
Expand Down
4 changes: 4 additions & 0 deletions src/runtime/server/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -997,6 +997,7 @@ impl<const SSL: bool, const DEBUG: bool> NewServer<SSL, DEBUG> {

// SAFETY: `this` is the live server backref for this request.
let server = unsafe { &*this };
let _entered = server.vm().enter_event_loop_scope_without_checkpoint();
let global = server.global_this();
let response_value = match callback.call(global, server_js, &args) {
Ok(v) => v,
Expand Down Expand Up @@ -1143,6 +1144,7 @@ impl<const SSL: bool, const DEBUG: bool> NewServer<SSL, DEBUG> {

// SAFETY: `this` is the live server backref for this request.
let server = unsafe { &*this };
let _entered = server.vm().enter_event_loop_scope_without_checkpoint();
let on_request = server.config.on_request;
debug_assert!(!on_request.is_empty());

Expand Down Expand Up @@ -1194,6 +1196,7 @@ impl<const SSL: bool, const DEBUG: bool> NewServer<SSL, DEBUG> {

// 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();
let global = server_ref.global_this();
let server_request_list =
Self::js_route_list_get_cached(server_js).expect("routeList cached value missing");
Expand Down Expand Up @@ -1275,6 +1278,7 @@ impl<const SSL: bool, const DEBUG: bool> NewServer<SSL, DEBUG> {
core::ptr::NonNull::new(this).expect("on_node_http_request: this non-null"),
);
let vm = this_ref.vm_mut();
let _entered = this_ref.vm().enter_event_loop_scope_without_checkpoint();
req.set_yield(false);
resp.timeout(this_ref.config.idle_timeout);

Expand Down
4 changes: 4 additions & 0 deletions src/runtime/server/server_body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2947,6 +2947,7 @@ where
return;
};

let _entered = server_ref.vm().enter_event_loop_scope_without_checkpoint();
let server_request_list = Self::js_route_list_get_cached(server_js).unwrap();
let call_route = if Ctx::IS_H3 {
Bun__ServerRouteList__callRouteH3
Expand Down Expand Up @@ -3045,6 +3046,7 @@ where
// SAFETY: `self_ptr` is `self`, live for this frame. Shared — the
// handler call below re-enters JS, so no `&mut` may span it.
let server = unsafe { &*self_ptr };
let _entered = server.vm().enter_event_loop_scope_without_checkpoint();
let on_request_fn = server.config.on_request;
debug_assert!(!on_request_fn.is_empty());

Expand Down Expand Up @@ -3359,6 +3361,7 @@ where
.upgrade_context
.set(UpgradeState::Pending(NonNull::from(upgrade_ctx)))
};
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.
let global = bun_opaque::opaque_deref(server_ref.global_this);
Expand Down Expand Up @@ -3441,6 +3444,7 @@ where
resp.end_without_body(true);
return;
}
let _entered = this.vm().enter_event_loop_scope_without_checkpoint();
this.on_pending_request();
req.set_yield(false);
// SAFETY: `request_pool` is non-null while the server is alive; `claim()`
Expand Down
66 changes: 66 additions & 0 deletions test/js/bun/http/serve-http3.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import type { ServerWebSocket } from "bun";
import { describe, expect, test } from "bun:test";
import { createHash, createPrivateKey, randomBytes } from "crypto";
import { readFileSync } from "fs";
Expand Down Expand Up @@ -1449,3 +1450,68 @@ describe("Bun.serve HTTP/3 request validation", () => {
expect({ selfSigned, chained }).toEqual({ selfSigned: "closed", chained: "200 1" });
});
});

// The HTTP/3 twin of the HTTP/1 cases in websocket-server.test.ts: ws.close()
// runs close() before it returns, and a request handler that calls it must still
// run to completion before the nextTick and promise callbacks it queued. The
// socket being closed lives on a plain HTTP/1 server, since HTTP/3 carries no
// WebSockets; any handler can close it.
describe("Bun.serve HTTP/3 request handlers run to completion before the callbacks they queued", () => {
async function openHeldSocket() {
const order: string[] = [];
const opened = Promise.withResolvers<ServerWebSocket<unknown>>();
const closed = Promise.withResolvers<void>();
const wsServer = Bun.serve({
port: 0,
fetch: (req, srv) => (srv.upgrade(req) ? undefined : new Response("upgrade() failed", { status: 500 })),
websocket: {
open: ws => opened.resolve(ws),
message() {},
close() {
order.push("close()");
},
},
});
const client = new WebSocket(wsServer.url.href.replace(/^http/, "ws"));
client.onerror = () => closed.resolve();
client.onclose = () => closed.resolve();
const held = await opened.promise;
return {
order,
closed: closed.promise,
handler() {
process.nextTick(() => order.push("nextTick"));
Promise.resolve().then(() => order.push("microtask"));
held.close();
order.push("rest of handler");
return new Response("ok");
},
[Symbol.dispose]: () => wsServer.stop(true),
};
}

test("fetch() and a route handler closing an open ServerWebSocket", async () => {
using viaFetch = await openHeldSocket();
using viaRoute = await openHeldSocket();
await using server = Bun.serve({
port: 0,
tls,
http3: true,
routes: { "/route": viaRoute.handler },
fetch: viaFetch.handler,
});

const responses = {
fetch: await h3Exchange(server.port, requestHeaders("/")),
route: await h3Exchange(server.port, requestHeaders("/route")),
};
await Promise.all([viaFetch.closed, viaRoute.closed]);

const expectedOrder = ["close()", "rest of handler", "nextTick", "microtask"];
expect({ responses, fetch: viaFetch.order, route: viaRoute.order }).toEqual({
responses: { fetch: "200 ok", route: "200 ok" },
fetch: expectedOrder,
route: expectedOrder,
});
});
});
Loading
Loading