use crate::modules::input::token;
use crate::{
Runtime, RuntimeError,
constants::{
DEAD_KQUEUE_ID, MANAGER_TICK, MANAGER_TICK_IDENT, MAX_TASK_ID, NO_SELECT, NO_TASK,
PARK_TIMER, RESTART_BACKOFF, RESTART_LIMIT, RESTART_WINDOW, SCHEDULE_IDENT_BASE,
SELECT_IDENT, SELECT_POLL, SHUTDOWN_POLL, TIMEOUT_IDENT_BASE, UNSTARTED_TASK_ID,
WAKE_IDENT,
},
futures::task::{Task, sealed::Park},
modules::{
address_lock,
erased_task::ErasedTask,
event_desc::EventDesc,
extras::Extras,
faults,
forward::Forward,
gate::{AfterRun, Gate, Trigger},
gated::Gated,
gather::{Gather, access},
handle_kind::Waiting,
handle_set::HandleSet,
help::{self, Patience},
input::Receives,
int_check::IntCheck,
kevent::{KEvent, eventlist},
kqueue,
mailbox::Mailbox,
merge_set::MergeSet,
series::SeriesTask,
task_data::TaskData,
task_data::deadline_epoch,
task_handle::TaskHandle,
task_setup::TaskSetup,
task_state::TaskState,
task_table::TaskTable,
tuning::{self, Tuning},
worker_pool::POOL,
},
};
use libc::c_void;
use std::{
cell::Cell,
mem,
panic::{self, AssertUnwindSafe},
ptr,
sync::{
Arc, Mutex,
atomic::{AtomicI32, AtomicU8, AtomicU32, AtomicU64, Ordering},
},
thread::{self, JoinHandle},
time::{Duration, Instant},
};
static DATA: TaskTable = TaskTable::new();
static EXECUTOR_KQUEUE_ID: AtomicI32 = AtomicI32::new(DEAD_KQUEUE_ID);
static SEQUENCE: AtomicU64 = AtomicU64::new(0);
thread_local! {
static CURRENT: Cell<usize> = const { Cell::new(NO_TASK) };
}
const STOPPED: u8 = 0;
const STARTING: u8 = 1;
const RUNNING: u8 = 2;
const STOPPING: u8 = 3;
static LIFECYCLE: AtomicU8 = AtomicU8::new(STOPPED);
static SUPERVISOR: Mutex<Option<JoinHandle<()>>> = Mutex::new(None);
#[inline(always)]
pub(crate) fn shutting_down() -> bool {
matches!(LIFECYCLE.load(Ordering::SeqCst), STOPPING | STOPPED)
}
pub(crate) fn shutdown_now() {
loop {
match LIFECYCLE.compare_exchange(RUNNING, STOPPING, Ordering::SeqCst, Ordering::SeqCst) {
Ok(_) => break,
Err(STOPPED) => return,
Err(STOPPING) => {
while LIFECYCLE.load(Ordering::SeqCst) == STOPPING {
thread::sleep(SHUTDOWN_POLL);
}
return;
}
Err(_) => thread::sleep(SHUTDOWN_POLL),
}
}
POOL.close();
let manager = EXECUTOR_KQUEUE_ID.swap(DEAD_KQUEUE_ID, Ordering::SeqCst);
if manager != DEAD_KQUEUE_ID {
let _ = unsafe {
KEvent::register(
manager,
WAKE_IDENT,
0,
ptr::null_mut(),
EventDesc::new_user_trigger(),
)
}
.check();
}
let supervisor = SUPERVISOR
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.take();
if let Some(supervisor) = supervisor {
let _ = supervisor.join();
}
while POOL.stats().has_any_task() {
thread::sleep(SHUTDOWN_POLL);
}
loop {
POOL.stop_all();
POOL.abandon();
if POOL.live() == 0 && POOL.sleeps_live() == 0 {
break;
}
thread::sleep(SHUTDOWN_POLL);
}
for task in POOL
.injector()
.drain()
.into_iter()
.chain(POOL.blocking().drain())
{
fail(task);
}
for task in 0..DATA.high_water() {
let Some(data) = slot(task) else {
continue;
};
data.disarm();
if data.parked() {
unpark(task, data);
continue;
}
let state = data.state();
if state == TaskState::Running && !data.kind().schedules() {
continue;
}
if !state.terminal() {
data.set_state(TaskState::Failed);
wake(data);
}
close_extras(data);
release(task);
}
LIFECYCLE.store(STOPPED, Ordering::SeqCst);
}
static INJECTED_FAULTS: AtomicU32 = AtomicU32::new(0);
#[cfg(feature = "fault-injection")]
pub(crate) fn inject_manager_faults(count: u32) {
INJECTED_FAULTS.store(count, Ordering::SeqCst);
}
fn injected_fault() -> bool {
INJECTED_FAULTS
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |left| match left {
0 => None,
_ => Some(left - 1),
})
.is_ok()
}
pub(crate) struct Executor;
impl Executor {
pub(crate) fn init(tuning: Tuning) -> Result<(), RuntimeError> {
loop {
match LIFECYCLE.compare_exchange(STOPPED, STARTING, Ordering::SeqCst, Ordering::SeqCst)
{
Ok(_) => break,
Err(RUNNING) => return Err(RuntimeError::AlreadyInit),
Err(_) => thread::sleep(SHUTDOWN_POLL),
}
}
let id = match unsafe { libc::kqueue() }.check() {
Ok(id) => id,
Err(error) => {
LIFECYCLE.store(STOPPED, Ordering::SeqCst);
return Err(error);
}
};
EXECUTOR_KQUEUE_ID.store(id, Ordering::SeqCst);
tuning::apply(tuning);
let _ = deadline_epoch();
POOL.open();
POOL.ensure_floor();
supervise(id);
LIFECYCLE.store(RUNNING, Ordering::SeqCst);
Ok(())
}
pub(crate) fn new_task<F>(task: F, setup: TaskSetup) -> TaskHandle<F::Output>
where
F: Task,
{
create(task, setup, None).0
}
pub(crate) fn new_series<F>(task: F, setup: TaskSetup) -> TaskHandle<F::Output>
where
F: Task + Clone,
{
create_series(task, setup, None)
}
pub(crate) fn new_waiting<F, T, M>(
task: F,
setup: TaskSetup,
) -> TaskHandle<F::Output, Waiting<T>>
where
F: Task,
F::Input: Receives<T, M>,
T: Send + 'static,
M: 'static,
{
let (gate, mailbox, setup) = waiting_parts::<T>(setup);
let gated = Gated::<F, T, M>::new(task, Arc::clone(&mailbox));
create(gated, setup, Some(gate)).0.into_kind(mailbox)
}
pub(crate) fn new_waiting_series<F, T, M>(
task: F,
setup: TaskSetup,
) -> TaskHandle<F::Output, Waiting<T>>
where
F: Task + Clone,
F::Input: Receives<T, M>,
T: Send + 'static,
M: 'static,
{
let (gate, mailbox, setup) = waiting_parts::<T>(setup);
let gated = Gated::<F, T, M>::new(task, Arc::clone(&mailbox));
create_series(gated, setup, Some(gate)).into_kind(mailbox)
}
pub(crate) fn give<T>(id: usize, mailbox: &Mailbox<T>, value: T) -> Result<(), RuntimeError> {
let Some(data) = slot(id) else {
return Err(missing(id));
};
let gate = mailbox.gate();
if gate.finished()
|| matches!(
data.state(),
TaskState::Cancelled | TaskState::TimedOut | TaskState::Failed
)
{
return Err(closed(data));
}
drop(mailbox.replace(value));
match gate.trigger() {
Trigger::Replaced => Ok(()),
Trigger::Closed => Err(closed(data)),
Trigger::Start => match start_waiting(id, data, gate) {
true => Ok(()),
false => Err(RuntimeError::TaskFailed),
},
}
}
pub(crate) fn forward(upstream: usize, forward: Box<dyn Forward>) {
let Some(data) = slot(upstream) else {
drop(forward);
return;
};
data.extras_or_attach().receivers().register(forward, data);
}
pub(crate) fn receive_all<O, H>(
waiting: TaskHandle<O, Waiting<H::Output>>,
set: H,
) -> TaskHandle<O>
where
H: HandleSet,
{
let gather = Arc::new(Gather::<H>::new(set.slots(), waiting.retyped()));
let mut held = Vec::new();
set.link(&gather, access(|root: &mut H::Slots| root), &mut held);
hold(waiting.id(), held);
gather.settle();
waiting.into_plain()
}
pub(crate) fn receive_any<O, V, M, H>(
waiting: TaskHandle<O, Waiting<H::Given>>,
set: H,
) -> TaskHandle<O>
where
H: MergeSet<V, M>,
{
let target = waiting.retyped();
let mut held = Vec::new();
set.link(&target, &mut held);
hold(waiting.id(), held);
drop(target);
waiting.into_plain()
}
pub(crate) fn add_listener(id: usize) {
let Some(data) = slot(id) else {
return;
};
data.add_listener();
}
pub(crate) fn drop_listener(id: usize) {
let Some(data) = slot(id) else {
return;
};
if !data.drop_listener() {
return;
}
unsafe { data.destroy() };
DATA.free(id);
}
pub(crate) fn finished(id: usize) -> bool {
let Some(data) = slot(id) else {
return true;
};
let state = data.state();
match state {
TaskState::Cancelled | TaskState::TimedOut | TaskState::Failed => true,
_ => state.terminal() && !data.kind().repeats() && !data.open_for_gives(),
}
}
pub(crate) fn state(id: usize) -> TaskState {
match slot(id) {
Some(data) => data.state(),
None => TaskState::Failed,
}
}
pub(crate) fn wait(id: usize) -> Result<TaskState, RuntimeError> {
Self::wait_until(id, None)
}
pub(crate) fn wait_until(
id: usize,
deadline: Option<Instant>,
) -> Result<TaskState, RuntimeError> {
let Some(data) = slot(id) else {
return Err(missing(id));
};
let mut patience = Patience::new();
loop {
let state = data.state();
if state.terminal() {
return Ok(state);
}
if help::run_awaited(id) || help::help_once() {
patience.helped();
continue;
}
let slice = patience.slice();
let sleep = match deadline {
None => slice,
Some(deadline) => {
let left = deadline.saturating_duration_since(Instant::now());
if left.is_zero() {
return Err(RuntimeError::NotReady);
}
Some(slice.map_or(left, |slice| slice.min(left)))
}
};
let Some(sleep) = sleep else {
address_lock::wait(data.wait_address(), state as u32)?;
continue;
};
if address_lock::wait_until(data.wait_address(), state as u32, sleep)? {
continue;
}
if deadline.is_none_or(|deadline| Instant::now() < deadline) {
continue;
}
let state = data.state();
if state.terminal() {
return Ok(state);
}
return Err(RuntimeError::NotReady);
}
}
pub(crate) fn clone_result<T>(id: usize) -> Result<T, RuntimeError>
where
T: Clone,
{
Self::clone_result_until(id, None)
}
pub(crate) fn clone_result_until<T>(
id: usize,
deadline: Option<Instant>,
) -> Result<T, RuntimeError>
where
T: Clone,
{
let Some(data) = slot(id) else {
return Err(missing(id));
};
loop {
settled(data, Self::wait_until(id, deadline)?)?;
if data.enter_read() {
break;
}
let error = lost(data, data.state());
if error != RuntimeError::NotReady {
return Err(error);
}
if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
return Err(RuntimeError::NotReady);
}
}
debug_assert_eq!(data.size(), mem::size_of::<T>());
let value = unsafe { (*data.payload().cast::<T>()).clone() };
data.leave_read();
Ok(value)
}
pub(crate) fn poll_result<T>(id: usize) -> Result<T, RuntimeError>
where
T: Clone,
{
let Some(data) = slot(id) else {
return Err(missing(id));
};
if !data.enter_read() {
return Err(lost(data, data.state()));
}
debug_assert_eq!(data.size(), mem::size_of::<T>());
let value = unsafe { (*data.payload().cast::<T>()).clone() };
data.leave_read();
Ok(value)
}
pub(crate) fn poll_take<T>(id: usize) -> Result<T, RuntimeError> {
let Some(data) = slot(id) else {
return Err(missing(id));
};
if !data.claim_result() {
return Err(lost(data, data.state()));
}
debug_assert_eq!(data.size(), mem::size_of::<T>());
data.empty();
Ok(unsafe { ptr::read(data.payload().cast::<T>()) })
}
pub(crate) fn take_result<T>(id: usize) -> Result<T, RuntimeError> {
Self::take_result_until(id, None)
}
pub(crate) fn take_result_until<T>(
id: usize,
deadline: Option<Instant>,
) -> Result<T, RuntimeError> {
let Some(data) = slot(id) else {
return Err(missing(id));
};
loop {
settled(data, Self::wait_until(id, deadline)?)?;
if data.claim_result() {
break;
}
let error = lost(data, data.state());
if error != RuntimeError::NotReady {
return Err(error);
}
if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
return Err(RuntimeError::NotReady);
}
}
debug_assert_eq!(data.size(), mem::size_of::<T>());
data.empty();
Ok(unsafe { ptr::read(data.payload().cast::<T>()) })
}
pub(crate) fn cancel(id: usize) {
if let Some(data) = slot(id) {
stop(id, data, TaskState::Cancelled);
}
}
}
fn stop(id: usize, data: &TaskData, to: TaskState) {
loop {
let state = data.state();
let repeats = data.kind().repeats();
match state {
TaskState::Cancelled | TaskState::TimedOut | TaskState::Failed | TaskState::Free => {
return;
}
TaskState::Taken if !repeats => return,
_ => {}
}
if data.try_state(state, to) {
stopped(id, data);
return;
}
}
}
fn stopped(id: usize, data: &TaskData) {
wake(data);
interrupt(data, id);
unpark(id, data);
if data.takes_input() {
let_go_waiting(id, data);
}
}
fn create_series<F>(task: F, setup: TaskSetup, gate: Option<Arc<Gate>>) -> TaskHandle<F::Output>
where
F: Task + Clone,
{
let boxed: Box<dyn SeriesTask> = Box::new(task);
let prototype = Box::into_raw(Box::new(boxed)).cast::<c_void>();
if !Runtime::initialised() {
drop(unsafe { Box::from_raw(prototype.cast::<Box<dyn SeriesTask>>()) });
return TaskHandle::new(UNSTARTED_TASK_ID);
}
let Some(id) = DATA.alloc() else {
return abandoned(prototype);
};
let Some(entry) = DATA.slot(id) else {
DATA.free(id);
return abandoned(prototype);
};
let ready = unsafe {
TaskData::init::<F::Output>(
entry as *const TaskData as *mut TaskData,
ptr::null_mut(),
TaskState::Pending,
setup,
SEQUENCE.fetch_add(1, Ordering::Relaxed),
)
};
if !ready {
DATA.free(id);
return abandoned(prototype);
}
let waits = gate.is_some();
entry.attach_extras(Box::new(Extras::new(gate).timed(setup.timeout, NO_TASK)));
entry.set_prototype(prototype);
let handle = TaskHandle::new(id);
if waits {
return handle;
}
let started = match setup.start_delay.as_nanos() as u64 {
0 => start_series(id, entry),
delay => wait_for(entry, id, delay),
};
if started {
return handle;
}
unschedule(id);
if entry.try_state(TaskState::Pending, TaskState::Failed) {
wake(entry);
}
release(id);
handle
}
fn create<F>(task: F, setup: TaskSetup, gate: Option<Arc<Gate>>) -> (TaskHandle<F::Output>, bool)
where
F: Task,
{
let setup = setup.blocking(task.blocking(token()));
let boxed: Box<dyn ErasedTask> = Box::new(task);
let erased = Box::into_raw(Box::new(boxed)).cast::<c_void>();
if !Runtime::initialised() {
return (unstarted(erased), false);
}
let Some(id) = DATA.alloc() else {
return (failed(erased), false);
};
let Some(entry) = DATA.slot(id) else {
DATA.free(id);
return (failed(erased), false);
};
let ready = unsafe {
TaskData::init::<F::Output>(
entry as *const TaskData as *mut TaskData,
erased,
TaskState::Pending,
setup,
SEQUENCE.fetch_add(1, Ordering::Relaxed),
)
};
if !ready {
DATA.free(id);
return (failed(erased), false);
}
let handle = TaskHandle::new(id);
if gate.is_some() || setup.timeout.is_some() {
let waits = gate.is_some();
entry.attach_extras(Box::new(
Extras::new(gate).timed(setup.timeout, setup.parent),
));
if waits {
return (handle, true);
}
}
let started = match setup.start_delay.as_nanos() as u64 {
0 => queue_spawned(id, setup.blocking),
delay => wait_for(entry, id, delay),
};
if started {
return (handle, true);
}
if let Some(published) = slot(id) {
published.set_state(TaskState::Failed);
wake(published);
}
release(id);
(handle, false)
}
pub(crate) fn spawn_run<F>(task: F, priority: u8, timeout: Option<Duration>, series: usize) -> bool
where
F: Task,
{
let setup = TaskSetup {
timeout,
parent: series,
..TaskSetup::once(priority)
};
create(task, setup, None).1
}
pub(crate) fn publish<T>(id: usize, value: T) {
let Some(data) = slot(id) else {
return;
};
if !data.begin() {
close_finished_series(data);
return;
}
debug_assert_eq!(data.size(), mem::size_of::<T>());
unsafe { data.payload().cast::<T>().write(value) };
data.fill();
note_output(data);
if !data.try_state(TaskState::Running, TaskState::Ready) {
return;
}
wake(data);
forward_output(data);
close_finished_series(data);
}
#[inline(always)]
pub(crate) fn slot(id: usize) -> Option<&'static TaskData> {
let data = DATA.slot(id)?;
if data.state() == TaskState::Free {
return None;
}
Some(data)
}
#[inline(always)]
pub(crate) fn queue_link(id: usize) -> Option<usize> {
Some(DATA.slot(id)?.queue_next())
}
#[inline(always)]
pub(crate) fn set_queue_link(id: usize, next: usize) {
if let Some(data) = DATA.slot(id) {
data.set_queue_next(next);
}
}
#[inline(always)]
pub(crate) fn manager_alive() -> bool {
EXECUTOR_KQUEUE_ID.load(Ordering::Relaxed) != DEAD_KQUEUE_ID
}
#[inline(always)]
pub(crate) fn peak_slots() -> usize {
DATA.peak()
}
#[inline(always)]
pub(crate) fn slots() -> usize {
DATA.high_water()
}
#[inline(always)]
pub(crate) fn live() -> usize {
DATA.live()
}
#[inline(always)]
pub(crate) fn trim() -> Result<usize, RuntimeError> {
DATA.trim()
}
#[inline(always)]
pub(crate) fn sequence() -> u64 {
SEQUENCE.load(Ordering::Relaxed)
}
pub(crate) fn run(id: usize) {
let Some(data) = slot(id) else {
return;
};
let raw = data.claim();
if raw.is_null() {
return;
}
let mut task = unsafe { Box::from_raw(raw.cast::<Box<dyn ErasedTask>>()) };
let resumed = match data.begin() {
true => false,
false if data.state() == TaskState::Running => true,
false => {
close_extras(data);
drop(task);
release(id);
return;
}
};
data.clear_start_delay();
if !resumed {
time_run(id, data);
}
let reactor = Runtime::reactor_id();
let payload = data.payload();
let outer = CURRENT.with(|current| current.replace(id));
let stepped = panic::catch_unwind(AssertUnwindSafe(|| unsafe {
task.run(reactor, id, payload, resumed)
}));
CURRENT.with(|current| current.set(outer));
let finished = match stepped {
Ok(Some(park)) => {
park_task(id, data, task, park);
return;
}
Ok(None) => true,
Err(payload) if faults::is_thread_death(payload.as_ref()) => {
drop(task);
panic::resume_unwind(payload);
}
Err(_) => false,
};
untime_run(id, data);
if !finished {
close_extras(data);
drop(task);
if data.try_state(TaskState::Running, TaskState::Failed) {
wake(data);
}
release(id);
return;
}
data.fill();
note_output(data);
if !data.try_state(TaskState::Running, TaskState::Ready) {
close_extras(data);
drop(task);
release(id);
return;
}
wake(data);
forward_output(data);
if !data.kind().repeats() {
if data.takes_input() {
rewait(id, data, task);
return;
}
close_extras(data);
drop(task);
release(id);
return;
}
let gap = match data.kind().waits() {
true => Duration::from_nanos(data.interval()),
false => Duration::ZERO,
};
if data.count_run() || data.past_deadline(gap) {
if data.takes_input() {
rewait(id, data, task);
return;
}
data.finish_series();
close_extras(data);
drop(task);
release(id);
return;
}
data.rearm(Box::into_raw(task).cast::<c_void>());
let armed = match data.kind().waits() {
true => wait_out(data, id),
false => queue(id, data.blocking()),
};
if armed {
return;
}
if data.try_state(TaskState::Ready, TaskState::Failed) {
wake(data);
}
close_extras(data);
release(id);
}
fn time_run(id: usize, data: &TaskData) {
let Some(timeout) = data.extras().and_then(Extras::timeout) else {
return;
};
let word = data.time_run(timeout);
let _ = arm_timeout(id, timeout, word);
}
#[inline(always)]
fn untime_run(id: usize, data: &TaskData) {
if !data.untime_run() {
return;
}
let manager = EXECUTOR_KQUEUE_ID.load(Ordering::Relaxed);
if manager == DEAD_KQUEUE_ID {
return;
}
let _ = unsafe {
KEvent::register(
manager,
id + TIMEOUT_IDENT_BASE,
0,
ptr::null_mut(),
EventDesc::new_timer_delete(),
)
}
.check();
}
fn arm_timeout(id: usize, left: Duration, word: u64) -> bool {
let manager = EXECUTOR_KQUEUE_ID.load(Ordering::Relaxed);
if manager == DEAD_KQUEUE_ID {
return false;
}
let nanos = left.as_nanos().clamp(1, libc::intptr_t::MAX as u128);
unsafe {
KEvent::register(
manager,
id + TIMEOUT_IDENT_BASE,
nanos as libc::intptr_t,
word as *mut c_void,
EventDesc::new_timer(),
)
}
.check()
.is_ok()
}
fn timed_out(id: usize, word: u64) {
let Some(data) = slot(id) else {
return;
};
if word == 0 || data.run_deadline() != word {
return;
}
let parent = data.extras().map_or(NO_TASK, Extras::parent);
let held = slot(parent).is_some_and(|series| {
series.add_listener();
true
});
if let Some(parked) = data.claim_parked() {
unwatch_park(id, parked, Fired::Neither);
let raw = data.claim();
if !raw.is_null() {
drop(unsafe { Box::from_raw(raw.cast::<Box<dyn ErasedTask>>()) });
}
if data.try_state(TaskState::Running, TaskState::TimedOut) {
close_extras(data);
wake(data);
}
release(id);
} else if !data.try_state(TaskState::Running, TaskState::TimedOut) {
if held {
Executor::drop_listener(parent);
}
return;
} else {
stopped(id, data);
}
if !held {
return;
}
if let Some(series) = slot(parent) {
stop(parent, series, TaskState::TimedOut);
}
Executor::drop_listener(parent);
}
#[inline(always)]
fn queue(id: usize, blocking: bool) -> bool {
match blocking {
true => POOL.offload(id),
false => POOL.submit(id),
}
}
#[inline(always)]
fn queue_spawned(id: usize, blocking: bool) -> bool {
if !blocking && CURRENT.with(|current| current.get()) != NO_TASK {
if let Some(worker) = help::worker() {
return POOL.submit_local(worker, id);
}
}
queue(id, blocking)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Fired {
Event,
Deadline,
Neither,
}
fn park_task(id: usize, data: &TaskData, task: Box<Box<dyn ErasedTask>>, park: Park) {
data.add_listener();
data.rearm(Box::into_raw(task).cast::<c_void>());
data.park(park.ident, park.filter, park.deadline.is_some());
let watched = watch_park(id, park);
if !watched || data.state() != TaskState::Running {
unpark(id, data);
}
Executor::drop_listener(id);
}
fn watch_park(id: usize, park: Park) -> bool {
let manager = EXECUTOR_KQUEUE_ID.load(Ordering::Relaxed);
if manager == DEAD_KQUEUE_ID {
return false;
}
if let Some(deadline) = park.deadline {
let left = deadline
.saturating_duration_since(Instant::now())
.as_nanos()
.clamp(1, libc::intptr_t::MAX as u128);
let timed = unsafe {
KEvent::register(
manager,
id + SCHEDULE_IDENT_BASE,
left as libc::intptr_t,
PARK_TIMER as *mut c_void,
EventDesc::new_timer(),
)
}
.check();
if timed.is_err() {
return false;
}
}
unsafe {
KEvent::register(
manager,
park.ident as usize,
0,
id as *mut c_void,
EventDesc::new_park(park.filter, park.notes),
)
}
.check()
.is_ok()
}
fn unwatch_park(id: usize, parked: (i32, i16, bool), fired: Fired) {
let manager = EXECUTOR_KQUEUE_ID.load(Ordering::Relaxed);
if manager == DEAD_KQUEUE_ID {
return;
}
let (ident, filter, timed) = parked;
if timed && fired != Fired::Deadline {
let _ = unsafe {
KEvent::register(
manager,
id + SCHEDULE_IDENT_BASE,
0,
ptr::null_mut(),
EventDesc::new_timer_delete(),
)
}
.check();
}
if fired != Fired::Event {
let _ = unsafe {
KEvent::register(
manager,
ident as usize,
0,
id as *mut c_void,
EventDesc::new_park_delete(filter),
)
}
.check();
}
}
fn wake_parked(id: usize, fired: Fired) {
let Some(data) = slot(id) else {
return;
};
let Some(parked) = data.claim_parked() else {
return;
};
unwatch_park(id, parked, fired);
if queue(id, data.blocking()) {
return;
}
take_down(id, data);
}
fn unpark(id: usize, data: &TaskData) {
let Some(parked) = data.claim_parked() else {
return;
};
unwatch_park(id, parked, Fired::Neither);
take_down(id, data);
}
fn take_down(id: usize, data: &TaskData) {
let raw = data.claim();
if !raw.is_null() {
drop(unsafe { Box::from_raw(raw.cast::<Box<dyn ErasedTask>>()) });
}
if data.try_state(TaskState::Running, TaskState::Failed) {
wake(data);
}
release(id);
}
fn wait_out(data: &TaskData, id: usize) -> bool {
wait_for(data, id, data.interval())
}
fn wait_for(data: &TaskData, id: usize, nanos: u64) -> bool {
data.arm();
if arm_timer(id, nanos) {
return true;
}
data.disarm();
false
}
fn arm_timer(id: usize, interval: u64) -> bool {
let manager = EXECUTOR_KQUEUE_ID.load(Ordering::Relaxed);
if manager == DEAD_KQUEUE_ID {
return false;
}
unsafe {
KEvent::register(
manager,
id + SCHEDULE_IDENT_BASE,
interval as libc::intptr_t,
ptr::null_mut(),
EventDesc::new_timer(),
)
}
.check()
.is_ok()
}
fn fire(ident: usize) {
let Some(id) = ident.checked_sub(SCHEDULE_IDENT_BASE) else {
return;
};
let Some(data) = slot(id) else {
return;
};
if data.kind().schedules() {
if data.armed() && !data.claim_armed() {
return;
}
tick(id, data);
return;
}
if !data.claim_armed() {
return;
}
if queue(id, data.blocking()) {
return;
}
if data.try_state(TaskState::Ready, TaskState::Failed)
|| data.try_state(TaskState::Taken, TaskState::Failed)
|| data.try_state(TaskState::Pending, TaskState::Failed)
{
wake(data);
}
close_extras(data);
release(id);
}
fn start_series(id: usize, data: &TaskData) -> bool {
if !data.runs_remain() {
return false;
}
if !launch(id, data) {
return false;
}
if data.count_run() {
match data.takes_input() {
true => rewait_series(id, data),
false => {
data.finish_series();
close_finished_series(data);
release(id);
}
}
return true;
}
schedule(id, data.interval())
}
fn tick(id: usize, data: &TaskData) {
if data.start_delay() != 0 {
data.clear_start_delay();
if !schedule(id, data.interval()) {
end_schedule(id, data);
return;
}
}
if data.past_deadline(Duration::from_nanos(data.interval())) || !data.runs_remain() {
end_schedule(id, data);
return;
}
if !over(data.state()) && launch(id, data) {
if data.count_run() {
end_schedule(id, data);
}
return;
}
unschedule(id);
if !data.state().terminal() {
data.set_state(TaskState::Failed);
wake(data);
}
close_extras(data);
release(id);
}
fn end_schedule(id: usize, data: &TaskData) {
unschedule(id);
if data.takes_input() {
rewait_series(id, data);
return;
}
data.finish_series();
close_finished_series(data);
release(id);
}
#[inline(always)]
fn over(state: TaskState) -> bool {
matches!(
state,
TaskState::Cancelled | TaskState::TimedOut | TaskState::Failed
)
}
fn launch(id: usize, data: &TaskData) -> bool {
let prototype = data.prototype();
if prototype.is_null() {
return false;
}
let task = unsafe { &**prototype.cast::<Box<dyn SeriesTask>>() };
let timeout = data.extras().and_then(Extras::timeout);
task.launch(id, data.priority_class(), timeout)
}
fn schedule(id: usize, interval: u64) -> bool {
let manager = EXECUTOR_KQUEUE_ID.load(Ordering::Relaxed);
if manager == DEAD_KQUEUE_ID {
return false;
}
unsafe {
KEvent::register(
manager,
id + SCHEDULE_IDENT_BASE,
interval as libc::intptr_t,
ptr::null_mut(),
EventDesc::new_interval(),
)
}
.check()
.is_ok()
}
fn unschedule(id: usize) {
let manager = EXECUTOR_KQUEUE_ID.load(Ordering::Relaxed);
if manager == DEAD_KQUEUE_ID {
return;
}
let _ = unsafe {
KEvent::register(
manager,
id + SCHEDULE_IDENT_BASE,
0,
ptr::null_mut(),
EventDesc::new_timer_delete(),
)
}
.check();
}
pub(crate) fn waiting_on(queue: i32) -> bool {
let id = CURRENT.with(|current| current.get());
if id == NO_TASK {
return true;
}
let Some(data) = slot(id) else {
return true;
};
if data.state().stopped() {
return false;
}
data.set_waiting(queue);
if data.state().stopped() {
data.clear_waiting();
return false;
}
true
}
pub(crate) fn stopped_waiting() -> bool {
let id = CURRENT.with(|current| current.get());
if id == NO_TASK {
return true;
}
let Some(data) = slot(id) else {
return true;
};
data.clear_waiting();
!data.state().stopped()
}
pub(crate) fn cancelled() -> bool {
let id = CURRENT.with(|current| current.get());
if id == NO_TASK {
return false;
}
let Some(data) = slot(id) else {
return false;
};
data.state().stopped()
}
fn interrupt(data: &TaskData, id: usize) {
let Some(queue) = data.claim_waiting() else {
return;
};
let _ =
unsafe { KEvent::register(queue, id, 0, ptr::null_mut(), EventDesc::new_timer_delete()) }
.check();
let _ = unsafe {
KEvent::register(
queue,
WAKE_IDENT,
0,
ptr::null_mut(),
EventDesc::new_user_trigger(),
)
}
.check();
data.release_waiting();
}
pub(crate) fn fail(id: usize) {
let Some(data) = slot(id) else {
return;
};
if !data.state().terminal() {
data.set_state(TaskState::Failed);
wake(data);
}
close_extras(data);
release(id);
}
#[inline(always)]
fn settled(data: &TaskData, state: TaskState) -> Result<(), RuntimeError> {
if state == TaskState::Ready {
return Ok(());
}
Err(lost(data, state))
}
#[inline(always)]
fn lost(data: &TaskData, state: TaskState) -> RuntimeError {
match state {
TaskState::Taken => {
match !data.kind().repeats() && data.spent() && !data.open_for_gives() {
true => RuntimeError::Finished,
false => RuntimeError::AlreadyTaken,
}
}
TaskState::Cancelled => RuntimeError::Cancelled,
TaskState::TimedOut => RuntimeError::TimedOut,
TaskState::Failed => RuntimeError::TaskFailed,
TaskState::Pending | TaskState::Running | TaskState::Ready => RuntimeError::NotReady,
TaskState::Free => RuntimeError::NoSuchTask,
}
}
#[inline(always)]
fn wake(data: &TaskData) {
address_lock::wake(data.wait_address());
let queue = data.select_queue();
if queue == NO_SELECT {
return;
}
let _ = unsafe {
KEvent::register(
queue,
SELECT_IDENT,
0,
ptr::null_mut(),
EventDesc::new_user_trigger(),
)
}
.check();
}
pub(crate) fn join_first(ids: &[usize]) -> Option<usize> {
if ids.is_empty() {
return None;
}
if let Some(done) = settled_any(ids) {
return Some(done);
}
let Ok(queue) = kqueue::id() else {
loop {
if let Some(done) = settled_any(ids) {
return Some(done);
}
if help::help_once() {
continue;
}
thread::sleep(SELECT_POLL);
}
};
let registered: Vec<usize> = ids
.iter()
.filter(|id| slot(**id).is_some_and(|data| data.set_select(queue)))
.copied()
.collect();
let mut patience = Patience::new();
let winner = loop {
if let Some(done) = settled_any(ids) {
break Some(done);
}
if help::help_once() {
patience.helped();
continue;
}
let slice = patience
.slice()
.map_or(SELECT_POLL, |slice| slice.min(SELECT_POLL));
kqueue::wait_any(queue, slice);
};
for id in registered {
if let Some(data) = slot(id) {
data.clear_select(queue);
}
}
winner
}
fn settled_any(ids: &[usize]) -> Option<usize> {
ids.iter().copied().find(|id| match slot(*id) {
Some(data) => data.state().terminal(),
None => true,
})
}
fn failed<T>(erased: *mut c_void) -> TaskHandle<T> {
drop(unsafe { Box::from_raw(erased.cast::<Box<dyn ErasedTask>>()) });
TaskHandle::new(MAX_TASK_ID)
}
#[inline(always)]
fn missing(id: usize) -> RuntimeError {
match id == UNSTARTED_TASK_ID {
true => RuntimeError::NotInitialised,
false => RuntimeError::NoSuchTask,
}
}
fn unstarted<T>(erased: *mut c_void) -> TaskHandle<T> {
drop(unsafe { Box::from_raw(erased.cast::<Box<dyn ErasedTask>>()) });
TaskHandle::new(UNSTARTED_TASK_ID)
}
fn abandoned<T>(prototype: *mut c_void) -> TaskHandle<T> {
drop(unsafe { Box::from_raw(prototype.cast::<Box<dyn SeriesTask>>()) });
TaskHandle::new(MAX_TASK_ID)
}
#[inline(always)]
fn release(id: usize) {
let Some(data) = slot(id) else {
return;
};
if !data.claim_release() {
return;
}
Executor::drop_listener(id);
}
fn waiting_parts<T>(setup: TaskSetup) -> (Arc<Gate>, Arc<Mailbox<T>>, TaskSetup) {
let gate = Arc::new(Gate::new(
setup.gives,
setup.runs,
setup.deadline,
setup.start_delay,
setup.kind.repeats(),
));
let mailbox = Arc::new(Mailbox::new(Arc::clone(&gate)));
let setup = TaskSetup {
start_delay: Duration::ZERO,
waits: true,
..setup
};
(gate, mailbox, setup)
}
fn closed(data: &TaskData) -> RuntimeError {
match data.state() {
TaskState::Cancelled => RuntimeError::Cancelled,
TaskState::TimedOut => RuntimeError::TimedOut,
TaskState::Failed | TaskState::Free => RuntimeError::TaskFailed,
_ => RuntimeError::Finished,
}
}
fn start_waiting(id: usize, data: &TaskData, gate: &Gate) -> bool {
data.reset_series(gate.runs(), gate.until());
let delay = gate.delay();
if delay != 0 {
data.set_start_delay(delay);
}
let started = match (data.kind().schedules(), delay) {
(true, 0) => start_series(id, data),
(false, 0) => queue(id, data.blocking()),
(_, delay) => wait_for(data, id, delay),
};
if started {
return true;
}
gate.finish();
if data.try_state(TaskState::Pending, TaskState::Failed)
|| data.try_state(TaskState::Ready, TaskState::Failed)
|| data.try_state(TaskState::Taken, TaskState::Failed)
{
wake(data);
}
finish_waiting(id, data);
false
}
fn rewait(id: usize, data: &TaskData, task: Box<Box<dyn ErasedTask>>) {
data.rearm(Box::into_raw(task).cast::<c_void>());
rewait_series(id, data);
}
fn rewait_series(id: usize, data: &TaskData) {
data.add_listener();
match data.gate().map(Gate::after_run) {
Some(AfterRun::Wait) => {}
Some(AfterRun::Again) => {
if let Some(gate) = data.gate() {
let _ = start_waiting(id, data, gate);
}
}
Some(AfterRun::Finish) | None => finish_waiting(id, data),
}
Executor::drop_listener(id);
}
fn finish_waiting(id: usize, data: &TaskData) {
close_extras(data);
let task = data.claim();
if !task.is_null() {
drop(unsafe { Box::from_raw(task.cast::<Box<dyn ErasedTask>>()) });
}
if let Some(extras) = data.extras() {
extras.release_held();
}
data.finish_series();
release(id);
}
fn let_go_waiting(id: usize, data: &TaskData) {
let Some(gate) = data.gate() else {
return;
};
if !gate.close_waiting() {
return;
}
if data.try_state(TaskState::Pending, TaskState::Failed) {
wake(data);
}
finish_waiting(id, data);
}
pub(crate) fn abandon(id: usize) {
let Some(data) = slot(id) else {
return;
};
let_go_waiting(id, data);
}
#[inline(always)]
fn close_extras(data: &TaskData) {
let Some(extras) = data.extras() else {
return;
};
if let Some(gate) = extras.gate() {
gate.finish();
}
extras.receivers().close(data);
}
fn close_finished_series(data: &TaskData) {
if data.kind().repeats() || data.takes_input() || data.state() == TaskState::Pending {
return;
}
close_extras(data);
}
#[inline(always)]
fn note_output(data: &TaskData) {
if let Some(extras) = data.extras() {
extras.receivers().note_output();
}
}
#[inline(always)]
fn forward_output(data: &TaskData) {
if let Some(extras) = data.extras() {
extras.receivers().walk(data);
}
}
fn hold(id: usize, claims: Vec<Box<dyn Send>>) {
match slot(id) {
Some(data) => data.extras_or_attach().hold(claims),
None => drop(claims),
}
}
pub(crate) fn write_off(id: usize) {
let Some(data) = slot(id) else {
return;
};
loop {
let state = data.state();
if matches!(
state,
TaskState::Cancelled | TaskState::TimedOut | TaskState::Failed | TaskState::Free
) {
break;
}
if data.try_state(state, TaskState::Failed) {
wake(data);
break;
}
}
let Some(gate) = data.gate() else {
return;
};
if gate.close_waiting() {
finish_waiting(id, data);
return;
}
gate.finish();
}
fn supervise(id: i32) {
let supervisor = thread::spawn(move || {
let mut failures = 0;
let mut started = Instant::now();
loop {
let _ = thread::spawn(move || executor_loop(id)).join();
if shutting_down() {
let _ = unsafe { libc::close(id) };
break;
}
if started.elapsed() >= RESTART_WINDOW {
failures = 0;
}
failures += 1;
if failures > RESTART_LIMIT {
shutdown(id);
break;
}
thread::sleep(RESTART_BACKOFF * failures);
started = Instant::now();
}
});
*SUPERVISOR
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(supervisor);
}
fn executor_loop(id: i32) {
let mut events = eventlist();
let armed = unsafe {
KEvent::register(
id,
MANAGER_TICK_IDENT,
MANAGER_TICK.as_nanos() as libc::intptr_t,
ptr::null_mut(),
EventDesc::new_interval(),
)
}
.check();
if armed.is_err() {
return;
}
recover_waits();
loop {
if shutting_down() {
return;
}
POOL.tick();
let count = match unsafe { KEvent::listen(id, &mut events) }.check() {
Ok(count) => count as usize,
Err(RuntimeError::CheckError(Some(libc::EINTR))) => continue,
Err(_) => break,
};
if shutting_down() {
return;
}
if injected_fault() {
panic!("injected manager fault");
}
for event in events.iter().take(count) {
if event.flags & libc::EV_ERROR != 0 {
continue;
}
match event.filter {
libc::EVFILT_READ
| libc::EVFILT_WRITE
| libc::EVFILT_SIGNAL
| libc::EVFILT_VNODE => {
wake_parked(event.udata as usize, Fired::Event);
}
libc::EVFILT_TIMER if event.ident == MANAGER_TICK_IDENT => {}
libc::EVFILT_TIMER if event.ident >= TIMEOUT_IDENT_BASE => {
timed_out(event.ident - TIMEOUT_IDENT_BASE, event.udata as u64);
}
libc::EVFILT_TIMER if event.udata as usize == PARK_TIMER => {
if let Some(id) = event.ident.checked_sub(SCHEDULE_IDENT_BASE) {
wake_parked(id, Fired::Deadline);
}
}
libc::EVFILT_TIMER => fire(event.ident),
_ => {}
}
}
}
}
fn orphaned() {
for task in 0..DATA.high_water() {
let Some(data) = slot(task) else {
continue;
};
if data.parked() {
unpark(task, data);
continue;
}
let kind = data.kind();
if !kind.waits() && !kind.schedules() && !data.armed() {
continue;
}
let state = data.state();
if state == TaskState::Running && !kind.schedules() {
continue;
}
if !state.terminal() {
data.set_state(TaskState::Failed);
wake(data);
}
close_extras(data);
release(task);
}
}
fn recover_waits() {
for task in 0..DATA.high_water() {
let Some(data) = slot(task) else {
continue;
};
if let Some(deadline) = data.run_deadline_at() {
let left = deadline.saturating_duration_since(Instant::now());
let _ = arm_timeout(task, left, data.run_deadline());
}
if let Some(parked) = data.claim_parked() {
unwatch_park(task, parked, Fired::Neither);
if !queue(task, data.blocking()) {
take_down(task, data);
}
continue;
}
if !data.armed() {
continue;
}
let owed = match data.start_delay() {
0 => data.interval(),
delay => delay,
};
arm_timer(task, owed);
}
}
fn shutdown(id: i32) {
EXECUTOR_KQUEUE_ID.store(DEAD_KQUEUE_ID, Ordering::SeqCst);
let _ = unsafe { libc::close(id) };
orphaned();
POOL.sweep_all();
POOL.ensure_floor();
if POOL.live() > 0 {
return;
}
write_off_pool();
}
pub(crate) fn write_off_pool() {
POOL.close();
POOL.abandon();
for task in POOL
.injector()
.drain()
.into_iter()
.chain(POOL.blocking().drain())
{
fail(task);
}
for task in 0..DATA.high_water() {
let Some(data) = slot(task) else {
continue;
};
let state = data.state();
if state == TaskState::Running && !data.kind().schedules() {
continue;
}
if !state.terminal() {
data.set_state(TaskState::Failed);
wake(data);
}
close_extras(data);
release(task);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::modules::input::Token;
use crate::{Nothing, futures::task::sealed, sleep::Sleep};
use std::time::Duration;
struct Panics;
impl sealed::Sealed for Panics {}
impl Task for Panics {
type Output = usize;
type Input = Nothing;
fn execute(&self, _token: Token, _reactor_id: i32, _task_id: usize) -> Self::Output {
panic!("this task is meant to go down");
}
}
struct Liar;
impl sealed::Sealed for Liar {}
impl Task for Liar {
type Output = usize;
type Input = Nothing;
fn execute(&self, _token: Token, _reactor_id: i32, _task_id: usize) -> Self::Output {
thread::sleep(Duration::from_millis(200));
0
}
}
#[test]
fn pool_grows_when_tasks_hold_their_workers() {
let _ = crate::Runtime::init();
let cores = thread::available_parallelism()
.map(|count| count.get())
.unwrap_or(1);
let handles: Vec<_> = (0..cores * 2)
.map(|_| crate::Runtime::task(Liar).spawn())
.collect();
thread::sleep(Duration::from_millis(150));
let grown = crate::Runtime::pool().workers().len();
for handle in handles {
handle.join().expect("every task finishes");
}
assert!(
grown > cores,
"pool stayed at {} workers with {} tasks holding threads and a floor of {}",
grown,
cores * 2,
cores,
);
}
#[test]
fn a_stale_park_wake_cannot_start_a_delayed_task() {
let _ = crate::Runtime::init();
let handle = crate::Runtime::task(Sleep::sleep(Duration::from_micros(1)))
.after(Duration::from_millis(300))
.spawn();
wake_parked(handle.id(), Fired::Event);
wake_parked(handle.id(), Fired::Deadline);
thread::sleep(Duration::from_millis(50));
assert!(
handle.is_pending(),
"a stale park wake started a delayed task early, which is now {:?}",
handle.state(),
);
assert!(
slot(handle.id()).is_some_and(|data| data.armed()),
"a stale park wake took the delay's own wake",
);
handle
.join()
.expect("the delayed task still runs at its own time");
}
#[test]
fn panicking_task_does_not_lose_its_queue() {
let _ = crate::Runtime::init();
let quick = || Sleep::sleep(Duration::from_micros(50));
let before: Vec<_> = (0..256)
.map(|_| crate::Runtime::task(quick()).spawn())
.collect();
let doomed = crate::Runtime::task(Panics).spawn();
let after: Vec<_> = (0..256)
.map(|_| crate::Runtime::task(quick()).spawn())
.collect();
assert_eq!(
doomed.join(),
Err(RuntimeError::TaskFailed),
"the task that went down comes back as an error rather than blocking forever",
);
for handle in before.into_iter().chain(after) {
handle
.join()
.expect("every task the dead worker was holding still finishes");
}
}
}