use std::any::Any;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::{Duration, Instant};
use teksilo_async::{AsyncRuntimeHandle, spawn_blocking};
use teksilo_core::{AppEventPoster, Signal, SubscriptionId};
fn pump_until(rt: &AsyncRuntimeHandle, label: &str, mut cond: impl FnMut() -> bool) {
let start = Instant::now();
while !cond() {
rt.tick();
if start.elapsed() > Duration::from_secs(5) {
panic!("timed out waiting for: {label}");
}
std::thread::sleep(Duration::from_millis(1));
}
}
#[test]
fn spawn_local_awaiting_spawn_blocking_sets_signal() {
let rt = AsyncRuntimeHandle::new();
let result = Signal::new(0_i32);
let sink = result.clone();
rt.spawn_local(async move {
let value = spawn_blocking(|| 21 * 2).await.expect("worker ran ok");
sink.set(value); })
.detach();
pump_until(&rt, "spawn_blocking result", || result.get() == 42);
assert_eq!(result.get(), 42);
}
#[test]
fn dropping_task_handle_cancels_the_continuation() {
let rt = AsyncRuntimeHandle::new();
let result = Signal::new(0_i32);
let sink = result.clone();
let handle = rt.spawn_local(async move {
let value = spawn_blocking(|| {
std::thread::sleep(Duration::from_millis(40));
99
})
.await
.unwrap_or(0);
sink.set(value);
});
rt.tick();
drop(handle);
let start = Instant::now();
while start.elapsed() < Duration::from_millis(150) {
rt.tick();
std::thread::sleep(Duration::from_millis(2));
}
assert_eq!(
result.get(),
0,
"a cancelled task must not run its continuation"
);
}
#[test]
fn detached_task_runs_to_completion() {
let rt = AsyncRuntimeHandle::new();
let done = Signal::new(false);
let sink = done.clone();
rt.spawn_local(async move {
let _ = spawn_blocking(|| ()).await;
sink.set(true);
})
.detach();
pump_until(&rt, "detached task completion", || done.get());
assert!(done.get());
}
#[test]
fn idle_tick_reports_no_work() {
let rt = AsyncRuntimeHandle::new();
assert!(!rt.tick());
assert!(!rt.poll_source().get());
}
#[derive(Default)]
struct CountingPoster {
nudges: AtomicUsize,
}
impl AppEventPoster for CountingPoster {
fn post_subscription_event(&self, _sub_id: SubscriptionId, _event: Box<dyn Any + Send>) {}
fn post_external(&self, _payload: Box<dyn Any + Send>) {
self.nudges.fetch_add(1, Ordering::SeqCst);
}
}
#[test]
fn off_thread_completion_nudges_the_event_loop_poster() {
let rt = AsyncRuntimeHandle::new();
let poster = Arc::new(CountingPoster::default());
rt.set_poster(poster.clone());
let done = Signal::new(false);
let sink = done.clone();
rt.spawn_local(async move {
let _ = spawn_blocking(|| std::thread::sleep(Duration::from_millis(20))).await;
sink.set(true);
})
.detach();
pump_until(&rt, "off-thread completion", || done.get());
assert!(done.get());
assert!(
poster.nudges.load(Ordering::SeqCst) >= 1,
"off-thread worker completion must nudge the event-loop poster at least \
once (got {})",
poster.nudges.load(Ordering::SeqCst),
);
}
#[test]
fn spawn_blocking_panic_surfaces_as_error() {
let rt = AsyncRuntimeHandle::new();
let outcome = Signal::new(0_i32);
let sink = outcome.clone();
rt.spawn_local(async move {
let result: Result<i32, _> = spawn_blocking(|| panic!("boom")).await;
sink.set(if result.is_err() { 1 } else { 2 });
})
.detach();
pump_until(&rt, "spawn_blocking panic outcome", || outcome.get() != 0);
assert_eq!(
outcome.get(),
1,
"a panicking closure must resolve to Err, not crash the executor"
);
}