use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Mutex, PoisonError};
use std::time::{Duration, Instant};
use crate::macros::sdk_warn;
type Job = Box<dyn FnOnce() + Send + 'static>;
const BACKLOG_WARNING: usize = 10_000;
struct Queue {
jobs: Mutex<Vec<Job>>,
len: AtomicUsize,
}
impl Queue {
fn lock(&self) -> std::sync::MutexGuard<'_, Vec<Job>> {
self.jobs.lock().unwrap_or_else(PoisonError::into_inner)
}
fn push(&self, job: Job) -> usize {
let mut jobs = self.lock();
jobs.push(job);
self.len.store(jobs.len(), Ordering::Release);
jobs.len()
}
fn take(&self) -> Vec<Job> {
if self.len.load(Ordering::Acquire) == 0 {
return Vec::new();
}
let mut jobs = self.lock();
self.len.store(0, Ordering::Release);
std::mem::take(&mut *jobs)
}
fn put_back(&self, left: Vec<Job>) {
let mut jobs = self.lock();
let newer = std::mem::replace(&mut *jobs, left);
jobs.extend(newer);
self.len.store(jobs.len(), Ordering::Release);
}
}
fn queue() -> &'static Queue {
static QUEUE: Queue = Queue {
jobs: Mutex::new(Vec::new()),
len: AtomicUsize::new(0),
};
&QUEUE
}
pub fn post<F>(job: F)
where
F: FnOnce() + Send + 'static,
{
let backlog = queue().push(Box::new(job));
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().len.load(Ordering::Acquire)
}
static BUDGET_NANOS: AtomicU64 = AtomicU64::new(0);
pub fn set_budget(budget: Option<Duration>) {
let nanos = budget.map_or(0, |budget| {
u64::try_from(budget.as_nanos()).unwrap_or(u64::MAX).max(1)
});
BUDGET_NANOS.store(nanos, Ordering::Release);
}
#[must_use]
pub fn budget() -> Option<Duration> {
match BUDGET_NANOS.load(Ordering::Acquire) {
0 => None,
nanos => Some(Duration::from_nanos(nanos)),
}
}
pub fn run_pending() -> usize {
let jobs = queue().take();
if jobs.is_empty() {
return 0;
}
let deadline = budget().map(|budget| Instant::now() + budget);
let mut jobs = jobs.into_iter();
let mut ran = 0;
for job in jobs.by_ref() {
if crate::panic_guard::catch(job).is_err() {
sdk_warn!("a job posted to the main thread panicked; it was dropped");
}
ran += 1;
if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
break;
}
}
let left: Vec<Job> = jobs.collect();
if !left.is_empty() {
queue().put_back(left);
}
ran
}
#[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_budget_leaves_the_rest_for_the_next_drain_in_order() {
let _g = exclusive();
drain_quietly();
let seen = Arc::new(StdMutex::new(Vec::new()));
for i in 0..4 {
let seen = Arc::clone(&seen);
post(move || {
std::thread::sleep(Duration::from_millis(2));
seen.lock().unwrap().push(i);
});
}
set_budget(Some(Duration::from_millis(1)));
assert_eq!(run_pending(), 1);
let late = Arc::clone(&seen);
post(move || late.lock().unwrap().push(99));
set_budget(None);
assert_eq!(run_pending(), 4);
assert_eq!(*seen.lock().unwrap(), vec![0, 1, 2, 3, 99]);
assert_eq!(budget(), None);
}
#[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);
}
}