Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

17 changes: 13 additions & 4 deletions src/bun_core/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -950,21 +950,30 @@ pub unsafe trait IntrusiveField<F>: Sized {
///
/// ```ignore
/// bun_core::intrusive_field!(ShellCpTask, task: ShellTask);
/// bun_core::intrusive_field!(ShellCpTask, task.task: WorkPoolTask); // nested field
/// bun_core::intrusive_field!([T: Send] Wrapper<T>, inner: Mixin<Wrapper<T>>);
/// ```
#[macro_export]
macro_rules! intrusive_field {
// Bracketed-generics arm MUST come first: the bare `$T:ty` arm below would
// otherwise try to parse `['a]` as a slice type and hard-error on the
// lifetime before backtracking to this arm.
([$($gen:tt)*] $T:ty, $field:ident : $F:ty) => {
([$($gen:tt)*] $T:ty, $($field:ident).+ : $F:ty) => {
unsafe impl<$($gen)*> $crate::IntrusiveField<$F> for $T {
const OFFSET: usize = ::core::mem::offset_of!($T, $field);
const OFFSET: usize = {
// The named field must actually be an `$F`.
let _ = |s: &$T| -> *const $F { &raw const s.$($field).+ };
::core::mem::offset_of!($T, $($field).+)
};
}
};
($T:ty, $field:ident : $F:ty) => {
($T:ty, $($field:ident).+ : $F:ty) => {
unsafe impl $crate::IntrusiveField<$F> for $T {
const OFFSET: usize = ::core::mem::offset_of!($T, $field);
const OFFSET: usize = {
// The named field must actually be an `$F`.
let _ = |s: &$T| -> *const $F { &raw const s.$($field).+ };
::core::mem::offset_of!($T, $($field).+)
};
}
};
}
Expand Down
36 changes: 36 additions & 0 deletions src/event_loop/AnyEventLoop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -252,13 +252,42 @@ pub enum EventLoopTask {
Mini(AnyTaskWithExtraContext),
}

/// A leaked box's embedded node, armed and ready to post to its loop
/// ([`EventLoopTask::arm_boxed`]).
#[derive(Clone, Copy)]
pub enum ArmedLoopTask {
Js(NonNull<ConcurrentTask>),
Mini(NonNull<AnyTaskWithExtraContext>),
}

impl EventLoopTask {
pub fn from_event_loop(loop_: EventLoopHandle) -> EventLoopTask {
match loop_ {
EventLoopHandle::Js { .. } => EventLoopTask::Js(ConcurrentTask::default()),
EventLoopHandle::Mini(_) => EventLoopTask::Mini(AnyTaskWithExtraContext::default()),
}
}

/// Leak `owner` into the node it embeds (`node(owner)`). The JS arm queues
/// `Task { T::TAG, owner }` for `bun_runtime::dispatch` to rebox; the mini
/// arm hands the box to `R::run_from_loop_thread`. No allocation.
pub fn arm_boxed<T, R>(owner: Box<T>, node: fn(&mut T) -> &mut EventLoopTask) -> ArmedLoopTask
where
T: crate::Taskable,
R: crate::AnyTaskWithExtraContext::BoxedMiniTaskRunner<T>,
{
let raw: *mut T = bun_core::heap::into_raw(owner);
// SAFETY: `raw` was just leaked; this thread owns it exclusively until
// the node is queued.
match node(unsafe { &mut *raw }) {
EventLoopTask::Js(ct) => ArmedLoopTask::Js(NonNull::from(
ct.from(raw, crate::ConcurrentTask::AutoDeinit::ManualDeinit),
)),
EventLoopTask::Mini(at) => {
ArmedLoopTask::Mini(AnyTaskWithExtraContext::arm_at::<T, R>(at, raw))
}
}
}
}

/// RAII pairing for [`EventLoopHandle::enter`] / [`EventLoopHandle::exit`].
Expand Down Expand Up @@ -508,6 +537,13 @@ impl EventLoopHandle {
f(unsafe { (*env).get(key) })
}

/// `f(loader)` on the loop's dotenv loader; the borrow ends with `f` (the
/// map is mutable at runtime).
pub fn with_env<R>(self, f: impl FnOnce(&DotEnvLoader) -> R) -> R {
// SAFETY: as `with_env_var`.
f(unsafe { &*self.env() })
}

pub fn top_level_dir(self) -> &'static [u8] {
match self {
// SAFETY: slice borrowed for VM lifetime.
Expand Down
62 changes: 57 additions & 5 deletions src/event_loop/AnyTaskWithExtraContext.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,27 @@ impl AnyTaskWithExtraContext {
}
}

/// Leak `owner` and arm the node it embeds (`node(owner)`) so the mini
/// loop hands the box to `R::run_from_loop_thread`. No allocation.
pub fn arm_boxed<T, R: BoxedMiniTaskRunner<T>>(
owner: Box<T>,
node: fn(&mut T) -> &mut AnyTaskWithExtraContext,
) -> NonNull<AnyTaskWithExtraContext> {
let raw: *mut T = bun_core::heap::into_raw(owner);
// SAFETY: `raw` was just leaked; this thread owns it exclusively until
// the node is queued.
Self::arm_at::<T, R>(node(unsafe { &mut *raw }), raw)
}

/// Arm `at` (a node inside the leaked box `raw`) for `R`.
pub(crate) fn arm_at<T, R: BoxedMiniTaskRunner<T>>(
at: &mut AnyTaskWithExtraContext,
raw: *mut T,
) -> NonNull<AnyTaskWithExtraContext> {
*at = New::<T, ()>::init(raw, mini_boxed_trampoline::<T, R>);
NonNull::from(at)
}

/// Initializes `self` in place to call `callback(of, extra)`.
// The unit context means the callee is effectively `fn(*T)` only; mapped
// to `*mut ()` to keep the two-arg stored ABI uniform.
Expand All @@ -77,14 +98,45 @@ impl AnyTaskWithExtraContext {
std::ptr::from_mut::<Self>(self)
}

pub(crate) fn run(&mut self, extra: *mut c_void) {
let callback = self.callback;
let ctx = self.ctx;
// SAFETY: caller contract — `ctx` was set by `init`/`from*` to a live pointer.
callback(ctx.expect("ctx is non-null").as_ptr(), extra.cast::<()>());
/// Copy the node's two words out so running it borrows nothing from the
/// allocation it lives in (the callback may free that allocation).
///
/// # Safety
/// `this` is a queued node, live at the time of the call.
pub(crate) unsafe fn load(this: *const Self) -> Runnable {
// SAFETY: fn contract.
let (callback, ctx) = unsafe { ((*this).callback, (*this).ctx) };
Runnable { callback, ctx }
}
}

/// A dequeued [`AnyTaskWithExtraContext`], detached from its storage.
pub(crate) struct Runnable {
callback: fn(*mut (), *mut ()),
ctx: Option<NonNull<()>>,
}

impl Runnable {
pub(crate) fn run(self, extra: *mut c_void) {
(self.callback)(
self.ctx.expect("ctx is non-null").as_ptr(),
extra.cast::<()>(),
);
}
}

/// Receives a boxed payload back on the mini loop's thread
/// ([`AnyTaskWithExtraContext::arm_boxed`]).
pub trait BoxedMiniTaskRunner<T> {
fn run_from_loop_thread(owner: Box<T>);
}

fn mini_boxed_trampoline<T, R: BoxedMiniTaskRunner<T>>(this: *mut T, _extra: *mut ()) {
// SAFETY: `this` is the box `arm_boxed` leaked into its own node; the loop
// copied the node out (`load`) and fires it once.
R::run_from_loop_thread(unsafe { bun_core::heap::take(this) });
}

/// Stable Rust cannot take a fn value as a const generic, so `Callback` moves to
/// a runtime argument on `init` and is type-erased (ABI-identical: both forms
/// are thin fn pointers taking two thin data pointers).
Expand Down
27 changes: 27 additions & 0 deletions src/event_loop/ConcurrentTask.rs
Original file line number Diff line number Diff line change
Expand Up @@ -164,6 +164,33 @@ pub trait Taskable {
unsafe fn release_unrun(this: *mut Self);
}

/// [`Taskable`] for a type that is only ever queued as a leaked `Box<Self>`
/// (`heap::into_raw` / `Box::into_raw` at every post site, `Box::from_raw` in
/// its `bun_runtime::dispatch` arm). Released unrun by reclaiming the box and
/// handing it to `|this| $release` (default: drop it).
///
/// ```ignore
/// bun_event_loop::boxed_taskable!(ShellGlobTask, ShellGlobTask, |this| this.task.unref_unrun());
/// ```
#[macro_export]
macro_rules! boxed_taskable {
($ty:ty, $tag:ident) => {
$crate::boxed_taskable!($ty, $tag, |this| ());
};
($ty:ty, $tag:ident, |$this:ident| $release:expr) => {
impl $crate::Taskable for $ty {
const TAG: $crate::TaskTag = $crate::task_tag::$tag;
unsafe fn release_unrun(this: *mut Self) {
// SAFETY: `release_unrun` contract — `this` is the queued
// `Task::ptr` under `TAG`, which for this type is a leaked `Box<Self>`.
#[allow(unused_mut)]
let mut $this: ::std::boxed::Box<Self> = unsafe { ::bun_core::heap::take(this) };
$release;
}
}
};
}

impl TaskTag {
/// The tag's identifier, for diagnostics.
pub fn name(self) -> &'static str {
Expand Down
8 changes: 5 additions & 3 deletions src/event_loop/MiniEventLoop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -315,8 +315,9 @@ impl MiniEventLoop {
}

while let Some(task) = self.tasks.read_item() {
// SAFETY: tasks are pushed by enqueue_task* and remain valid until run() consumes them.
unsafe { (*task).run(context) };
// SAFETY: tasks are pushed by enqueue_task* and remain valid until loaded here.
let task = unsafe { AnyTaskWithExtraContext::load(task) };
task.run(context);
}
}

Expand All @@ -325,7 +326,8 @@ impl MiniEventLoop {
let _ = self.tick_concurrent_with_count();
while let Some(task) = self.tasks.read_item() {
// SAFETY: see tick_once.
unsafe { (*task).run(context) };
let task = unsafe { AnyTaskWithExtraContext::load(task) };
task.run(context);
}

// SAFETY: see `loop_ptr()` invariant.
Expand Down
2 changes: 1 addition & 1 deletion src/event_loop/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ pub use ConcurrentTask::{Task, TaskTag, Taskable, task_tag};
pub use DeferredTaskQueue as deferred_task_queue;

pub use any_event_loop::{
AnyEventLoop, EventLoopHandle, EventLoopTask, JsPoster, JsPosterVTable, Posted,
AnyEventLoop, ArmedLoopTask, EventLoopHandle, EventLoopTask, JsPoster, JsPosterVTable, Posted,
};
pub use bun_io::PipeReadScratch;

Expand Down
Loading
Loading