use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Mutex, OnceLock, PoisonError};
use crate::macros::sdk_warn;
type Job = Box<dyn FnOnce() + Send + 'static>;
const BACKLOG_WARNING: usize = 10_000;
fn queue() -> &'static Mutex<Vec<Job>> {
static QUEUE: OnceLock<Mutex<Vec<Job>>> = OnceLock::new();
QUEUE.get_or_init(|| Mutex::new(Vec::new()))
}
pub fn post<F>(job: F)
where
F: FnOnce() + Send + 'static,
{
let mut guard = queue().lock().unwrap_or_else(PoisonError::into_inner);
guard.push(Box::new(job));
if guard.len() >= BACKLOG_WARNING {
static WARNED: AtomicBool = AtomicBool::new(false);
if !WARNED.swap(true, Ordering::Relaxed) {
sdk_warn!(
"{} jobs queued for the main thread and nothing is draining them \
— is `samp::plugin::enable_tick()` missing from `initialize_plugin!`?",
guard.len()
);
}
}
}
#[must_use]
pub fn pending() -> usize {
queue().lock().unwrap_or_else(PoisonError::into_inner).len()
}
pub fn run_pending() -> usize {
let jobs: Vec<Job> = {
let mut guard = queue().lock().unwrap_or_else(PoisonError::into_inner);
std::mem::take(&mut *guard)
};
let count = jobs.len();
for job in jobs {
if std::panic::catch_unwind(std::panic::AssertUnwindSafe(job)).is_err() {
sdk_warn!("a job posted to the main thread panicked; it was dropped");
}
}
count
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicUsize;
use std::sync::{Arc, Mutex as StdMutex};
static TEST_LOCK: StdMutex<()> = StdMutex::new(());
fn drain_quietly() {
run_pending();
}
#[test]
fn jobs_run_in_the_order_they_were_posted() {
let _g = TEST_LOCK.lock().unwrap();
drain_quietly();
let seen = Arc::new(StdMutex::new(Vec::new()));
for i in 0..3 {
let seen = Arc::clone(&seen);
post(move || seen.lock().unwrap().push(i));
}
assert_eq!(run_pending(), 3);
assert_eq!(*seen.lock().unwrap(), vec![0, 1, 2]);
}
#[test]
fn pending_counts_and_draining_clears() {
let _g = TEST_LOCK.lock().unwrap();
drain_quietly();
post(|| {});
post(|| {});
assert_eq!(pending(), 2);
assert_eq!(run_pending(), 2);
assert_eq!(pending(), 0);
assert_eq!(run_pending(), 0);
}
#[test]
fn a_job_posted_while_draining_waits_for_the_next_drain() {
let _g = TEST_LOCK.lock().unwrap();
drain_quietly();
let runs = Arc::new(AtomicUsize::new(0));
let inner = Arc::clone(&runs);
post(move || {
let deeper = Arc::clone(&inner);
post(move || {
deeper.fetch_add(1, Ordering::Relaxed);
});
});
assert_eq!(run_pending(), 1, "only the outer job runs in this drain");
assert_eq!(runs.load(Ordering::Relaxed), 0);
assert_eq!(run_pending(), 1, "the re-posted job runs in the next one");
assert_eq!(runs.load(Ordering::Relaxed), 1);
}
#[test]
fn a_panicking_job_does_not_stop_the_ones_after_it() {
let _g = TEST_LOCK.lock().unwrap();
drain_quietly();
let ran = Arc::new(AtomicUsize::new(0));
let after = Arc::clone(&ran);
post(|| panic!("job blew up"));
post(move || {
after.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_pending(), 2);
assert_eq!(ran.load(Ordering::Relaxed), 1);
let later = Arc::clone(&ran);
post(move || {
later.fetch_add(1, Ordering::Relaxed);
});
assert_eq!(run_pending(), 1);
assert_eq!(ran.load(Ordering::Relaxed), 2);
}
#[test]
fn a_worker_thread_can_post() {
let _g = TEST_LOCK.lock().unwrap();
drain_quietly();
let ran = Arc::new(AtomicUsize::new(0));
let from_worker = Arc::clone(&ran);
std::thread::spawn(move || {
post(move || {
from_worker.fetch_add(1, Ordering::Relaxed);
});
})
.join()
.unwrap();
assert_eq!(run_pending(), 1);
assert_eq!(ran.load(Ordering::Relaxed), 1);
}
}