Skip to main content

samp/
mainthread.rs

1//! Handing work back to the server's main thread.
2//!
3//! Both servers are single-threaded: every native, callback and tick runs on
4//! one thread, and the AMX VM must only ever be touched from it. A plugin that
5//! does I/O — HTTP, a database, SMTP — has to do that work elsewhere and bring
6//! the result back, because blocking the main thread freezes the server for
7//! every player.
8//!
9//! This module is that return path. A worker thread calls [`post`] with a
10//! closure; the closure runs on the main thread on the next tick.
11//!
12//! ```rust,no_run
13//! use samp::exec_public;
14//! use samp::prelude::*;
15//! # fn example(amx_ident: samp::amx::AmxIdent) {
16//! std::thread::spawn(move || {
17//!     let answer = 42; // ... the slow work ...
18//!
19//!     samp::mainthread::post(move || {
20//!         // Back on the main thread: safe to touch the VM.
21//!         if let Some(amx) = samp::amx::get(amx_ident) {
22//!             let _ = exec_public!(amx, "OnWorkDone", answer);
23//!         }
24//!     });
25//! });
26//! # }
27//! ```
28//!
29//! An [`AmxIdent`] is a plain address wrapper, so it crosses thread boundaries;
30//! the `&Amx` it resolves to does not, which is why the job looks it up again
31//! after arriving. [`samp::amx::get`] returns `None` if the script was unloaded
32//! meanwhile — the case this pattern makes easy to handle instead of dangling.
33//!
34//! ## Draining
35//!
36//! Jobs run from the same place as [`SampPlugin::on_tick`], so the plugin needs
37//! `samp::plugin::enable_tick()` in `initialize_plugin!`. Without it nothing
38//! drains the queue and jobs pile up; the SDK logs a warning once the backlog
39//! is large enough to be a mistake rather than a burst. A plugin that does not
40//! want the tick can call [`run_pending`] from wherever it prefers, such as
41//! inside a native.
42//!
43//! A job posted while the queue is draining runs on the **next** tick, not the
44//! current one. That keeps a job that re-posts itself from spinning forever
45//! inside one tick.
46//!
47//! [`AmxIdent`]: crate::amx::AmxIdent
48//! [`samp::amx::get`]: crate::amx::get
49//! [`SampPlugin::on_tick`]: crate::plugin::SampPlugin::on_tick
50
51use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
52use std::sync::{Mutex, PoisonError};
53use std::time::{Duration, Instant};
54
55use crate::macros::sdk_warn;
56
57/// Work queued from another thread, to run on the main thread.
58type Job = Box<dyn FnOnce() + Send + 'static>;
59
60/// Backlog at which the SDK warns once. A plugin bursting a few hundred jobs
61/// between ticks is normal; five figures means nothing is draining them.
62const BACKLOG_WARNING: usize = 10_000;
63
64/// The jobs, and their count mirrored outside the lock: the main thread asks
65/// on every tick, and almost always the answer is "none", which then costs a
66/// load instead of locking and unlocking the mutex.
67struct Queue {
68    jobs: Mutex<Vec<Job>>,
69    /// `jobs.len()`, written only while `jobs` is locked.
70    len: AtomicUsize,
71}
72
73impl Queue {
74    fn lock(&self) -> std::sync::MutexGuard<'_, Vec<Job>> {
75        // A poisoned lock means some earlier holder panicked. The queue itself
76        // is still consistent — a `Vec` of jobs — so recovering beats refusing
77        // every later post for the lifetime of the process.
78        self.jobs.lock().unwrap_or_else(PoisonError::into_inner)
79    }
80
81    /// Appends `job`; returns the new length.
82    fn push(&self, job: Job) -> usize {
83        let mut jobs = self.lock();
84        jobs.push(job);
85        self.len.store(jobs.len(), Ordering::Release);
86        jobs.len()
87    }
88
89    /// Takes every queued job, without locking when there is none.
90    fn take(&self) -> Vec<Job> {
91        if self.len.load(Ordering::Acquire) == 0 {
92            return Vec::new();
93        }
94        let mut jobs = self.lock();
95        self.len.store(0, Ordering::Release);
96        std::mem::take(&mut *jobs)
97    }
98
99    /// Puts `left` back ahead of whatever was posted since it was taken.
100    fn put_back(&self, left: Vec<Job>) {
101        let mut jobs = self.lock();
102        let newer = std::mem::replace(&mut *jobs, left);
103        jobs.extend(newer);
104        self.len.store(jobs.len(), Ordering::Release);
105    }
106}
107
108fn queue() -> &'static Queue {
109    static QUEUE: Queue = Queue {
110        jobs: Mutex::new(Vec::new()),
111        len: AtomicUsize::new(0),
112    };
113    &QUEUE
114}
115
116/// Queues `job` to run on the main thread at the next tick.
117///
118/// Callable from any thread. The job runs once, in the order it was posted.
119///
120/// A panic inside the job is caught and logged: it must not unwind into the
121/// server's C++ frames, which would abort the process.
122pub fn post<F>(job: F)
123where
124    F: FnOnce() + Send + 'static,
125{
126    let backlog = queue().push(Box::new(job));
127
128    // Outside the lock: the warning goes through the logger, which may post.
129    if backlog >= BACKLOG_WARNING {
130        static WARNED: AtomicBool = AtomicBool::new(false);
131        if !WARNED.swap(true, Ordering::Relaxed) {
132            sdk_warn!(
133                "{} jobs queued for the main thread and nothing is draining them \
134                 — is `samp::plugin::enable_tick()` missing from `initialize_plugin!`?",
135                backlog
136            );
137        }
138    }
139}
140
141/// Queues `job` to run on the main thread with the plugin instance.
142///
143/// The plugin's own state is where the result of background work usually has to
144/// land, and a worker thread cannot reach it: the plugin is not `Sync`, and only
145/// the main thread may touch it. This posts the job and hands it `&mut T` at the
146/// moment it runs.
147///
148/// ```rust,no_run
149/// # use samp::prelude::*;
150/// # struct Mailer { sent: u32 }
151/// # impl SampPlugin for Mailer {}
152/// # fn example() {
153/// std::thread::spawn(move || {
154///     let delivered = true; // ... the slow work ...
155///
156///     samp::mainthread::post_with::<Mailer>(move |plugin| {
157///         if delivered {
158///             plugin.sent += 1;
159///         }
160///     });
161/// });
162/// # }
163/// ```
164///
165/// `T` must be the type `initialize_plugin!` creates; naming another logs a
166/// warning and skips the job, rather than reinterpreting the plugin's bytes.
167///
168/// **Do not call into Pawn from inside the closure.** A `public` that re-enters
169/// one of this plugin's natives would take a second `&mut` to the same plugin.
170/// Collect what to send, let the closure end, and send it after — or use
171/// [`post_with_amx`], which is that shape already.
172pub fn post_with<T>(job: impl FnOnce(&mut T) + Send + 'static)
173where
174    T: crate::plugin::SampPlugin + 'static,
175{
176    post(move || {
177        crate::plugin::with_instance::<T, _>(job);
178    });
179}
180
181/// Queues `job` to run on the main thread with the plugin and a resolved
182/// [`Amx`], in that order.
183///
184/// The common shape of background work coming back: record the result in the
185/// plugin, then tell the script. The closure returns what to do with the script,
186/// and the SDK runs it after the plugin borrow has ended, so calling a `public`
187/// from there is safe.
188///
189/// ```rust,no_run
190/// # use samp::prelude::*;
191/// # use samp::exec_public;
192/// # struct Mailer { sent: u32 }
193/// # impl SampPlugin for Mailer {}
194/// # fn example(script: samp::amx::AmxIdent) {
195/// samp::mainthread::post_with_amx::<Mailer, _>(script, |plugin| {
196///     plugin.sent += 1;
197///     let total = plugin.sent;
198///
199///     // Runs next, with no borrow of the plugin alive.
200///     move |amx: &Amx| {
201///         let _ = exec_public!(amx, "OnMailSent", total);
202///     }
203/// });
204/// # }
205/// ```
206///
207/// The job is skipped, quietly, if the script was unloaded before it ran: that
208/// is the normal end of a gamemode restart, not a fault.
209pub fn post_with_amx<T, R>(
210    script: crate::amx::AmxIdent,
211    job: impl FnOnce(&mut T) -> R + Send + 'static,
212) where
213    T: crate::plugin::SampPlugin + 'static,
214    R: FnOnce(&samp_sdk::amx::Amx),
215{
216    post(move || {
217        // Resolved first: with the script gone there is nothing to report, so
218        // the plugin is not disturbed either.
219        if crate::amx::get(script).is_none() {
220            return;
221        }
222        let Some(reply) = crate::plugin::with_instance::<T, _>(job) else {
223            return;
224        };
225        if let Some(amx) = crate::amx::get(script) {
226            reply(amx);
227        }
228    });
229}
230
231/// Number of jobs waiting to run.
232#[must_use]
233pub fn pending() -> usize {
234    queue().len.load(Ordering::Acquire)
235}
236
237/// Time one drain may spend running jobs, in nanoseconds; `0` for no limit.
238static BUDGET_NANOS: AtomicU64 = AtomicU64::new(0);
239
240/// Limits how long one drain spends running jobs; `None` (the default) runs
241/// every queued job.
242///
243/// The server is frozen while jobs run, so a burst of ten thousand replies
244/// arriving at once becomes one long stall. With a budget, a drain stops at the
245/// first job that ends past it and leaves the rest, in order, for the next
246/// tick — the stall is spread over several ticks instead. A job is never cut
247/// short, and at least one runs per drain, so the queue always advances.
248///
249/// ```rust,no_run
250/// // At most ~2 ms of a 5 ms SA-MP tick goes to background replies.
251/// samp::mainthread::set_budget(Some(std::time::Duration::from_millis(2)));
252/// ```
253pub fn set_budget(budget: Option<Duration>) {
254    let nanos = budget.map_or(0, |budget| {
255        u64::try_from(budget.as_nanos()).unwrap_or(u64::MAX).max(1)
256    });
257    BUDGET_NANOS.store(nanos, Ordering::Release);
258}
259
260/// The limit [`set_budget`] set, if any.
261#[must_use]
262pub fn budget() -> Option<Duration> {
263    match BUDGET_NANOS.load(Ordering::Acquire) {
264        0 => None,
265        nanos => Some(Duration::from_nanos(nanos)),
266    }
267}
268
269/// Runs the queued jobs and returns how many ran: all of them, or as many as
270/// fit in the [`budget`] when one is set.
271///
272/// Called by the SDK on each tick. A plugin only needs it when it drives the
273/// queue itself — with the tick disabled, for instance.
274///
275/// Must be called from the main thread: the jobs assume they are on it.
276pub fn run_pending() -> usize {
277    // Take the jobs out under the lock and run them with it released: a job is
278    // allowed to post more work (which lands on the next drain), and running
279    // while holding the lock would deadlock on that.
280    let jobs = queue().take();
281    if jobs.is_empty() {
282        return 0;
283    }
284
285    let deadline = budget().map(|budget| Instant::now() + budget);
286    let mut jobs = jobs.into_iter();
287    let mut ran = 0;
288    for job in jobs.by_ref() {
289        if crate::panic_guard::catch(job).is_err() {
290            sdk_warn!("a job posted to the main thread panicked; it was dropped");
291        }
292        ran += 1;
293        if deadline.is_some_and(|deadline| Instant::now() >= deadline) {
294            break;
295        }
296    }
297
298    // Over budget: what is left goes back ahead of anything posted meanwhile,
299    // keeping the order jobs were posted in.
300    let left: Vec<Job> = jobs.collect();
301    if !left.is_empty() {
302        queue().put_back(left);
303    }
304    ran
305}
306
307#[cfg(test)]
308mod tests {
309    use super::*;
310    use std::sync::atomic::AtomicUsize;
311    use std::sync::{Arc, Mutex as StdMutex};
312
313    use crate::test_support::{TestPlugin, exclusive, value};
314
315    /// A second plugin type, to check that naming it reaches nothing.
316    struct OtherPlugin;
317    impl crate::plugin::SampPlugin for OtherPlugin {}
318
319    /// Leaves the queue empty, so a test starts from a known state.
320    fn drain_quietly() {
321        run_pending();
322    }
323
324    #[test]
325    fn jobs_run_in_the_order_they_were_posted() {
326        let _g = exclusive();
327        drain_quietly();
328
329        let seen = Arc::new(StdMutex::new(Vec::new()));
330        for i in 0..3 {
331            let seen = Arc::clone(&seen);
332            post(move || seen.lock().unwrap().push(i));
333        }
334
335        assert_eq!(run_pending(), 3);
336        assert_eq!(*seen.lock().unwrap(), vec![0, 1, 2]);
337    }
338
339    #[test]
340    fn pending_counts_and_draining_clears() {
341        let _g = exclusive();
342        drain_quietly();
343
344        post(|| {});
345        post(|| {});
346        assert_eq!(pending(), 2);
347
348        assert_eq!(run_pending(), 2);
349        assert_eq!(pending(), 0);
350        assert_eq!(run_pending(), 0);
351    }
352
353    #[test]
354    fn a_job_posted_while_draining_waits_for_the_next_drain() {
355        let _g = exclusive();
356        drain_quietly();
357
358        let runs = Arc::new(AtomicUsize::new(0));
359        let inner = Arc::clone(&runs);
360        post(move || {
361            let deeper = Arc::clone(&inner);
362            post(move || {
363                deeper.fetch_add(1, Ordering::Relaxed);
364            });
365        });
366
367        assert_eq!(run_pending(), 1, "only the outer job runs in this drain");
368        assert_eq!(runs.load(Ordering::Relaxed), 0);
369        assert_eq!(run_pending(), 1, "the re-posted job runs in the next one");
370        assert_eq!(runs.load(Ordering::Relaxed), 1);
371    }
372
373    #[test]
374    fn a_panicking_job_does_not_stop_the_ones_after_it() {
375        let _g = exclusive();
376        drain_quietly();
377
378        let ran = Arc::new(AtomicUsize::new(0));
379        let after = Arc::clone(&ran);
380        post(|| panic!("job blew up"));
381        post(move || {
382            after.fetch_add(1, Ordering::Relaxed);
383        });
384
385        assert_eq!(run_pending(), 2);
386        assert_eq!(ran.load(Ordering::Relaxed), 1);
387
388        // And the queue still works afterwards.
389        let later = Arc::clone(&ran);
390        post(move || {
391            later.fetch_add(1, Ordering::Relaxed);
392        });
393        assert_eq!(run_pending(), 1);
394        assert_eq!(ran.load(Ordering::Relaxed), 2);
395    }
396
397    #[test]
398    fn a_job_reaches_the_plugin_from_a_worker_thread() {
399        let _g = exclusive();
400        drain_quietly();
401
402        let before = value();
403
404        std::thread::spawn(|| {
405            post_with::<TestPlugin>(|plugin| plugin.value += 5);
406        })
407        .join()
408        .unwrap();
409
410        assert_eq!(run_pending(), 1);
411        assert_eq!(value(), before + 5);
412    }
413
414    #[test]
415    fn a_job_naming_the_wrong_plugin_type_is_skipped() {
416        let _g = exclusive();
417        drain_quietly();
418
419        let before = value();
420        post_with::<OtherPlugin>(|_| unreachable!("must not run"));
421
422        assert_eq!(run_pending(), 1, "the job ran and refused itself");
423        assert_eq!(value(), before, "the plugin was left alone");
424    }
425
426    #[test]
427    fn a_reply_to_a_script_that_went_away_is_dropped() {
428        let _g = exclusive();
429        drain_quietly();
430
431        let before = value();
432        // An ident no AMX was ever registered under: the script is gone.
433        let gone =
434            crate::amx::AmxIdent::from(std::ptr::without_provenance_mut(0xDEAD_BEEF_u32 as _));
435
436        post_with_amx::<TestPlugin, _>(gone, |plugin| {
437            plugin.value += 1;
438            |_amx: &samp_sdk::amx::Amx| unreachable!("there is no script to answer")
439        });
440
441        assert_eq!(run_pending(), 1);
442        assert_eq!(
443            value(),
444            before,
445            "the plugin is not disturbed when there is nothing to report to"
446        );
447    }
448
449    #[test]
450    fn posting_from_many_threads_while_draining_loses_nothing() {
451        let _g = exclusive();
452        drain_quietly();
453
454        const THREADS: usize = 8;
455        const JOBS: usize = 5_000;
456        let ran = Arc::new(AtomicUsize::new(0));
457        let workers: Vec<_> = (0..THREADS)
458            .map(|t| {
459                let ran = Arc::clone(&ran);
460                std::thread::spawn(move || {
461                    for i in 0..JOBS {
462                        let ran = Arc::clone(&ran);
463                        if i % 997 == t {
464                            post(|| panic!("a job blew up"));
465                        }
466                        post(move || {
467                            ran.fetch_add(1, Ordering::Relaxed);
468                        });
469                    }
470                })
471            })
472            .collect();
473
474        // The main thread drains while the workers are still posting.
475        while workers.iter().any(|w| !w.is_finished()) {
476            run_pending();
477        }
478        for worker in workers {
479            worker.join().unwrap();
480        }
481        run_pending();
482        assert_eq!(ran.load(Ordering::Relaxed), THREADS * JOBS);
483        assert_eq!(pending(), 0);
484    }
485
486    #[test]
487    fn a_budget_leaves_the_rest_for_the_next_drain_in_order() {
488        let _g = exclusive();
489        drain_quietly();
490
491        let seen = Arc::new(StdMutex::new(Vec::new()));
492        for i in 0..4 {
493            let seen = Arc::clone(&seen);
494            post(move || {
495                std::thread::sleep(Duration::from_millis(2));
496                seen.lock().unwrap().push(i);
497            });
498        }
499        set_budget(Some(Duration::from_millis(1)));
500        // One job runs past the budget: the drain stops after it.
501        assert_eq!(run_pending(), 1);
502        let late = Arc::clone(&seen);
503        post(move || late.lock().unwrap().push(99));
504        set_budget(None);
505        assert_eq!(run_pending(), 4);
506        assert_eq!(*seen.lock().unwrap(), vec![0, 1, 2, 3, 99]);
507        assert_eq!(budget(), None);
508    }
509
510    #[test]
511    fn a_worker_thread_can_post() {
512        let _g = exclusive();
513        drain_quietly();
514
515        let ran = Arc::new(AtomicUsize::new(0));
516        let from_worker = Arc::clone(&ran);
517        std::thread::spawn(move || {
518            post(move || {
519                from_worker.fetch_add(1, Ordering::Relaxed);
520            });
521        })
522        .join()
523        .unwrap();
524
525        assert_eq!(run_pending(), 1);
526        assert_eq!(ran.load(Ordering::Relaxed), 1);
527    }
528}