use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use super::runtime::with_rt;
pub use crate::base::FrameRequester;
type PostedJob = Box<dyn FnOnce() + Send>;
pub(crate) struct RemoteShared {
posted: Mutex<Vec<PostedJob>>,
woken: AtomicBool,
waker: Mutex<Option<Box<dyn Fn() + Send + Sync>>>,
}
impl RemoteShared {
pub(crate) fn new() -> Self {
RemoteShared {
posted: Mutex::new(Vec::new()),
woken: AtomicBool::new(false),
waker: Mutex::new(None),
}
}
fn notify(&self) {
self.woken.store(true, Ordering::Release);
if let Some(waker) = self.waker.lock().expect("waker lock").as_ref() {
waker();
}
}
}
#[derive(Clone)]
pub struct WakeHandle {
shared: Arc<RemoteShared>,
}
impl WakeHandle {
pub fn wake(&self) {
self.shared.notify();
}
pub fn post(&self, f: impl FnOnce() + Send + 'static) {
self.shared
.posted
.lock()
.expect("posted lock")
.push(Box::new(f));
self.shared.notify();
}
}
pub fn wake_handle() -> WakeHandle {
WakeHandle {
shared: with_rt(|rt| rt.remote.clone()),
}
}
pub fn set_wake_callback(f: impl Fn() + Send + Sync + 'static) {
let shared = with_rt(|rt| rt.remote.clone());
*shared.waker.lock().expect("waker lock") = Some(Box::new(f));
}
pub fn wake_pending() -> bool {
with_rt(|rt| rt.remote.clone())
.woken
.load(Ordering::Acquire)
}
pub fn drain_posted() -> usize {
let shared = with_rt(|rt| rt.remote.clone());
shared.woken.store(false, Ordering::Release);
let jobs: Vec<PostedJob> = {
let mut posted = shared.posted.lock().expect("posted lock");
std::mem::take(&mut *posted)
};
let count = jobs.len();
for job in jobs {
job(); }
count
}
pub fn set_frame_requester(requester: std::rc::Rc<dyn FrameRequester>) {
with_rt(|rt| rt.frame_requester = Some(requester));
}
pub fn request_frame() {
let requester = with_rt(|rt| {
if rt.frame_requested {
None
} else {
rt.frame_requested = true;
rt.frame_requester.clone()
}
});
if let Some(r) = requester {
r.request_frame(); }
}
pub fn take_frame_request() -> bool {
with_rt(|rt| std::mem::take(&mut rt.frame_requested))
}
pub fn spawn_worker(
label: &'static str,
f: impl FnOnce() + Send + 'static,
) -> std::thread::JoinHandle<()> {
let handle = wake_handle();
std::thread::Builder::new()
.name(format!("abstracttui-worker-{label}"))
.spawn(move || {
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(f));
if let Err(payload) = result {
let msg = panic_text(payload.as_ref());
let text = format!("background worker '{label}' died: {msg}");
handle.post(move || super::diag::record_worker_failure(text));
}
})
.expect("spawn worker thread")
}
fn panic_text(payload: &(dyn std::any::Any + Send)) -> String {
if let Some(s) = payload.downcast_ref::<&str>() {
(*s).to_string()
} else if let Some(s) = payload.downcast_ref::<String>() {
s.clone()
} else {
"non-string panic payload".to_string()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicUsize;
#[test]
fn posted_work_runs_on_draining_thread() {
let handle = wake_handle();
let hits = Arc::new(AtomicUsize::new(0));
let h2 = hits.clone();
std::thread::spawn(move || {
h2.fetch_add(1, Ordering::SeqCst);
handle.post(move || {
});
})
.join()
.expect("thread");
assert!(wake_pending());
assert_eq!(drain_posted(), 1);
assert!(!wake_pending());
assert_eq!(hits.load(Ordering::SeqCst), 1);
}
#[test]
fn frame_requests_coalesce() {
struct Counter(Arc<AtomicUsize>);
impl FrameRequester for Counter {
fn request_frame(&self) {
self.0.fetch_add(1, Ordering::SeqCst);
}
}
let calls = Arc::new(AtomicUsize::new(0));
set_frame_requester(std::rc::Rc::new(Counter(calls.clone())));
let _ = take_frame_request(); request_frame();
request_frame();
request_frame();
assert_eq!(calls.load(Ordering::SeqCst), 1, "requests must coalesce");
assert!(take_frame_request());
assert!(!take_frame_request());
request_frame();
assert_eq!(calls.load(Ordering::SeqCst), 2);
assert!(take_frame_request());
}
}