#[cfg(unix)]
use core::ffi::c_int;
use core::ffi::c_void;
use core::fmt;
#[cfg(unix)]
use core::ptr;
#[cfg(not(windows))]
use bun_sys::{self as sys, Fd};
use bun_uws_sys::Loop as UwsLoop;
pub type Loop = UwsLoop;
#[cfg(not(windows))]
#[inline]
fn loop_add_active(loop_: &mut Loop, value: u32) {
loop_.active = loop_.active.saturating_add(value);
}
#[cfg(not(windows))]
#[inline]
fn loop_sub_active(loop_: &mut Loop, value: u32) {
loop_.active = loop_.active.saturating_sub(value);
}
bun_core::declare_scope!(KeepAlive, visible);
#[cfg(not(windows))]
macro_rules! syslog {
($($arg:tt)*) => {{ let _ = ::core::format_args!($($arg)*); }};
}
#[cfg(not(windows))]
#[inline]
fn errno_sys<R>(rc: R, syscall: sys::Tag) -> Option<sys::Result<()>>
where
R: sys::GetErrno,
{
match sys::get_errno(rc) {
sys::E::SUCCESS => None,
e => Some(sys::Result::Err(sys::Error::from_code(e, syscall))),
}
}
pub use crate::{EventLoopCtx, EventLoopCtxKind, EventLoopKind, OpaqueCallback};
unsafe extern "Rust" {
safe fn __bun_get_vm_ctx(kind: AllocatorType) -> EventLoopCtx;
}
#[derive(Copy, Clone, Eq, PartialEq)]
pub enum FileType {
File,
Pipe,
NonblockingPipe,
Socket,
}
impl FileType {
pub fn is_pollable(self) -> bool {
matches!(
self,
FileType::Pipe | FileType::NonblockingPipe | FileType::Socket
)
}
pub fn is_blocking(self) -> bool {
self == FileType::Pipe
}
}
#[inline]
pub fn get_vm_ctx(kind: AllocatorType) -> EventLoopCtx {
__bun_get_vm_ctx(kind)
}
#[inline]
pub fn js_vm_ctx() -> EventLoopCtx {
get_vm_ctx(AllocatorType::Js)
}
#[cfg(all(target_os = "macos", debug_assertions))]
type KQueueGenerationNumber = usize;
#[cfg(all(unix, not(all(target_os = "macos", debug_assertions))))]
type KQueueGenerationNumber = u8;
#[cfg(all(target_os = "macos", debug_assertions))]
static MAX_GENERATION_NUMBER: core::sync::atomic::AtomicUsize =
core::sync::atomic::AtomicUsize::new(0);
#[cfg(target_os = "macos")]
type KQueueEvent = bun_sys::darwin::kevent64_s;
#[cfg(target_os = "freebsd")]
type KQueueEvent = bun_sys::freebsd::Kevent;
#[cfg(target_os = "freebsd")]
#[inline]
fn make_kevent(
ident: usize,
filter: i16,
flags: u16,
fflags: u32,
udata: *mut core::ffi::c_void,
) -> KQueueEvent {
let mut ev: KQueueEvent = bun_core::ffi::zeroed();
ev.ident = ident;
ev.filter = filter;
ev.flags = flags;
ev.fflags = fflags;
ev.data = 0;
ev.udata = udata;
ev
}
#[cfg(any(target_os = "macos", target_os = "freebsd"))]
const EV_EOF: u16 = 0x8000;
#[repr(u8)]
#[derive(Copy, Clone, PartialEq, Eq)]
pub enum PollTag {
Null = 0,
FileSink,
StaticPipeWriter,
ShellStaticPipeWriter,
SecurityScanStaticPipeWriter,
BufferedReader,
DnsResolver,
GetAddrInfoRequest,
Request,
Process,
ShellBufferedWriter,
TerminalPoll,
ParentDeathWatchdog,
LifecycleScriptSubprocessOutputReader,
}
pub mod poll_tag {
use super::PollTag;
pub const NULL: PollTag = PollTag::Null;
pub const FILE_SINK: PollTag = PollTag::FileSink;
pub const STATIC_PIPE_WRITER: PollTag = PollTag::StaticPipeWriter;
pub const SHELL_STATIC_PIPE_WRITER: PollTag = PollTag::ShellStaticPipeWriter;
pub const SECURITY_SCAN_STATIC_PIPE_WRITER: PollTag = PollTag::SecurityScanStaticPipeWriter;
pub const BUFFERED_READER: PollTag = PollTag::BufferedReader;
pub const DNS_RESOLVER: PollTag = PollTag::DnsResolver;
pub const GET_ADDR_INFO_REQUEST: PollTag = PollTag::GetAddrInfoRequest;
pub const REQUEST: PollTag = PollTag::Request;
pub const PROCESS: PollTag = PollTag::Process;
pub const SHELL_BUFFERED_WRITER: PollTag = PollTag::ShellBufferedWriter;
pub const TERMINAL_POLL: PollTag = PollTag::TerminalPoll;
pub const PARENT_DEATH_WATCHDOG: PollTag = PollTag::ParentDeathWatchdog;
pub const LIFECYCLE_SCRIPT_SUBPROCESS_OUTPUT_READER: PollTag =
PollTag::LifecycleScriptSubprocessOutputReader;
}
#[derive(Copy, Clone)]
pub struct Owner {
pub tag: PollTag,
pub ptr: *mut (),
}
impl Owner {
pub const NULL: Owner = Owner {
tag: PollTag::Null,
ptr: core::ptr::null_mut(),
};
#[inline]
pub const fn new(tag: PollTag, ptr: *mut ()) -> Owner {
Owner { tag, ptr }
}
#[inline]
pub fn is_null(&self) -> bool {
self.ptr.is_null()
}
#[inline]
pub fn clear(&mut self) {
*self = Self::NULL;
}
#[inline]
pub fn tag(&self) -> PollTag {
self.tag
}
}
unsafe extern "Rust" {
fn __bun_run_file_poll(poll: *mut crate::FilePoll, size_or_offset: i64);
}
#[repr(u8)]
#[derive(Copy, Clone, PartialEq, Eq, Default)]
pub enum AllocatorType {
#[default]
Js,
Mini,
}
#[cfg(not(windows))]
pub struct FilePoll {
pub fd: Fd,
pub flags: FlagsSet,
pub owner: Owner,
pub generation_number: KQueueGenerationNumber,
pub next_to_free: *mut FilePoll,
pub allocator_type: AllocatorType,
}
#[cfg(not(windows))]
impl Default for FilePoll {
fn default() -> Self {
Self {
fd: INVALID_FD,
flags: FlagsSet::empty(),
owner: Owner::NULL,
generation_number: 0,
next_to_free: ptr::null_mut(),
allocator_type: AllocatorType::Js,
}
}
}
#[cfg(not(windows))]
impl FilePoll {
fn update_flags(&mut self, updated: FlagsSet) {
let mut flags = self.flags;
flags.remove(Flags::Readable);
flags.remove(Flags::Writable);
flags.remove(Flags::Process);
flags.remove(Flags::Machport);
flags.remove(Flags::Eof);
flags.remove(Flags::Hup);
flags |= updated;
self.flags = flags;
}
pub fn file_type(&self) -> FileType {
let flags = self.flags;
if flags.contains(Flags::Socket) {
return FileType::Socket;
}
if flags.contains(Flags::Nonblocking) {
return FileType::NonblockingPipe;
}
FileType::Pipe
}
#[cfg(any(target_os = "macos", target_os = "freebsd"))]
pub fn on_kqueue_event(&mut self, kqueue_event: &KQueueEvent) {
self.update_flags(Flags::from_kqueue_event(kqueue_event));
syslog!("onKQueueEvent: {}", self);
#[cfg(all(target_os = "macos", debug_assertions))]
debug_assert!(self.generation_number == kqueue_event.ext[0] as usize);
self.on_update(kqueue_event.data as i64);
}
#[cfg(any(target_os = "linux", target_os = "android"))]
pub fn on_epoll_event(&mut self, epoll_event: &bun_sys::linux::epoll_event) {
self.update_flags(Flags::from_epoll_event(epoll_event));
self.on_update(0);
}
pub fn clear_event(&mut self, flag: Flags) {
self.flags.remove(flag);
}
pub fn is_readable(&mut self) -> bool {
let readable = self.flags.contains(Flags::Readable);
self.flags.remove(Flags::Readable);
readable
}
pub fn is_hup(&mut self) -> bool {
let readable = self.flags.contains(Flags::Hup);
self.flags.remove(Flags::Hup);
readable
}
pub fn is_eof(&mut self) -> bool {
let readable = self.flags.contains(Flags::Eof);
self.flags.remove(Flags::Eof);
readable
}
pub fn is_writable(&mut self) -> bool {
let readable = self.flags.contains(Flags::Writable);
self.flags.remove(Flags::Writable);
readable
}
pub fn deinit(&mut self) {
let ctx = get_vm_ctx(self.allocator_type);
self.deinit_possibly_defer(ctx, false);
}
pub fn deinit_force_unregister(&mut self) {
let ctx = get_vm_ctx(self.allocator_type);
self.deinit_possibly_defer(ctx, true);
}
fn deinit_possibly_defer(&mut self, vm: EventLoopCtx, force_unregister: bool) {
let _ = self.unregister(vm.loop_mut(), force_unregister);
self.owner.clear();
let was_ever_registered = self.flags.contains(Flags::WasEverRegistered);
self.flags = FlagsSet::empty();
self.fd = INVALID_FD;
let this = ptr::NonNull::from(self);
vm.file_polls_mut().put(this, vm, was_ever_registered);
}
pub fn deinit_with_vm(&mut self, vm: EventLoopCtx) {
self.deinit_possibly_defer(vm, false);
}
pub fn is_registered(&self) -> bool {
self.flags.contains(Flags::PollWritable)
|| self.flags.contains(Flags::PollReadable)
|| self.flags.contains(Flags::PollProcess)
|| self.flags.contains(Flags::PollMachport)
}
pub fn on_update(&mut self, size_or_offset: i64) {
if self.flags.contains(Flags::OneShot) && !self.flags.contains(Flags::NeedsRearm) {
self.flags.insert(Flags::NeedsRearm);
}
debug_assert!(!self.owner.is_null());
unsafe { __bun_run_file_poll(self, size_or_offset) };
}
#[inline]
pub fn is_active(&self) -> bool {
self.flags.contains(Flags::HasIncrementedPollCount)
}
#[inline]
pub fn is_watching(&self) -> bool {
!self.flags.contains(Flags::NeedsRearm)
&& (self.flags.contains(Flags::PollReadable)
|| self.flags.contains(Flags::PollWritable)
|| self.flags.contains(Flags::PollProcess))
}
pub fn disable_keeping_process_alive(&mut self, event_loop_ctx: EventLoopCtx) {
event_loop_ctx
.loop_sub_active(self.flags.contains(Flags::HasIncrementedActiveCount) as u32);
self.flags.remove(Flags::KeepsEventLoopAlive);
self.flags.remove(Flags::HasIncrementedActiveCount);
}
#[inline]
pub fn can_enable_keeping_process_alive(&self) -> bool {
self.flags.contains(Flags::KeepsEventLoopAlive)
&& self.flags.contains(Flags::HasIncrementedPollCount)
}
pub fn set_keeping_process_alive(&mut self, event_loop_ctx: EventLoopCtx, value: bool) {
if value {
self.enable_keeping_process_alive(event_loop_ctx);
} else {
self.disable_keeping_process_alive(event_loop_ctx);
}
}
pub fn enable_keeping_process_alive(&mut self, event_loop_ctx: EventLoopCtx) {
if self.flags.contains(Flags::Closed) {
return;
}
event_loop_ctx
.loop_add_active((!self.flags.contains(Flags::HasIncrementedActiveCount)) as u32);
self.flags.insert(Flags::KeepsEventLoopAlive);
self.flags.insert(Flags::HasIncrementedActiveCount);
}
fn deactivate(&mut self, loop_: &mut Loop) {
if self.flags.contains(Flags::HasIncrementedPollCount) {
loop_.dec();
}
self.flags.remove(Flags::HasIncrementedPollCount);
loop_sub_active(
loop_,
self.flags.contains(Flags::HasIncrementedActiveCount) as u32,
);
self.flags.remove(Flags::KeepsEventLoopAlive);
self.flags.remove(Flags::HasIncrementedActiveCount);
}
fn activate(&mut self, loop_: &mut Loop) {
self.flags.remove(Flags::Closed);
if !self.flags.contains(Flags::HasIncrementedPollCount) {
loop_.inc();
}
self.flags.insert(Flags::HasIncrementedPollCount);
if self.flags.contains(Flags::KeepsEventLoopAlive) {
loop_add_active(
loop_,
(!self.flags.contains(Flags::HasIncrementedActiveCount)) as u32,
);
self.flags.insert(Flags::HasIncrementedActiveCount);
}
}
#[inline]
fn new_value(vm: EventLoopCtx, fd: Fd, flags: FlagsSet, owner: Owner) -> FilePoll {
FilePoll {
fd,
flags,
owner,
next_to_free: ptr::null_mut(),
allocator_type: if vm.is_js() { AllocatorType::Js } else { AllocatorType::Mini },
#[cfg(all(target_os = "macos", debug_assertions))]
generation_number: MAX_GENERATION_NUMBER
.fetch_add(1, core::sync::atomic::Ordering::Relaxed)
.wrapping_add(1),
#[cfg(not(all(target_os = "macos", debug_assertions)))]
generation_number: 0,
}
}
pub fn init(vm: EventLoopCtx, fd: Fd, flags: FlagsSet, owner: Owner) -> *mut FilePoll {
let value = Self::new_value(vm, fd, flags, owner);
let generation_number = value.generation_number;
let poll = vm.alloc_file_poll(value).as_ptr();
syslog!(
"FilePoll.init(0x{:x}, generation_number={}, fd={})",
poll as usize,
generation_number,
fd
);
poll
}
#[inline]
pub fn can_ref(&self) -> bool {
!self.flags.contains(Flags::HasIncrementedPollCount)
}
#[inline]
pub fn can_unref(&self) -> bool {
self.flags.contains(Flags::HasIncrementedPollCount)
}
pub fn unref(&mut self, event_loop_ctx: EventLoopCtx) {
syslog!("unref");
self.disable_keeping_process_alive(event_loop_ctx);
}
pub fn ref_(&mut self, event_loop_ctx: EventLoopCtx) {
if self.flags.contains(Flags::Closed) {
return;
}
syslog!("ref");
self.enable_keeping_process_alive(event_loop_ctx);
}
pub fn on_ended(&mut self, event_loop_ctx: EventLoopCtx) {
self.flags.remove(Flags::KeepsEventLoopAlive);
self.flags.insert(Flags::Closed);
self.deactivate(event_loop_ctx.loop_mut());
}
#[inline]
pub fn file_descriptor(&self) -> Fd {
self.fd
}
pub fn register(&mut self, loop_: &mut Loop, flag: Flags, one_shot: bool) -> sys::Result<()> {
self.register_with_fd(
loop_,
flag,
if one_shot {
OneShotFlag::OneShot
} else {
OneShotFlag::None
},
self.fd,
)
}
pub fn register_with_fd(
&mut self,
loop_: &mut Loop,
flag: Flags,
one_shot: OneShotFlag,
fd: Fd,
) -> sys::Result<()> {
#[cfg(any(
target_os = "linux",
target_os = "android",
target_os = "macos",
target_os = "freebsd"
))]
return self.register_with_fd_impl(loop_, flag, one_shot, fd);
#[cfg(not(any(
target_os = "linux",
target_os = "android",
target_os = "macos",
target_os = "freebsd"
)))]
{
let _ = (loop_, flag, one_shot, fd);
sys::Result::Ok(())
}
}
#[cfg(any(
target_os = "linux",
target_os = "android",
target_os = "macos",
target_os = "freebsd"
))]
fn register_with_fd_impl(
&mut self,
loop_: &mut Loop,
flag: Flags,
one_shot: OneShotFlag,
fd: Fd,
) -> sys::Result<()> {
let watcher_fd = loop_.fd;
syslog!(
"register: FilePoll(0x{:x}, generation_number={}) {} ({})",
std::ptr::from_mut(self) as usize,
self.generation_number,
<&'static str>::from(flag),
fd
);
debug_assert!(fd != INVALID_FD);
if one_shot != OneShotFlag::None {
self.flags.insert(Flags::OneShot);
}
#[cfg(any(target_os = "linux", target_os = "android"))]
{
use bun_sys::linux::{self, EPOLL};
let one_shot_flag: u32 = if !self.flags.contains(Flags::OneShot) {
0
} else {
EPOLL::ONESHOT
};
let mut flags: u32 = match flag {
Flags::Process | Flags::Readable => EPOLL::IN | EPOLL::HUP | one_shot_flag,
Flags::Writable => EPOLL::OUT | EPOLL::HUP | EPOLL::ERR | one_shot_flag,
_ => unreachable!(),
};
if flag == Flags::Readable && self.flags.contains(Flags::PollWritable) {
debug_assert!(!self.flags.contains(Flags::OneShot));
flags |= EPOLL::OUT | EPOLL::ERR;
}
if flag == Flags::Writable && self.flags.contains(Flags::PollReadable) {
debug_assert!(!self.flags.contains(Flags::OneShot));
flags |= EPOLL::IN;
}
let mut event = linux::epoll_event {
events: flags,
u64: Pollable::init(self).ptr() as u64,
};
let op: c_int = if self.is_registered() || self.flags.contains(Flags::NeedsRearm) {
EPOLL::CTL_MOD
} else {
EPOLL::CTL_ADD
};
let ctl = unsafe { linux::epoll_ctl(watcher_fd, op, fd.native(), &raw mut event) };
self.flags.insert(Flags::WasEverRegistered);
if let Some(errno) = errno_sys(ctl, sys::Tag::epoll_ctl) {
self.deactivate(loop_);
return errno;
}
}
#[cfg(target_os = "macos")]
{
use bun_sys::darwin::{EV, EVFILT, NOTE, kevent64_s};
let mut changelist: [kevent64_s; 2] = bun_core::ffi::zeroed();
let one_shot_flag: u16 = if !self.flags.contains(Flags::OneShot) {
0
} else if one_shot == OneShotFlag::Dispatch {
EV::DISPATCH | EV::ENABLE
} else {
EV::ONESHOT
};
changelist[0] = match flag {
Flags::Readable => kevent64_s {
ident: u64::try_from(fd.native()).expect("int cast"),
filter: EVFILT::READ,
data: 0,
fflags: 0,
udata: Pollable::init(self).ptr() as u64,
flags: EV::ADD | one_shot_flag,
ext: [self.generation_number as u64, 0],
},
Flags::Writable => kevent64_s {
ident: u64::try_from(fd.native()).expect("int cast"),
filter: EVFILT::WRITE,
data: 0,
fflags: 0,
udata: Pollable::init(self).ptr() as u64,
flags: EV::ADD | one_shot_flag,
ext: [self.generation_number as u64, 0],
},
Flags::Process => kevent64_s {
ident: u64::try_from(fd.native()).expect("int cast"),
filter: EVFILT::PROC,
data: 0,
fflags: NOTE::EXIT,
udata: Pollable::init(self).ptr() as u64,
flags: EV::ADD | one_shot_flag,
ext: [self.generation_number as u64, 0],
},
Flags::Machport => kevent64_s {
ident: u64::try_from(fd.native()).expect("int cast"),
filter: EVFILT::MACHPORT,
data: 0,
fflags: 0,
udata: Pollable::init(self).ptr() as u64,
flags: EV::ADD | one_shot_flag,
ext: [self.generation_number as u64, 0],
},
_ => unreachable!(),
};
const KEVENT_FLAG_ERROR_EVENTS: u32 = 0x000002;
let rc = 'rc: loop {
let rc = unsafe {
bun_sys::darwin::kevent64(
watcher_fd,
changelist.as_ptr(),
1,
changelist.as_mut_ptr(),
0,
KEVENT_FLAG_ERROR_EVENTS,
&raw const TIMEOUT,
)
};
if sys::get_errno(rc) == sys::E::EINTR {
continue;
}
break 'rc rc;
};
self.flags.insert(Flags::WasEverRegistered);
if (changelist[0].flags & EV::ERROR) != 0 && changelist[0].data != 0 {
return errno_sys(changelist[0].data, sys::Tag::kevent).unwrap();
}
let errno = sys::get_errno(rc);
if errno != sys::E::SUCCESS {
self.deactivate(loop_);
return sys::Result::Err(sys::Error::from_code(errno, sys::Tag::kqueue));
}
}
#[cfg(target_os = "freebsd")]
{
use bun_sys::freebsd::{EV, EVFILT, Kevent, NOTE, kevent};
let mut changelist: [Kevent; 1] = bun_core::ffi::zeroed();
let one_shot_flag: u16 = if !self.flags.contains(Flags::OneShot) {
0
} else if one_shot == OneShotFlag::Dispatch {
EV::DISPATCH | EV::ENABLE
} else {
EV::ONESHOT
};
let ident = usize::try_from(fd.native()).expect("int cast");
let udata = Pollable::init(self).ptr();
changelist[0] = match flag {
Flags::Readable => {
make_kevent(ident, EVFILT::READ, EV::ADD | one_shot_flag, 0, udata)
}
Flags::Writable => {
make_kevent(ident, EVFILT::WRITE, EV::ADD | one_shot_flag, 0, udata)
}
Flags::Process => make_kevent(
ident,
EVFILT::PROC,
EV::ADD | one_shot_flag,
NOTE::EXIT,
udata,
),
Flags::Machport => {
return sys::Result::Err(sys::Error::from_code(
sys::E::EOPNOTSUPP,
sys::Tag::kevent,
));
}
_ => unreachable!(),
};
let rc = 'rc: loop {
let rc = unsafe {
kevent(
watcher_fd,
changelist.as_ptr(),
1,
changelist.as_mut_ptr(),
0,
ptr::null(),
)
};
if sys::get_errno(rc) == sys::E::EINTR {
continue;
}
break 'rc rc;
};
self.flags.insert(Flags::WasEverRegistered);
if let Some(err) = errno_sys(rc, sys::Tag::kevent) {
self.deactivate(loop_);
return err;
}
}
self.activate(loop_);
self.flags.insert(match flag {
Flags::Readable => Flags::PollReadable,
Flags::Process => {
#[cfg(any(target_os = "linux", target_os = "android"))]
{
Flags::PollReadable
}
#[cfg(not(any(target_os = "linux", target_os = "android")))]
{
Flags::PollProcess
}
}
Flags::Writable => Flags::PollWritable,
Flags::Machport => Flags::PollMachport,
_ => unreachable!(),
});
self.flags.remove(Flags::NeedsRearm);
sys::Result::Ok(())
}
pub fn unregister(&mut self, loop_: &mut Loop, force_unregister: bool) -> sys::Result<()> {
self.unregister_with_fd(loop_, self.fd, force_unregister)
}
pub fn unregister_with_fd(
&mut self,
loop_: &mut Loop,
fd: Fd,
force_unregister: bool,
) -> sys::Result<()> {
#[cfg(any(
target_os = "linux",
target_os = "android",
target_os = "macos",
target_os = "freebsd"
))]
let result = self.unregister_with_fd_impl(loop_, fd, force_unregister);
#[cfg(not(any(
target_os = "linux",
target_os = "android",
target_os = "macos",
target_os = "freebsd"
)))]
let result: sys::Result<()> = {
let _ = (fd, force_unregister);
sys::Result::Ok(())
};
self.deactivate(loop_);
result
}
#[cfg(any(
target_os = "linux",
target_os = "android",
target_os = "macos",
target_os = "freebsd"
))]
fn unregister_with_fd_impl(
&mut self,
loop_: &mut Loop,
fd: Fd,
force_unregister: bool,
) -> sys::Result<()> {
#[cfg(debug_assertions)]
debug_assert!(fd.native() >= 0 && fd != INVALID_FD);
if !(self.flags.contains(Flags::PollReadable)
|| self.flags.contains(Flags::PollWritable)
|| self.flags.contains(Flags::PollProcess)
|| self.flags.contains(Flags::PollMachport))
{
return sys::Result::Ok(());
}
debug_assert!(fd != INVALID_FD);
let watcher_fd = loop_.fd;
let both_directions =
self.flags.contains(Flags::PollReadable) && self.flags.contains(Flags::PollWritable);
let flag: Flags = 'brk: {
if self.flags.contains(Flags::PollReadable) {
break 'brk Flags::Readable;
}
if self.flags.contains(Flags::PollWritable) {
break 'brk Flags::Writable;
}
if self.flags.contains(Flags::PollProcess) {
break 'brk Flags::Process;
}
if self.flags.contains(Flags::PollMachport) {
break 'brk Flags::Machport;
}
return sys::Result::Ok(());
};
if self.flags.contains(Flags::NeedsRearm) && !force_unregister {
syslog!(
"unregister: {} ({}) skipped due to needs_rearm",
<&'static str>::from(flag),
fd
);
self.flags.remove(Flags::PollProcess);
self.flags.remove(Flags::PollReadable);
self.flags.remove(Flags::PollWritable);
self.flags.remove(Flags::PollMachport);
return sys::Result::Ok(());
}
syslog!(
"unregister: FilePoll(0x{:x}, generation_number={}) {}{} ({})",
std::ptr::from_mut(self) as usize,
self.generation_number,
<&'static str>::from(flag),
if both_directions { "+writable" } else { "" },
fd
);
#[cfg(any(target_os = "linux", target_os = "android"))]
{
use bun_sys::linux::{self, EPOLL};
let ctl = unsafe {
linux::epoll_ctl(watcher_fd, EPOLL::CTL_DEL, fd.native(), ptr::null_mut())
};
if let Some(errno) = errno_sys(ctl, sys::Tag::epoll_ctl) {
return errno;
}
}
#[cfg(target_os = "macos")]
{
use bun_sys::darwin::{EV, EVFILT, NOTE, kevent64, kevent64_s};
let mut changelist: [kevent64_s; 2] = bun_core::ffi::zeroed();
changelist[0] = match flag {
Flags::Readable => kevent64_s {
ident: u64::try_from(fd.native()).expect("int cast"),
filter: EVFILT::READ,
data: 0,
fflags: 0,
udata: Pollable::init(self).ptr() as u64,
flags: EV::DELETE,
ext: [0, 0],
},
Flags::Machport => kevent64_s {
ident: u64::try_from(fd.native()).expect("int cast"),
filter: EVFILT::MACHPORT,
data: 0,
fflags: 0,
udata: Pollable::init(self).ptr() as u64,
flags: EV::DELETE,
ext: [0, 0],
},
Flags::Writable => kevent64_s {
ident: u64::try_from(fd.native()).expect("int cast"),
filter: EVFILT::WRITE,
data: 0,
fflags: 0,
udata: Pollable::init(self).ptr() as u64,
flags: EV::DELETE,
ext: [0, 0],
},
Flags::Process => kevent64_s {
ident: u64::try_from(fd.native()).expect("int cast"),
filter: EVFILT::PROC,
data: 0,
fflags: NOTE::EXIT,
udata: Pollable::init(self).ptr() as u64,
flags: EV::DELETE,
ext: [0, 0],
},
_ => unreachable!(),
};
let mut nchanges: c_int = 1;
if both_directions {
changelist[1] = kevent64_s {
ident: u64::try_from(fd.native()).expect("int cast"),
filter: EVFILT::WRITE,
data: 0,
fflags: 0,
udata: Pollable::init(self).ptr() as u64,
flags: EV::DELETE,
ext: [0, 0],
};
nchanges = 2;
}
const KEVENT_FLAG_ERROR_EVENTS: u32 = 0x000002;
let rc = unsafe {
kevent64(
watcher_fd,
changelist.as_ptr(),
nchanges,
changelist.as_mut_ptr(),
nchanges,
KEVENT_FLAG_ERROR_EVENTS,
&raw const TIMEOUT,
)
};
let errno = sys::get_errno(rc);
if rc < 0 {
return sys::Result::Err(sys::Error::from_code(errno, sys::Tag::kevent));
}
if rc >= 1 && (changelist[0].flags & EV::ERROR) != 0 && changelist[0].data != 0 {
return errno_sys(changelist[0].data, sys::Tag::kevent).unwrap();
}
if rc >= 2 && (changelist[1].flags & EV::ERROR) != 0 && changelist[1].data != 0 {
return errno_sys(changelist[1].data, sys::Tag::kevent).unwrap();
}
}
#[cfg(target_os = "freebsd")]
{
use bun_sys::freebsd::{EV, EVFILT, Kevent, NOTE, kevent};
let mut changelist: [Kevent; 2] = bun_core::ffi::zeroed();
let ident = usize::try_from(fd.native()).expect("int cast");
let udata = Pollable::init(self).ptr();
changelist[0] = match flag {
Flags::Readable => make_kevent(ident, EVFILT::READ, EV::DELETE, 0, udata),
Flags::Writable => make_kevent(ident, EVFILT::WRITE, EV::DELETE, 0, udata),
Flags::Process => make_kevent(ident, EVFILT::PROC, EV::DELETE, NOTE::EXIT, udata),
Flags::Machport => {
return sys::Result::Err(sys::Error::from_code(
sys::E::EOPNOTSUPP,
sys::Tag::kevent,
));
}
_ => unreachable!(),
};
let mut nchanges: c_int = 1;
if both_directions {
changelist[1] = make_kevent(ident, EVFILT::WRITE, EV::DELETE, 0, udata);
nchanges = 2;
}
let rc = unsafe {
kevent(
watcher_fd,
changelist.as_ptr(),
nchanges,
changelist.as_mut_ptr(),
0,
ptr::null(),
)
};
if let Some(err) = errno_sys(rc, sys::Tag::kevent) {
return err;
}
}
self.flags.remove(Flags::NeedsRearm);
self.flags.remove(Flags::OneShot);
self.flags.remove(Flags::PollReadable);
self.flags.remove(Flags::PollWritable);
self.flags.remove(Flags::PollProcess);
self.flags.remove(Flags::PollMachport);
sys::Result::Ok(())
}
}
#[cfg(not(windows))]
impl fmt::Display for FilePoll {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
f,
"FilePoll(fd={}, generation_number={}) = {}",
self.fd,
self.generation_number,
FlagsFormatter(self.flags)
)
}
}
#[derive(enumset::EnumSetType, strum::IntoStaticStr, Debug)]
#[strum(serialize_all = "snake_case")]
pub enum Flags {
PollReadable,
PollWritable,
PollProcess,
PollMachport,
Readable,
Writable,
Process,
Eof,
Hup,
Machport,
Fifo,
Tty,
OneShot,
NeedsRearm,
HasIncrementedPollCount,
HasIncrementedActiveCount,
Closed,
KeepsEventLoopAlive,
Nonblocking,
WasEverRegistered,
IgnoreUpdates,
Nonblock,
Socket,
}
pub type FlagsSet = enumset::EnumSet<Flags>;
#[allow(dead_code)]
pub type FlagsStruct = FlagsSet;
impl Flags {
pub fn poll(self) -> Flags {
match self {
Flags::Readable => Flags::PollReadable,
Flags::Writable => Flags::PollWritable,
Flags::Process => Flags::PollProcess,
Flags::Machport => Flags::PollMachport,
other => other,
}
}
#[cfg(any(target_os = "macos", target_os = "freebsd"))]
pub fn from_kqueue_event(kqueue_event: &KQueueEvent) -> FlagsSet {
#[cfg(target_os = "macos")]
use bun_sys::darwin::EVFILT;
#[cfg(target_os = "freebsd")]
use bun_sys::freebsd::EVFILT;
let mut flags = FlagsSet::empty();
if kqueue_event.filter == EVFILT::READ {
flags.insert(Flags::Readable);
if kqueue_event.flags & EV_EOF != 0 {
flags.insert(Flags::Hup);
}
} else if kqueue_event.filter == EVFILT::WRITE {
flags.insert(Flags::Writable);
if kqueue_event.flags & EV_EOF != 0 {
flags.insert(Flags::Hup);
}
} else if kqueue_event.filter == EVFILT::PROC {
flags.insert(Flags::Process);
} else {
#[cfg(target_os = "macos")]
if kqueue_event.filter == EVFILT::MACHPORT {
flags.insert(Flags::Machport);
}
}
flags
}
#[cfg(any(target_os = "linux", target_os = "android"))]
pub fn from_epoll_event(epoll: &bun_sys::linux::epoll_event) -> FlagsSet {
use bun_sys::linux::EPOLL;
let mut flags = FlagsSet::empty();
if epoll.events & EPOLL::IN != 0 {
flags.insert(Flags::Readable);
}
if epoll.events & EPOLL::OUT != 0 {
flags.insert(Flags::Writable);
}
if epoll.events & EPOLL::ERR != 0 {
flags.insert(Flags::Eof);
}
if epoll.events & EPOLL::HUP != 0 {
flags.insert(Flags::Hup);
}
flags
}
}
#[allow(dead_code)]
pub(crate) struct FlagsFormatter(pub FlagsSet);
impl fmt::Display for FlagsFormatter {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let mut is_first = true;
for flag in self.0.iter() {
if !is_first {
write!(f, " | ")?;
}
f.write_str(<&'static str>::from(flag))?;
is_first = false;
}
Ok(())
}
}
#[cfg(not(windows))]
const HIVE_SIZE: usize = 128;
#[cfg(not(windows))]
type FilePollHive = bun_collections::hive_array::Fallback<FilePoll, HIVE_SIZE>;
#[cfg(not(windows))]
pub struct Store {
hive: FilePollHive,
pending_free_head: *mut FilePoll,
pending_free_tail: *mut FilePoll,
}
#[cfg(not(windows))]
impl Store {
pub fn init() -> Store {
Store {
hive: FilePollHive::init(),
pending_free_head: ptr::null_mut(),
pending_free_tail: ptr::null_mut(),
}
}
#[inline]
pub fn get_init(&mut self, value: FilePoll) -> ptr::NonNull<FilePoll> {
self.hive.get_init(value)
}
pub fn process_deferred_frees(&mut self) {
let mut next = self.pending_free_head;
while !next.is_null() {
let current = next;
unsafe {
next = (*current).next_to_free;
(*current).next_to_free = ptr::null_mut();
self.hive.put(current);
}
}
self.pending_free_head = ptr::null_mut();
self.pending_free_tail = ptr::null_mut();
}
pub fn put(&mut self, poll: ptr::NonNull<FilePoll>, vm: EventLoopCtx, ever_registered: bool) {
let poll = poll.as_ptr();
if !ever_registered {
unsafe { self.hive.put(poll) };
return;
}
debug_assert!(unsafe { (*poll).next_to_free }.is_null());
if !self.pending_free_tail.is_null() {
debug_assert!(!self.pending_free_head.is_null());
unsafe {
debug_assert!((*self.pending_free_tail).next_to_free.is_null());
(*self.pending_free_tail).next_to_free = poll;
}
}
if self.pending_free_head.is_null() {
self.pending_free_head = poll;
debug_assert!(self.pending_free_tail.is_null());
}
unsafe { (*poll).flags.insert(Flags::IgnoreUpdates) };
self.pending_free_tail = poll;
let callback: OpaqueCallback = Self::process_deferred_frees_thunk;
debug_assert!(
vm.after_event_loop_callback().is_none()
|| vm.after_event_loop_callback().map(|f| f as usize) == Some(callback as usize)
);
vm.set_after_event_loop_callback(
Some(callback),
core::ptr::NonNull::new(std::ptr::from_mut::<Store>(self).cast::<c_void>()),
);
}
extern "C" fn process_deferred_frees_thunk(ctx: *mut c_void) {
let this = unsafe { bun_ptr::callback_ctx::<Store>(ctx) };
this.process_deferred_frees();
}
}
#[derive(Copy, Clone)]
#[allow(dead_code)]
pub(crate) struct Pollable {
repr: bun_collections::TaggedPointer,
}
impl Pollable {
#[allow(dead_code)]
pub(crate) const FILE_POLL_TAG: u16 = 1024;
#[inline]
#[allow(dead_code)]
pub(crate) fn init(ptr: *const crate::FilePoll) -> Self {
Self {
repr: bun_collections::TaggedPointer::init(ptr, Self::FILE_POLL_TAG),
}
}
#[inline]
#[allow(dead_code)]
pub(crate) fn from(val: *mut c_void) -> Self {
Self {
repr: bun_collections::TaggedPointer::from(val),
}
}
#[inline]
#[allow(dead_code)]
pub(crate) fn tag(self) -> u16 {
self.repr.data()
}
#[inline]
#[allow(dead_code)]
pub(crate) fn as_file_poll(self) -> *mut crate::FilePoll {
self.repr.get::<crate::FilePoll>()
}
#[inline]
#[allow(dead_code)]
pub(crate) fn ptr(self) -> *mut c_void {
self.repr.to()
}
}
#[cfg(any(
target_os = "linux",
target_os = "android",
target_os = "macos",
target_os = "freebsd"
))]
#[unsafe(no_mangle)]
pub(crate) unsafe extern "C" fn Bun__internal_dispatch_ready_poll(
loop_: *mut Loop,
tagged_pointer: *mut c_void,
) {
let tag = Pollable::from(tagged_pointer);
if tag.tag() != Pollable::FILE_POLL_TAG {
return;
}
let file_poll: &mut FilePoll = unsafe { &mut *tag.as_file_poll() };
if file_poll.flags.contains(Flags::IgnoreUpdates) {
return;
}
let ev = unsafe { &*loop_ }.current_ready_event();
#[cfg(any(target_os = "macos", target_os = "freebsd"))]
file_poll.on_kqueue_event(&ev);
#[cfg(any(target_os = "linux", target_os = "android"))]
file_poll.on_epoll_event(&ev);
}
#[cfg(target_os = "macos")]
static TIMEOUT: bun_sys::posix::timespec = bun_sys::posix::timespec {
tv_sec: 0,
tv_nsec: 0,
};
#[repr(u8)]
#[derive(Copy, Clone, PartialEq, Eq)]
pub enum OneShotFlag {
Dispatch,
OneShot,
None,
}
#[cfg(not(windows))]
const INVALID_FD: Fd = Fd::INVALID;
pub use crate::closer::Closer;
#[cfg(target_os = "macos")]
pub use crate::waker::KEventWaker;
#[cfg(any(target_os = "linux", target_os = "android", target_os = "freebsd"))]
pub use crate::waker::Waker;