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}