Skip to content
Merged
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
34 changes: 23 additions & 11 deletions src/bundler/ParseTask.rs
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,7 @@ pub enum ParseTaskStage {
// ───────────────────────────────────────────────────────────────────────────

/// The information returned to the Bundler thread when a parse finishes.
pub(crate) struct Result {
pub struct Result {
pub(crate) task: EventLoop::Task,
pub(crate) ctx: bun_ptr::ParentRef<BundleV2<'static>, bun_ptr::Mut>,
pub(crate) value: ResultValue,
Expand All @@ -138,6 +138,24 @@ pub(crate) struct Result {
/// returned source code by the plugin.
pub(crate) external: ExternalFreeFunction,
}
impl bun_event_loop::Taskable for Result {
const TAG: bun_event_loop::TaskTag = bun_event_loop::task_tag::BundleV2ParseTaskResult;
/// The VM that runs the bundle is going: the result goes to nobody.
unsafe fn release_unrun(this: *mut Self) {
// SAFETY: fn contract — the box `run_from_thread_pool_impl` (or
// `ServerComponentParseTask`) leaked.
let mut result = unsafe { bun_core::heap::take(this) };
// A native plugin's source buffer: `on_parse_task_complete` would have handed this to
// the bundle's finalizers. The source may borrow the buffer, so it goes first.
let external = core::mem::take(&mut result.external);
drop(result);
external.call();
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
/// A step of the bundle.
unsafe fn context(_: *const Self) -> bun_event_loop::ContextId {
bun_event_loop::ContextId::NONE
}
}
// `Result` lives in a bump arena (no Drop on free); boxing the large arm
// would leak the heap allocation. The size diff is acceptable.
#[allow(clippy::large_enum_variant)]
Expand Down Expand Up @@ -2889,13 +2907,7 @@ pub mod parse_worker {
.expect("BundleV2.linker.loop must be set before scheduling ParseTask")
{
bun_event_loop::AnyEventLoop::Js { .. } => {
let ct =
bun_event_loop::ConcurrentTask::ConcurrentTask::from_callback(result, |p| {
// SAFETY: `p` is the `result` Box leaked above; ownership
// transfers to `on_complete`, which deallocates it.
unsafe { on_complete(p) };
Ok(())
});
let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::create_from(result);
let poster = worker
.ctx
.js_poster
Expand All @@ -2906,7 +2918,7 @@ pub mod parse_worker {
// SAFETY: refused ⇒ we own the task box and the leaked result.
unsafe {
bun_event_loop::ConcurrentTask::ConcurrentTask::release_refused(ct);
drop(bun_core::heap::take(result));
<Result as bun_event_loop::Taskable>::release_unrun(result);
}
}
}
Expand Down Expand Up @@ -2961,7 +2973,7 @@ pub mod parse_worker {
/// (or `ServerComponentParseTask`'s equivalent). Ownership transfers to
/// this fn, which deallocates `result` before returning. Must run on the
/// main/bundler thread (it dereferences `result.ctx` mutably).
pub(crate) unsafe fn on_complete(result: *mut Result) {
pub unsafe fn on_complete(result: *mut Result) {
// SAFETY: result allocated via heap::alloc above; uniquely owned here.
let r = unsafe { &mut *result };
let ctx = r.ctx;
Expand All @@ -2978,4 +2990,4 @@ pub mod parse_worker {
}
} // end mod parse_worker

pub(crate) use parse_worker::on_complete;
pub use parse_worker::on_complete;
9 changes: 2 additions & 7 deletions src/bundler/ServerComponentParseTask.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,12 +117,7 @@ fn task_callback_wrap(thread_pool_task: *mut ThreadPoolTask) {
.expect("BundleV2.linker.loop must be set before scheduling ServerComponentParseTask")
{
bun_event_loop::AnyEventLoop::Js { .. } => {
let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::from_callback(result, |p| {
// SAFETY: `p` is the `result` Box leaked above; ownership
// transfers to `on_complete`, which deallocates it.
unsafe { on_complete(p) };
Ok(())
});
let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::create_from(result);
let poster = worker
.ctx
.js_poster
Expand All @@ -133,7 +128,7 @@ fn task_callback_wrap(thread_pool_task: *mut ThreadPoolTask) {
// SAFETY: refused ⇒ we own the task box and the leaked result.
unsafe {
bun_event_loop::ConcurrentTask::ConcurrentTask::release_refused(ct);
drop(bun_core::heap::take(result));
<parse_task::Result as bun_event_loop::Taskable>::release_unrun(result);
}
}
}
Expand Down
100 changes: 66 additions & 34 deletions src/bundler/bundle_v2.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1143,6 +1143,27 @@ pub mod bv2_impl {
unsafe { (*(*this).bv2).plugin_context }
}
}
/// The plugins answered a [`Resolve`]: the hop back to the loop that runs the bundle, when
/// that is a JS loop. Same pointer as the request, its own tag.
#[repr(transparent)]
pub struct ResolveAnswered(Resolve);
impl bun_event_loop::Taskable for ResolveAnswered {
const TAG: bun_event_loop::TaskTag =
bun_event_loop::task_tag::BundleV2PluginResolveAnswered;
/// The VM that runs the bundle is going: the answer goes to nobody.
unsafe fn release_unrun(_: *mut Self) {}
/// A step of the bundle.
unsafe fn context(_: *const Self) -> bun_event_loop::ContextId {
bun_event_loop::ContextId::NONE
}
}
impl ResolveAnswered {
pub fn run(&mut self) {
// SAFETY: `bv2` is a live backref set in `Resolve::init`.
let bv2 = unsafe { &mut *self.0.bv2 };
BundleV2::on_resolve(&mut self.0, bv2);
}
}
impl Resolve {
pub(crate) fn init(bv2: &mut BundleV2<'_>, record: MiniImportRecord) -> Self {
Self {
Expand Down Expand Up @@ -1352,6 +1373,47 @@ pub mod bv2_impl {
unsafe { (*(*this).bv2).plugin_context }
}
}
/// The plugins answered a [`Load`]: as [`ResolveAnswered`].
#[repr(transparent)]
pub struct LoadAnswered(Load);
impl bun_event_loop::Taskable for LoadAnswered {
const TAG: bun_event_loop::TaskTag =
bun_event_loop::task_tag::BundleV2PluginLoadAnswered;
/// As `ResolveAnswered::release_unrun`.
unsafe fn release_unrun(_: *mut Self) {}
/// A step of the bundle.
unsafe fn context(_: *const Self) -> bun_event_loop::ContextId {
bun_event_loop::ContextId::NONE
}
}
impl LoadAnswered {
pub fn run(&mut self) {
// SAFETY: `bv2` is a live backref set in `Load::init`.
let bv2 = unsafe { &mut *self.0.bv2 };
BundleV2::on_load(&mut self.0, bv2);
}
}
/// A plugin `.defer()`red a [`Load`]: the notice to the loop that runs the bundle, when
/// that is a JS loop. Same pointer as the request, its own tag.
#[repr(transparent)]
pub struct LoadDeferred(Load);
impl bun_event_loop::Taskable for LoadDeferred {
const TAG: bun_event_loop::TaskTag =
bun_event_loop::task_tag::BundleV2PluginLoadDeferred;
/// As `ResolveAnswered::release_unrun`.
unsafe fn release_unrun(_: *mut Self) {}
/// A step of the bundle.
unsafe fn context(_: *const Self) -> bun_event_loop::ContextId {
bun_event_loop::ContextId::NONE
}
}
impl LoadDeferred {
pub fn run(&mut self) {
// SAFETY: `bv2` is a live backref set in `Load::init`.
let bv2 = unsafe { &mut *self.0.bv2 };
BundleV2::on_notify_defer(&mut self.0, bv2);
}
}
}
}

Expand Down Expand Up @@ -4509,9 +4571,8 @@ pub mod bv2_impl {
// mutate `graph` / allocate from `graph.heap` off-thread.
match self.any_loop_mut() {
bun_event_loop::AnyEventLoop::Js { .. } => {
let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::from_callback(
std::ptr::from_mut(load),
on_load_from_js_loop_raw,
let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::create_from(
std::ptr::from_mut(load).cast::<jsc_api::JSBundler::LoadAnswered>(),
);
let poster = self
.js_poster
Expand Down Expand Up @@ -4543,9 +4604,8 @@ pub mod bv2_impl {
// See `on_load_async` — must dispatch on the bundler's own loop.
match self.any_loop_mut() {
bun_event_loop::AnyEventLoop::Js { .. } => {
let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::from_callback(
std::ptr::from_mut(resolve),
on_resolve_from_js_loop_raw,
let ct = bun_event_loop::ConcurrentTask::ConcurrentTask::create_from(
std::ptr::from_mut(resolve).cast::<jsc_api::JSBundler::ResolveAnswered>(),
);
let poster = self
.js_poster
Expand Down Expand Up @@ -4586,20 +4646,6 @@ pub mod bv2_impl {
BundleV2::on_resolve(unsafe { &mut *resolve }, unsafe { &mut *this });
}

fn on_load_from_js_loop(load: &mut jsc_api::JSBundler::Load) {
// SAFETY: `bv2` is a live backref set in `Load::init`.
let bv2 = unsafe { &mut *load.bv2 };
BundleV2::on_load(load, bv2);
}

fn on_load_from_js_loop_raw(
load: *mut jsc_api::JSBundler::Load,
) -> bun_event_loop::JsResult<()> {
// SAFETY: `load` is a valid pointer set up by `from_callback`.
on_load_from_js_loop(unsafe { &mut *load });
Ok(())
}

impl<'a> BundleV2<'a> {
pub(crate) fn on_load(load: &mut jsc_api::JSBundler::Load, this: &mut BundleV2) {
if load.deferred_in.take() == Some(this.graph.defer_epoch) {
Expand Down Expand Up @@ -4787,20 +4833,6 @@ pub mod bv2_impl {
}
}

fn on_resolve_from_js_loop(resolve: &mut jsc_api::JSBundler::Resolve) {
// SAFETY: `bv2` is a live backref set in `Resolve::init`.
let bv2 = unsafe { &mut *resolve.bv2 };
BundleV2::on_resolve(resolve, bv2);
}

fn on_resolve_from_js_loop_raw(
resolve: *mut jsc_api::JSBundler::Resolve,
) -> bun_event_loop::JsResult<()> {
// SAFETY: `resolve` is a valid pointer set up by `from_callback`.
on_resolve_from_js_loop(unsafe { &mut *resolve });
Ok(())
}

impl<'a> BundleV2<'a> {
/// Re-run the idempotent barrel seeding pass after a plugin `onResolve` result patches a record that the importer's parse-completion pass saw unresolved and skipped (#40606).
fn schedule_barrel_imports_after_plugin_resolve(
Expand Down
59 changes: 25 additions & 34 deletions src/event_loop/ConcurrentTask.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
//! If `auto_delete` is true, the task is automatically deallocated when it's finished.
//! Otherwise, it's expected that the containing struct will deallocate the task.

use crate::ManagedTask;
use bun_threading::UnboundedQueue;
use bun_threading::unbounded_queue::{Link, Linked};

Expand Down Expand Up @@ -68,20 +67,33 @@ pub mod task_tag {
BundleV2DeferredBatchTask, // bun.bundle_v2.DeferredBatchTask
BundleV2PluginResolve, // bun.bundle_v2.Resolve (JS-thread hop)
BundleV2PluginLoad, // bun.bundle_v2.Load (JS-thread hop)
BundleV2PluginResolveAnswered, // bun.bundle_v2.Resolve (hop back to the bundle's loop)
BundleV2PluginLoadAnswered, // bun.bundle_v2.Load (hop back to the bundle's loop)
BundleV2PluginLoadDeferred, // bun.bundle_v2.Load (`.defer()` notice to the bundle's loop)
BundleV2ParseTaskResult, // bun.bundle_v2.ParseTask.Result
ChromePipeEvent,
CopyFileWindowsMkdirp,
WriteFileWindowsMkdirp,
DnsErrorDeferred,
ShellYesTask, // shell.Interpreter.Builtin.Yes.YesTask
Close,
CppTask,
DuplexUpgradeContext,
FetchTasklet,
FetchTaskletDeinit,
FetchTaskletPromiseSettle,
FetchTaskletRequestDrain,
FSWatchTask,
GetAddrInfoLibuvComplete,
GraphContextStopAgain,
GraphContextStopAndFree,
DeadContextStopAgain,
HandledPromise,
HotReloadTask,
HTMLRewriterBackgroundPull,
WatchReloadTask,
JSBundleCompletionTask,
JSCDeferredWorkTask,
ManagedTask,
NapiAsyncWork, // napi_async_work
NapiFinalizerTask,
NativePromiseContextDeferredDerefTask,
Expand All @@ -96,11 +108,18 @@ pub mod task_tag {
Read,
Readv,
FlushPendingFileSinkTask,
RunTestsTask,
RuntimeTranspilerStore,
S3HttpDownloadStreamingTask,
S3HttpSimpleTask,
SendQueueDeferred, // bun_runtime::ipc::SendQueue (close / after-close hop)
ServerAllConnectionsClosedTask,
HTTPServerDeinit,
HTTPSServerDeinit,
DebugHTTPServerDeinit,
DebugHTTPSServerDeinit,
HTTPAppClose,
HTTPSAppClose,
ShellAsync,
ShellCondExprStatTask,
ShellCpTask,
Expand All @@ -120,6 +139,7 @@ pub mod task_tag {
StreamPending,
ThreadSafeFunction,
ValkeyDeferredClose,
ValkeyDeferredFailure,
WindowsNamedPipeContext,
Write,
Writev,
Expand Down Expand Up @@ -250,19 +270,6 @@ impl Task {
}
}

// Taskable impls for the low-tier task wrappers defined in this crate.
impl Taskable for crate::ManagedTask::ManagedTask {
const TAG: TaskTag = task_tag::ManagedTask;
unsafe fn release_unrun(this: *mut Self) {
// SAFETY: fn contract — a queued ManagedTask is the heap box `new*` made.
unsafe { crate::ManagedTask::ManagedTask::release(this) }
}
/// A callback task always runs: a callback that continues some script enters that script's
/// context itself, so what it reports goes to nobody once the context has stopped.
unsafe fn context(_this: *const Self) -> ContextId {
ContextId::NONE
}
}
// ────────────────────────────────────────────────────────────────────────────

#[repr(C)]
Expand Down Expand Up @@ -347,16 +354,6 @@ impl ConcurrentTask {
Self::create(Task::init(task))
}

// callback returns `JsResult<()>` to match `ManagedTask::new`'s stored ABI;
// callers that have a `fn(*mut T)` should wrap it as `|p| { f(p); Ok(()) }` at the call site.
pub fn from_callback<T>(
ptr: *mut T,
callback: fn(*mut T) -> crate::JsResult<()>,
) -> core::ptr::NonNull<ConcurrentTask> {
bun_core::mark_binding!();
Self::create(ManagedTask::ManagedTask::new(ptr, callback))
}

pub fn from<T: Taskable>(
&mut self,
of: *mut T,
Expand Down Expand Up @@ -388,20 +385,14 @@ impl ConcurrentTask {
}

/// A weak poster got `task` back because the target VM has closed: free
/// it if it is a heap task (`create*`); an intrusive one belongs to its
/// container.
/// the carrier if it is a heap one (`create*`); an intrusive one belongs to
/// its container. What the task points at stays the poster's.
///
/// # Safety
/// `task` was just refused and is not queued anywhere.
pub unsafe fn release_refused(task: core::ptr::NonNull<ConcurrentTask>) {
// SAFETY: fn contract.
let inner = unsafe { Self::into_task(task) };
// A callback task (`from_callback`, `ManagedTask::new*`) owns a heap
// `ManagedTask` behind `task.ptr` as well.
if inner.tag == crate::task_tag::ManagedTask {
// SAFETY: as above; refused ⇒ ours.
unsafe { crate::ManagedTask::ManagedTask::release(inner.ptr.cast()) };
}
let _ = unsafe { Self::into_task(task) };
}

/// Returns whether this task should be automatically deallocated after execution.
Expand Down
Loading
Loading