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}