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, Ordering};
52use std::sync::{Mutex, OnceLock, PoisonError};
53
54use crate::macros::sdk_warn;
55
56/// Work queued from another thread, to run on the main thread.
57type Job = Box<dyn FnOnce() + Send + 'static>;
58
59/// Backlog at which the SDK warns once. A plugin bursting a few hundred jobs
60/// between ticks is normal; five figures means nothing is draining them.
61const BACKLOG_WARNING: usize = 10_000;
62
63fn queue() -> &'static Mutex<Vec<Job>> {
64    static QUEUE: OnceLock<Mutex<Vec<Job>>> = OnceLock::new();
65    QUEUE.get_or_init(|| Mutex::new(Vec::new()))
66}
67
68/// Queues `job` to run on the main thread at the next tick.
69///
70/// Callable from any thread. The job runs once, in the order it was posted.
71///
72/// A panic inside the job is caught and logged: it must not unwind into the
73/// server's C++ frames, which would abort the process.
74pub fn post<F>(job: F)
75where
76    F: FnOnce() + Send + 'static,
77{
78    // A poisoned lock means some earlier holder panicked. The queue itself is
79    // still consistent — a `Vec` of jobs — so recovering beats refusing every
80    // later post for the lifetime of the process.
81    let mut guard = queue().lock().unwrap_or_else(PoisonError::into_inner);
82    guard.push(Box::new(job));
83
84    if guard.len() >= BACKLOG_WARNING {
85        static WARNED: AtomicBool = AtomicBool::new(false);
86        if !WARNED.swap(true, Ordering::Relaxed) {
87            sdk_warn!(
88                "{} jobs queued for the main thread and nothing is draining them \
89                 — is `samp::plugin::enable_tick()` missing from `initialize_plugin!`?",
90                guard.len()
91            );
92        }
93    }
94}
95
96/// Number of jobs waiting to run.
97#[must_use]
98pub fn pending() -> usize {
99    queue().lock().unwrap_or_else(PoisonError::into_inner).len()
100}
101
102/// Runs every queued job and returns how many ran.
103///
104/// Called by the SDK on each tick. A plugin only needs it when it drives the
105/// queue itself — with the tick disabled, for instance.
106///
107/// Must be called from the main thread: the jobs assume they are on it.
108pub fn run_pending() -> usize {
109    // Take the jobs out under the lock and run them with it released: a job is
110    // allowed to post more work (which lands on the next drain), and running
111    // while holding the lock would deadlock on that.
112    let jobs: Vec<Job> = {
113        let mut guard = queue().lock().unwrap_or_else(PoisonError::into_inner);
114        std::mem::take(&mut *guard)
115    };
116
117    let count = jobs.len();
118    for job in jobs {
119        if std::panic::catch_unwind(std::panic::AssertUnwindSafe(job)).is_err() {
120            sdk_warn!("a job posted to the main thread panicked; it was dropped");
121        }
122    }
123    count
124}
125
126#[cfg(test)]
127mod tests {
128    use super::*;
129    use std::sync::atomic::AtomicUsize;
130    use std::sync::{Arc, Mutex as StdMutex};
131
132    /// The queue is process-wide, so the tests take turns on it.
133    static TEST_LOCK: StdMutex<()> = StdMutex::new(());
134
135    fn drain_quietly() {
136        run_pending();
137    }
138
139    #[test]
140    fn jobs_run_in_the_order_they_were_posted() {
141        let _g = TEST_LOCK.lock().unwrap();
142        drain_quietly();
143
144        let seen = Arc::new(StdMutex::new(Vec::new()));
145        for i in 0..3 {
146            let seen = Arc::clone(&seen);
147            post(move || seen.lock().unwrap().push(i));
148        }
149
150        assert_eq!(run_pending(), 3);
151        assert_eq!(*seen.lock().unwrap(), vec![0, 1, 2]);
152    }
153
154    #[test]
155    fn pending_counts_and_draining_clears() {
156        let _g = TEST_LOCK.lock().unwrap();
157        drain_quietly();
158
159        post(|| {});
160        post(|| {});
161        assert_eq!(pending(), 2);
162
163        assert_eq!(run_pending(), 2);
164        assert_eq!(pending(), 0);
165        assert_eq!(run_pending(), 0);
166    }
167
168    #[test]
169    fn a_job_posted_while_draining_waits_for_the_next_drain() {
170        let _g = TEST_LOCK.lock().unwrap();
171        drain_quietly();
172
173        let runs = Arc::new(AtomicUsize::new(0));
174        let inner = Arc::clone(&runs);
175        post(move || {
176            let deeper = Arc::clone(&inner);
177            post(move || {
178                deeper.fetch_add(1, Ordering::Relaxed);
179            });
180        });
181
182        assert_eq!(run_pending(), 1, "only the outer job runs in this drain");
183        assert_eq!(runs.load(Ordering::Relaxed), 0);
184        assert_eq!(run_pending(), 1, "the re-posted job runs in the next one");
185        assert_eq!(runs.load(Ordering::Relaxed), 1);
186    }
187
188    #[test]
189    fn a_panicking_job_does_not_stop_the_ones_after_it() {
190        let _g = TEST_LOCK.lock().unwrap();
191        drain_quietly();
192
193        let ran = Arc::new(AtomicUsize::new(0));
194        let after = Arc::clone(&ran);
195        post(|| panic!("job blew up"));
196        post(move || {
197            after.fetch_add(1, Ordering::Relaxed);
198        });
199
200        assert_eq!(run_pending(), 2);
201        assert_eq!(ran.load(Ordering::Relaxed), 1);
202
203        // And the queue still works afterwards.
204        let later = Arc::clone(&ran);
205        post(move || {
206            later.fetch_add(1, Ordering::Relaxed);
207        });
208        assert_eq!(run_pending(), 1);
209        assert_eq!(ran.load(Ordering::Relaxed), 2);
210    }
211
212    #[test]
213    fn a_worker_thread_can_post() {
214        let _g = TEST_LOCK.lock().unwrap();
215        drain_quietly();
216
217        let ran = Arc::new(AtomicUsize::new(0));
218        let from_worker = Arc::clone(&ran);
219        std::thread::spawn(move || {
220            post(move || {
221                from_worker.fetch_add(1, Ordering::Relaxed);
222            });
223        })
224        .join()
225        .unwrap();
226
227        assert_eq!(run_pending(), 1);
228        assert_eq!(ran.load(Ordering::Relaxed), 1);
229    }
230}