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
31 changes: 31 additions & 0 deletions src/boringssl_sys/boringssl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -305,6 +305,37 @@ impl Drop for OwnedSslCtx {
}
}

/// Owns one `X509`; `X509_free`s it on drop.
#[repr(transparent)]
pub struct OwnedX509(core::ptr::NonNull<X509>);

impl OwnedX509 {
/// Parses one DER-encoded certificate; `None` if it does not parse.
pub fn from_der(der: &[u8]) -> Option<Self> {
let len = c_long::try_from(der.len()).ok()?;
let mut ptr = der.as_ptr();
// SAFETY: `d2i_X509` reads at most `len` bytes from `ptr`, which `der` covers.
let raw = unsafe { d2i_X509(core::ptr::null_mut(), &raw mut ptr, len) };
core::ptr::NonNull::new(raw).map(Self)
}

pub fn as_ptr(&self) -> *mut X509 {
self.0.as_ptr()
}

pub fn as_mut(&mut self) -> &mut X509 {
// SAFETY: we own the live `X509` for as long as `self` is borrowed.
unsafe { self.0.as_mut() }
}
}

impl Drop for OwnedX509 {
fn drop(&mut self) {
// SAFETY: we own exactly one reference, released once.
unsafe { X509_free(self.0.as_ptr()) }
}
}

/// Owns the `STACK_OF(GENERAL_NAME)` that `X509V3_EXT_d2i` returns for a
/// subjectAltName extension. Frees every `GENERAL_NAME` and then the stack.
pub struct GeneralNames(core::ptr::NonNull<struct_stack_st_GENERAL_NAME>);
Expand Down
92 changes: 89 additions & 3 deletions src/event_loop/ConcurrentTask.rs
Original file line number Diff line number Diff line change
Expand Up @@ -72,9 +72,10 @@ pub mod task_tag {
Close,
CppTask,
DuplexUpgradeContext,
FetchTasklet,
FetchTaskletDeinit,
FetchTasklet, // progress update (bun_runtime FetchTasklet)
FetchTaskletHandBack, // the HTTP thread is done with the fetch
FetchTaskletPromiseSettle,
FetchTaskletRequestDataDrain, // the streaming request body buffer drained
FSWatchTask,
GetAddrInfoLibuvComplete,
HotReloadTask,
Expand Down Expand Up @@ -132,6 +133,25 @@ pub struct Task {
pub ptr: *mut (),
}

/// A task whose queued pointer is a [`ThisPtr`](bun_ptr::ThisPtr) to
/// `Target`, which keeps itself alive for the task; the dispatcher runs it or
/// releases it through here. Implemented by a zero-sized hop type per tag so
/// the ref protocol lives next to `Target`.
pub trait TaskHop {
type Target;
/// The tag constant from [`task_tag`]; the `bun_runtime::dispatch` match
/// arms MUST agree.
const TAG: TaskTag;
fn run(this: bun_ptr::ThisPtr<Self::Target>) -> crate::JsResult<()>;
/// As [`Taskable::release_unrun`].
fn release_unrun(this: bun_ptr::ThisPtr<Self::Target>);

#[inline]
fn task(this: bun_ptr::ThisPtr<Self::Target>) -> Task {
Task::new(Self::TAG, this.as_ptr().cast::<()>())
}
}

/// What it takes to be queued as a [`Task`]: a tag, and how the task is
/// freed when it will never run. Implement on every type that can be
/// enqueued; the impl lives in whatever crate owns the type.
Expand Down Expand Up @@ -216,6 +236,10 @@ pub struct ConcurrentTask {
/// dispatch. Immutable after construction; read only on the consumer thread,
/// so it does not need to share a word with the contended `next` link.
pub auto_delete: bool,
/// A [`ReusableConcurrentTask`] is armed (posted and not yet consumed).
/// Cleared by the consumer after its last read of the node; `false` for a
/// freshly built node.
pub queued: core::sync::atomic::AtomicBool,
}

impl Default for ConcurrentTask {
Expand All @@ -226,6 +250,52 @@ impl Default for ConcurrentTask {
task: unsafe { bun_core::ffi::zeroed_unchecked() },
next: Link::new(),
auto_delete: false,
queued: core::sync::atomic::AtomicBool::new(false),
}
}
}

/// An intrusive [`ConcurrentTask`] its owner posts again and again, at most one
/// post in flight: [`arm`](Self::arm) hands the node out only once the consumer
/// is done with the previous post.
pub struct ReusableConcurrentTask(core::cell::UnsafeCell<ConcurrentTask>);

// SAFETY: the node is written only by the thread that wins `queued`
// (false -> true) and read only by the consumer, which clears `queued` after
// its last read; `queued` itself is atomic.
unsafe impl Sync for ReusableConcurrentTask {}
// SAFETY: as above; the payload is a tag and an address.
unsafe impl Send for ReusableConcurrentTask {}

impl Default for ReusableConcurrentTask {
fn default() -> Self {
Self(core::cell::UnsafeCell::new(ConcurrentTask::default()))
}
}

impl ReusableConcurrentTask {
/// The node loaded with `task`, ready to post; `None` while the previous
/// post has not been consumed yet.
pub fn arm(&self, task: Task) -> Option<core::ptr::NonNull<ConcurrentTask>> {
let node = self.0.get();
// SAFETY: `queued` is atomic; winning false -> true makes this thread the
// node's only writer until the consumer clears it (`consumed`).
unsafe {
if (*node)
.queued
.compare_exchange(
false,
true,
core::sync::atomic::Ordering::AcqRel,
core::sync::atomic::Ordering::Acquire,
)
.is_err()
{
return None;
}
(*node).task = task;
(*node).auto_delete = false;
Some(core::ptr::NonNull::new_unchecked(node))
}
}
}
Expand All @@ -241,7 +311,7 @@ impl Default for ConcurrentTask {
const _: () = assert!(
core::mem::size_of::<ConcurrentTask>()
== core::mem::size_of::<Task>() + 2 * core::mem::size_of::<usize>(),
"ConcurrentTask = Task + next ptr + auto_delete (padded)"
"ConcurrentTask = Task + next ptr + auto_delete/queued (padded)"
);

// SAFETY: `link()` always projects to the same embedded `next: Link<Self>`
Expand Down Expand Up @@ -277,6 +347,7 @@ impl ConcurrentTask {
task,
next: Link::new(),
auto_delete: true,
queued: core::sync::atomic::AtomicBool::new(false),
});
// SAFETY: `new` heap-allocates via `heap::into_raw` — never null.
unsafe { core::ptr::NonNull::new_unchecked(raw) }
Expand Down Expand Up @@ -307,6 +378,7 @@ impl ConcurrentTask {
task: Task::init(of),
next: Link::new(),
auto_delete: auto_deinit == AutoDeinit::AutoDeinit,
queued: core::sync::atomic::AtomicBool::new(false),
};
self
}
Expand All @@ -322,11 +394,25 @@ impl ConcurrentTask {
let (task, auto_delete) = (this.as_ref().task, this.as_ref().auto_delete());
if auto_delete {
drop(bun_core::heap::take(this.as_ptr()));
} else {
Self::consumed(this);
}
task
}
}

/// Consuming thread: done reading an intrusive node (its `task` copied,
/// the batch iterator past it); a [`ReusableConcurrentTask`] may be armed again.
///
/// # Safety
/// `this` is live and not read again by the consumer for this post.
pub unsafe fn consumed(this: core::ptr::NonNull<ConcurrentTask>) {
// SAFETY: fn contract.
unsafe { this.as_ref() }
.queued
.store(false, core::sync::atomic::Ordering::Release);
}

/// 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.
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 @@ -29,7 +29,7 @@ pub mod any_event_loop;
// ─── public surface ─────────────────────────────────────────────────────────

pub type JsResult<T> = core::result::Result<T, bun_core::JsError>;
pub use ConcurrentTask::{Task, TaskTag, Taskable, task_tag};
pub use ConcurrentTask::{ReusableConcurrentTask, Task, TaskHop, TaskTag, Taskable, task_tag};

// snake_case alias for the file-level-struct module so higher tiers avoid
// the type/module namespace collision on the PascalCase form.
Expand Down
Loading
Loading