Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions src/io/io_darwin.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,15 @@ extern "C" bool io_darwin_schedule_wakeup(mach_port_t waker)
}
}

extern "C" void io_darwin_close_machport(mach_port_t 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.
Comment thread
robobun marked this conversation as resolved.
mach_port_deallocate(self, port);
mach_port_mod_refs(self, port, MACH_PORT_RIGHT_RECEIVE, -1);
}

#else

// stub out these symbols
Expand Down
31 changes: 31 additions & 0 deletions src/sys/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7356,6 +7356,37 @@
}
}

/// `kevent64()` — slice-wrapped Maybe form of [`kevent`]. Retries on EINTR.
/// XNU allows one kevent flavor per kqueue, so a kqueue that saw `kevent64()` once uses this for every call.
Comment thread
robobun marked this conversation as resolved.
#[cfg(target_os = "macos")]
pub fn kevent64(
fd: Fd,
changelist: &[libc::kevent64_s],
eventlist: &mut [libc::kevent64_s],
timeout: Option<&libc::timespec>,
) -> Maybe<usize> {
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) {

Check warning on line 7382 in src/sys/lib.rs

View workflow job for this annotation

GitHub Actions / mordant

this `match` on `bun_errno::SystemErrno` repeats an earlier one arm for arm. The mapping exists twice, and a change to one copy will miss the other
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"
Expand Down
7 changes: 7 additions & 0 deletions src/sys/windows/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
88 changes: 69 additions & 19 deletions src/watcher/INotifyWatcher.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -41,6 +39,8 @@ pub(crate) type Platform = INotifyWatcher;

pub struct INotifyWatcher {
pub(crate) fd: Fd,
/// eventfd that `wake()` writes; `read()` polls it next to `fd`.
pub(crate) wake_fd: Fd,
pub(crate) loaded: bool,

// Avoid statically allocating because it increases the binary size.
Expand All @@ -53,7 +53,6 @@ pub struct INotifyWatcher {
/// see `test-fs-watch-recursive-linux-parallel-remove.js`
read_ptr: Option<ReadPtr>,

pub(crate) watch_count: AtomicU32,
/// nanoseconds
pub(crate) coalesce_interval: isize,
}
Expand All @@ -62,11 +61,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,
}
}
Expand Down Expand Up @@ -129,32 +128,26 @@ impl INotifyWatcher {
pub(crate) fn watch_path(&mut self, pathname: &ZStr) -> bun_sys::Result<EventListIndex> {
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.
let rc = unsafe {
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<EventListIndex> {
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
Expand All @@ -168,18 +161,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<Self> {
Expand All @@ -193,9 +182,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()
Expand All @@ -222,12 +219,53 @@ 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 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.
Comment thread
robobun marked this conversation as resolved.
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: let `watch_loop` re-check `running`.
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 {
Expand Down Expand Up @@ -365,6 +403,18 @@ 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;
}
}

/// 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;
}
let _ = bun_sys::write(self.wake_fd, &1u64.to_ne_bytes());
}
}

Expand Down
Loading
Loading