#![forbid(unsafe_code)]
use std::error::Error as StdError;
use std::fmt;
use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};
#[cfg(not(target_arch = "wasm32"))]
use std::sync::Arc;
#[cfg(not(target_arch = "wasm32"))]
use std::sync::atomic::{AtomicBool, Ordering};
#[cfg(not(target_arch = "wasm32"))]
use std::task::{Wake, Waker};
#[allow(unused_imports)]
pub(crate) use std::future::poll_fn;
#[allow(unused_imports)]
pub(crate) use std::future::pending;
#[allow(clippy::manual_async_fn)]
pub(crate) fn poll_once<F>(future: F) -> impl Future<Output = Option<F::Output>>
where
F: Future,
{
async move {
let mut future = std::pin::pin!(future);
poll_fn(|context| {
Poll::Ready(match future.as_mut().poll(context) {
Poll::Ready(output) => Some(output),
Poll::Pending => None,
})
})
.await
}
}
#[derive(Debug, Default)]
pub(crate) struct YieldNow {
yielded: bool,
}
impl Future for YieldNow {
type Output = ();
fn poll(mut self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
if self.yielded {
Poll::Ready(())
} else {
self.yielded = true;
context.waker().wake_by_ref();
Poll::Pending
}
}
}
#[must_use]
pub(crate) const fn yield_now() -> YieldNow {
YieldNow { yielded: false }
}
pub(crate) fn zip<F1, F2>(future1: F1, future2: F2) -> Zip<F1, F2>
where
F1: Future,
F2: Future,
{
Zip {
future1: Some(future1),
output1: None,
future2: Some(future2),
output2: None,
}
}
#[pin_project::pin_project]
#[derive(Debug)]
#[must_use = "futures do nothing unless polled or awaited"]
pub(crate) struct Zip<F1, F2>
where
F1: Future,
F2: Future,
{
#[pin]
future1: Option<F1>,
output1: Option<F1::Output>,
#[pin]
future2: Option<F2>,
output2: Option<F2::Output>,
}
impl<F1, F2> Future for Zip<F1, F2>
where
F1: Future,
F2: Future,
{
type Output = (F1::Output, F2::Output);
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
let mut this = self.project();
if let Some(future) = this.future1.as_mut().as_pin_mut() {
if let Poll::Ready(output) = future.poll(context) {
*this.output1 = Some(output);
this.future1.set(None);
}
}
if let Some(future) = this.future2.as_mut().as_pin_mut() {
if let Poll::Ready(output) = future.poll(context) {
*this.output2 = Some(output);
this.future2.set(None);
}
}
take_zip_outputs(this.output1, this.output2)
}
}
fn take_zip_outputs<T1, T2>(output1: &mut Option<T1>, output2: &mut Option<T2>) -> Poll<(T1, T2)> {
match (output1.take(), output2.take()) {
(Some(output1), Some(output2)) => Poll::Ready((output1, output2)),
(remaining1, remaining2) => {
*output1 = remaining1;
*output2 = remaining2;
Poll::Pending
}
}
}
pub(crate) fn or<T, F1, F2>(future1: F1, future2: F2) -> Or<F1, F2>
where
F1: Future<Output = T>,
F2: Future<Output = T>,
{
Or { future1, future2 }
}
#[pin_project::pin_project]
#[derive(Debug)]
#[must_use = "futures do nothing unless polled or awaited"]
pub(crate) struct Or<F1, F2> {
#[pin]
future1: F1,
#[pin]
future2: F2,
}
impl<T, F1, F2> Future for Or<F1, F2>
where
F1: Future<Output = T>,
F2: Future<Output = T>,
{
type Output = T;
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
let this = self.project();
if let Poll::Ready(output) = this.future1.poll(context) {
return Poll::Ready(output);
}
if let Poll::Ready(output) = this.future2.poll(context) {
return Poll::Ready(output);
}
Poll::Pending
}
}
pub(crate) fn catch_unwind<F>(future: F) -> CatchUnwind<F>
where
F: Future + std::panic::UnwindSafe,
{
CatchUnwind {
inner: future,
completed: false,
}
}
#[pin_project::pin_project]
#[derive(Debug)]
#[must_use = "futures do nothing unless polled or awaited"]
pub(crate) struct CatchUnwind<F> {
#[pin]
inner: F,
completed: bool,
}
impl<F> Future for CatchUnwind<F>
where
F: Future + std::panic::UnwindSafe,
{
type Output = Result<F::Output, Box<dyn std::any::Any + Send>>;
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
let mut this = self.project();
if *this.completed {
return Poll::Pending;
}
match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
this.inner.as_mut().poll(context)
})) {
Ok(Poll::Ready(output)) => {
*this.completed = true;
Poll::Ready(Ok(output))
}
Ok(Poll::Pending) => Poll::Pending,
Err(payload) => {
*this.completed = true;
Poll::Ready(Err(payload))
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum BlockOnError {
RuntimeContext,
UnsupportedPlatform,
}
impl fmt::Display for BlockOnError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::RuntimeContext => formatter.write_str(
"cannot block the current thread while an Asupersync runtime handle is installed",
),
Self::UnsupportedPlatform => {
formatter.write_str("blocking future execution is unsupported on this target")
}
}
}
}
impl StdError for BlockOnError {}
pub(crate) fn block_on<F>(future: F) -> F::Output
where
F: Future,
{
try_block_on(future).unwrap_or_else(|error| panic!("{error}"))
}
pub(crate) fn try_block_on<F>(future: F) -> Result<F::Output, BlockOnError>
where
F: Future,
{
#[cfg(target_arch = "wasm32")]
{
drop(future);
Err(BlockOnError::UnsupportedPlatform)
}
#[cfg(not(target_arch = "wasm32"))]
{
if crate::runtime::Runtime::current_handle().is_some() {
return Err(BlockOnError::RuntimeContext);
}
Ok(block_on_with_park(future, |_| std::thread::park()))
}
}
#[cfg(not(target_arch = "wasm32"))]
struct ThreadNotification {
thread: std::thread::Thread,
notified: AtomicBool,
}
#[cfg(not(target_arch = "wasm32"))]
impl ThreadNotification {
fn new() -> Self {
Self {
thread: std::thread::current(),
notified: AtomicBool::new(false),
}
}
fn prepare_for_poll(&self) {
self.notified.swap(false, Ordering::AcqRel);
}
fn record_notification(&self) {
self.notified.store(true, Ordering::Release);
}
fn notify(&self) {
self.record_notification();
self.thread.unpark();
}
fn wait<P>(&self, park: &mut P)
where
P: FnMut(&Self),
{
while !self.notified.swap(false, Ordering::AcqRel) {
park(self);
}
}
}
#[cfg(not(target_arch = "wasm32"))]
impl Wake for ThreadNotification {
fn wake(self: Arc<Self>) {
self.notify();
}
fn wake_by_ref(self: &Arc<Self>) {
self.notify();
}
}
#[cfg(not(target_arch = "wasm32"))]
fn block_on_with_park<F, P>(future: F, mut park: P) -> F::Output
where
F: Future,
P: FnMut(&ThreadNotification),
{
let notification = Arc::new(ThreadNotification::new());
let waker = Waker::from(Arc::clone(¬ification));
let mut context = Context::from_waker(&waker);
let mut future = std::pin::pin!(future);
loop {
notification.prepare_for_poll();
match future.as_mut().poll(&mut context) {
Poll::Ready(output) => return output,
Poll::Pending => notification.wait(&mut park),
}
}
}
#[cfg(all(test, not(target_arch = "wasm32")))]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, clippy::future_not_send)]
use super::*;
use crate::runtime::{BlockingPool, RuntimeBuilder};
use std::cell::{Cell, RefCell};
use std::rc::Rc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::{Duration, Instant};
#[cfg(target_os = "linux")]
fn process_cpu_ticks() -> u64 {
let stat = std::fs::read_to_string("/proc/self/stat").expect("read process stat");
let close = stat.rfind(')').expect("process stat comm terminator");
let fields: Vec<&str> = stat[close + 1..].split_whitespace().collect();
let user: u64 = fields[11].parse().expect("process user ticks");
let system: u64 = fields[12].parse().expect("process system ticks");
user + system
}
#[cfg(target_os = "linux")]
fn delayed_wake(delay: Duration) -> (impl Future<Output = ()>, std::thread::JoinHandle<()>) {
let ready = Arc::new(AtomicBool::new(false));
let parked_waker = Arc::new(Mutex::new(None::<Waker>));
let (polled_tx, polled_rx) = std::sync::mpsc::sync_channel(1);
let helper_ready = Arc::clone(&ready);
let helper_waker = Arc::clone(&parked_waker);
let helper = std::thread::spawn(move || {
polled_rx.recv().expect("future reports its pending poll");
std::thread::sleep(delay);
helper_ready.store(true, Ordering::Release);
helper_waker
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
.expect("pending future installed its waker")
.wake();
});
let mut reported_pending = false;
let future = poll_fn(move |context| {
if ready.load(Ordering::Acquire) {
return Poll::Ready(());
}
parked_waker
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.replace(context.waker().clone());
if !reported_pending {
reported_pending = true;
polled_tx.send(()).expect("wake helper remains available");
}
Poll::Pending
});
(future, helper)
}
#[cfg(target_os = "linux")]
fn ready_batch<D>(driver: &mut D, iterations: u64) -> (u128, u64)
where
D: FnMut(u64) -> u64,
{
let start = Instant::now();
let mut checksum = 0_u64;
for value in 0..iterations {
checksum =
checksum.wrapping_add(std::hint::black_box(driver(std::hint::black_box(value))));
}
(start.elapsed().as_nanos(), checksum)
}
#[derive(Default)]
struct CountingWaker {
wakes: AtomicUsize,
}
impl Wake for CountingWaker {
fn wake(self: Arc<Self>) {
self.wakes.fetch_add(1, Ordering::Relaxed);
}
fn wake_by_ref(self: &Arc<Self>) {
self.wakes.fetch_add(1, Ordering::Relaxed);
}
}
struct PendingDrop {
polls: Rc<Cell<usize>>,
drops: Rc<Cell<usize>>,
}
impl Future for PendingDrop {
type Output = u8;
fn poll(self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll<Self::Output> {
self.polls.set(self.polls.get() + 1);
Poll::Pending
}
}
impl Drop for PendingDrop {
fn drop(&mut self) {
self.drops.set(self.drops.get() + 1);
}
}
struct ReadyDrop {
output: u8,
polls: Rc<Cell<usize>>,
drops: Rc<Cell<usize>>,
}
impl Future for ReadyDrop {
type Output = u8;
fn poll(self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll<Self::Output> {
self.polls.set(self.polls.get() + 1);
Poll::Ready(self.output)
}
}
impl Drop for ReadyDrop {
fn drop(&mut self) {
self.drops.set(self.drops.get() + 1);
}
}
struct DropCounter {
drops: Rc<Cell<usize>>,
}
impl Drop for DropCounter {
fn drop(&mut self) {
self.drops.set(self.drops.get() + 1);
}
}
fn ready_after(
label: &'static str,
pending_polls: usize,
output: u8,
trace: Rc<RefCell<Vec<&'static str>>>,
) -> impl Future<Output = u8> {
let mut remaining = pending_polls;
poll_fn(move |context| {
trace.borrow_mut().push(label);
if remaining == 0 {
Poll::Ready(output)
} else {
remaining -= 1;
context.waker().wake_by_ref();
Poll::Pending
}
})
}
#[test]
fn poll_fn_forwards_context_and_calls_once_per_wrapper_poll() {
let wake_state = Arc::new(CountingWaker::default());
let waker = Waker::from(Arc::clone(&wake_state));
let mut context = Context::from_waker(&waker);
let polls = Cell::new(0_usize);
let mut future = std::pin::pin!(poll_fn(|received_context| {
assert!(received_context.waker().will_wake(&waker));
let current = polls.get();
polls.set(current + 1);
if current == 0 {
Poll::Pending
} else {
Poll::Ready(41_u8)
}
}));
assert_eq!(future.as_mut().poll(&mut context), Poll::Pending);
assert_eq!(polls.get(), 1);
assert_eq!(future.as_mut().poll(&mut context), Poll::Ready(41));
assert_eq!(polls.get(), 2);
assert_eq!(wake_state.wakes.load(Ordering::Relaxed), 0);
}
#[test]
fn poll_once_observes_ready_and_pending_without_waiting() {
assert_eq!(block_on(poll_once(async { 43_u8 })), Some(43));
let polls = Rc::new(Cell::new(0_usize));
let drops = Rc::new(Cell::new(0_usize));
let observed = block_on(poll_once(PendingDrop {
polls: Rc::clone(&polls),
drops: Rc::clone(&drops),
}));
assert_eq!(observed, None);
assert_eq!(polls.get(), 1);
assert_eq!(drops.get(), 1);
}
#[test]
fn yield_now_wakes_once_then_remains_ready() {
let wake_state = Arc::new(CountingWaker::default());
let waker = Waker::from(Arc::clone(&wake_state));
let mut context = Context::from_waker(&waker);
let mut future = std::pin::pin!(yield_now());
assert_eq!(future.as_mut().poll(&mut context), Poll::Pending);
assert_eq!(wake_state.wakes.load(Ordering::Relaxed), 1);
assert_eq!(future.as_mut().poll(&mut context), Poll::Ready(()));
assert_eq!(future.as_mut().poll(&mut context), Poll::Ready(()));
assert_eq!(wake_state.wakes.load(Ordering::Relaxed), 1);
}
#[test]
fn pending_never_completes_or_schedules_a_wake() {
let wake_state = Arc::new(CountingWaker::default());
let waker = Waker::from(Arc::clone(&wake_state));
let mut context = Context::from_waker(&waker);
let mut future = std::pin::pin!(pending::<u8>());
assert_eq!(future.as_mut().poll(&mut context), Poll::Pending);
assert_eq!(future.as_mut().poll(&mut context), Poll::Pending);
assert_eq!(wake_state.wakes.load(Ordering::Relaxed), 0);
}
#[test]
fn zip_polls_left_then_right_and_stops_polling_completed_children() {
let log = Rc::new(RefCell::new(Vec::new()));
let left_log = Rc::clone(&log);
let left = poll_fn(move |_| {
left_log.borrow_mut().push("left");
Poll::Ready(47_u8)
});
let right_log = Rc::clone(&log);
let right_polls = Rc::new(Cell::new(0_usize));
let observed_right_polls = Rc::clone(&right_polls);
let right = poll_fn(move |context| {
right_log.borrow_mut().push("right");
let current = right_polls.get();
right_polls.set(current + 1);
if current == 0 {
context.waker().wake_by_ref();
Poll::Pending
} else {
Poll::Ready(53_u8)
}
});
assert_eq!(block_on(zip(left, right)), (47, 53));
assert_eq!(&*log.borrow(), &["left", "right", "right"]);
assert_eq!(observed_right_polls.get(), 2);
}
#[test]
fn dropping_pending_zip_drops_retained_output_and_unfinished_child() {
let output_drops = Rc::new(Cell::new(0_usize));
let pending_polls = Rc::new(Cell::new(0_usize));
let pending_drops = Rc::new(Cell::new(0_usize));
let wake_state = Arc::new(CountingWaker::default());
let waker = Waker::from(Arc::clone(&wake_state));
let mut context = Context::from_waker(&waker);
{
let left_output_drops = Rc::clone(&output_drops);
let joined = zip(
async move {
DropCounter {
drops: left_output_drops,
}
},
PendingDrop {
polls: Rc::clone(&pending_polls),
drops: Rc::clone(&pending_drops),
},
);
let mut joined = std::pin::pin!(joined);
assert!(joined.as_mut().poll(&mut context).is_pending());
assert_eq!(output_drops.get(), 0);
assert_eq!(pending_drops.get(), 0);
}
assert_eq!(output_drops.get(), 1);
assert_eq!(pending_polls.get(), 1);
assert_eq!(pending_drops.get(), 1);
assert_eq!(wake_state.wakes.load(Ordering::Relaxed), 0);
}
#[test]
fn or_is_left_biased_and_drops_the_loser_with_the_wrapper() {
let left_polls = Rc::new(Cell::new(0_usize));
let left_drops = Rc::new(Cell::new(0_usize));
let skipped_right_polls = Rc::new(Cell::new(0_usize));
let skipped_right_drops = Rc::new(Cell::new(0_usize));
let result = block_on(or(
ReadyDrop {
output: 59,
polls: Rc::clone(&left_polls),
drops: Rc::clone(&left_drops),
},
ReadyDrop {
output: 61,
polls: Rc::clone(&skipped_right_polls),
drops: Rc::clone(&skipped_right_drops),
},
));
assert_eq!(result, 59);
assert_eq!(left_polls.get(), 1);
assert_eq!(left_drops.get(), 1);
assert_eq!(skipped_right_polls.get(), 0);
assert_eq!(skipped_right_drops.get(), 1);
let pending_polls = Rc::new(Cell::new(0_usize));
let pending_drops = Rc::new(Cell::new(0_usize));
let right_polls = Rc::new(Cell::new(0_usize));
let right_drops = Rc::new(Cell::new(0_usize));
let result = block_on(or(
PendingDrop {
polls: Rc::clone(&pending_polls),
drops: Rc::clone(&pending_drops),
},
ReadyDrop {
output: 67,
polls: Rc::clone(&right_polls),
drops: Rc::clone(&right_drops),
},
));
assert_eq!(result, 67);
assert_eq!(pending_polls.get(), 1);
assert_eq!(pending_drops.get(), 1);
assert_eq!(right_polls.get(), 1);
assert_eq!(right_drops.get(), 1);
}
#[test]
fn zip_and_or_readiness_matrix_is_deterministic() {
fn run_zip(left_pending: usize, right_pending: usize) -> ((u8, u8), Vec<&'static str>) {
let trace = Rc::new(RefCell::new(Vec::new()));
let result = block_on(zip(
ready_after("left", left_pending, 71, Rc::clone(&trace)),
ready_after("right", right_pending, 73, Rc::clone(&trace)),
));
let events = trace.borrow().clone();
(result, events)
}
fn run_or(left_pending: usize, right_pending: usize) -> (u8, Vec<&'static str>) {
let trace = Rc::new(RefCell::new(Vec::new()));
let result = block_on(or(
ready_after("left", left_pending, 79, Rc::clone(&trace)),
ready_after("right", right_pending, 83, Rc::clone(&trace)),
));
let events = trace.borrow().clone();
(result, events)
}
for left_pending in 0..=3 {
for right_pending in 0..=3 {
let first_zip = run_zip(left_pending, right_pending);
let second_zip = run_zip(left_pending, right_pending);
assert_eq!(first_zip, second_zip);
assert_eq!(first_zip.0, (71, 73));
assert_eq!(
first_zip.1.iter().filter(|event| **event == "left").count(),
left_pending + 1
);
assert_eq!(
first_zip
.1
.iter()
.filter(|event| **event == "right")
.count(),
right_pending + 1
);
let first_or = run_or(left_pending, right_pending);
let second_or = run_or(left_pending, right_pending);
assert_eq!(first_or, second_or);
assert_eq!(
first_or.0,
if left_pending <= right_pending {
79
} else {
83
}
);
}
}
}
#[test]
fn helper_futures_quiesce_under_lab_dpor_exploration() {
use crate::lab::{DporExplorer, ExplorerConfig};
use crate::types::Budget;
let mut explorer = DporExplorer::new(
ExplorerConfig::new(0x00F0_74A4, 8)
.worker_count(1)
.max_steps(2_000),
);
let report = explorer.explore(|runtime| {
let region = runtime.state.create_root_region(Budget::INFINITE);
let (zip_task, _) = runtime
.state
.create_task(region, Budget::INFINITE, async {
assert_eq!(poll_once(async { 89_u8 }).await, Some(89));
assert_eq!(poll_once(pending::<u8>()).await, None);
assert_eq!(zip(async { 97_u8 }, async { 101_u8 }).await, (97, 101));
})
.expect("create zip helper task");
let (or_task, _) = runtime
.state
.create_task(region, Budget::INFINITE, async {
yield_now().await;
assert_eq!(or(async { 103_u8 }, async { 107_u8 }).await, 103);
})
.expect("create or helper task");
{
let mut scheduler = runtime.scheduler.lock();
scheduler.schedule(zip_task, 0);
scheduler.schedule(or_task, 0);
}
runtime.run_until_quiescent();
assert!(runtime.is_quiescent());
assert_eq!(runtime.state.pending_obligation_count(), 0);
});
assert!(!report.has_violations());
assert!(report.unique_classes >= 1);
}
#[test]
fn catch_unwind_forwards_pending_wake_and_ready() {
let polls = Rc::new(Cell::new(0_usize));
let observed_polls = Rc::clone(&polls);
let future = poll_fn(move |context| {
let current = polls.get();
polls.set(current + 1);
if current == 0 {
context.waker().wake_by_ref();
Poll::Pending
} else {
Poll::Ready(71_u8)
}
});
let observed = block_on(catch_unwind(std::panic::AssertUnwindSafe(future)));
assert!(matches!(observed, Ok(71)));
assert_eq!(observed_polls.get(), 2);
}
#[test]
fn catch_unwind_preserves_payload_and_refuses_repoll() {
let polls = Rc::new(Cell::new(0_usize));
let observed_polls = Rc::clone(&polls);
let future = poll_fn(move |_| -> Poll<u8> {
polls.set(polls.get() + 1);
std::panic::panic_any(String::from("owned poll panic"));
});
let mut caught = std::pin::pin!(catch_unwind(std::panic::AssertUnwindSafe(future)));
let wake_state = Arc::new(CountingWaker::default());
let waker = Waker::from(Arc::clone(&wake_state));
let mut context = Context::from_waker(&waker);
let Poll::Ready(Err(payload)) = caught.as_mut().poll(&mut context) else {
panic!("poll panic must become a ready error");
};
assert_eq!(
payload.downcast_ref::<String>().map(String::as_str),
Some("owned poll panic")
);
assert!(caught.as_mut().poll(&mut context).is_pending());
assert_eq!(observed_polls.get(), 1);
}
#[test]
fn dropping_unpolled_catch_unwind_drops_inner() {
let polls = Rc::new(Cell::new(0_usize));
let drops = Rc::new(Cell::new(0_usize));
{
let _caught = catch_unwind(std::panic::AssertUnwindSafe(PendingDrop {
polls: Rc::clone(&polls),
drops: Rc::clone(&drops),
}));
}
assert_eq!(polls.get(), 0);
assert_eq!(drops.get(), 1);
}
#[test]
fn ready_future_completes_without_parking() {
let park_calls = Cell::new(0_usize);
let output = block_on_with_park(async { 42_u8 }, |_| {
park_calls.set(park_calls.get() + 1);
});
assert_eq!(output, 42);
assert_eq!(park_calls.get(), 0);
}
#[test]
fn borrowed_non_send_future_and_recursive_call_are_admitted() {
let value = Rc::new(Cell::new(1_u8));
let borrowed = &value;
let output = block_on(async {
borrowed.set(2);
let inner_polls = Cell::new(0_usize);
let inner = block_on(poll_fn(|context| {
let current = inner_polls.get();
inner_polls.set(current + 1);
if current == 0 {
context.waker().wake_by_ref();
Poll::Pending
} else {
borrowed.set(3);
Poll::Ready(borrowed.get())
}
}));
assert_eq!(inner_polls.get(), 2);
inner + borrowed.get()
});
assert_eq!(output, 6);
assert_eq!(value.get(), 3);
}
#[test]
fn wakes_during_poll_are_coalesced_without_parking() {
let polls = Cell::new(0_usize);
let park_calls = Cell::new(0_usize);
let output = block_on_with_park(
poll_fn(|context| {
let current = polls.get();
polls.set(current + 1);
if current == 0 {
for _ in 0..4 {
context.waker().wake_by_ref();
}
Poll::Pending
} else {
Poll::Ready(7_u8)
}
}),
|_| park_calls.set(park_calls.get() + 1),
);
assert_eq!(output, 7);
assert_eq!(polls.get(), 2);
assert_eq!(park_calls.get(), 0);
}
#[test]
fn spurious_park_return_does_not_trigger_an_unnotified_poll() {
let polls = Cell::new(0_usize);
let park_calls = Cell::new(0_usize);
let output = block_on_with_park(
poll_fn(|_| {
let current = polls.get();
polls.set(current + 1);
if current == 0 {
Poll::Pending
} else {
Poll::Ready(11_u8)
}
}),
|notification| {
let current = park_calls.get();
park_calls.set(current + 1);
if current == 1 {
notification.record_notification();
}
},
);
assert_eq!(output, 11);
assert_eq!(polls.get(), 2);
assert_eq!(park_calls.get(), 2);
}
#[test]
fn notification_state_model_exhausts_spurious_and_coalesced_wakes() {
for spurious_returns in 0..=4 {
for coalesced_wakes in 1..=4 {
let ready = Cell::new(false);
let polls = Cell::new(0_usize);
let park_calls = Cell::new(0_usize);
let output = block_on_with_park(
poll_fn(|_| {
polls.set(polls.get() + 1);
if ready.get() {
Poll::Ready(13_u8)
} else {
Poll::Pending
}
}),
|notification| {
let current = park_calls.get();
park_calls.set(current + 1);
if current == spurious_returns {
ready.set(true);
for _ in 0..coalesced_wakes {
notification.record_notification();
}
}
},
);
assert_eq!(output, 13);
assert_eq!(polls.get(), 2, "spurious returns must not repoll");
assert_eq!(
park_calls.get(),
spurious_returns + 1,
"the model must park until its first recorded wake"
);
}
}
}
#[cfg(target_os = "linux")]
#[test]
#[ignore = "explicit FUT A3 measurement receipt"]
fn block_on_incumbent_idle_cpu_and_ready_latency_receipt() {
const IDLE_MILLIS: u64 = 750;
const READY_ITERATIONS: u64 = 50_000;
const READY_SAMPLES: usize = 7;
let idle_delay = Duration::from_millis(IDLE_MILLIS);
let (owned_idle_future, owned_helper) = delayed_wake(idle_delay);
let owned_cpu_before = process_cpu_ticks();
let owned_idle_started = Instant::now();
block_on(owned_idle_future);
let owned_idle_elapsed = owned_idle_started.elapsed();
let owned_idle_cpu_ticks = process_cpu_ticks().saturating_sub(owned_cpu_before);
owned_helper.join().expect("owned wake helper completes");
let (incumbent_idle_future, incumbent_helper) = delayed_wake(idle_delay);
let incumbent_cpu_before = process_cpu_ticks();
let incumbent_idle_started = Instant::now();
futures_lite::future::block_on(incumbent_idle_future);
let incumbent_idle_elapsed = incumbent_idle_started.elapsed();
let incumbent_idle_cpu_ticks = process_cpu_ticks().saturating_sub(incumbent_cpu_before);
incumbent_helper
.join()
.expect("incumbent wake helper completes");
let mut owned_driver = |value| block_on(std::hint::black_box(async move { value }));
let mut incumbent_driver =
|value| futures_lite::future::block_on(std::hint::black_box(async move { value }));
let mut owned_ready_samples = Vec::with_capacity(READY_SAMPLES);
let mut incumbent_ready_samples = Vec::with_capacity(READY_SAMPLES);
let mut owned_checksum = 0_u64;
let mut incumbent_checksum = 0_u64;
for sample in 0..READY_SAMPLES {
let (first_ns, first_checksum);
let (second_ns, second_checksum);
if sample % 2 == 0 {
(first_ns, first_checksum) = ready_batch(&mut owned_driver, READY_ITERATIONS);
(second_ns, second_checksum) = ready_batch(&mut incumbent_driver, READY_ITERATIONS);
owned_ready_samples.push(first_ns);
incumbent_ready_samples.push(second_ns);
owned_checksum = first_checksum;
incumbent_checksum = second_checksum;
} else {
(first_ns, first_checksum) = ready_batch(&mut incumbent_driver, READY_ITERATIONS);
(second_ns, second_checksum) = ready_batch(&mut owned_driver, READY_ITERATIONS);
incumbent_ready_samples.push(first_ns);
owned_ready_samples.push(second_ns);
incumbent_checksum = first_checksum;
owned_checksum = second_checksum;
}
}
owned_ready_samples.sort_unstable();
incumbent_ready_samples.sort_unstable();
let owned_ready_median_ns = owned_ready_samples[READY_SAMPLES / 2];
let incumbent_ready_median_ns = incumbent_ready_samples[READY_SAMPLES / 2];
assert!(owned_idle_elapsed >= idle_delay);
assert!(incumbent_idle_elapsed >= idle_delay);
assert!(
owned_idle_cpu_ticks <= incumbent_idle_cpu_ticks.saturating_add(3),
"owned idle process CPU exceeded the incumbent by more than three clock ticks"
);
assert_eq!(owned_checksum, incumbent_checksum);
println!(
"FUT_A3_PERF_RECEIPT platform=linux idle_wait_ms={IDLE_MILLIS} \
owned_idle_cpu_ticks={owned_idle_cpu_ticks} \
incumbent_idle_cpu_ticks={incumbent_idle_cpu_ticks} \
owned_idle_wall_ns={} incumbent_idle_wall_ns={} \
ready_iterations_per_sample={READY_ITERATIONS} ready_samples={READY_SAMPLES} \
owned_ready_median_batch_ns={owned_ready_median_ns} \
incumbent_ready_median_batch_ns={incumbent_ready_median_ns}",
owned_idle_elapsed.as_nanos(),
incumbent_idle_elapsed.as_nanos(),
);
}
#[test]
fn repeated_polls_receive_the_same_waker_identity() {
let first_waker = RefCell::new(None::<Waker>);
let polls = Cell::new(0_usize);
block_on(poll_fn(|context| {
let current = polls.get();
polls.set(current + 1);
if current == 0 {
first_waker.replace(Some(context.waker().clone()));
context.waker().wake_by_ref();
Poll::Pending
} else {
assert!(
first_waker
.borrow()
.as_ref()
.is_some_and(|waker| waker.will_wake(context.waker()))
);
Poll::Ready(())
}
}));
assert_eq!(polls.get(), 2);
}
#[test]
fn wake_after_pending_makes_progress() {
let ready = Arc::new(AtomicBool::new(false));
let parked_waker = Arc::new(Mutex::new(None::<Waker>));
let (polled_tx, polled_rx) = std::sync::mpsc::channel();
let helper_ready = Arc::clone(&ready);
let helper_waker = Arc::clone(&parked_waker);
let helper = std::thread::spawn(move || {
polled_rx.recv().expect("future reports its pending poll");
helper_ready.store(true, Ordering::Release);
helper_waker
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
.expect("pending future installed its waker")
.wake();
});
let mut reported_pending = false;
let output = block_on(poll_fn(|context| {
if ready.load(Ordering::Acquire) {
return Poll::Ready(19_u8);
}
parked_waker
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.replace(context.waker().clone());
if !reported_pending {
reported_pending = true;
polled_tx.send(()).expect("wake helper remains available");
}
Poll::Pending
}));
helper.join().expect("wake helper does not panic");
assert_eq!(output, 19);
}
#[test]
fn explicit_cancellation_wake_makes_progress() {
let cancelled = Arc::new(AtomicBool::new(false));
let parked_waker = Arc::new(Mutex::new(None::<Waker>));
let (polled_tx, polled_rx) = std::sync::mpsc::channel();
let helper_cancelled = Arc::clone(&cancelled);
let helper_waker = Arc::clone(&parked_waker);
let helper = std::thread::spawn(move || {
polled_rx.recv().expect("future reports its pending poll");
helper_cancelled.store(true, Ordering::Release);
helper_waker
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
.expect("pending future installed its waker")
.wake();
});
let mut reported_pending = false;
let observed = block_on(poll_fn(|context| {
if cancelled.load(Ordering::Acquire) {
return Poll::Ready("cancelled");
}
parked_waker
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.replace(context.waker().clone());
if !reported_pending {
reported_pending = true;
polled_tx.send(()).expect("cancel helper remains available");
}
Poll::Pending
}));
helper.join().expect("cancel helper does not panic");
assert_eq!(observed, "cancelled");
}
#[test]
fn future_panic_propagates_without_poisoning_kernel_state() {
let panic = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
block_on(async { panic!("future panic sentinel") });
}));
assert!(panic.is_err());
assert_eq!(block_on(async { 23_u8 }), 23);
}
#[test]
fn installed_runtime_context_is_refused_before_poll() {
let runtime = RuntimeBuilder::new()
.worker_threads(1)
.build()
.expect("runtime build");
let polls = Cell::new(0_usize);
let result = runtime.block_on(async {
try_block_on(poll_fn(|_| {
polls.set(polls.get() + 1);
Poll::Ready(31_u8)
}))
});
assert_eq!(result, Err(BlockOnError::RuntimeContext));
assert_eq!(polls.get(), 0);
}
#[test]
fn scheduler_worker_context_is_refused_before_poll() {
let runtime = RuntimeBuilder::new()
.worker_threads(1)
.build()
.expect("runtime build");
let polls = Arc::new(AtomicUsize::new(0));
let task_polls = Arc::clone(&polls);
let task = runtime.handle().spawn(async move {
try_block_on(poll_fn(|_| {
task_polls.fetch_add(1, Ordering::Relaxed);
Poll::Ready(41_u8)
}))
});
assert_eq!(runtime.block_on(task), Err(BlockOnError::RuntimeContext));
assert_eq!(polls.load(Ordering::Relaxed), 0);
}
#[test]
fn foreign_executor_admits_self_contained_nested_future() {
let wake_state = Arc::new(CountingWaker::default());
let waker = Waker::from(Arc::clone(&wake_state));
let mut context = Context::from_waker(&waker);
let polls = Cell::new(0_usize);
let mut outer = std::pin::pin!(async {
try_block_on(poll_fn(|context| {
let current = polls.get();
polls.set(current + 1);
if current == 0 {
context.waker().wake_by_ref();
Poll::Pending
} else {
Poll::Ready(43_u8)
}
}))
});
assert_eq!(outer.as_mut().poll(&mut context), Poll::Ready(Ok(43)));
assert_eq!(polls.get(), 2);
assert_eq!(wake_state.wakes.load(Ordering::Relaxed), 0);
}
#[test]
fn owned_kernel_leaves_lab_runtime_quiescent() {
let mut lab = crate::lab::LabRuntime::new(crate::lab::LabConfig::new(0xF07A_0003));
assert!(lab.is_quiescent());
assert_eq!(block_on(async { 47_u8 }), 47);
assert_eq!(lab.run_until_quiescent(), 0);
assert!(lab.is_quiescent());
}
#[test]
fn blocking_pool_thread_is_admitted() {
let pool = BlockingPool::new(1, 1);
let output = Arc::new(AtomicUsize::new(0));
let task_output = Arc::clone(&output);
let task = pool.spawn(move || {
task_output.store(block_on(async { 37_usize }), Ordering::Release);
});
assert!(task.wait_timeout(Duration::from_secs(2)));
assert_eq!(output.load(Ordering::Acquire), 37);
assert!(pool.shutdown_and_wait(Duration::from_secs(2)));
}
}