From 411995db526f1a75241c481079889e82e7f522eb Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 28 Jul 2026 20:12:42 +0000 Subject: [PATCH 01/16] watcher: wake the blocked read()/kevent() on shutdown so each dev server releases its inotify/kqueue instance --- src/watcher/INotifyWatcher.rs | 98 ++++++++++++++---- src/watcher/KEventWatcher.rs | 51 ++++++--- src/watcher/Watcher.rs | 40 ++++++- src/watcher/WindowsWatcher.rs | 27 +++++ test/bake/dev-server-watcher-release.test.ts | 103 +++++++++++++++++++ 5 files changed, 284 insertions(+), 35 deletions(-) create mode 100644 test/bake/dev-server-watcher-release.test.ts diff --git a/src/watcher/INotifyWatcher.rs b/src/watcher/INotifyWatcher.rs index 23637c21ca76..7e0facaaa5cf 100644 --- a/src/watcher/INotifyWatcher.rs +++ b/src/watcher/INotifyWatcher.rs @@ -3,12 +3,10 @@ use core::ffi::c_int; use core::mem::{align_of, size_of}; -use core::sync::atomic::{AtomicU32, Ordering}; use bun_core::{ZStr, env_var, output as Output}; use bun_paths::MAX_PATH_BYTES; use bun_sys::{self, Fd}; -use bun_threading::Futex; use crate::watcher_impl::{MAX_COUNT as max_count, Op, WatchEvent, WatchItemIndex, Watcher}; use bun_collections::index_sort; @@ -41,6 +39,11 @@ pub(crate) type Platform = INotifyWatcher; pub struct INotifyWatcher { pub(crate) fd: Fd, + /// eventfd used by `wake()` to unblock the watcher thread's `ppoll()` so + /// it can observe `Watcher.running == false` and exit. Without this the + /// thread is parked in a blocking `read()` on `fd` and holds the inotify + /// instance (and its kernel `max_user_instances` slot) until process exit. + pub(crate) wake_fd: Fd, pub(crate) loaded: bool, // Avoid statically allocating because it increases the binary size. @@ -53,7 +56,6 @@ pub struct INotifyWatcher { /// see `test-fs-watch-recursive-linux-parallel-remove.js` read_ptr: Option, - pub(crate) watch_count: AtomicU32, /// nanoseconds pub(crate) coalesce_interval: isize, } @@ -62,11 +64,11 @@ impl Default for INotifyWatcher { fn default() -> Self { Self { fd: Fd::INVALID, + wake_fd: Fd::INVALID, loaded: false, eventlist_bytes: bun_core::boxed_zeroed(), eventlist_ptrs: [core::ptr::null(); max_count], read_ptr: None, - watch_count: AtomicU32::new(0), coalesce_interval: 100_000, } } @@ -129,7 +131,6 @@ impl INotifyWatcher { pub(crate) fn watch_path(&mut self, pathname: &ZStr) -> bun_sys::Result { use bun_sys::linux::IN; debug_assert!(self.loaded); - let old_count = self.watch_count.fetch_add(1, Ordering::Release); let watch_file_mask = IN::EXCL_UNLINK | IN::MOVE_SELF | IN::DELETE_SELF | IN::MOVED_TO | IN::MODIFY; // SAFETY: fd is a valid inotify fd (loaded == true), pathname is NUL-terminated. @@ -137,24 +138,19 @@ impl INotifyWatcher { bun_sys::linux::inotify_add_watch(self.fd.native(), pathname.as_ptr(), watch_file_mask) }; bun_core::scoped_log!(watcher, "inotify_add_watch({}) = {}", self.fd, rc); - let result = if rc < 0 { + if rc < 0 { Err( bun_sys::Error::from_code_int(bun_sys::last_errno(), bun_sys::Tag::watch) .with_path(pathname.as_bytes()), ) } else { Ok(rc) - }; - if old_count == 0 { - Futex::wake(&self.watch_count, 10); } - result } pub(crate) fn watch_dir(&mut self, pathname: &ZStr) -> bun_sys::Result { use bun_sys::linux::IN; debug_assert!(self.loaded); - let old_count = self.watch_count.fetch_add(1, Ordering::Release); let watch_dir_mask = IN::EXCL_UNLINK | IN::DELETE | IN::DELETE_SELF @@ -168,18 +164,14 @@ impl INotifyWatcher { bun_sys::linux::inotify_add_watch(self.fd.native(), pathname.as_ptr(), watch_dir_mask) }; bun_core::scoped_log!(watcher, "inotify_add_watch({}) = {}", self.fd, rc); - let result = if rc < 0 { + if rc < 0 { Err( bun_sys::Error::from_code_int(bun_sys::last_errno(), bun_sys::Tag::watch) .with_path(pathname.as_bytes()), ) } else { Ok(rc) - }; - if old_count == 0 { - Futex::wake(&self.watch_count, 10); } - result } pub(crate) fn new(_root: &[u8]) -> crate::Result { @@ -193,9 +185,17 @@ impl INotifyWatcher { return Err(crate::Error::Sys(errno)); } let fd = Fd::from_native(raw); - bun_core::scoped_log!(watcher, "{} init", fd); + let wake_fd = match bun_sys::eventfd(0, libc::EFD_CLOEXEC | libc::EFD_NONBLOCK) { + Ok(fd) => fd, + Err(err) => { + let _ = bun_sys::close(fd); + return Err(crate::Error::Sys(err.get_errno())); + } + }; + bun_core::scoped_log!(watcher, "{} init (wake_fd {})", fd, wake_fd); Ok(Self { fd, + wake_fd, loaded: true, coalesce_interval: env_var::BUN_INOTIFY_COALESCE_INTERVAL .get() @@ -222,12 +222,58 @@ impl INotifyWatcher { // reshaped for borrowck — track length instead of borrowing a sub-slice // of self.eventlist_bytes across the whole function. let read_len: usize = if let Some(ptr) = self.read_ptr { - Futex::wait_forever(&self.watch_count, 0); i = ptr.i; ptr.len as usize } else { 'outer: loop { - Futex::wait_forever(&self.watch_count, 0); + // Block until either the inotify fd has events or `wake()` has + // signalled the eventfd. With no watches registered the inotify + // fd simply never becomes readable, so this replaces the + // previous `Futex::wait_forever(&watch_count, 0)` and also lets + // `Watcher::shutdown` unpark the thread while it is waiting for + // a filesystem event. + let mut fds = [ + system::pollfd { + fd: self.fd.native(), + events: libc::POLLIN, + revents: 0, + }, + system::pollfd { + fd: self.wake_fd.native(), + events: libc::POLLIN, + revents: 0, + }, + ]; + // SAFETY: fds is a valid stack array of `fds.len()` entries; + // timeout/sigmask are null (block indefinitely). + let poll_rc = unsafe { + system::ppoll( + fds.as_mut_ptr(), + fds.len(), + core::ptr::null(), + core::ptr::null(), + ) + }; + if poll_rc < 0 { + let e = get_errno(poll_rc); + if matches!(e, E::EAGAIN | E::EINTR) { + continue 'outer; + } + return Err(bun_sys::Error { + errno: e as u32 as _, + syscall: bun_sys::Tag::poll, + ..Default::default() + }); + } + if fds[1].revents != 0 { + // `wake()` fired. Yield an empty batch so `watch_loop` can + // re-check `running` and exit; the eventfd is closed by + // `stop()` on the way out. + return Ok(&[]); + } + if fds[0].revents & libc::POLLIN == 0 { + continue 'outer; + } // SAFETY: fd is a valid inotify fd; buffer is valid for eventlist_bytes.len() bytes. let rc = unsafe { @@ -365,6 +411,20 @@ impl INotifyWatcher { let _ = bun_sys::close(self.fd); self.fd = Fd::INVALID; } + if self.wake_fd != Fd::INVALID { + let _ = bun_sys::close(self.wake_fd); + self.wake_fd = Fd::INVALID; + } + } + + /// Unblock the watcher thread's `ppoll()` in `read()` so it can observe + /// `Watcher.running == false` and exit. Called from `Watcher::shutdown` + /// on the main thread while holding `Watcher.mutex`. + pub(crate) fn wake(&self) { + if self.wake_fd == Fd::INVALID { + return; + } + let _ = bun_sys::write(self.wake_fd, &1u64.to_ne_bytes()); } } diff --git a/src/watcher/KEventWatcher.rs b/src/watcher/KEventWatcher.rs index c784b8636584..846e3e13300c 100644 --- a/src/watcher/KEventWatcher.rs +++ b/src/watcher/KEventWatcher.rs @@ -11,12 +11,23 @@ pub struct KEventWatcher { const CHANGELIST_COUNT: usize = 128; +/// Arbitrary non-zero `ident` for the EVFILT_USER wakeup event registered in +/// `new()` and triggered by `wake()`. +const WAKE_EVENT_IDENT: usize = 0x2307; + impl KEventWatcher { pub(crate) fn new(_root: &[u8]) -> crate::Result { let fd = bun_sys::kqueue()?; if fd.native() == 0 { return Err(crate::Error::KQueueError); } + // Register a user-triggered event so `wake()` can unblock the + // blocking `kevent()` in `watch_loop_cycle` during shutdown. + let mut ev: libc::kevent = bun_core::ffi::zeroed(); + ev.ident = WAKE_EVENT_IDENT; + ev.filter = libc::EVFILT_USER; + ev.flags = (libc::EV_ADD | libc::EV_CLEAR) as _; + let _ = bun_sys::kevent(fd, core::slice::from_ref(&ev), &mut [], None); Ok(Self { fd }) } @@ -26,6 +37,20 @@ impl KEventWatcher { self.fd = Fd::INVALID; } } + + /// Unblock the watcher thread's blocking `kevent()` so it can observe + /// `Watcher.running == false` and exit. Called from `Watcher::shutdown` + /// on the main thread while holding `Watcher.mutex`. + pub(crate) fn wake(&self) { + if !self.fd.is_valid() { + return; + } + let mut ev: libc::kevent = bun_core::ffi::zeroed(); + ev.ident = WAKE_EVENT_IDENT; + ev.filter = libc::EVFILT_USER; + ev.fflags = libc::NOTE_TRIGGER; + let _ = bun_sys::kevent(self.fd, core::slice::from_ref(&ev), &mut [], None); + } } fn watch_event_from_kevent(kevent: &libc::kevent) -> WatchEvent { @@ -70,21 +95,23 @@ pub(crate) fn watch_loop_cycle(this: &mut Watcher) -> bun_sys::Result<()> { let changes = &changelist[..count]; let watchevents = &mut this.watch_events[..count]; let mut out_len: usize = 0; - if let [first, rest @ ..] = changes { - watchevents[0] = watch_event_from_kevent(first); - out_len = 1; - let mut prev_event = first; - for event in rest { - if prev_event.udata == event.udata { - let new = watch_event_from_kevent(event); - watchevents[out_len - 1].merge(new); + let mut prev_event: Option<&libc::kevent> = None; + for event in changes { + // Only VNODE events map to watch items; skip the EVFILT_USER wakeup + // posted by `wake()`. + if event.filter != libc::EVFILT_VNODE { + continue; + } + if let Some(prev) = prev_event { + if prev.udata == event.udata { + watchevents[out_len - 1].merge(watch_event_from_kevent(event)); + prev_event = Some(event); continue; } - - watchevents[out_len] = watch_event_from_kevent(event); - prev_event = event; - out_len += 1; } + watchevents[out_len] = watch_event_from_kevent(event); + prev_event = Some(event); + out_len += 1; } this.dispatch_file_updates(out_len, out_len); diff --git a/src/watcher/Watcher.rs b/src/watcher/Watcher.rs index 59cdd4020637..28b5a7c528ec 100644 --- a/src/watcher/Watcher.rs +++ b/src/watcher/Watcher.rs @@ -294,13 +294,24 @@ impl Watcher { me.mutex.lock(); me.close_descriptors.store(close_descriptors); me.running.store(false); + // Wake the watcher thread out of its blocking wait so it can + // observe `running == false`, run `platform.stop()`, and free + // `*this`. Without this the thread stays parked (in inotify + // `read()` / `kevent()`) and every disposed dev server leaks its + // inotify/kqueue instance until process exit. + me.platform.wake(); me.mutex.unlock(); + // `*this` may be freed by the watcher thread any time after this + // point; `thread_main` takes/releases `mutex` as a barrier before + // `heap::take(this)` so it cannot proceed until the unlock above. false } else { if close_descriptors && me.running.load() { let fds = me.watchlist.items_fd(); for &fd in fds { - let _ = bun_sys::close(fd); + if fd.is_valid() { + let _ = bun_sys::close(fd); + } } } true @@ -310,7 +321,13 @@ impl Watcher { // watchlist freed by Drop on Box // SAFETY: this was heap-allocated by caller of init(); no borrow of it // is live here. - drop(unsafe { bun_core::heap::take(this) }); + let mut me = unsafe { bun_core::heap::take(this) }; + // A spawned thread runs `platform.stop()` itself in `thread_body`, + // also when it hands `*this` back after a watch error. + if me.thread.is_none() { + me.platform.stop(); + } + drop(me); } } @@ -357,7 +374,6 @@ impl Watcher { let owner_still_alive = match self.watch_loop() { Err(err) => { self.watchloop_handle.store(false); - self.platform.stop(); let running = self.running.load(); if running { (self.on_error)(self.ctx, err); @@ -367,11 +383,27 @@ impl Watcher { Ok(()) => false, }; + // Barrier: `shutdown()` holds `self.mutex` across + // `running.store(false)` and `platform.wake()`. This thread can + // observe `running == false` at the unlocked `while` check and + // fall through here before `shutdown()` has unlocked, so without + // this pair `platform.stop()` below could race `wake()` touching + // the same platform fds, and `heap::take(this)` could free + // `self.mutex` out from under `shutdown()`'s pending unlock. + self.mutex.lock(); + self.mutex.unlock(); + + // Release platform resources. `wake()` makes the loop exit via + // `Ok(())`, so this must run on both arms (previously only `Err`). + self.platform.stop(); + // deinit and close descriptors if needed if self.close_descriptors.load() { let fds = self.watchlist.items_fd(); for &fd in fds { - let _ = bun_sys::close(fd); + if fd.is_valid() { + let _ = bun_sys::close(fd); + } } } owner_still_alive diff --git a/src/watcher/WindowsWatcher.rs b/src/watcher/WindowsWatcher.rs index e39adde6e990..f16b80f20f89 100644 --- a/src/watcher/WindowsWatcher.rs +++ b/src/watcher/WindowsWatcher.rs @@ -22,6 +22,12 @@ pub struct WindowsWatcher { pub(crate) watcher: DirWatcher, pub(crate) buf: PathBuffer, pub(crate) base_idx: usize, + /// Latched true once `next()` has armed a `ReadDirectoryChangesW` on + /// `self.watcher.overlapped`. While set, `stop()` must not close the + /// handles from `thread_main`: the kernel may still write the + /// cancellation status into `overlapped` after `CloseHandle`, which lands + /// in freed memory once `heap::take(this)` runs. + pub(crate) armed: bool, } impl Default for WindowsWatcher { @@ -35,6 +41,7 @@ impl Default for WindowsWatcher { }, buf: PathBuffer::ZEROED, base_idx: 0, + armed: false, } } } @@ -298,6 +305,7 @@ impl WindowsWatcher { bun_core::scoped_log!(watcher, "prepare() returned error"); return Err(err); } + self.armed = true; let mut nbytes: w::DWORD = 0; let mut key: w::ULONG_PTR = 0; @@ -368,12 +376,31 @@ impl WindowsWatcher { } pub(crate) fn stop(&mut self) { + if self.armed { + // A `ReadDirectoryChangesW` is (or may still be) pending on + // `self.watcher.overlapped`; `CloseHandle(dir_handle)` cancels it + // asynchronously and `thread_main` frees `*self` right after this + // returns, so the kernel's cancellation write can land in freed + // memory. Leak the two handles until process exit; the proper fix + // is a `CancelIoEx` + IOCP drain before `heap::take`. + return; + } // SAFETY: handles were opened in init() and are valid until stop() is called once. unsafe { w::CloseHandle(self.watcher.dir_handle); w::CloseHandle(self.iocp); } } + + /// On Linux/macOS `wake()` unblocks the watcher thread so it can observe + /// `Watcher.running == false`, run `stop()`, and exit. On Windows that + /// teardown is not yet safe: `next()` keeps a `ReadDirectoryChangesW` + /// pending on `self.watcher.overlapped`, and freeing `self` after `stop()` + /// races the kernel's cancellation write into that buffer. A proper fix + /// needs `CancelIoEx` + draining the IOCP before `heap::take`; until then + /// the thread stays parked in `GetQueuedCompletionStatus` until process + /// exit (unchanged from before). + pub(crate) fn wake(&self) {} } #[repr(u32)] diff --git a/test/bake/dev-server-watcher-release.test.ts b/test/bake/dev-server-watcher-release.test.ts new file mode 100644 index 000000000000..df17b53e671d --- /dev/null +++ b/test/bake/dev-server-watcher-release.test.ts @@ -0,0 +1,103 @@ +// Each stopped `Bun.serve({ development: true })` must release its file +// watcher. On Linux the watcher thread was parked in a blocking `read()` on +// the inotify fd after `server.stop()`, so every disposed dev server leaked +// one inotify instance (and one thread) until process exit. In a container +// with a tight `fs.inotify.max_user_instances` budget this surfaced as +// `EMFILE while initializing file watcher for development server`. + +import { test, expect } from "bun:test"; +import { bunEnv, bunExe, tempDir, isLinux, isWindows } from "harness"; + +// Windows watcher `wake()` is intentionally a no-op (see +// `src/watcher/WindowsWatcher.rs`), so the thread/handle are held until +// process exit there; the inotify-instance budget problem this test covers +// is POSIX-specific. +test.skipIf(isWindows)("dev server releases its file watcher on stop()", async () => { + const fixture = /* ts */ ` + import { readdirSync, readlinkSync, readFileSync } from "node:fs"; + import html from "./index.html"; + + function scan() { + if (process.platform !== "linux") return { inotify: 0, threads: 0 }; + let inotify = 0; + for (const name of readdirSync("/proc/self/fd")) { + try { + if (readlinkSync("/proc/self/fd/" + name) === "anon_inode:inotify") inotify++; + } catch {} + } + const status = readFileSync("/proc/self/status", "utf8"); + const threads = Number(/^Threads:\\s+(\\d+)/m.exec(status)?.[1] ?? 0); + return { inotify, threads }; + } + + const ITER = 10; + + // warm-up: the first server initialises process-global state + { + const s = Bun.serve({ port: 0, development: true, static: { "/": html }, fetch: () => new Response("") }); + await (await fetch(s.url)).text(); + s.stop(true); + } + // wait for the warm-up watcher to release so it isn't counted + for (let i = 0; i < 40 && scan().inotify > 0; i++) { + Bun.gc(true); + await Bun.sleep(50); + } + const before = scan(); + + for (let i = 0; i < ITER; i++) { + const s = Bun.serve({ port: 0, development: true, static: { "/": html }, fetch: () => new Response("") }); + await (await fetch(s.url)).text(); + s.stop(true); + } + + // poll until watcher threads have observed running=false and exited + for (let i = 0; i < 40; i++) { + Bun.gc(true); + const now = scan(); + if (now.inotify <= before.inotify && now.threads <= before.threads) break; + await Bun.sleep(50); + } + + const after = scan(); + console.log(JSON.stringify({ + iterations: ITER, + inotifyDelta: after.inotify - before.inotify, + threadDelta: after.threads - before.threads, + })); + `; + + using dir = tempDir("dev-server-watcher-release", { + "index.html": "hi", + "fixture.ts": fixture, + }); + + await using proc = Bun.spawn({ + cmd: [bunExe(), "run", "fixture.ts"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + + const line = stdout + .split("\n") + .reverse() + .find(l => l.startsWith("{")); + if (!line) { + throw new Error(`no JSON summary in stdout.\nstdout:\n${stdout}\nstderr:\n${stderr}`); + } + const { iterations, inotifyDelta, threadDelta } = JSON.parse(line); + + expect(stderr).not.toContain("error:"); + expect(exitCode).toBe(0); + + if (isLinux) { + // Without the fix every iteration leaks one inotify instance + // (inotifyDelta == iterations). With the fix all of them are released. + expect(inotifyDelta).toBeLessThanOrEqual(1); + expect(inotifyDelta).toBeLessThan(iterations); + expect(threadDelta).toBeLessThan(iterations); + } +}); From 3c4ca0d2e3cad4110d9d2b762d05d777cfb5643f Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Tue, 28 Jul 2026 20:15:25 +0000 Subject: [PATCH 02/16] [autofix.ci] apply automated fixes --- test/bake/dev-server-watcher-release.test.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/test/bake/dev-server-watcher-release.test.ts b/test/bake/dev-server-watcher-release.test.ts index df17b53e671d..6befd3cd5b2b 100644 --- a/test/bake/dev-server-watcher-release.test.ts +++ b/test/bake/dev-server-watcher-release.test.ts @@ -5,8 +5,8 @@ // with a tight `fs.inotify.max_user_instances` budget this surfaced as // `EMFILE while initializing file watcher for development server`. -import { test, expect } from "bun:test"; -import { bunEnv, bunExe, tempDir, isLinux, isWindows } from "harness"; +import { expect, test } from "bun:test"; +import { bunEnv, bunExe, isLinux, isWindows, tempDir } from "harness"; // Windows watcher `wake()` is intentionally a no-op (see // `src/watcher/WindowsWatcher.rs`), so the thread/handle are held until From 50a4ad57de034927a2cf4da0ebb6005c1b502c35 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 28 Jul 2026 21:01:46 +0000 Subject: [PATCH 03/16] address review: tighten comments, reset WindowsWatcher.armed after GQCS dequeues, skip test outside Linux, gate poll on inotify only, drop per-watcher WatcherTrace::deinit --- src/watcher/INotifyWatcher.rs | 24 +++++-------- src/watcher/KEventWatcher.rs | 14 +++----- src/watcher/Watcher.rs | 23 +++---------- src/watcher/WatcherTrace.rs | 7 ---- src/watcher/WindowsWatcher.rs | 30 +++++++--------- test/bake/dev-server-watcher-release.test.ts | 36 +++++++++----------- 6 files changed, 45 insertions(+), 89 deletions(-) diff --git a/src/watcher/INotifyWatcher.rs b/src/watcher/INotifyWatcher.rs index 7e0facaaa5cf..a0b98c79809c 100644 --- a/src/watcher/INotifyWatcher.rs +++ b/src/watcher/INotifyWatcher.rs @@ -39,10 +39,8 @@ pub(crate) type Platform = INotifyWatcher; pub struct INotifyWatcher { pub(crate) fd: Fd, - /// eventfd used by `wake()` to unblock the watcher thread's `ppoll()` so - /// it can observe `Watcher.running == false` and exit. Without this the - /// thread is parked in a blocking `read()` on `fd` and holds the inotify - /// instance (and its kernel `max_user_instances` slot) until process exit. + /// eventfd written by `wake()`; `read()` ppolls on `[fd, wake_fd]` so + /// `Watcher::shutdown` can unpark the thread. pub(crate) wake_fd: Fd, pub(crate) loaded: bool, @@ -226,12 +224,9 @@ impl INotifyWatcher { ptr.len as usize } else { 'outer: loop { - // Block until either the inotify fd has events or `wake()` has - // signalled the eventfd. With no watches registered the inotify - // fd simply never becomes readable, so this replaces the - // previous `Futex::wait_forever(&watch_count, 0)` and also lets - // `Watcher::shutdown` unpark the thread while it is waiting for - // a filesystem event. + // Block until the inotify fd has events or `wake()` has + // signalled the eventfd; an inotify fd with no watches never + // becomes readable, so `wake_fd` is the only way out then. let mut fds = [ system::pollfd { fd: self.fd.native(), @@ -266,9 +261,7 @@ impl INotifyWatcher { }); } if fds[1].revents != 0 { - // `wake()` fired. Yield an empty batch so `watch_loop` can - // re-check `running` and exit; the eventfd is closed by - // `stop()` on the way out. + // `wake()` fired: let `watch_loop` re-check `running`. return Ok(&[]); } if fds[0].revents & libc::POLLIN == 0 { @@ -417,9 +410,8 @@ impl INotifyWatcher { } } - /// Unblock the watcher thread's `ppoll()` in `read()` so it can observe - /// `Watcher.running == false` and exit. Called from `Watcher::shutdown` - /// on the main thread while holding `Watcher.mutex`. + /// Unblock the watcher thread's `ppoll()` so it re-checks `running`. + /// Called from `Watcher::shutdown` under `Watcher.mutex`. pub(crate) fn wake(&self) { if self.wake_fd == Fd::INVALID { return; diff --git a/src/watcher/KEventWatcher.rs b/src/watcher/KEventWatcher.rs index 846e3e13300c..cc2a41d5158c 100644 --- a/src/watcher/KEventWatcher.rs +++ b/src/watcher/KEventWatcher.rs @@ -11,8 +11,8 @@ pub struct KEventWatcher { const CHANGELIST_COUNT: usize = 128; -/// Arbitrary non-zero `ident` for the EVFILT_USER wakeup event registered in -/// `new()` and triggered by `wake()`. +/// `ident` for the EVFILT_USER wakeup event registered in `new()` and +/// triggered by `wake()`. const WAKE_EVENT_IDENT: usize = 0x2307; impl KEventWatcher { @@ -21,8 +21,6 @@ impl KEventWatcher { if fd.native() == 0 { return Err(crate::Error::KQueueError); } - // Register a user-triggered event so `wake()` can unblock the - // blocking `kevent()` in `watch_loop_cycle` during shutdown. let mut ev: libc::kevent = bun_core::ffi::zeroed(); ev.ident = WAKE_EVENT_IDENT; ev.filter = libc::EVFILT_USER; @@ -38,9 +36,8 @@ impl KEventWatcher { } } - /// Unblock the watcher thread's blocking `kevent()` so it can observe - /// `Watcher.running == false` and exit. Called from `Watcher::shutdown` - /// on the main thread while holding `Watcher.mutex`. + /// Unblock the watcher thread's `kevent()` so it re-checks `running`. + /// Called from `Watcher::shutdown` under `Watcher.mutex`. pub(crate) fn wake(&self) { if !self.fd.is_valid() { return; @@ -97,8 +94,7 @@ pub(crate) fn watch_loop_cycle(this: &mut Watcher) -> bun_sys::Result<()> { let mut out_len: usize = 0; let mut prev_event: Option<&libc::kevent> = None; for event in changes { - // Only VNODE events map to watch items; skip the EVFILT_USER wakeup - // posted by `wake()`. + // Only VNODE events map to watch items (filters out `wake()`'s EVFILT_USER). if event.filter != libc::EVFILT_VNODE { continue; } diff --git a/src/watcher/Watcher.rs b/src/watcher/Watcher.rs index 28b5a7c528ec..30d85e09ceeb 100644 --- a/src/watcher/Watcher.rs +++ b/src/watcher/Watcher.rs @@ -294,16 +294,10 @@ impl Watcher { me.mutex.lock(); me.close_descriptors.store(close_descriptors); me.running.store(false); - // Wake the watcher thread out of its blocking wait so it can - // observe `running == false`, run `platform.stop()`, and free - // `*this`. Without this the thread stays parked (in inotify - // `read()` / `kevent()`) and every disposed dev server leaks its - // inotify/kqueue instance until process exit. me.platform.wake(); me.mutex.unlock(); // `*this` may be freed by the watcher thread any time after this - // point; `thread_main` takes/releases `mutex` as a barrier before - // `heap::take(this)` so it cannot proceed until the unlock above. + // unlock; `thread_main` lock/unlocks `mutex` before `heap::take`. false } else { if close_descriptors && me.running.load() { @@ -350,9 +344,6 @@ impl Watcher { // argument is still protected is UB under Stacked Borrows / Tree Borrows. let owner_still_alive = unsafe { (*this).thread_body() }; - // Close trace file if open - WatcherTrace::deinit(); - Output::flush(); if !owner_still_alive { @@ -383,18 +374,12 @@ impl Watcher { Ok(()) => false, }; - // Barrier: `shutdown()` holds `self.mutex` across - // `running.store(false)` and `platform.wake()`. This thread can - // observe `running == false` at the unlocked `while` check and - // fall through here before `shutdown()` has unlocked, so without - // this pair `platform.stop()` below could race `wake()` touching - // the same platform fds, and `heap::take(this)` could free - // `self.mutex` out from under `shutdown()`'s pending unlock. + // Barrier: `shutdown()` holds `mutex` across `running.store(false)` + // and `platform.wake()`; `stop()` and `heap::take(this)` below + // must not run until `shutdown()` has unlocked. self.mutex.lock(); self.mutex.unlock(); - // Release platform resources. `wake()` makes the loop exit via - // `Ok(())`, so this must run on both arms (previously only `Err`). self.platform.stop(); // deinit and close descriptors if needed diff --git a/src/watcher/WatcherTrace.rs b/src/watcher/WatcherTrace.rs index 3301de3b2521..c0b50c942ae0 100644 --- a/src/watcher/WatcherTrace.rs +++ b/src/watcher/WatcherTrace.rs @@ -172,10 +172,3 @@ pub(crate) fn write_events( return; } } - -/// Close the trace file if open -// free-function `deinit` (no `self`), so this stays a plain fn -// rather than `impl Drop`. -pub(crate) fn deinit() { - let _ = TRACE_FILE.lock().take(); -} diff --git a/src/watcher/WindowsWatcher.rs b/src/watcher/WindowsWatcher.rs index f16b80f20f89..f9136784d132 100644 --- a/src/watcher/WindowsWatcher.rs +++ b/src/watcher/WindowsWatcher.rs @@ -23,10 +23,8 @@ pub struct WindowsWatcher { pub(crate) buf: PathBuffer, pub(crate) base_idx: usize, /// Latched true once `next()` has armed a `ReadDirectoryChangesW` on - /// `self.watcher.overlapped`. While set, `stop()` must not close the - /// handles from `thread_main`: the kernel may still write the - /// cancellation status into `overlapped` after `CloseHandle`, which lands - /// in freed memory once `heap::take(this)` runs. + /// `self.watcher.overlapped`; while set, `stop()` must leak the handles + /// (the kernel's cancellation write would land in freed memory). pub(crate) armed: bool, } @@ -337,6 +335,9 @@ impl WindowsWatcher { if overlapped != &mut self.watcher.overlapped as *mut w::OVERLAPPED { continue; } + // Our completion was dequeued; nothing is pending on + // `overlapped` until the next successful `prepare()`. + self.armed = false; if nbytes == 0 { // ReadDirectoryChangesW internal change-buffer overflow — too many // events arrived between drain and re-arm. This is NOT a shutdown @@ -354,6 +355,7 @@ impl WindowsWatcher { if let Err(err) = self.watcher.prepare() { return Err(err); } + self.armed = true; continue; } return Ok(Some(EventIterator { @@ -377,12 +379,8 @@ impl WindowsWatcher { pub(crate) fn stop(&mut self) { if self.armed { - // A `ReadDirectoryChangesW` is (or may still be) pending on - // `self.watcher.overlapped`; `CloseHandle(dir_handle)` cancels it - // asynchronously and `thread_main` frees `*self` right after this - // returns, so the kernel's cancellation write can land in freed - // memory. Leak the two handles until process exit; the proper fix - // is a `CancelIoEx` + IOCP drain before `heap::take`. + // See `armed`. Proper fix: `CancelIoEx` + IOCP drain before + // `heap::take`; until then leak the two handles. return; } // SAFETY: handles were opened in init() and are valid until stop() is called once. @@ -392,14 +390,10 @@ impl WindowsWatcher { } } - /// On Linux/macOS `wake()` unblocks the watcher thread so it can observe - /// `Watcher.running == false`, run `stop()`, and exit. On Windows that - /// teardown is not yet safe: `next()` keeps a `ReadDirectoryChangesW` - /// pending on `self.watcher.overlapped`, and freeing `self` after `stop()` - /// races the kernel's cancellation write into that buffer. A proper fix - /// needs `CancelIoEx` + draining the IOCP before `heap::take`; until then - /// the thread stays parked in `GetQueuedCompletionStatus` until process - /// exit (unchanged from before). + /// No-op: `next()` keeps a `ReadDirectoryChangesW` pending on + /// `self.watcher.overlapped`, so freeing `self` after waking would race + /// the kernel's cancellation write. The thread stays parked in + /// `GetQueuedCompletionStatus` until process exit (see `armed`). pub(crate) fn wake(&self) {} } diff --git a/test/bake/dev-server-watcher-release.test.ts b/test/bake/dev-server-watcher-release.test.ts index 6befd3cd5b2b..82bdfc44c279 100644 --- a/test/bake/dev-server-watcher-release.test.ts +++ b/test/bake/dev-server-watcher-release.test.ts @@ -5,20 +5,18 @@ // with a tight `fs.inotify.max_user_instances` budget this surfaced as // `EMFILE while initializing file watcher for development server`. -import { expect, test } from "bun:test"; -import { bunEnv, bunExe, isLinux, isWindows, tempDir } from "harness"; +import { test, expect } from "bun:test"; +import { bunEnv, bunExe, tempDir, isLinux } from "harness"; -// Windows watcher `wake()` is intentionally a no-op (see -// `src/watcher/WindowsWatcher.rs`), so the thread/handle are held until -// process exit there; the inotify-instance budget problem this test covers -// is POSIX-specific. -test.skipIf(isWindows)("dev server releases its file watcher on stop()", async () => { +// The `/proc/self/{fd,status}` probes are Linux-specific; on macOS the kqueue +// leak is silent (no observable assertion here) and on Windows `wake()` is an +// intentional no-op (see `src/watcher/WindowsWatcher.rs`). +test.skipIf(!isLinux)("dev server releases its file watcher on stop()", async () => { const fixture = /* ts */ ` import { readdirSync, readlinkSync, readFileSync } from "node:fs"; import html from "./index.html"; function scan() { - if (process.platform !== "linux") return { inotify: 0, threads: 0 }; let inotify = 0; for (const name of readdirSync("/proc/self/fd")) { try { @@ -51,11 +49,10 @@ test.skipIf(isWindows)("dev server releases its file watcher on stop()", async ( s.stop(true); } - // poll until watcher threads have observed running=false and exited - for (let i = 0; i < 40; i++) { + // poll until the watcher threads have closed their inotify instances + // (Threads: is not a reliable gate; JSC may spawn a collector thread) + for (let i = 0; i < 40 && scan().inotify > before.inotify; i++) { Bun.gc(true); - const now = scan(); - if (now.inotify <= before.inotify && now.threads <= before.threads) break; await Bun.sleep(50); } @@ -91,13 +88,12 @@ test.skipIf(isWindows)("dev server releases its file watcher on stop()", async ( const { iterations, inotifyDelta, threadDelta } = JSON.parse(line); expect(stderr).not.toContain("error:"); - expect(exitCode).toBe(0); - if (isLinux) { - // Without the fix every iteration leaks one inotify instance - // (inotifyDelta == iterations). With the fix all of them are released. - expect(inotifyDelta).toBeLessThanOrEqual(1); - expect(inotifyDelta).toBeLessThan(iterations); - expect(threadDelta).toBeLessThan(iterations); - } + // Without the fix every iteration leaks one inotify instance + // (inotifyDelta == iterations). With the fix all of them are released. + expect(inotifyDelta).toBeLessThanOrEqual(1); + expect(inotifyDelta).toBeLessThan(iterations); + expect(threadDelta).toBeLessThan(iterations); + + expect(exitCode).toBe(0); }); From 1394fb11bbfd2a272fe60ed6df453e64d2105146 Mon Sep 17 00:00:00 2001 From: "autofix-ci[bot]" <114827586+autofix-ci[bot]@users.noreply.github.com> Date: Tue, 28 Jul 2026 21:04:47 +0000 Subject: [PATCH 04/16] [autofix.ci] apply automated fixes --- test/bake/dev-server-watcher-release.test.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/test/bake/dev-server-watcher-release.test.ts b/test/bake/dev-server-watcher-release.test.ts index 82bdfc44c279..1af9f6e1a684 100644 --- a/test/bake/dev-server-watcher-release.test.ts +++ b/test/bake/dev-server-watcher-release.test.ts @@ -5,8 +5,8 @@ // with a tight `fs.inotify.max_user_instances` budget this surfaced as // `EMFILE while initializing file watcher for development server`. -import { test, expect } from "bun:test"; -import { bunEnv, bunExe, tempDir, isLinux } from "harness"; +import { expect, test } from "bun:test"; +import { bunEnv, bunExe, isLinux, tempDir } from "harness"; // The `/proc/self/{fd,status}` probes are Linux-specific; on macOS the kqueue // leak is silent (no observable assertion here) and on Windows `wake()` is an From 265f8c01cc252732fd0f8f92088ca23b5ff5606f Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 28 Jul 2026 21:49:37 +0000 Subject: [PATCH 05/16] windows: clear armed when GQCS dequeues a failed-I/O completion for our overlapped --- src/watcher/WindowsWatcher.rs | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/src/watcher/WindowsWatcher.rs b/src/watcher/WindowsWatcher.rs index f9136784d132..3e9b60c4db56 100644 --- a/src/watcher/WindowsWatcher.rs +++ b/src/watcher/WindowsWatcher.rs @@ -325,6 +325,12 @@ impl WindowsWatcher { if err == w::Win32Error::TIMEOUT || err == w::Win32Error(258) { return Ok(None); } else { + // GQCS returning FALSE with `*lpOverlapped != NULL` + // dequeued a failed-I/O completion; nothing remains + // outstanding on our OVERLAPPED in that case. + if overlapped == &mut self.watcher.overlapped as *mut w::OVERLAPPED { + self.armed = false; + } bun_core::scoped_log!(watcher, "GetQueuedCompletionStatus failed: {}", err.0); return Err(bun_sys::Error::from_win32(err, bun_sys::Tag::watch)); } From 2f0dd2ceafdc1f11eb714c9a802da031613aabeb Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 28 Jul 2026 22:08:21 +0000 Subject: [PATCH 06/16] watcher: key shutdown() off thread.is_some() instead of a thread-written flag; drop flaky threadDelta assertion Under ASAN the spawned thread can be scheduled after shutdown() runs, so a thread-written flag is still false and shutdown() frees *this from under thread_main. Use the JoinHandle written synchronously by start() instead. threadDelta in the test is racy (thread closes its inotify fd before exiting, and JSC/bundler threads inflate the count); the inotifyDelta==0 assertion is the proof and is kept. --- test/bake/dev-server-watcher-release.test.ts | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/test/bake/dev-server-watcher-release.test.ts b/test/bake/dev-server-watcher-release.test.ts index 1af9f6e1a684..cbff086c02c2 100644 --- a/test/bake/dev-server-watcher-release.test.ts +++ b/test/bake/dev-server-watcher-release.test.ts @@ -91,9 +91,14 @@ test.skipIf(!isLinux)("dev server releases its file watcher on stop()", async () // Without the fix every iteration leaks one inotify instance // (inotifyDelta == iterations). With the fix all of them are released. - expect(inotifyDelta).toBeLessThanOrEqual(1); + // `threadDelta` is reported for diagnostics only: `Threads:` also counts + // JSC/bundler threads and can transiently read high right after `stop()` + // has closed the inotify fd but before the watcher thread has exited. + expect({ inotifyDelta, threadDelta }).toEqual({ + inotifyDelta: 0, + threadDelta: expect.any(Number), + }); expect(inotifyDelta).toBeLessThan(iterations); - expect(threadDelta).toBeLessThan(iterations); expect(exitCode).toBe(0); }); From 905d292aab911bbb27e520d2ce2fb41636b16a13 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 28 Jul 2026 22:33:11 +0000 Subject: [PATCH 07/16] test: drop vacuous inotifyDelta < iterations assertion --- test/bake/dev-server-watcher-release.test.ts | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/test/bake/dev-server-watcher-release.test.ts b/test/bake/dev-server-watcher-release.test.ts index cbff086c02c2..1ccda628c535 100644 --- a/test/bake/dev-server-watcher-release.test.ts +++ b/test/bake/dev-server-watcher-release.test.ts @@ -85,7 +85,7 @@ test.skipIf(!isLinux)("dev server releases its file watcher on stop()", async () if (!line) { throw new Error(`no JSON summary in stdout.\nstdout:\n${stdout}\nstderr:\n${stderr}`); } - const { iterations, inotifyDelta, threadDelta } = JSON.parse(line); + const { inotifyDelta, threadDelta } = JSON.parse(line); expect(stderr).not.toContain("error:"); @@ -98,7 +98,6 @@ test.skipIf(!isLinux)("dev server releases its file watcher on stop()", async () inotifyDelta: 0, threadDelta: expect.any(Number), }); - expect(inotifyDelta).toBeLessThan(iterations); expect(exitCode).toBe(0); }); From 082c6a6eaed3ef5faaa5feb340173a380bdd2116 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Tue, 28 Jul 2026 23:10:13 +0000 Subject: [PATCH 08/16] test: add dev-server-watcher-release to no-validate-leaksan The subprocess creates 11 Bun.serve({development:true}) instances; each leaks its NewServer/ServerConfig init allocations through the pre-existing NewServer<->HTMLBundle::Route cycle that survives server.stop(). On main the parked watcher threads' stacks kept those allocations LSAN-reachable; with the watcher threads now exiting they surface as leaks. --- test/no-validate-leaksan.txt | 3 +++ 1 file changed, 3 insertions(+) diff --git a/test/no-validate-leaksan.txt b/test/no-validate-leaksan.txt index 32a677be1b1f..cdc4780e6d63 100644 --- a/test/no-validate-leaksan.txt +++ b/test/no-validate-leaksan.txt @@ -373,6 +373,9 @@ test/bake/dev/react-spa.test.ts test/bake/dev/sourcemap.test.ts test/bake/dev/ssg-pages-router.test.ts test/bake/dev/deinitialization.test.ts +# NewServer <-> HTMLBundle::Route back-pointer cycle survives server.stop(); +# each disposed dev server leaks its ServerConfig/NewServer init allocations. +test/bake/dev-server-watcher-release.test.ts # Need to terminate HTTP thread. test/js/bun/test/parallel/test-http-should-not-accept-untrusted-certificates.ts From fab20d3c9d6e85d69de27e7015a2c07f5b1be5a4 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 29 Jul 2026 00:30:37 +0000 Subject: [PATCH 09/16] watcher(macOS): use the io_darwin machport waker instead of EVFILT_USER Matches bun_io::waker::KEventWaker: io_darwin_create_machport registers an EVFILT_MACHPORT on the kqueue, io_darwin_schedule_wakeup sends to it, and io_darwin_close_machport deallocates in stop(). FreeBSD has no mach ports so it keeps the kqueue-native EVFILT_USER wakeup. --- src/io/io_darwin.cpp | 5 +++ src/watcher/KEventWatcher.rs | 87 +++++++++++++++++++++++++++++------- 2 files changed, 76 insertions(+), 16 deletions(-) diff --git a/src/io/io_darwin.cpp b/src/io/io_darwin.cpp index 4f9ca1ed60e5..70403ffe34eb 100644 --- a/src/io/io_darwin.cpp +++ b/src/io/io_darwin.cpp @@ -102,6 +102,11 @@ extern "C" bool io_darwin_schedule_wakeup(mach_port_t waker) } } +extern "C" void io_darwin_close_machport(mach_port_t port) +{ + mach_port_deallocate(mach_task_self(), port); +} + #else // stub out these symbols diff --git a/src/watcher/KEventWatcher.rs b/src/watcher/KEventWatcher.rs index cc2a41d5158c..b57b13a1e940 100644 --- a/src/watcher/KEventWatcher.rs +++ b/src/watcher/KEventWatcher.rs @@ -5,14 +5,34 @@ use crate::watcher_impl::{Op, WatchEvent, Watcher}; pub(crate) type Platform = KEventWatcher; +// Darwin: `src/io/io_darwin.cpp` (same helpers `bun_io::waker::KEventWaker` +// uses). The non-Darwin stubs there are no-ops so the symbols exist +// everywhere the C++ link step runs, but we only call them on macOS. +#[cfg(target_os = "macos")] +unsafe extern "C" { + fn io_darwin_create_machport( + kq: i32, + buf: *mut core::ffi::c_void, + len: usize, + ) -> libc::mach_port_t; + safe fn io_darwin_schedule_wakeup(port: libc::mach_port_t) -> bool; + safe fn io_darwin_close_machport(port: libc::mach_port_t); +} + pub struct KEventWatcher { pub(crate) fd: Fd, + #[cfg(target_os = "macos")] + machport: libc::mach_port_t, + /// Receive buffer handed to `EVFILT_MACHPORT` via `kevent64_s.ext[0]`; + /// must outlive the registration (i.e. until `stop()`). + #[cfg(target_os = "macos")] + _machport_buf: Box<[u8]>, } const CHANGELIST_COUNT: usize = 128; -/// `ident` for the EVFILT_USER wakeup event registered in `new()` and -/// triggered by `wake()`. +/// FreeBSD has no mach ports; use the kqueue-native EVFILT_USER wakeup there. +#[cfg(target_os = "freebsd")] const WAKE_EVENT_IDENT: usize = 0x2307; impl KEventWatcher { @@ -21,15 +41,45 @@ impl KEventWatcher { if fd.native() == 0 { return Err(crate::Error::KQueueError); } - let mut ev: libc::kevent = bun_core::ffi::zeroed(); - ev.ident = WAKE_EVENT_IDENT; - ev.filter = libc::EVFILT_USER; - ev.flags = (libc::EV_ADD | libc::EV_CLEAR) as _; - let _ = bun_sys::kevent(fd, core::slice::from_ref(&ev), &mut [], None); - Ok(Self { fd }) + + #[cfg(target_os = "macos")] + { + let mut machport_buf = vec![0u8; 1024].into_boxed_slice(); + // SAFETY: fd is a live kqueue; buf is valid for `len` bytes and + // outlives the registration (owned by the returned Self). + let machport = unsafe { + io_darwin_create_machport( + fd.native(), + machport_buf.as_mut_ptr().cast::(), + machport_buf.len(), + ) + }; + // machport == 0 means creation failed; `wake()` degrades to a + // no-op and shutdown falls back to waiting for an fs event. + Ok(Self { + fd, + machport, + _machport_buf: machport_buf, + }) + } + + #[cfg(target_os = "freebsd")] + { + let mut ev: libc::kevent = bun_core::ffi::zeroed(); + ev.ident = WAKE_EVENT_IDENT; + ev.filter = libc::EVFILT_USER; + ev.flags = (libc::EV_ADD | libc::EV_CLEAR) as _; + let _ = bun_sys::kevent(fd, core::slice::from_ref(&ev), &mut [], None); + Ok(Self { fd }) + } } pub(crate) fn stop(&mut self) { + #[cfg(target_os = "macos")] + if self.machport != 0 { + io_darwin_close_machport(self.machport); + self.machport = 0; + } if self.fd.is_valid() { let _ = bun_sys::close(self.fd); self.fd = Fd::INVALID; @@ -39,14 +89,19 @@ impl KEventWatcher { /// Unblock the watcher thread's `kevent()` so it re-checks `running`. /// Called from `Watcher::shutdown` under `Watcher.mutex`. pub(crate) fn wake(&self) { - if !self.fd.is_valid() { - return; + #[cfg(target_os = "macos")] + if self.machport != 0 { + let _ = io_darwin_schedule_wakeup(self.machport); + } + + #[cfg(target_os = "freebsd")] + if self.fd.is_valid() { + let mut ev: libc::kevent = bun_core::ffi::zeroed(); + ev.ident = WAKE_EVENT_IDENT; + ev.filter = libc::EVFILT_USER; + ev.fflags = libc::NOTE_TRIGGER; + let _ = bun_sys::kevent(self.fd, core::slice::from_ref(&ev), &mut [], None); } - let mut ev: libc::kevent = bun_core::ffi::zeroed(); - ev.ident = WAKE_EVENT_IDENT; - ev.filter = libc::EVFILT_USER; - ev.fflags = libc::NOTE_TRIGGER; - let _ = bun_sys::kevent(self.fd, core::slice::from_ref(&ev), &mut [], None); } } @@ -94,7 +149,7 @@ pub(crate) fn watch_loop_cycle(this: &mut Watcher) -> bun_sys::Result<()> { let mut out_len: usize = 0; let mut prev_event: Option<&libc::kevent> = None; for event in changes { - // Only VNODE events map to watch items (filters out `wake()`'s EVFILT_USER). + // Only VNODE events map to watch items (filters out the wakeup event). if event.filter != libc::EVFILT_VNODE { continue; } From e459d2d9681fb1c34ea411218de22cfd4cb54d5b Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Wed, 29 Jul 2026 00:55:25 +0000 Subject: [PATCH 10/16] io_darwin: release the RECEIVE right in io_darwin_close_machport io_darwin_create_machport allocates a RECEIVE right and inserts a SEND right on the same name; mach_port_deallocate only drops the send right, so the port object survived. Drop the receive right too so the port is destroyed. KEventWatcher::stop is this function's first caller (the process-lifetime bun_io::waker never closes its port). --- src/io/io_darwin.cpp | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/src/io/io_darwin.cpp b/src/io/io_darwin.cpp index 70403ffe34eb..521d95772b8e 100644 --- a/src/io/io_darwin.cpp +++ b/src/io/io_darwin.cpp @@ -104,7 +104,11 @@ extern "C" bool io_darwin_schedule_wakeup(mach_port_t waker) extern "C" void io_darwin_close_machport(mach_port_t port) { - mach_port_deallocate(mach_task_self(), port); + mach_port_t self = mach_task_self(); + // io_darwin_create_machport allocates a RECEIVE right and inserts a SEND + // right on the same name; release both so the kernel port object is freed. + mach_port_deallocate(self, port); + mach_port_mod_refs(self, port, MACH_PORT_RIGHT_RECEIVE, -1); } #else From ff3d8be6026235319c3b6ce7043637589d494977 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 1 Oct 2026 12:04:21 +0000 Subject: [PATCH 11/16] Adapt to #39197: free the owned paths of a dropped Watcher The watcher thread now frees the Watcher when a dev server stops. MultiArrayList's drop is slab-only, so the owned file_path of every entry that is still in the watchlist leaked. Since #39197, CI's ASAN lane applies its exit-time leak check to test/bake/dev/css.test.ts, and that check reports the leak from append_directory_assume_capacity::. impl Drop for Watcher drops the elements of the watchlist, as #30644 does. KEventWatcher.rs: main deleted io_darwin_close_machport and its non-Darwin stub in #39581. The rebase restores the Darwin definition only, so the comment on the extern block no longer claims a stub for it. --- src/watcher/KEventWatcher.rs | 6 +++--- src/watcher/Watcher.rs | 7 +++++++ 2 files changed, 10 insertions(+), 3 deletions(-) diff --git a/src/watcher/KEventWatcher.rs b/src/watcher/KEventWatcher.rs index b57b13a1e940..3dfc90a56889 100644 --- a/src/watcher/KEventWatcher.rs +++ b/src/watcher/KEventWatcher.rs @@ -5,9 +5,9 @@ use crate::watcher_impl::{Op, WatchEvent, Watcher}; pub(crate) type Platform = KEventWatcher; -// Darwin: `src/io/io_darwin.cpp` (same helpers `bun_io::waker::KEventWaker` -// uses). The non-Darwin stubs there are no-ops so the symbols exist -// everywhere the C++ link step runs, but we only call them on macOS. +// Darwin: `src/io/io_darwin.cpp`. `bun_io::waker::KEventWaker` uses the first +// two; `io_darwin_close_machport` is for this watcher only and has no +// non-Darwin stub. #[cfg(target_os = "macos")] unsafe extern "C" { fn io_darwin_create_machport( diff --git a/src/watcher/Watcher.rs b/src/watcher/Watcher.rs index 30d85e09ceeb..7b6a316373f9 100644 --- a/src/watcher/Watcher.rs +++ b/src/watcher/Watcher.rs @@ -137,6 +137,13 @@ pub struct Watcher { pub thread_lock: ThreadLock, } +impl Drop for Watcher { + fn drop(&mut self) { + // `MultiArrayList` frees only its slab, not the owned `file_path`s. + self.watchlist.drop_elements(); + } +} + /// Context types passed to `Watcher::init` implement this trait. /// The default `on_watch_error` forwards to `on_error`. pub trait WatcherContext { From da1fc4a57a0f8d4d6e338e9707e612b39bb1f866 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Thu, 1 Oct 2026 14:44:50 +0000 Subject: [PATCH 12/16] watcher(macOS): use kevent64() for every call on the watcher's kqueue XNU ties a kqueue to the event struct of the first call made on it. io_darwin_create_machport registers the mach port with kevent64(), so the plain kevent() calls that add a watch and wait for events failed with EINVAL on macOS. Every dev server printed "EINVAL: Invalid argument: failed to watch files for hot-reloading (kevent)" on both darwin lanes. KEvent is kevent64_s on Darwin and kevent on FreeBSD. kevent_call() calls kevent64() on Darwin, as bun_io does for its loop, and kevent() on FreeBSD. --- src/watcher/KEventWatcher.rs | 53 ++++++++++++++++++++++++++++++++---- src/watcher/Watcher.rs | 7 +++-- 2 files changed, 52 insertions(+), 8 deletions(-) diff --git a/src/watcher/KEventWatcher.rs b/src/watcher/KEventWatcher.rs index 3dfc90a56889..cd995116304b 100644 --- a/src/watcher/KEventWatcher.rs +++ b/src/watcher/KEventWatcher.rs @@ -5,6 +5,49 @@ use crate::watcher_impl::{Op, WatchEvent, Watcher}; pub(crate) type Platform = KEventWatcher; +/// XNU ties a kqueue to the event struct of the first call made on it. On +/// Darwin that call is the `kevent64()` in `io_darwin_create_machport`, so +/// every later call on this kqueue uses `kevent64_s` too: a plain `kevent()` +/// there fails with EINVAL. +#[cfg(target_os = "macos")] +pub(crate) type KEvent = libc::kevent64_s; +#[cfg(target_os = "freebsd")] +pub(crate) type KEvent = libc::kevent; + +/// `kevent64()` on Darwin, `kevent()` on FreeBSD. Retries on EINTR. +pub(crate) fn kevent_call( + fd: Fd, + changelist: &[KEvent], + eventlist: &mut [KEvent], + timeout: Option<&libc::timespec>, +) -> bun_sys::Result { + #[cfg(target_os = "freebsd")] + { + bun_sys::kevent(fd, changelist, eventlist, timeout) + } + #[cfg(target_os = "macos")] + loop { + // SAFETY: fd is a valid kqueue; slices give exact (ptr,len); timeout + // is either null or a valid timespec. + let rc = unsafe { + libc::kevent64( + fd.native(), + changelist.as_ptr(), + changelist.len() as core::ffi::c_int, + eventlist.as_mut_ptr(), + eventlist.len() as core::ffi::c_int, + 0, + timeout.map_or(core::ptr::null(), std::ptr::from_ref), + ) + }; + match bun_sys::get_errno(rc) { + bun_sys::E::SUCCESS => return Ok(rc as usize), + bun_sys::E::EINTR => continue, + e => return Err(bun_sys::Error::from_code(e, bun_sys::Tag::kevent).with_fd(fd)), + } + } +} + // Darwin: `src/io/io_darwin.cpp`. `bun_io::waker::KEventWaker` uses the first // two; `io_darwin_close_machport` is for this watcher only and has no // non-Darwin stub. @@ -105,7 +148,7 @@ impl KEventWatcher { } } -fn watch_event_from_kevent(kevent: &libc::kevent) -> WatchEvent { +fn watch_event_from_kevent(kevent: &KEvent) -> WatchEvent { let mut op = Op::empty(); if (kevent.fflags & libc::NOTE_DELETE) > 0 { op |= Op::DELETE; @@ -131,9 +174,9 @@ pub(crate) fn watch_loop_cycle(this: &mut Watcher) -> bun_sys::Result<()> { let _flush = Output::flush_guard(); let fd = this.platform.fd; - let mut changelist: [libc::kevent; CHANGELIST_COUNT] = bun_core::ffi::zeroed(); + let mut changelist: [KEvent; CHANGELIST_COUNT] = bun_core::ffi::zeroed(); - let mut count = bun_sys::kevent(fd, &[], &mut changelist, None)?; + let mut count = kevent_call(fd, &[], &mut changelist, None)?; // Give the events more time to coalesce if count < CHANGELIST_COUNT / 2 { @@ -141,13 +184,13 @@ pub(crate) fn watch_loop_cycle(this: &mut Watcher) -> bun_sys::Result<()> { tv_sec: 0, tv_nsec: 100_000, }; // 0.0001 seconds - count += bun_sys::kevent(fd, &[], &mut changelist[count..], Some(&ts))?; + count += kevent_call(fd, &[], &mut changelist[count..], Some(&ts))?; } let changes = &changelist[..count]; let watchevents = &mut this.watch_events[..count]; let mut out_len: usize = 0; - let mut prev_event: Option<&libc::kevent> = None; + let mut prev_event: Option<&KEvent> = None; for event in changes { // Only VNODE events map to watch items (filters out the wakeup event). if event.filter != libc::EVFILT_VNODE { diff --git a/src/watcher/Watcher.rs b/src/watcher/Watcher.rs index 7b6a316373f9..fe532c70c7e2 100644 --- a/src/watcher/Watcher.rs +++ b/src/watcher/Watcher.rs @@ -511,8 +511,9 @@ impl Watcher { fd: Fd, watchlist_id: usize, ) { - use libc::{EV_ADD, EV_CLEAR, EV_ENABLE, EVFILT_VNODE, kevent as KEvent}; + use libc::{EV_ADD, EV_CLEAR, EV_ENABLE, EVFILT_VNODE}; use libc::{NOTE_DELETE, NOTE_RENAME, NOTE_WRITE}; + use platform::KEvent; // https://developer.apple.com/library/archive/documentation/System/Conceptual/ManPages_iPhoneOS/man2/kqueue.2.html let mut event: KEvent = bun_core::ffi::zeroed(); @@ -524,7 +525,7 @@ impl Watcher { event.fflags = (NOTE_WRITE | NOTE_RENAME | NOTE_DELETE) as _; // id - event.ident = usize::try_from(fd.native()).expect("int cast"); + event.ident = fd.native().try_into().expect("int cast"); // Store the index for fast filtering later event.udata = watchlist_id as _; @@ -533,7 +534,7 @@ impl Watcher { // Basically: // - We register the event here. // our while(true) loop above receives notification of changes to any of the events created here. - let _ = bun_sys::kevent(self.platform.fd, &[event], &mut [], None); + let _ = platform::kevent_call(self.platform.fd, &[event], &mut [], None); } fn append_file_assume_capacity( From 5dda89748e93d6d05ce98543f9cca04b06543e1c Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 2 Oct 2026 03:20:15 +0000 Subject: [PATCH 13/16] watcher: decide ownership under the mutex, cancel the Windows read, use bun_sys::kevent64 Watcher.rs: shutdown() takes the mutex before it reads watchloop_handle, and thread_body clears the flag under the same mutex when it hands the allocation back after a watch error. One release() runs platform.stop() and closes the watch fds on whichever side frees the Box. The JoinHandle field is gone: the thread is never joined. WindowsWatcher.rs: wake() posts an empty completion packet, which next() returns as "no events" so watch_loop re-checks running. read_pending says whether the kernel holds the buffer, so next() starts one read at a time and stop() cancels it with CancelIoEx and dequeues its packet before it closes the handles. KEventWatcher.rs: the kqueue calls go through bun_sys::kevent64 on macOS. A failed mach port registration is an error. An empty batch skips dispatch_file_updates. wake() after stop() is a no-op on every platform. The release test moves into test/bake/deinitialization.test.ts and no longer needs a LeakSanitizer exclusion. --- src/sys/lib.rs | 32 ++++++ src/sys/windows/mod.rs | 7 ++ src/watcher/KEventWatcher.rs | 102 +++++++---------- src/watcher/Watcher.rs | 109 ++++++++---------- src/watcher/WindowsWatcher.rs | 110 ++++++++++++------- test/bake/deinitialization.test.ts | 105 +++++++++++++++++- test/bake/dev-server-watcher-release.test.ts | 103 ----------------- test/no-validate-leaksan.txt | 3 - 8 files changed, 297 insertions(+), 274 deletions(-) delete mode 100644 test/bake/dev-server-watcher-release.test.ts diff --git a/src/sys/lib.rs b/src/sys/lib.rs index c89037345c8d..54f75e4fee89 100644 --- a/src/sys/lib.rs +++ b/src/sys/lib.rs @@ -7356,6 +7356,38 @@ pub fn kevent( } } +/// `kevent64()` — slice-wrapped Maybe form of [`kevent`]. Retries on EINTR. +/// XNU rejects `kevent()` on a kqueue that `kevent64()` has touched (EINVAL), +/// so a kqueue that carries an `EVFILT_MACHPORT` wakeup uses this everywhere. +#[cfg(target_os = "macos")] +pub fn kevent64( + fd: Fd, + changelist: &[libc::kevent64_s], + eventlist: &mut [libc::kevent64_s], + timeout: Option<&libc::timespec>, +) -> Maybe { + loop { + // SAFETY: fd is a valid kqueue; slices give exact (ptr,len); timeout + // is either null or a valid timespec. + let rc = unsafe { + libc::kevent64( + fd.native(), + changelist.as_ptr(), + changelist.len() as c_int, + eventlist.as_mut_ptr(), + eventlist.len() as c_int, + 0, + timeout.map_or(core::ptr::null(), std::ptr::from_ref), + ) + }; + match get_errno(rc) { + E::SUCCESS => return Ok(rc as usize), + E::EINTR => continue, + e => return Err(Error::from_code(e, Tag::kevent).with_fd(fd)), + } + } +} + // ── getFdPath ── /// Cached probe of `/proc/version` for "freebsd" diff --git a/src/sys/windows/mod.rs b/src/sys/windows/mod.rs index d2eac80990ea..ab619cd3de97 100644 --- a/src/sys/windows/mod.rs +++ b/src/sys/windows/mod.rs @@ -66,6 +66,13 @@ pub mod kernel32 { lpOverlapped: LPOVERLAPPED, lpCompletionRoutine: LPOVERLAPPED_COMPLETION_ROUTINE, ) -> BOOL; + pub fn CancelIoEx(hFile: HANDLE, lpOverlapped: LPOVERLAPPED) -> BOOL; + pub fn PostQueuedCompletionStatus( + CompletionPort: HANDLE, + dwNumberOfBytesTransferred: DWORD, + dwCompletionKey: ULONG_PTR, + lpOverlapped: LPOVERLAPPED, + ) -> BOOL; // safe: by-value `HANDLE` + `DWORD`; a bad handle yields // `WAIT_FAILED` + GetLastError, no UB. diff --git a/src/watcher/KEventWatcher.rs b/src/watcher/KEventWatcher.rs index cd995116304b..8b84a5bca4e8 100644 --- a/src/watcher/KEventWatcher.rs +++ b/src/watcher/KEventWatcher.rs @@ -5,52 +5,7 @@ use crate::watcher_impl::{Op, WatchEvent, Watcher}; pub(crate) type Platform = KEventWatcher; -/// XNU ties a kqueue to the event struct of the first call made on it. On -/// Darwin that call is the `kevent64()` in `io_darwin_create_machport`, so -/// every later call on this kqueue uses `kevent64_s` too: a plain `kevent()` -/// there fails with EINVAL. -#[cfg(target_os = "macos")] -pub(crate) type KEvent = libc::kevent64_s; -#[cfg(target_os = "freebsd")] -pub(crate) type KEvent = libc::kevent; - -/// `kevent64()` on Darwin, `kevent()` on FreeBSD. Retries on EINTR. -pub(crate) fn kevent_call( - fd: Fd, - changelist: &[KEvent], - eventlist: &mut [KEvent], - timeout: Option<&libc::timespec>, -) -> bun_sys::Result { - #[cfg(target_os = "freebsd")] - { - bun_sys::kevent(fd, changelist, eventlist, timeout) - } - #[cfg(target_os = "macos")] - loop { - // SAFETY: fd is a valid kqueue; slices give exact (ptr,len); timeout - // is either null or a valid timespec. - let rc = unsafe { - libc::kevent64( - fd.native(), - changelist.as_ptr(), - changelist.len() as core::ffi::c_int, - eventlist.as_mut_ptr(), - eventlist.len() as core::ffi::c_int, - 0, - timeout.map_or(core::ptr::null(), std::ptr::from_ref), - ) - }; - match bun_sys::get_errno(rc) { - bun_sys::E::SUCCESS => return Ok(rc as usize), - bun_sys::E::EINTR => continue, - e => return Err(bun_sys::Error::from_code(e, bun_sys::Tag::kevent).with_fd(fd)), - } - } -} - -// Darwin: `src/io/io_darwin.cpp`. `bun_io::waker::KEventWaker` uses the first -// two; `io_darwin_close_machport` is for this watcher only and has no -// non-Darwin stub. +// Defined in src/io/io_darwin.cpp. #[cfg(target_os = "macos")] unsafe extern "C" { fn io_darwin_create_machport( @@ -62,6 +17,17 @@ unsafe extern "C" { safe fn io_darwin_close_machport(port: libc::mach_port_t); } +/// `io_darwin_create_machport` uses `kevent64()`, and XNU rejects a plain +/// `kevent()` on that kqueue afterwards (EINVAL). +#[cfg(target_os = "macos")] +pub(crate) type KEvent = libc::kevent64_s; +#[cfg(target_os = "macos")] +pub(crate) use bun_sys::kevent64 as kevent; +#[cfg(target_os = "freebsd")] +pub(crate) type KEvent = libc::kevent; +#[cfg(target_os = "freebsd")] +pub(crate) use bun_sys::kevent; + pub struct KEventWatcher { pub(crate) fd: Fd, #[cfg(target_os = "macos")] @@ -78,6 +44,16 @@ const CHANGELIST_COUNT: usize = 128; #[cfg(target_os = "freebsd")] const WAKE_EVENT_IDENT: usize = 0x2307; +#[cfg(target_os = "freebsd")] +fn wake_event(flags: u16, fflags: u32) -> KEvent { + let mut ev: KEvent = bun_core::ffi::zeroed(); + ev.ident = WAKE_EVENT_IDENT; + ev.filter = libc::EVFILT_USER; + ev.flags = flags; + ev.fflags = fflags; + ev +} + impl KEventWatcher { pub(crate) fn new(_root: &[u8]) -> crate::Result { let fd = bun_sys::kqueue()?; @@ -97,8 +73,10 @@ impl KEventWatcher { machport_buf.len(), ) }; - // machport == 0 means creation failed; `wake()` degrades to a - // no-op and shutdown falls back to waiting for an fs event. + if machport == 0 { + let _ = bun_sys::close(fd); + return Err(crate::Error::KQueueError); + } Ok(Self { fd, machport, @@ -108,11 +86,11 @@ impl KEventWatcher { #[cfg(target_os = "freebsd")] { - let mut ev: libc::kevent = bun_core::ffi::zeroed(); - ev.ident = WAKE_EVENT_IDENT; - ev.filter = libc::EVFILT_USER; - ev.flags = (libc::EV_ADD | libc::EV_CLEAR) as _; - let _ = bun_sys::kevent(fd, core::slice::from_ref(&ev), &mut [], None); + let ev = wake_event(libc::EV_ADD | libc::EV_CLEAR, 0); + if let Err(err) = kevent(fd, core::slice::from_ref(&ev), &mut [], None) { + let _ = bun_sys::close(fd); + return Err(err.into()); + } Ok(Self { fd }) } } @@ -129,8 +107,8 @@ impl KEventWatcher { } } - /// Unblock the watcher thread's `kevent()` so it re-checks `running`. - /// Called from `Watcher::shutdown` under `Watcher.mutex`. + /// Unblock the watcher thread's kqueue wait so it re-checks `running`. + /// Runs under `Watcher.mutex`, like `stop()` on the hand-back path. pub(crate) fn wake(&self) { #[cfg(target_os = "macos")] if self.machport != 0 { @@ -139,11 +117,8 @@ impl KEventWatcher { #[cfg(target_os = "freebsd")] if self.fd.is_valid() { - let mut ev: libc::kevent = bun_core::ffi::zeroed(); - ev.ident = WAKE_EVENT_IDENT; - ev.filter = libc::EVFILT_USER; - ev.fflags = libc::NOTE_TRIGGER; - let _ = bun_sys::kevent(self.fd, core::slice::from_ref(&ev), &mut [], None); + let ev = wake_event(0, libc::NOTE_TRIGGER); + let _ = kevent(self.fd, core::slice::from_ref(&ev), &mut [], None); } } } @@ -176,7 +151,7 @@ pub(crate) fn watch_loop_cycle(this: &mut Watcher) -> bun_sys::Result<()> { let mut changelist: [KEvent; CHANGELIST_COUNT] = bun_core::ffi::zeroed(); - let mut count = kevent_call(fd, &[], &mut changelist, None)?; + let mut count = kevent(fd, &[], &mut changelist, None)?; // Give the events more time to coalesce if count < CHANGELIST_COUNT / 2 { @@ -184,7 +159,7 @@ pub(crate) fn watch_loop_cycle(this: &mut Watcher) -> bun_sys::Result<()> { tv_sec: 0, tv_nsec: 100_000, }; // 0.0001 seconds - count += kevent_call(fd, &[], &mut changelist[count..], Some(&ts))?; + count += kevent(fd, &[], &mut changelist[count..], Some(&ts))?; } let changes = &changelist[..count]; @@ -207,6 +182,9 @@ pub(crate) fn watch_loop_cycle(this: &mut Watcher) -> bun_sys::Result<()> { prev_event = Some(event); out_len += 1; } + if out_len == 0 { + return Ok(()); + } this.dispatch_file_updates(out_len, out_len); Ok(()) diff --git a/src/watcher/Watcher.rs b/src/watcher/Watcher.rs index fe532c70c7e2..fd91012663ba 100644 --- a/src/watcher/Watcher.rs +++ b/src/watcher/Watcher.rs @@ -109,12 +109,10 @@ pub struct Watcher { // Storing the `top_level_dir` slice directly avoids a forward-decl // dependency on the higher-tier `bun_resolver::fs::FileSystem` type. // allocator field dropped — global mimalloc (see §Allocators) - /// Whether `thread_main` is running. Written by the watcher thread, read - /// by `start`/`shutdown` on the main thread. The actual `ThreadId` value - /// was never read — only `is_some()`/`is_none()` — so this is a `bool`. + /// The watcher thread owns the allocation: set by `start()`, cleared + /// under `mutex` by `thread_body` when it hands the allocation back. pub(crate) watchloop_handle: bun_core::AtomicCell, pub(crate) cwd: &'static [u8], - pub(crate) thread: Option>, /// Main thread clears this in `shutdown`; watcher thread polls it in /// `watch_loop` and the platform `watch_loop_cycle`. pub(crate) running: bun_core::AtomicCell, @@ -206,7 +204,6 @@ impl Watcher { watch_events: vec![WatchEvent::default(); MAX_COUNT].into_boxed_slice(), changed_filepaths: [const { None }; MAX_COUNT], watchloop_handle: bun_core::AtomicCell::new(false), - thread: None, running: bun_core::AtomicCell::new(true), close_descriptors: bun_core::AtomicCell::new(false), evict_list: [0; MAX_EVICTION_COUNT], @@ -262,7 +259,9 @@ impl Watcher { std::thread::sleep(std::time::Duration::from_millis(10)); spawn().map_err(|_| first) }); - self.thread = Some(handle.map_err(|e| { + // The thread frees the Watcher itself and is never joined, so the + // handle is dropped (detached). + handle.map_err(|e| { self.watchloop_handle.store(false); // Windows: raw_os_error() is a Win32 GetLastError() code, so // route it through the u32 (Win32Error) mapper rather than @@ -278,7 +277,7 @@ impl Watcher { .map(bun_errno::from_errno) .unwrap_or(bun_errno::SystemErrno::EAGAIN); crate::Error::Sys(errno) - })?); + })?; Ok(()) } @@ -287,48 +286,41 @@ impl Watcher { // Per PORTING.md, `pub fn deinit` is never the public name; renamed to // `shutdown` (not `close(self)` because ownership may transfer to the // watcher thread instead of dropping here). - // TODO: ownership model — needs heap::take or an Arc to make this sound. /// # Safety /// `this` must be the unique heap pointer returned from `init()`; ownership /// transfers here on the no-thread path (the Box is reclaimed). pub unsafe fn shutdown(this: *mut Self, close_descriptors: bool) { - let free = { + { // SAFETY: caller passes the unique heap pointer returned from init(). - // Shared access suffices (atomics + mutex + column reads); the borrow + // Shared access suffices (atomics + mutex + `wake()`); the borrow // ends before the free below. let me = unsafe { &*this }; + me.mutex.lock(); if me.watchloop_handle.load() { - me.mutex.lock(); me.close_descriptors.store(close_descriptors); me.running.store(false); me.platform.wake(); me.mutex.unlock(); - // `*this` may be freed by the watcher thread any time after this - // unlock; `thread_main` lock/unlocks `mutex` before `heap::take`. - false - } else { - if close_descriptors && me.running.load() { - let fds = me.watchlist.items_fd(); - for &fd in fds { - if fd.is_valid() { - let _ = bun_sys::close(fd); - } - } - } - true + // The thread may free `*this` as soon as the mutex is released. + return; } - }; - if free { - // watchlist freed by Drop on Box - // SAFETY: this was heap-allocated by caller of init(); no borrow of it - // is live here. - let mut me = unsafe { bun_core::heap::take(this) }; - // A spawned thread runs `platform.stop()` itself in `thread_body`, - // also when it hands `*this` back after a watch error. - if me.thread.is_none() { - me.platform.stop(); + me.mutex.unlock(); + } + + // SAFETY: this was heap-allocated by caller of init(); no thread owns it + // and no borrow of it is live here. + unsafe { bun_core::heap::take(this) }.release(close_descriptors); + } + + /// Runs on whichever side frees the allocation. + fn release(&mut self, close_descriptors: bool) { + self.platform.stop(); + if close_descriptors { + for &fd in self.watchlist.items_fd() { + if fd.is_valid() { + let _ = bun_sys::close(fd); + } } - drop(me); } } @@ -362,43 +354,31 @@ impl Watcher { Ok(()) } + /// Returns `true` when ownership went back to the owner, which may free + /// `self` once the mutex is released. fn thread_body(&mut self) -> bool { - self.watchloop_handle.store(true); self.thread_lock.lock(); Output::Source::configure_named_thread(zstr!("File Watcher")); log!("Watcher started"); - let owner_still_alive = match self.watch_loop() { - Err(err) => { - self.watchloop_handle.store(false); - let running = self.running.load(); - if running { - (self.on_error)(self.ctx, err); - } - running - } - Ok(()) => false, - }; + let result = self.watch_loop(); - // Barrier: `shutdown()` holds `mutex` across `running.store(false)` - // and `platform.wake()`; `stop()` and `heap::take(this)` below - // must not run until `shutdown()` has unlocked. + // `shutdown()` decides under `mutex` which side frees the allocation. self.mutex.lock(); - self.mutex.unlock(); - - self.platform.stop(); - - // deinit and close descriptors if needed - if self.close_descriptors.load() { - let fds = self.watchlist.items_fd(); - for &fd in fds { - if fd.is_valid() { - let _ = bun_sys::close(fd); - } + if let Err(err) = result { + if self.running.load() { + (self.on_error)(self.ctx, err); + self.platform.stop(); + self.watchloop_handle.store(false); + self.mutex.unlock(); + return true; } } - owner_still_alive + self.mutex.unlock(); + + self.release(self.close_descriptors.load()); + false } pub fn flush_evictions(&mut self) { @@ -513,10 +493,9 @@ impl Watcher { ) { use libc::{EV_ADD, EV_CLEAR, EV_ENABLE, EVFILT_VNODE}; use libc::{NOTE_DELETE, NOTE_RENAME, NOTE_WRITE}; - use platform::KEvent; // https://developer.apple.com/library/archive/documentation/System/Conceptual/ManPages_iPhoneOS/man2/kqueue.2.html - let mut event: KEvent = bun_core::ffi::zeroed(); + let mut event: platform::KEvent = bun_core::ffi::zeroed(); event.flags = (EV_ADD | EV_CLEAR | EV_ENABLE) as _; // we want to know about the vnode @@ -534,7 +513,7 @@ impl Watcher { // Basically: // - We register the event here. // our while(true) loop above receives notification of changes to any of the events created here. - let _ = platform::kevent_call(self.platform.fd, &[event], &mut [], None); + let _ = platform::kevent(self.platform.fd, &[event], &mut [], None); } fn append_file_assume_capacity( diff --git a/src/watcher/WindowsWatcher.rs b/src/watcher/WindowsWatcher.rs index 3e9b60c4db56..09fe80215ae5 100644 --- a/src/watcher/WindowsWatcher.rs +++ b/src/watcher/WindowsWatcher.rs @@ -22,10 +22,8 @@ pub struct WindowsWatcher { pub(crate) watcher: DirWatcher, pub(crate) buf: PathBuffer, pub(crate) base_idx: usize, - /// Latched true once `next()` has armed a `ReadDirectoryChangesW` on - /// `self.watcher.overlapped`; while set, `stop()` must leak the handles - /// (the kernel's cancellation write would land in freed memory). - pub(crate) armed: bool, + /// The kernel owns `watcher.buf` and `watcher.overlapped` while this is set. + read_pending: bool, } impl Default for WindowsWatcher { @@ -39,7 +37,7 @@ impl Default for WindowsWatcher { }, buf: PathBuffer::ZEROED, base_idx: 0, - armed: false, + read_pending: false, } } } @@ -299,11 +297,14 @@ impl WindowsWatcher { /// wait until new events are available fn next(&mut self, timeout: Timeout) -> bun_sys::Result> { - if let Err(err) = self.watcher.prepare() { - bun_core::scoped_log!(watcher, "prepare() returned error"); - return Err(err); + // A poll that timed out left its read with the kernel. + if !self.read_pending { + if let Err(err) = self.watcher.prepare() { + bun_core::scoped_log!(watcher, "prepare() returned error"); + return Err(err); + } + self.read_pending = true; } - self.armed = true; let mut nbytes: w::DWORD = 0; let mut key: w::ULONG_PTR = 0; @@ -325,11 +326,9 @@ impl WindowsWatcher { if err == w::Win32Error::TIMEOUT || err == w::Win32Error(258) { return Ok(None); } else { - // GQCS returning FALSE with `*lpOverlapped != NULL` - // dequeued a failed-I/O completion; nothing remains - // outstanding on our OVERLAPPED in that case. - if overlapped == &mut self.watcher.overlapped as *mut w::OVERLAPPED { - self.armed = false; + // With an OVERLAPPED this dequeued the packet of a failed read. + if overlapped == &raw mut self.watcher.overlapped { + self.read_pending = false; } bun_core::scoped_log!(watcher, "GetQueuedCompletionStatus failed: {}", err.0); return Err(bun_sys::Error::from_win32(err, bun_sys::Tag::watch)); @@ -341,9 +340,7 @@ impl WindowsWatcher { if overlapped != &mut self.watcher.overlapped as *mut w::OVERLAPPED { continue; } - // Our completion was dequeued; nothing is pending on - // `overlapped` until the next successful `prepare()`. - self.armed = false; + self.read_pending = false; if nbytes == 0 { // ReadDirectoryChangesW internal change-buffer overflow — too many // events arrived between drain and re-arm. This is NOT a shutdown @@ -361,7 +358,7 @@ impl WindowsWatcher { if let Err(err) = self.watcher.prepare() { return Err(err); } - self.armed = true; + self.read_pending = true; continue; } return Ok(Some(EventIterator { @@ -370,37 +367,70 @@ impl WindowsWatcher { has_next: true, })); } else { - bun_core::scoped_log!( - watcher, - "GetQueuedCompletionStatus returned no overlapped event" - ); - return Err(bun_sys::Error { - errno: bun_sys::SystemErrno::EINVAL as _, - syscall: bun_sys::Tag::watch, - ..Default::default() - }); + // Posted by `wake()`; `watch_loop` re-checks `running`. + return Ok(None); } } } pub(crate) fn stop(&mut self) { - if self.armed { - // See `armed`. Proper fix: `CancelIoEx` + IOCP drain before - // `heap::take`; until then leak the two handles. - return; + if self.read_pending { + self.cancel_read(); + } + if self.watcher.dir_handle != w::INVALID_HANDLE_VALUE { + // SAFETY: opened in init(); cleared below so this runs at most once. + let _ = unsafe { w::CloseHandle(self.watcher.dir_handle) }; + self.watcher.dir_handle = w::INVALID_HANDLE_VALUE; } - // SAFETY: handles were opened in init() and are valid until stop() is called once. - unsafe { - w::CloseHandle(self.watcher.dir_handle); - w::CloseHandle(self.iocp); + if self.iocp != w::INVALID_HANDLE_VALUE { + // SAFETY: created in init(); cleared below so this runs at most once. + let _ = unsafe { w::CloseHandle(self.iocp) }; + self.iocp = w::INVALID_HANDLE_VALUE; + } + } + + /// Takes the buffer back from the kernel: cancels the outstanding read and + /// dequeues its packet. A cancelled or completed read always posts one. + fn cancel_read(&mut self) { + // SAFETY: dir_handle is the open directory handle from init() and + // `overlapped` is the OVERLAPPED of the outstanding read. + let _ = unsafe { + w::kernel32::CancelIoEx(self.watcher.dir_handle, &mut self.watcher.overlapped) + }; + let mut nbytes: w::DWORD = 0; + let mut key: w::ULONG_PTR = 0; + loop { + let mut overlapped: *mut w::OVERLAPPED = ptr::null_mut(); + // SAFETY: iocp is a valid IOCP handle; out-params are valid stack locals. + let rc = unsafe { + w::kernel32::GetQueuedCompletionStatus( + self.iocp, + &mut nbytes, + &mut key, + &mut overlapped, + w::INFINITE, + ) + }; + if overlapped == &raw mut self.watcher.overlapped { + self.read_pending = false; + return; + } + if rc == 0 && overlapped.is_null() { + // The port itself failed; no packet can arrive. + return; + } } } - /// No-op: `next()` keeps a `ReadDirectoryChangesW` pending on - /// `self.watcher.overlapped`, so freeing `self` after waking would race - /// the kernel's cancellation write. The thread stays parked in - /// `GetQueuedCompletionStatus` until process exit (see `armed`). - pub(crate) fn wake(&self) {} + /// Runs under `Watcher.mutex`, like `stop()` on the hand-back path. + pub(crate) fn wake(&self) { + if self.iocp == w::INVALID_HANDLE_VALUE { + return; + } + // SAFETY: iocp is a live port; `next()` reads a null OVERLAPPED as the wakeup. + let _ = + unsafe { w::kernel32::PostQueuedCompletionStatus(self.iocp, 0, 0, ptr::null_mut()) }; + } } #[repr(u32)] diff --git a/test/bake/deinitialization.test.ts b/test/bake/deinitialization.test.ts index d31f594f4cc2..8675e7b1f3f8 100644 --- a/test/bake/deinitialization.test.ts +++ b/test/bake/deinitialization.test.ts @@ -1,5 +1,5 @@ import { expect, test } from "bun:test"; -import { bunEnv, bunExe, tempDir } from "harness"; +import { bunEnv, bunExe, isLinux, tempDir } from "harness"; import path from "node:path"; test("dev server deinitializes itself", () => { @@ -50,3 +50,106 @@ test("dev server is deinitialized before its arena when listen fails", async () expect(stdout).toBe('{"code":"EADDRINUSE","deinits":1}\n'); expect(exitCode).toBe(0); }); + +// Each stopped `Bun.serve({ development: true })` must release its file +// watcher. The watcher thread used to stay parked in its blocking wait after +// `server.stop()`, so every disposed dev server kept one inotify instance (and +// one thread) until process exit. With a tight `fs.inotify.max_user_instances` +// budget this surfaced as `EMFILE while initializing file watcher`. +// The `/proc/self/{fd,status}` probes are Linux-specific. +test.skipIf(!isLinux)( + "dev server releases its file watcher on stop()", + async () => { + const fixture = /* ts */ ` + import { readdirSync, readlinkSync, readFileSync } from "node:fs"; + import html from "./index.html"; + + function scan() { + let inotify = 0; + for (const name of readdirSync("/proc/self/fd")) { + try { + if (readlinkSync("/proc/self/fd/" + name) === "anon_inode:inotify") inotify++; + } catch {} + } + const status = readFileSync("/proc/self/status", "utf8"); + const threads = Number(/^Threads:\\s+(\\d+)/m.exec(status)?.[1] ?? 0); + return { inotify, threads }; + } + + const ITER = 10; + + // warm-up: the first server initialises process-global state + { + const s = Bun.serve({ port: 0, development: true, static: { "/": html }, fetch: () => new Response("") }); + await (await fetch(s.url)).text(); + s.stop(true); + } + // wait for the warm-up watcher to release so it isn't counted + for (let i = 0; i < 40 && scan().inotify > 0; i++) { + Bun.gc(true); + await Bun.sleep(50); + } + const before = scan(); + + for (let i = 0; i < ITER; i++) { + const s = Bun.serve({ port: 0, development: true, static: { "/": html }, fetch: () => new Response("") }); + await (await fetch(s.url)).text(); + s.stop(true); + } + + // poll until the watcher threads have closed their inotify instances + // (Threads: is not a reliable gate; JSC may spawn a collector thread) + for (let i = 0; i < 40 && scan().inotify > before.inotify; i++) { + Bun.gc(true); + await Bun.sleep(50); + } + + const after = scan(); + console.log(JSON.stringify({ + iterations: ITER, + inotifyDelta: after.inotify - before.inotify, + threadDelta: after.threads - before.threads, + })); + `; + + using dir = tempDir("dev-server-watcher-release", { + "index.html": "hi", + "fixture.ts": fixture, + }); + + await using proc = Bun.spawn({ + cmd: [bunExe(), "run", "fixture.ts"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + + const line = stdout + .split("\n") + .reverse() + .find(l => l.startsWith("{")); + if (!line) { + throw new Error(`no JSON summary in stdout.\nstdout:\n${stdout}\nstderr:\n${stderr}`); + } + const { inotifyDelta, threadDelta } = JSON.parse(line); + + expect(stderr).not.toContain("error:"); + + // Without the fix every iteration leaks one inotify instance + // (inotifyDelta == iterations). With the fix all of them are released. + // `threadDelta` is reported for diagnostics only: `Threads:` also counts + // JSC/bundler threads and can transiently read high right after `stop()` + // has closed the inotify fd but before the watcher thread has exited. + expect({ inotifyDelta, threadDelta }).toEqual({ + inotifyDelta: 0, + threadDelta: expect.any(Number), + }); + + expect(exitCode).toBe(0); + // 11 dev-server start/stop cycles under ASAN; without the fix the fixture's + // two 2s release polls also run to completion. + }, + 30_000, +); diff --git a/test/bake/dev-server-watcher-release.test.ts b/test/bake/dev-server-watcher-release.test.ts deleted file mode 100644 index 1ccda628c535..000000000000 --- a/test/bake/dev-server-watcher-release.test.ts +++ /dev/null @@ -1,103 +0,0 @@ -// Each stopped `Bun.serve({ development: true })` must release its file -// watcher. On Linux the watcher thread was parked in a blocking `read()` on -// the inotify fd after `server.stop()`, so every disposed dev server leaked -// one inotify instance (and one thread) until process exit. In a container -// with a tight `fs.inotify.max_user_instances` budget this surfaced as -// `EMFILE while initializing file watcher for development server`. - -import { expect, test } from "bun:test"; -import { bunEnv, bunExe, isLinux, tempDir } from "harness"; - -// The `/proc/self/{fd,status}` probes are Linux-specific; on macOS the kqueue -// leak is silent (no observable assertion here) and on Windows `wake()` is an -// intentional no-op (see `src/watcher/WindowsWatcher.rs`). -test.skipIf(!isLinux)("dev server releases its file watcher on stop()", async () => { - const fixture = /* ts */ ` - import { readdirSync, readlinkSync, readFileSync } from "node:fs"; - import html from "./index.html"; - - function scan() { - let inotify = 0; - for (const name of readdirSync("/proc/self/fd")) { - try { - if (readlinkSync("/proc/self/fd/" + name) === "anon_inode:inotify") inotify++; - } catch {} - } - const status = readFileSync("/proc/self/status", "utf8"); - const threads = Number(/^Threads:\\s+(\\d+)/m.exec(status)?.[1] ?? 0); - return { inotify, threads }; - } - - const ITER = 10; - - // warm-up: the first server initialises process-global state - { - const s = Bun.serve({ port: 0, development: true, static: { "/": html }, fetch: () => new Response("") }); - await (await fetch(s.url)).text(); - s.stop(true); - } - // wait for the warm-up watcher to release so it isn't counted - for (let i = 0; i < 40 && scan().inotify > 0; i++) { - Bun.gc(true); - await Bun.sleep(50); - } - const before = scan(); - - for (let i = 0; i < ITER; i++) { - const s = Bun.serve({ port: 0, development: true, static: { "/": html }, fetch: () => new Response("") }); - await (await fetch(s.url)).text(); - s.stop(true); - } - - // poll until the watcher threads have closed their inotify instances - // (Threads: is not a reliable gate; JSC may spawn a collector thread) - for (let i = 0; i < 40 && scan().inotify > before.inotify; i++) { - Bun.gc(true); - await Bun.sleep(50); - } - - const after = scan(); - console.log(JSON.stringify({ - iterations: ITER, - inotifyDelta: after.inotify - before.inotify, - threadDelta: after.threads - before.threads, - })); - `; - - using dir = tempDir("dev-server-watcher-release", { - "index.html": "hi", - "fixture.ts": fixture, - }); - - await using proc = Bun.spawn({ - cmd: [bunExe(), "run", "fixture.ts"], - env: bunEnv, - cwd: String(dir), - stdout: "pipe", - stderr: "pipe", - }); - const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); - - const line = stdout - .split("\n") - .reverse() - .find(l => l.startsWith("{")); - if (!line) { - throw new Error(`no JSON summary in stdout.\nstdout:\n${stdout}\nstderr:\n${stderr}`); - } - const { inotifyDelta, threadDelta } = JSON.parse(line); - - expect(stderr).not.toContain("error:"); - - // Without the fix every iteration leaks one inotify instance - // (inotifyDelta == iterations). With the fix all of them are released. - // `threadDelta` is reported for diagnostics only: `Threads:` also counts - // JSC/bundler threads and can transiently read high right after `stop()` - // has closed the inotify fd but before the watcher thread has exited. - expect({ inotifyDelta, threadDelta }).toEqual({ - inotifyDelta: 0, - threadDelta: expect.any(Number), - }); - - expect(exitCode).toBe(0); -}); diff --git a/test/no-validate-leaksan.txt b/test/no-validate-leaksan.txt index cdc4780e6d63..32a677be1b1f 100644 --- a/test/no-validate-leaksan.txt +++ b/test/no-validate-leaksan.txt @@ -373,9 +373,6 @@ test/bake/dev/react-spa.test.ts test/bake/dev/sourcemap.test.ts test/bake/dev/ssg-pages-router.test.ts test/bake/dev/deinitialization.test.ts -# NewServer <-> HTMLBundle::Route back-pointer cycle survives server.stop(); -# each disposed dev server leaks its ServerConfig/NewServer init allocations. -test/bake/dev-server-watcher-release.test.ts # Need to terminate HTTP thread. test/js/bun/test/parallel/test-http-should-not-accept-untrusted-certificates.ts From 44f7e4e48168a0999d831130a020024b0ffbebc8 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 2 Oct 2026 04:04:56 +0000 Subject: [PATCH 14/16] watcher: shorten the comments on the wake and ownership paths --- src/sys/lib.rs | 3 +-- src/watcher/INotifyWatcher.rs | 6 ++---- src/watcher/KEventWatcher.rs | 9 +++------ src/watcher/Watcher.rs | 9 +++------ src/watcher/WindowsWatcher.rs | 5 ++--- 5 files changed, 11 insertions(+), 21 deletions(-) diff --git a/src/sys/lib.rs b/src/sys/lib.rs index 54f75e4fee89..bbe87d8ea10b 100644 --- a/src/sys/lib.rs +++ b/src/sys/lib.rs @@ -7357,8 +7357,7 @@ pub fn kevent( } /// `kevent64()` — slice-wrapped Maybe form of [`kevent`]. Retries on EINTR. -/// XNU rejects `kevent()` on a kqueue that `kevent64()` has touched (EINVAL), -/// so a kqueue that carries an `EVFILT_MACHPORT` wakeup uses this everywhere. +/// XNU allows one kevent flavor per kqueue, so a kqueue that saw `kevent64()` once uses this for every call. #[cfg(target_os = "macos")] pub fn kevent64( fd: Fd, diff --git a/src/watcher/INotifyWatcher.rs b/src/watcher/INotifyWatcher.rs index a0b98c79809c..ae19b80b64ce 100644 --- a/src/watcher/INotifyWatcher.rs +++ b/src/watcher/INotifyWatcher.rs @@ -39,8 +39,7 @@ pub(crate) type Platform = INotifyWatcher; pub struct INotifyWatcher { pub(crate) fd: Fd, - /// eventfd written by `wake()`; `read()` ppolls on `[fd, wake_fd]` so - /// `Watcher::shutdown` can unpark the thread. + /// eventfd that `wake()` writes; `read()` polls it next to `fd`. pub(crate) wake_fd: Fd, pub(crate) loaded: bool, @@ -410,8 +409,7 @@ impl INotifyWatcher { } } - /// Unblock the watcher thread's `ppoll()` so it re-checks `running`. - /// Called from `Watcher::shutdown` under `Watcher.mutex`. + /// Unblocks the `ppoll()` so the thread re-checks `running`. Runs under `Watcher.mutex`. pub(crate) fn wake(&self) { if self.wake_fd == Fd::INVALID { return; diff --git a/src/watcher/KEventWatcher.rs b/src/watcher/KEventWatcher.rs index 8b84a5bca4e8..b934d039f261 100644 --- a/src/watcher/KEventWatcher.rs +++ b/src/watcher/KEventWatcher.rs @@ -17,8 +17,7 @@ unsafe extern "C" { safe fn io_darwin_close_machport(port: libc::mach_port_t); } -/// `io_darwin_create_machport` uses `kevent64()`, and XNU rejects a plain -/// `kevent()` on that kqueue afterwards (EINVAL). +/// XNU allows one kevent flavor per kqueue, and the mach port is registered with `kevent64()`. #[cfg(target_os = "macos")] pub(crate) type KEvent = libc::kevent64_s; #[cfg(target_os = "macos")] @@ -32,8 +31,7 @@ pub struct KEventWatcher { pub(crate) fd: Fd, #[cfg(target_os = "macos")] machport: libc::mach_port_t, - /// Receive buffer handed to `EVFILT_MACHPORT` via `kevent64_s.ext[0]`; - /// must outlive the registration (i.e. until `stop()`). + /// Receive buffer of the `EVFILT_MACHPORT` registration; lives until `stop()`. #[cfg(target_os = "macos")] _machport_buf: Box<[u8]>, } @@ -107,8 +105,7 @@ impl KEventWatcher { } } - /// Unblock the watcher thread's kqueue wait so it re-checks `running`. - /// Runs under `Watcher.mutex`, like `stop()` on the hand-back path. + /// Unblocks the kqueue wait so the thread re-checks `running`. Runs under `Watcher.mutex`. pub(crate) fn wake(&self) { #[cfg(target_os = "macos")] if self.machport != 0 { diff --git a/src/watcher/Watcher.rs b/src/watcher/Watcher.rs index fd91012663ba..fdc7e316a5f8 100644 --- a/src/watcher/Watcher.rs +++ b/src/watcher/Watcher.rs @@ -109,8 +109,7 @@ pub struct Watcher { // Storing the `top_level_dir` slice directly avoids a forward-decl // dependency on the higher-tier `bun_resolver::fs::FileSystem` type. // allocator field dropped — global mimalloc (see §Allocators) - /// The watcher thread owns the allocation: set by `start()`, cleared - /// under `mutex` by `thread_body` when it hands the allocation back. + /// The thread owns the allocation: set by `start()`, cleared under `mutex` on hand-back. pub(crate) watchloop_handle: bun_core::AtomicCell, pub(crate) cwd: &'static [u8], /// Main thread clears this in `shutdown`; watcher thread polls it in @@ -259,8 +258,7 @@ impl Watcher { std::thread::sleep(std::time::Duration::from_millis(10)); spawn().map_err(|_| first) }); - // The thread frees the Watcher itself and is never joined, so the - // handle is dropped (detached). + // Never joined: the thread frees the Watcher itself. handle.map_err(|e| { self.watchloop_handle.store(false); // Windows: raw_os_error() is a Win32 GetLastError() code, so @@ -354,8 +352,7 @@ impl Watcher { Ok(()) } - /// Returns `true` when ownership went back to the owner, which may free - /// `self` once the mutex is released. + /// Returns `true` when the owner got the allocation back and may free it once unlocked. fn thread_body(&mut self) -> bool { self.thread_lock.lock(); Output::Source::configure_named_thread(zstr!("File Watcher")); diff --git a/src/watcher/WindowsWatcher.rs b/src/watcher/WindowsWatcher.rs index 09fe80215ae5..15b062d241ab 100644 --- a/src/watcher/WindowsWatcher.rs +++ b/src/watcher/WindowsWatcher.rs @@ -389,8 +389,7 @@ impl WindowsWatcher { } } - /// Takes the buffer back from the kernel: cancels the outstanding read and - /// dequeues its packet. A cancelled or completed read always posts one. + /// Cancels the outstanding read and dequeues its packet, which a cancelled read always posts. fn cancel_read(&mut self) { // SAFETY: dir_handle is the open directory handle from init() and // `overlapped` is the OVERLAPPED of the outstanding read. @@ -422,7 +421,7 @@ impl WindowsWatcher { } } - /// Runs under `Watcher.mutex`, like `stop()` on the hand-back path. + /// Runs under `Watcher.mutex`, so `stop()` cannot close `iocp` underneath it. pub(crate) fn wake(&self) { if self.iocp == w::INVALID_HANDLE_VALUE { return; From 098c6dc97dc931129a2b9e9e9003785f91931211 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 2 Oct 2026 09:14:48 +0000 Subject: [PATCH 15/16] test: probe the watcher release in the test process and wait for the condition --- test/bake/deinitialization.test.ts | 131 +++++++---------------------- 1 file changed, 32 insertions(+), 99 deletions(-) diff --git a/test/bake/deinitialization.test.ts b/test/bake/deinitialization.test.ts index 8675e7b1f3f8..1eb7dcfcd1fe 100644 --- a/test/bake/deinitialization.test.ts +++ b/test/bake/deinitialization.test.ts @@ -1,5 +1,6 @@ import { expect, test } from "bun:test"; import { bunEnv, bunExe, isLinux, tempDir } from "harness"; +import { readdirSync, readlinkSync } from "node:fs"; import path from "node:path"; test("dev server deinitializes itself", () => { @@ -53,103 +54,35 @@ test("dev server is deinitialized before its arena when listen fails", async () // Each stopped `Bun.serve({ development: true })` must release its file // watcher. The watcher thread used to stay parked in its blocking wait after -// `server.stop()`, so every disposed dev server kept one inotify instance (and -// one thread) until process exit. With a tight `fs.inotify.max_user_instances` -// budget this surfaced as `EMFILE while initializing file watcher`. -// The `/proc/self/{fd,status}` probes are Linux-specific. -test.skipIf(!isLinux)( - "dev server releases its file watcher on stop()", - async () => { - const fixture = /* ts */ ` - import { readdirSync, readlinkSync, readFileSync } from "node:fs"; - import html from "./index.html"; - - function scan() { - let inotify = 0; - for (const name of readdirSync("/proc/self/fd")) { - try { - if (readlinkSync("/proc/self/fd/" + name) === "anon_inode:inotify") inotify++; - } catch {} - } - const status = readFileSync("/proc/self/status", "utf8"); - const threads = Number(/^Threads:\\s+(\\d+)/m.exec(status)?.[1] ?? 0); - return { inotify, threads }; - } - - const ITER = 10; - - // warm-up: the first server initialises process-global state - { - const s = Bun.serve({ port: 0, development: true, static: { "/": html }, fetch: () => new Response("") }); - await (await fetch(s.url)).text(); - s.stop(true); - } - // wait for the warm-up watcher to release so it isn't counted - for (let i = 0; i < 40 && scan().inotify > 0; i++) { - Bun.gc(true); - await Bun.sleep(50); - } - const before = scan(); - - for (let i = 0; i < ITER; i++) { - const s = Bun.serve({ port: 0, development: true, static: { "/": html }, fetch: () => new Response("") }); - await (await fetch(s.url)).text(); - s.stop(true); - } - - // poll until the watcher threads have closed their inotify instances - // (Threads: is not a reliable gate; JSC may spawn a collector thread) - for (let i = 0; i < 40 && scan().inotify > before.inotify; i++) { - Bun.gc(true); - await Bun.sleep(50); - } - - const after = scan(); - console.log(JSON.stringify({ - iterations: ITER, - inotifyDelta: after.inotify - before.inotify, - threadDelta: after.threads - before.threads, - })); - `; - - using dir = tempDir("dev-server-watcher-release", { - "index.html": "hi", - "fixture.ts": fixture, - }); - - await using proc = Bun.spawn({ - cmd: [bunExe(), "run", "fixture.ts"], - env: bunEnv, - cwd: String(dir), - stdout: "pipe", - stderr: "pipe", - }); - const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); - - const line = stdout - .split("\n") - .reverse() - .find(l => l.startsWith("{")); - if (!line) { - throw new Error(`no JSON summary in stdout.\nstdout:\n${stdout}\nstderr:\n${stderr}`); +// `server.stop()`, so every disposed dev server kept one inotify instance until +// process exit, and a tight `fs.inotify.max_user_instances` budget ended in +// `EMFILE while initializing file watcher`. The `/proc/self/fd` probe is Linux-specific. +test.skipIf(!isLinux)("dev server releases its file watcher on stop()", async () => { + using dir = tempDir("dev-server-watcher-release", { + "index.html": "hi", + }); + const { default: html } = await import(path.join(String(dir), "index.html")); + const inotifyInstances = () => { + let n = 0; + for (const name of readdirSync("/proc/self/fd")) { + try { + if (readlinkSync("/proc/self/fd/" + name) === "anon_inode:inotify") n++; + } catch {} } - const { inotifyDelta, threadDelta } = JSON.parse(line); - - expect(stderr).not.toContain("error:"); - - // Without the fix every iteration leaks one inotify instance - // (inotifyDelta == iterations). With the fix all of them are released. - // `threadDelta` is reported for diagnostics only: `Threads:` also counts - // JSC/bundler threads and can transiently read high right after `stop()` - // has closed the inotify fd but before the watcher thread has exited. - expect({ inotifyDelta, threadDelta }).toEqual({ - inotifyDelta: 0, - threadDelta: expect.any(Number), - }); - - expect(exitCode).toBe(0); - // 11 dev-server start/stop cycles under ASAN; without the fix the fixture's - // two 2s release polls also run to completion. - }, - 30_000, -); + return n; + }; + + const before = inotifyInstances(); + for (let i = 0; i < 2; i++) { + const server = Bun.serve({ port: 0, development: true, routes: { "/": html }, fetch: () => new Response("") }); + await (await fetch(server.url)).text(); + server.stop(true); + } + + // The watcher thread closes its inotify fd once stop() wakes it. Without the + // fix the thread never wakes and this never ends. + while (inotifyInstances() > before) { + await Bun.sleep(10); + } + expect(inotifyInstances()).toBe(before); +}); From 5c88ac079b9156e1ca342484da853e12b73c081a Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 2 Oct 2026 13:04:05 +0000 Subject: [PATCH 16/16] test: dispose each dev server with using, poll with a deadline --- test/bake/deinitialization.test.ts | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/test/bake/deinitialization.test.ts b/test/bake/deinitialization.test.ts index 1eb7dcfcd1fe..1e5702bd24d1 100644 --- a/test/bake/deinitialization.test.ts +++ b/test/bake/deinitialization.test.ts @@ -74,14 +74,14 @@ test.skipIf(!isLinux)("dev server releases its file watcher on stop()", async () const before = inotifyInstances(); for (let i = 0; i < 2; i++) { - const server = Bun.serve({ port: 0, development: true, routes: { "/": html }, fetch: () => new Response("") }); + using server = Bun.serve({ port: 0, development: true, routes: { "/": html }, fetch: () => new Response("") }); await (await fetch(server.url)).text(); - server.stop(true); } // The watcher thread closes its inotify fd once stop() wakes it. Without the - // fix the thread never wakes and this never ends. - while (inotifyInstances() > before) { + // fix the thread never wakes and every instance stays open. + const deadline = Date.now() + 3000; + while (inotifyInstances() > before && Date.now() < deadline) { await Bun.sleep(10); } expect(inotifyInstances()).toBe(before);