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 backlog = {
let mut guard = queue().lock().unwrap_or_else(PoisonError::into_inner);
guard.push(Box::new(job));
guard.len()
};
if backlog >= 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!`?",
backlog
);
}
}
}
pub fn post_with<T>(job: impl FnOnce(&mut T) + Send + 'static)
where
T: crate::plugin::SampPlugin + 'static,
{
post(move || {
crate::plugin::with_instance::<T, _>(job);
});
}
pub fn post_with_amx<T, R>(
script: crate::amx::AmxIdent,
job: impl FnOnce(&mut T) -> R + Send + 'static,
) where
T: crate::plugin::SampPlugin + 'static,
R: FnOnce(&samp_sdk::amx::Amx),
{
post(move || {
if crate::amx::get(script).is_none() {
return;
}
let Some(reply) = crate::plugin::with_instance::<T, _>(job) else {
return;
};
if let Some(amx) = crate::amx::get(script) {
reply(amx);
}
});
}
#[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 crate::panic_guard::catch(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};
use crate::test_support::{TestPlugin, exclusive, value};
struct OtherPlugin;
impl crate::plugin::SampPlugin for OtherPlugin {}
fn drain_quietly() {
run_pending();
}
#[test]
fn jobs_run_in_the_order_they_were_posted() {
let _g = exclusive();
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 = exclusive();
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 = exclusive();
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 = exclusive();
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_job_reaches_the_plugin_from_a_worker_thread() {
let _g = exclusive();
drain_quietly();
let before = value();
std::thread::spawn(|| {
post_with::<TestPlugin>(|plugin| plugin.value += 5);
})
.join()
.unwrap();
assert_eq!(run_pending(), 1);
assert_eq!(value(), before + 5);
}
#[test]
fn a_job_naming_the_wrong_plugin_type_is_skipped() {
let _g = exclusive();
drain_quietly();
let before = value();
post_with::<OtherPlugin>(|_| unreachable!("must not run"));
assert_eq!(run_pending(), 1, "the job ran and refused itself");
assert_eq!(value(), before, "the plugin was left alone");
}
#[test]
fn a_reply_to_a_script_that_went_away_is_dropped() {
let _g = exclusive();
drain_quietly();
let before = value();
let gone =
crate::amx::AmxIdent::from(std::ptr::without_provenance_mut(0xDEAD_BEEF_u32 as _));
post_with_amx::<TestPlugin, _>(gone, |plugin| {
plugin.value += 1;
|_amx: &samp_sdk::amx::Amx| unreachable!("there is no script to answer")
});
assert_eq!(run_pending(), 1);
assert_eq!(
value(),
before,
"the plugin is not disturbed when there is nothing to report to"
);
}
#[test]
fn posting_from_many_threads_while_draining_loses_nothing() {
let _g = exclusive();
drain_quietly();
const THREADS: usize = 8;
const JOBS: usize = 5_000;
let ran = Arc::new(AtomicUsize::new(0));
let workers: Vec<_> = (0..THREADS)
.map(|t| {
let ran = Arc::clone(&ran);
std::thread::spawn(move || {
for i in 0..JOBS {
let ran = Arc::clone(&ran);
if i % 997 == t {
post(|| panic!("a job blew up"));
}
post(move || {
ran.fetch_add(1, Ordering::Relaxed);
});
}
})
})
.collect();
while workers.iter().any(|w| !w.is_finished()) {
run_pending();
}
for worker in workers {
worker.join().unwrap();
}
run_pending();
assert_eq!(ran.load(Ordering::Relaxed), THREADS * JOBS);
assert_eq!(pending(), 0);
}
#[test]
fn a_worker_thread_can_post() {
let _g = exclusive();
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);
}
}