Skip to content
Closed
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
31 changes: 22 additions & 9 deletions src/runtime/server/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1522,9 +1522,11 @@ impl<const SSL: bool, const DEBUG: bool> NewServer<SSL, DEBUG> {
if Self::HAS_H3 && self.h3_app.is_some() {
self.unref();
self.notify_inspector_server_stopped();
if abrupt {
self.flags.insert(ServerFlags::TERMINATED);
}
}
// A previous graceful stop already took the listener. An abrupt stop
// still needs to close the app so in-flight connections are torn down.
if abrupt {
self.terminate_app();
}
return;
};
Expand Down Expand Up @@ -1552,13 +1554,24 @@ impl<const SSL: bool, const DEBUG: bool> NewServer<SSL, DEBUG> {
if !abrupt {
// S012: `app::ListenSocket<SSL>` is a ZST opaque — safe deref.
bun_opaque::opaque_deref_mut(listener).close();
} else if !self.flags.contains(ServerFlags::TERMINATED) {
if let Some(ws) = self.config.websocket.as_mut() {
ws.handler.app = None;
}
self.flags.insert(ServerFlags::TERMINATED);
} else {
self.terminate_app();
}
}

/// Force-close every connection on the uws app and mark the server
/// terminated. Guarded by `TERMINATED` so repeated abrupt stops are no-ops.
fn terminate_app(&mut self) {
if self.flags.contains(ServerFlags::TERMINATED) {
return;
}
if let Some(ws) = self.config.websocket.as_mut() {
ws.handler.app = None;
}
self.flags.insert(ServerFlags::TERMINATED);
if let Some(app) = self.app {
// S012: `NewApp<SSL>` is a ZST opaque — safe `*mut → &mut` deref.
bun_opaque::opaque_deref_mut(self.app.unwrap()).close();
bun_opaque::opaque_deref_mut(app).close();
}
}

Expand Down
18 changes: 9 additions & 9 deletions src/runtime/server/server_body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2502,23 +2502,23 @@ where
pub fn stop_from_js(&mut self, abruptly: Option<JSValue>) -> JSValue {
let rc = self.get_all_closed_promise(&self.global());

if self.has_listener() {
let abrupt = 'brk: {
if let Some(val) = abruptly {
if val.is_boolean() && val.to_boolean() {
break 'brk true;
}
let abrupt = 'brk: {
if let Some(val) = abruptly {
if val.is_boolean() && val.to_boolean() {
break 'brk true;
}
false
};
}
false
};
if self.has_listener() || (abrupt && !self.flags.contains(ServerFlags::TERMINATED)) {
self.stop(abrupt);
}

rc
}

pub fn dispose_from_js(&mut self) -> JSValue {
if self.has_listener() {
if self.has_listener() || !self.flags.contains(ServerFlags::TERMINATED) {
self.stop(true);
}
JSValue::UNDEFINED
Expand Down
86 changes: 86 additions & 0 deletions test/js/bun/http/serve.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2196,6 +2196,92 @@ it("should be able to abrupt stop the server", async () => {
}
});

describe("server.stop(true) after a prior graceful stop", () => {
async function setup(stopper: (server: Server) => void) {
const enc = new TextEncoder();
const firstChunk = Promise.withResolvers<void>();
const aborted = mock(() => {});
const cancelled = mock(() => {});
const server = Bun.serve({
port: 0,
idleTimeout: 0,
fetch(req) {
req.signal.addEventListener("abort", aborted);
return new Response(
new ReadableStream({
async pull(controller) {
controller.enqueue(enc.encode("data: x\n\n"));
await Bun.sleep(20);
},
cancel: cancelled,
}),
{ headers: { "Content-Type": "text/event-stream" } },
);
},
});
const closed = Promise.withResolvers<void>();
const sock = net.connect(server.port, "127.0.0.1", () => {
sock.write(`GET /sse HTTP/1.1\r\nHost: x\r\n\r\n`);
});
sock.once("data", () => firstChunk.resolve());
sock.on("error", () => {});
sock.on("close", () => closed.resolve());
try {
await firstChunk.promise;

server.stop(false);
expect(server.pendingRequests).toBe(1);

stopper(server);
await closed.promise;

expect({
pendingRequests: server.pendingRequests,
cancelled: cancelled.mock.calls.length,
aborted: aborted.mock.calls.length,
}).toEqual({ pendingRequests: 0, cancelled: 1, aborted: 1 });
} finally {
sock.destroy();
server.stop(true);
}
}

it("force-closes in-flight connections", async () => {
await setup(server => server.stop(true));
});

it("force-closes in-flight connections via [Symbol.dispose]", async () => {
await setup(server => server[Symbol.dispose]());
});

it("resolves the stop() promise once connections are gone", async () => {
const server = Bun.serve({
port: 0,
idleTimeout: 0,
fetch() {
return new Response(new ReadableStream({ pull: () => Bun.sleep(1000) }));
},
});
const gotData = Promise.withResolvers<void>();
const sock = net.connect(server.port, "127.0.0.1", () => {
sock.write(`GET / HTTP/1.1\r\nHost: x\r\n\r\n`);
});
sock.once("data", () => gotData.resolve());
sock.on("error", () => {});
try {
await gotData.promise;
const gracefulPromise = server.stop(false);
expect(server.pendingRequests).toBe(1);
const forcePromise = server.stop(true);
await Promise.all([gracefulPromise, forcePromise]);
expect(server.pendingRequests).toBe(0);
} finally {
sock.destroy();
server.stop(true);
}
});
});

it.concurrent("should not instanciate error instances in each request", async () => {
const startErrorCount = heapStats().objectTypeCounts.Error || 0;
using server = Bun.serve({
Expand Down
Loading