diff --git a/src/runtime/bake/DevServer.rs b/src/runtime/bake/DevServer.rs index f26d5e23b697..57ecce18909f 100644 --- a/src/runtime/bake/DevServer.rs +++ b/src/runtime/bake/DevServer.rs @@ -198,6 +198,8 @@ pub enum TestingBatchEvents { /// a message saying that new files have been seen. Once DevServer receives /// that signal, or times out, it will "release" this batch. Enabled(TestingBatch), + /// Released while a bundle ran; `finalize_bundle_cleanup` starts it after. + ReleaseAfterBundle(TestingBatch), } /// There is only ever one bundle executing at the same time, since all bundles @@ -1151,7 +1153,9 @@ impl Drop for DevServer { } } - if let TestingBatchEvents::Enabled(batch) = &mut self.testing_batch_events { + if let TestingBatchEvents::Enabled(batch) | TestingBatchEvents::ReleaseAfterBundle(batch) = + &mut self.testing_batch_events + { drop(std::mem::replace( &mut batch.entry_points, EntryPointList::empty(), @@ -3778,6 +3782,20 @@ fn finalize_bundle_cleanup(dev: &mut DevServer, bv2: &mut BundleV2, had_sent_hmr dev.start_next_bundle_if_present(); + // If the call above started another bundle, its cleanup releases the batch. + if matches!( + dev.testing_batch_events, + TestingBatchEvents::ReleaseAfterBundle(_) + ) && dev.current_bundle.is_none() + { + let TestingBatchEvents::ReleaseAfterBundle(batch) = + core::mem::replace(&mut dev.testing_batch_events, TestingBatchEvents::Disabled) + else { + unreachable!() + }; + dev.release_testing_batch(batch); + } + // Unref the ref added in `start_async_bundle` if let Some(server) = dev.server.as_mut() { server.on_static_request_complete(); @@ -4873,6 +4891,23 @@ pub(super) fn finalize_bundle( } impl DevServer { + /// Bundle the files a testing batch collected, or report an empty batch. + pub(crate) fn release_testing_batch(&mut self, batch: TestingBatch) { + debug_assert!(self.current_bundle.is_none()); + if batch.entry_points.set.count() == 0 { + self.publish( + HmrTopic::TestingWatchSynchronization, + &[MessageId::TestingWatchSynchronization.char(), 2], + Opcode::BINARY, + ); + return; + } + + self.start_async_bundle(batch.entry_points, true, Instant::now()) + // bun.handleOom(err) — Rust aborts on OOM by default + .expect("OOM"); + } + fn start_next_bundle_if_present(&mut self) { debug_assert!(self.magic == Magic::Valid); // Clear the current bundle diff --git a/src/runtime/bake/dev_server/hmr_socket.rs b/src/runtime/bake/dev_server/hmr_socket.rs index 671104b4c823..5a2b92de177f 100644 --- a/src/runtime/bake/dev_server/hmr_socket.rs +++ b/src/runtime/bake/dev_server/hmr_socket.rs @@ -158,6 +158,10 @@ impl HmrSocket { } x if x == IncomingMessageId::SetUrl as u8 => { let pattern = &msg[1..]; + // `match_slow` requires an absolute path; these are peer bytes. + if pattern.first() != Some(&b'/') { + return ws.close(); + } // SAFETY: JS-thread only; sole `&mut DevServer` for this scope. let dev = unsafe { self.dev() }; let maybe_rbi = dev.route_to_bundle_index_slow(pattern); @@ -204,36 +208,29 @@ impl HmrSocket { ); } } - super::TestingBatchEvents::EnableAfterBundle => { - // do not expose a websocket event that panics a release build - debug_assert!(false); + super::TestingBatchEvents::EnableAfterBundle + | super::TestingBatchEvents::ReleaseAfterBundle(_) => { + // A duplicate `H` is a protocol violation, not an invariant. ws.close(); } super::TestingBatchEvents::Enabled(_event_const) => { // Replace-and-extract to satisfy borrowck. - let super::TestingBatchEvents::Enabled(mut event) = core::mem::replace( + let super::TestingBatchEvents::Enabled(batch) = core::mem::replace( &mut dev.testing_batch_events, super::TestingBatchEvents::Disabled, ) else { unreachable!() }; - let _ = &mut event; - if event.entry_points.set.count() == 0 { - dev.publish( - HmrTopic::TestingWatchSynchronization, - &[MessageId::TestingWatchSynchronization.char(), 2], - bun_uws::Opcode::BINARY, - ); + // An unbundled route's request can start a bundle; + // `start_async_bundle` requires none in flight. + if dev.current_bundle.is_some() { + dev.testing_batch_events = + super::TestingBatchEvents::ReleaseAfterBundle(batch); return; } - let timer = std::time::Instant::now(); - dev.start_async_bundle(event.entry_points, true, timer) - // bun.handleOom(err) — Rust aborts on OOM by default - .expect("OOM"); - - // `event.entry_points.deinit(allocator)` → Drop handles this + dev.release_testing_batch(batch); } } } @@ -307,7 +304,7 @@ impl HmrSocket { } if field.contains(HmrTopic::MemoryVisualizer.as_bit()) { dev.emit_memory_visualizer_events -= 1; - if dev.emit_incremental_visualizer_events == 0 + if dev.emit_memory_visualizer_events == 0 && dev.memory_visualizer_timer.state == EventLoopTimerState::ACTIVE { // Note (jsc/runtime crate cycle): `vm.timer` is `()` on the low-tier diff --git a/src/runtime/bake/dev_server/memory_cost.rs b/src/runtime/bake/dev_server/memory_cost.rs index 766c627921cd..5c83b070d69e 100644 --- a/src/runtime/bake/dev_server/memory_cost.rs +++ b/src/runtime/bake/dev_server/memory_cost.rs @@ -234,7 +234,7 @@ pub(crate) fn memory_cost_detailed(dev: &DevServer) -> MemoryCost { // .testing_batch_events match &dev.testing_batch_events { TestingBatchEvents::Disabled => {} - TestingBatchEvents::Enabled(batch) => { + TestingBatchEvents::Enabled(batch) | TestingBatchEvents::ReleaseAfterBundle(batch) => { other_bytes += memory_cost_array_hash_map(&batch.entry_points.set); } TestingBatchEvents::EnableAfterBundle => {} diff --git a/src/runtime/bake/dev_server/mod.rs b/src/runtime/bake/dev_server/mod.rs index e8b18a04c1e0..bed3464472df 100644 --- a/src/runtime/bake/dev_server/mod.rs +++ b/src/runtime/bake/dev_server/mod.rs @@ -676,7 +676,7 @@ impl HotReloadEvent { match &mut dev_ref.testing_batch_events { TestingBatchEvents::Disabled => {} - TestingBatchEvents::Enabled(ev) => { + TestingBatchEvents::Enabled(ev) | TestingBatchEvents::ReleaseAfterBundle(ev) => { bun_core::handle_oom(ev.append(&entry_points)); dev_ref.publish( HmrTopic::TestingWatchSynchronization, diff --git a/test/bake/hmr-socket-protocol.test.ts b/test/bake/hmr-socket-protocol.test.ts new file mode 100644 index 000000000000..2c354f3d9ba7 --- /dev/null +++ b/test/bake/hmr-socket-protocol.test.ts @@ -0,0 +1,341 @@ +// The HMR websocket at `/_bun/hmr` accepts frames from any connected client +// (a browser tab, an extension, anything on the LAN with `--hostname 0.0.0.0`), +// so no frame may be able to reach an `assert`/`debug_assert` in the dev +// server. These tests drive the frames that could: +// - a second "H" (testing-batch) frame while a bundle is in flight hit the +// `TestingBatchEvents::EnableAfterBundle` arm's `debug_assert!(false)`. +// - "n" (SetUrl) with a pattern not starting with "/" reached +// `FrameworkRouter::match_slow`'s `debug_assert!(path[0] == b'/')`: an +// empty pattern indexes out of bounds inside it, any other fails it. +// - releasing a batch with "H" while an unrelated bundle was in flight +// reached `start_async_bundle`'s `debug_assert!(current_bundle.is_none())`. +// On a release build that assert is compiled out and the second +// `start_async_bundle` overwrites the in-flight `CurrentBundle`, freeing +// the arena its parse tasks are still reading: a multi-thread segfault. +import type { Subprocess } from "bun"; +import { expect, test } from "bun:test"; +import { bunEnv, bunExe, tempDir } from "harness"; + +const indexHtml = /* html */ ` + +`; + +const serverTs = /* ts */ ` + import html from "./index.html"; + const server = Bun.serve({ + port: 0, + development: { hmr: true, console: false }, + routes: { "/": html }, + fetch() { return new Response("fallback"); }, + }); + console.log("PORT=" + server.port); +`; + +/** Parks any import of `hold.block` on a fetch the test controls. */ +const holdPlugin = /* ts */ ` + export default { + name: "hold-bundle", + setup(build) { + build.onLoad({ filter: /hold\\.block$/ }, async () => { + await fetch(process.env.HMR_TEST_BUNDLE_GATE); + return { contents: "export default 1;", loader: "js" }; + }); + }, + }; +`; + +/** Drain stdout/stderr concurrently and resolve the port from the PORT= line. */ +function watchDevServer(proc: Subprocess<"ignore", "pipe", "pipe">) { + const port = Promise.withResolvers(); + let stdout = ""; + let stderr = ""; + (async () => { + for await (const chunk of proc.stdout) { + stdout += Buffer.from(chunk).toString(); + const m = stdout.match(/PORT=(\d+)/); + if (m) port.resolve(parseInt(m[1], 10)); + } + port.reject(new Error(`dev server exited before printing its port\n${stdout}${stderr}`)); + })().catch(() => {}); + (async () => { + for await (const chunk of proc.stderr) stderr += Buffer.from(chunk).toString(); + })().catch(() => {}); + return { port: port.promise, stderr: () => stderr }; +} + +/** + * Connect to `/_bun/hmr`. `onFrame` gets every server frame as its message-id + * byte plus the rest of the payload; `onClose` fires on any server-initiated + * close (an aborting dev server looks like an abrupt close to the client). + */ +async function connectHmr( + port: number, + onFrame: (id: string, body: Uint8Array) => void, + onClose: (err: Error) => void, +) { + const ws = new WebSocket(`ws://127.0.0.1:${port}/_bun/hmr`); + ws.binaryType = "arraybuffer"; + const received: string[] = []; + ws.onmessage = ev => { + const bytes = new Uint8Array(ev.data as ArrayBuffer); + const id = String.fromCharCode(bytes[0]); + received.push(id); + onFrame(id, bytes.subarray(1)); + }; + const opened = Promise.withResolvers(); + ws.onopen = () => opened.resolve(); + ws.onerror = () => opened.reject(new Error("hmr websocket failed to connect")); + ws.onclose = ev => onClose(new Error(`hmr websocket closed (code ${ev.code}, reason ${JSON.stringify(ev.reason)})`)); + await opened.promise; + return { + ws, + received, + [Symbol.dispose]() { + ws.onclose = null; + ws.close(); + }, + }; +} + +test.concurrent( + "a duplicate testing-batch frame (H) during an in-flight bundle closes the socket instead of aborting", + async () => { + // A bundler plugin parks the bundle on a fetch to this server, so "the + // bundle is in flight" is an awaited condition, not a sleep. + const bundleEntered = Promise.withResolvers(); + const bundleRelease = Promise.withResolvers(); + await using gate = Bun.serve({ + port: 0, + async fetch() { + bundleEntered.resolve(); + await bundleRelease.promise; + return new Response("go"); + }, + }); + + using dir = tempDir("hmr-socket-double-h", { + "bunfig.toml": `[serve.static]\nplugins = ["./plugin.ts"]\n`, + "plugin.ts": holdPlugin, + "index.html": indexHtml, + "entry.ts": `import "./hold.block";\nconsole.log("entry");`, + "hold.block": "", + "server.ts": serverTs, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "server.ts"], + env: { ...bunEnv, HMR_TEST_BUNDLE_GATE: String(gate.url) }, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const dev = watchDevServer(proc); + const port = await dev.port; + // If the bundle never reaches the plugin (so the gate fetch never comes), + // fail with the dev server's output instead of hanging on bundleEntered. + proc.exited.then(() => + bundleEntered.reject( + new Error(`dev server exited before the bundle reached the plugin\n--- dev server stderr ---\n${dev.stderr()}`), + ), + ); + + const closed = Promise.withResolvers(); + using hmr = await connectHmr( + port, + () => {}, + () => closed.resolve(), + ); + + // Kick off the bundle for `/`; the plugin holds it open on the gate fetch. + const pageFetch = fetch(`http://127.0.0.1:${port}/`); + pageFetch.catch(() => {}); + await bundleEntered.promise; + + // 1st "H" with a bundle in flight: Disabled -> EnableAfterBundle. + // 2nd "H": the EnableAfterBundle arm must close the socket, not assert. + hmr.ws.send("H"); + hmr.ws.send("H"); + await closed.promise; + + // Release the held bundle: the deferred `/` request completes only if the + // dev server survived the protocol violation. + bundleRelease.resolve(); + let pageStatus: string; + try { + pageStatus = String((await pageFetch).status); + } catch (e) { + pageStatus = `${(e as Error).message}\n--- dev server stderr ---\n${dev.stderr()}`; + } + expect(pageStatus).toBe("200"); + // And it is still accepting new requests. + const res = await fetch(`http://127.0.0.1:${port}/`); + expect(res.status).toBe(200); + }, +); + +// "n" is an empty pattern (the out-of-bounds index flavor); "nfoo" is a +// non-absolute pattern (the failed-assertion flavor). +test.concurrent.each(["n", "nfoo"])( + "a SetUrl frame without a leading slash (%j) closes the socket instead of aborting", + async frame => { + using dir = tempDir("hmr-socket-seturl", { + "index.html": indexHtml, + "entry.ts": `console.log("entry");`, + "server.ts": serverTs, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "server.ts"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const dev = watchDevServer(proc); + const port = await dev.port; + + // An aborting dev server also produces a close event, so the liveness + // check below is what distinguishes "closed the socket" from "died". + const closed = Promise.withResolvers(); + using hmr = await connectHmr( + port, + () => {}, + () => closed.resolve(), + ); + hmr.ws.send(frame); + await closed.promise; + + let status: string; + try { + status = String((await fetch(`http://127.0.0.1:${port}/`)).status); + } catch (e) { + status = `${(e as Error).message}\n--- dev server stderr ---\n${dev.stderr()}`; + } + expect(status).toBe("200"); + }, +); + +test.concurrent("releasing a testing batch while another bundle is in flight defers it", async () => { + // `/two` is not bundled at startup, and its entry imports `hold.block`, so + // the first request for it parks a bundle in the plugin until this server + // answers. That makes "a bundle is in flight" an awaited condition. + const bundleEntered = Promise.withResolvers(); + const bundleRelease = Promise.withResolvers(); + await using gate = Bun.serve({ + port: 0, + async fetch() { + bundleEntered.resolve(); + await bundleRelease.promise; + return new Response("go"); + }, + }); + + using dir = tempDir("hmr-socket-batch-defer", { + "bunfig.toml": `[serve.static]\nplugins = ["./plugin.ts"]\n`, + "plugin.ts": holdPlugin, + "index.html": indexHtml, + "entry.ts": `import { value } from "./dep.ts";\nconsole.log(value);`, + "dep.ts": `export const value = 0;`, + "two.html": /* html */ ` + +`, + "two.ts": `import "./hold.block";\nconsole.log("two");`, + "hold.block": "", + "server.ts": /* ts */ ` + import one from "./index.html"; + import two from "./two.html"; + const server = Bun.serve({ + port: 0, + development: { hmr: true, console: false }, + routes: { "/": one, "/two": two }, + fetch() { return new Response("fallback"); }, + }); + console.log("PORT=" + server.port); + `, + }); + await using proc = Bun.spawn({ + cmd: [bunExe(), "server.ts"], + env: { ...bunEnv, HMR_TEST_BUNDLE_GATE: String(gate.url) }, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const dev = watchDevServer(proc); + const port = await dev.port; + + // `r` frames are the testing-synchronization codes: 0 batching started, + // 1 the batch saw a file, 2 an empty batch was released, 3/4 a bundle + // finished. Waiting on them keeps every step below condition-based. + const seen = new Set(); + const waiters = new Map void>(); + const waitForSync = (code: number) => + new Promise(resolve => { + if (seen.delete(code)) return resolve(); + waiters.set(code, resolve); + }); + + // Only a process exit means the dev server died; a close is a close. + const died = Promise.withResolvers(); + died.promise.catch(() => {}); + proc.exited.then(() => died.reject(new Error(`dev server exited\n--- dev server stderr ---\n${dev.stderr()}`))); + const socketClosed = Promise.withResolvers(); + using hmr = await connectHmr( + port, + (id, body) => { + if (id !== "r") return; + const code = body[0]; + const waiter = waiters.get(code); + if (waiter) { + waiters.delete(code); + waiter(); + } else { + seen.add(code); + } + }, + () => socketClosed.resolve(), + ); + hmr.ws.send("sr"); + + // Bundle `/` so that `dep.ts` is watched and an edit to it reaches the batch. + expect((await fetch(`http://127.0.0.1:${port}/`)).status).toBe(200); + + // 1. "H" with no bundle running turns batching on (sync code 0). + const batchingOn = waitForSync(0); + hmr.ws.send("H"); + await Promise.race([batchingOn, died.promise]); + + // 2. An edit is parked in the batch instead of starting a bundle (code 1). + const sawFile = waitForSync(1); + await Bun.write(`${dir}/dep.ts`, `export const value = 1;`); + await Promise.race([sawFile, died.promise]); + + // 3. The first request for `/two` starts a bundle that the plugin holds. + const twoFetch = fetch(`http://127.0.0.1:${port}/two`); + twoFetch.catch(() => {}); + await Promise.race([bundleEntered.promise, died.promise]); + + // 4. "H" releases the batch while that bundle is still running. Starting a + // second bundle here trips start_async_bundle's assert on a debug build + // and corrupts the in-flight bundle on a release build, so the batch has + // to wait for the running bundle to finish. + hmr.ws.send("H"); + // A further "H" finds the batch pending release and closes the socket. + // Receiving that close proves the frame above was handled while the bundle + // was still held, which is the state this test is about. + hmr.ws.send("H"); + await Promise.race([socketClosed.promise, died.promise]); + + bundleRelease.resolve(); + expect( + String( + await twoFetch.then( + r => r.status, + e => `${e}\n${dev.stderr()}`, + ), + ), + ).toBe("200"); + + // The deferred batch now runs, so both routes still serve and the edit is + // live in the bundle for `/`. + const reloaded = await fetch(`http://127.0.0.1:${port}/`); + expect(reloaded.status).toBe(200); + expect(proc.killed).toBe(false); +});