Skip to main content

rumtk_core/
threading.rs

1/*
2 * rumtk attempts to implement HL7 and medical protocols for interoperability in medicine.
3 * This toolkit aims to be reliable, simple, performant, and standards compliant.
4 * Copyright (C) 2025  Luis M. Santos, M.D. <lsantos@medicalmasses.com>
5 * Copyright (C) 2025  MedicalMasses L.L.C. <contact@medicalmasses.com>
6 *
7 * This program is free software: you can redistribute it and/or modify
8 * it under the terms of the GNU General Public License as published by
9 * the Free Software Foundation, either version 3 of the License, or
10 * (at your option) any later version.
11 *
12 * This program is distributed in the hope that it will be useful,
13 * but WITHOUT ANY WARRANTY; without even the implied warranty of
14 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
15 * GNU General Public License for more details.
16 *
17 * You should have received a copy of the GNU General Public License
18 * along with this program.  If not, see <https://www.gnu.org/licenses/>.
19 */
20
21///
22/// This module provides all the primitives needed to build a multithreaded application.
23///
24pub mod thread_primitives {
25    pub use std::sync::Mutex as SyncMutex;
26    pub use std::sync::MutexGuard as SyncMutexGuard;
27    pub use std::sync::RwLock as SyncRwLock;
28    use std::sync::{Arc, OnceLock};
29    pub use tokio::io;
30    pub use tokio::io::{AsyncReadExt, AsyncWriteExt};
31    use tokio::runtime::Runtime as TokioRuntime;
32    pub use tokio::sync::{
33        Mutex as AsyncMutex, MutexGuard as AsyncMutexGuard,
34        OwnedRwLockReadGuard as AsyncOwnedRwLockMappedReadGuard,
35        OwnedRwLockReadGuard as AsyncOwnedRwLockReadGuard,
36        OwnedRwLockWriteGuard as AsyncOwnedRwLockWriteGuard, RwLock as AsyncRwLock,
37        RwLockMappedWriteGuard as AsyncRwLockMappedWriteGuard,
38        RwLockReadGuard as AsyncRwLockReadGuard, RwLockWriteGuard as AsyncRwLockWriteGuard,
39    };
40
41    /**************************** Types ***************************************/
42    pub type SafeLockReadGuard<T> = AsyncOwnedRwLockReadGuard<T>;
43    pub type MappedLockReadGuard<T> = AsyncOwnedRwLockReadGuard<T>;
44    pub type SafeLockWriteGuard<T> = AsyncOwnedRwLockWriteGuard<T>;
45    pub type SafeLock<T> = Arc<AsyncRwLock<T>>;
46    pub type SafeTokioRuntime = OnceLock<TokioRuntime>;
47}
48
49pub mod threading_manager {
50    use crate::base::{RUMResult, RUMVec};
51    use crate::strings::rumtk_format;
52    use crate::threading::thread_primitives::SafeLock;
53    use crate::threading::threading_functions::{async_sleep, sleep};
54    use crate::types::{RUMHashMap, RUMID};
55    use crate::{rumtk_init_threads, rumtk_resolve_task, threading};
56    use std::fmt::Debug;
57    use std::future::Future;
58    use std::sync::Arc;
59    pub use std::sync::RwLock as SyncRwLock;
60    use tokio::io::AsyncReadExt;
61    use tokio::task::JoinHandle;
62
63    const DEFAULT_SLEEP_DURATION: f32 = 0.001f32;
64    const DEFAULT_TASK_CAPACITY: usize = 100;
65
66    pub type AsyncHandle<T> = JoinHandle<T>;
67    pub type TaskItems<T> = RUMVec<T>;
68    /// This type aliases a vector of T elements that will be used for passing arguments to the task processor.
69    pub type TaskArgs<T> = TaskItems<T>;
70    /// Function signature defining the interface of task processing logic.
71    pub type SafeTaskArgs<T> = SafeLock<TaskItems<T>>;
72    pub type AsyncTaskHandle<R> = AsyncHandle<TaskResult<R>>;
73    pub type AsyncTaskHandles<R> = Vec<AsyncTaskHandle<R>>;
74    //pub type TaskProcessor<T, R, Fut: Future<Output = TaskResult<R>>> = impl FnOnce(&SafeTaskArgs<T>) -> Fut;
75    pub type TaskID = RUMID;
76
77    #[derive(Debug, Clone, Default)]
78    pub struct Task<R> {
79        pub id: TaskID,
80        pub finished: bool,
81        pub result: Option<R>,
82    }
83
84    pub type SafeTask<R> = Arc<Task<R>>;
85    type SafeInternalTask<R> = Arc<SyncRwLock<Task<R>>>;
86    pub type TaskTable<R> = RUMHashMap<TaskID, Task<R>>;
87    pub type SafeAsyncTaskTable<R> = SafeLock<TaskTable<R>>;
88    pub type SafeSyncTaskTable<R> = Arc<SyncRwLock<TaskTable<R>>>;
89    pub type TaskBatch = RUMVec<TaskID>;
90    /// Type to use to define how task results are expected to be returned.
91    pub type TaskResult<R> = RUMResult<Option<R>>;
92    pub type TaskResults<R> = TaskItems<TaskResult<R>>;
93
94    ///
95    /// Manages asynchronous tasks submitted as micro jobs from synchronous code. This type essentially
96    /// gives the multithreading, asynchronous superpowers to synchronous logic.
97    ///
98    /// ## Example Usage
99    ///
100    /// ```
101    /// use std::sync::{Arc};
102    /// use tokio::sync::RwLock as AsyncRwLock;
103    /// use rumtk_core::base::RUMResult;
104    /// use rumtk_core::strings::RUMString;
105    /// use rumtk_core::threading::threading_manager::{SafeTaskArgs, TaskItems, TaskManager};
106    /// use rumtk_core::{rumtk_create_task, };
107    ///
108    /// let expected = vec![
109    ///     RUMString::from("Hello"),
110    ///     RUMString::from("World!"),
111    ///     RUMString::from("Overcast"),
112    ///     RUMString::from("and"),
113    ///     RUMString::from("Sad"),
114    ///  ];
115    ///
116    /// type TestResult = RUMResult<Vec<RUMString>>;
117    /// let mut queue: TaskManager<TestResult> = TaskManager::new(&5).unwrap();
118    ///
119    /// let locked_args = AsyncRwLock::new(expected.clone());
120    /// let task_args = SafeTaskArgs::<RUMString>::new(locked_args);
121    /// let processor = rumtk_create_task!(
122    ///     async |args: &SafeTaskArgs<RUMString>| -> TestResult {
123    ///         let owned_args = Arc::clone(args);
124    ///         let locked_args = owned_args.read().await;
125    ///         let mut results = TaskItems::<RUMString>::with_capacity(locked_args.len());
126    ///
127    ///         for arg in locked_args.iter() {
128    ///             results.push(RUMString::from(arg));
129    ///         }
130    ///
131    ///         Ok(results)
132    ///     },
133    ///     task_args
134    /// );
135    ///
136    /// queue.add_task::<_>(processor);
137    /// let results = queue.wait();
138    ///
139    /// let mut result_data = Vec::<RUMString>::with_capacity(5);
140    /// for r in results {
141    ///     for v in r.unwrap().unwrap().iter() {
142    ///         for value in v.iter() {
143    ///             result_data.push(value.clone());
144    ///         }
145    ///     }
146    ///  }
147    ///
148    /// assert_eq!(result_data, expected, "Results do not match expected!");
149    ///
150    /// ```
151    ///
152    #[derive(Debug, Clone, Default)]
153    pub struct TaskManager<R> {
154        tasks: SafeSyncTaskTable<R>,
155        workers: usize,
156    }
157
158    impl<R> TaskManager<R>
159    where
160        R: Debug + Sync + Send + Clone + 'static,
161    {
162        ///
163        /// This method creates a [`TaskManager`] instance using sensible defaults.
164        ///
165        /// The `threads` field is computed from the number of cores present in system.
166        ///
167        pub fn default() -> RUMResult<TaskManager<R>> {
168            Self::new(&threading::threading_functions::get_default_system_thread_count())
169        }
170
171        ///
172        /// Creates an instance of [`TaskManager<R>`](TaskManager<R>).
173        /// Expects you to provide the count of threads to spawn and the microtask queue size
174        /// allocated by each thread.
175        ///
176        /// This method calls [`TaskTable::with_capacity()`](TaskTable::with_capacity) for the actual object creation.
177        /// The main queue capacity is pre-allocated to [`DEFAULT_TASK_CAPACITY`](DEFAULT_TASK_CAPACITY).
178        ///
179        pub fn new(worker_num: &usize) -> RUMResult<TaskManager<R>> {
180            let tasks = SafeSyncTaskTable::<R>::new(SyncRwLock::new(TaskTable::with_capacity(
181                DEFAULT_TASK_CAPACITY,
182            )));
183            Ok(TaskManager::<R> {
184                tasks,
185                workers: worker_num.to_owned(),
186            })
187        }
188
189        ///
190        /// Add a task to the processing queue. The idea is that you can queue a processor function
191        /// and list of args that will be picked up by one of the threads for processing.
192        ///
193        /// This is the async counterpart
194        ///
195        pub async fn add_task_async<F>(&mut self, task: F) -> TaskID
196        where
197            F: Future<Output = R> + Send + Sync + 'static,
198            F::Output: Send + 'static,
199        {
200            let id = TaskID::new_v4();
201            Self::_add_task_async(id.clone(), self.tasks.clone(), task).await
202        }
203
204        ///
205        /// See [`Self::add_task_async`]
206        ///
207        /// Unlike `add_task`, this method does not block which is key to avoiding panicking
208        /// the tokio runtim if trying to add task to queue from a normal function called from an
209        /// async environment.
210        ///
211        /// ## Example
212        ///
213        /// ```
214        /// use rumtk_core::threading::threading_manager::{TaskManager};
215        /// use rumtk_core::{rumtk_init_threads, strings::rumtk_format};
216        /// use std::sync::{Arc, LazyLock};
217        ///
218        /// type JobManager = LazyLock<TaskManager<usize>>;
219        /// static mut manager: JobManager = LazyLock::new( || TaskManager::new(&5).unwrap());
220        ///
221        /// async fn called_fn() -> usize {
222        ///     5
223        /// }
224        ///
225        /// fn push_job() -> usize {
226        ///     unsafe {(*manager).spawn_task(called_fn())};
227        ///     1
228        /// }
229        ///
230        /// async fn call_sync_fn() -> usize {
231        ///     push_job()
232        /// }
233        ///
234        /// unsafe {(*manager).spawn_task(call_sync_fn())};
235        ///
236        /// let result_raw = unsafe {(*manager).wait()};
237        ///
238        /// ```
239        ///
240        pub fn spawn_task<F>(&mut self, task: F) -> RUMResult<TaskID>
241        where
242            F: Future<Output = R> + Send + Sync + 'static,
243            F::Output: Send + Sized + 'static,
244        {
245            let id = TaskID::new_v4();
246            let tasks = self.tasks.clone();
247            rumtk_init_threads!(self.workers);
248            Ok(rumtk_resolve_task!(Self::_add_task_async(id.clone(), tasks, task)))
249        }
250
251        ///
252        /// See [add_task_async](Self::add_task_async)
253        ///
254        pub fn add_task<F>(&mut self, task: F) -> RUMResult<TaskID>
255        where
256            F: Future<Output = R> + Send + Sync + 'static,
257            F::Output: Send + Sized + 'static,
258        {
259            self.spawn_task(task)
260        }
261
262        async fn _add_task_async<F>(id: TaskID, tasks: SafeSyncTaskTable<R>, task: F) -> TaskID
263        where
264            F: Future<Output = R> + Send + Sync + 'static,
265            F::Output: Send + Sized + 'static,
266        {
267            let mut safe_task = Task::<R> {
268                id: id.clone(),
269                finished: false,
270                result: None,
271            };
272            tasks.write().unwrap().insert(id.clone(), safe_task.clone());
273
274            let task_wrapper = async move || {
275                // Run the task
276                let result = task.await;
277
278                // Cleanup task
279                let mut lock = tasks.write().unwrap();
280                if lock.contains_key(&id) {
281                    let mut task = lock.get_mut(&id).unwrap();
282                    task.result = Some(result);
283                    task.finished = true;
284                }
285            };
286
287            tokio::spawn(task_wrapper());
288
289            id
290        }
291
292        ///
293        /// See [wait_async](Self::wait_async)
294        ///
295        /// Duplicated here because we can't request the tokio runtime to do a quick exec for us if
296        /// this function happens to be called from the async context.
297        ///
298        pub fn wait(&mut self) -> TaskResults<R> {
299            let task_batch = self
300                .tasks
301                .read()
302                .unwrap()
303                .keys()
304                .cloned()
305                .collect::<Vec<_>>();
306            self.wait_on_batch(&task_batch)
307        }
308
309        ///
310        /// See [wait_on_batch_async](Self::wait_on_batch_async)
311        ///
312        /// Duplicated here because we can't request the tokio runtime to do a quick exec for us if
313        /// this function happens to be called from the async context.
314        ///
315        pub fn wait_on_batch(&mut self, tasks: &TaskBatch) -> TaskResults<R> {
316            let mut results = TaskResults::<R>::default();
317            for task in tasks {
318                results.push(self.wait_on(task));
319            }
320            results
321        }
322
323        ///
324        /// See [wait_on_async](Self::wait_on_async)
325        ///
326        /// Duplicated here because we can't request the tokio runtime to do a quick exec for us if
327        /// this function happens to be called from the async context.
328        ///
329        pub fn wait_on(&mut self, task_id: &TaskID) -> TaskResult<R> {
330            while !self.is_finished(task_id) {
331                sleep(DEFAULT_SLEEP_DURATION);
332            }
333
334            let task = match self.tasks.write().unwrap().remove(task_id) {
335                Some(task) => task.clone(),
336                None => return Err(rumtk_format!("No task with id {}", task_id)),
337            };
338
339            Ok(task.result)
340        }
341
342        ///
343        /// This method waits until a queued task with [TaskID](TaskID) has been processed from the main queue.
344        ///
345        /// We poll the status of the task every [DEFAULT_SLEEP_DURATION](DEFAULT_SLEEP_DURATION) ms.
346        ///
347        /// Upon completion,
348        ///
349        /// 2. Return the result ([TaskResults<R>](TaskResults)).
350        ///
351        /// This operation consumes the task.
352        ///
353        /// ### Note:
354        /// ```text
355        ///     Results returned here are not guaranteed to be in the same order as the order in which
356        ///     the tasks were queued for work. You will need to pass a type as T that automatically
357        ///     tracks its own id or has a way for you to resort results.
358        /// ```
359        pub async fn wait_on_async(&mut self, task_id: &TaskID) -> TaskResult<R> {
360            while !self.is_finished(task_id) {
361                async_sleep(DEFAULT_SLEEP_DURATION).await;
362            }
363
364            let task = match self.tasks.write().unwrap().remove(task_id) {
365                Some(task) => task.clone(),
366                None => return Err(rumtk_format!("No task with id {}", task_id)),
367            };
368
369            Ok(task.result)
370        }
371
372        ///
373        /// This method waits until a set of queued tasks with [TaskID](TaskID) has been processed from the main queue.
374        ///
375        /// We poll the status of the task every [DEFAULT_SLEEP_DURATION](DEFAULT_SLEEP_DURATION) ms.
376        ///
377        /// Upon completion,
378        ///
379        /// 1. We collect the results generated (if any).
380        /// 2. Return the list of results ([TaskResults<R>](TaskResults)).
381        ///
382        /// ### Note:
383        /// ```text
384        ///     Results returned here are not guaranteed to be in the same order as the order in which
385        ///     the tasks were queued for work. You will need to pass a type as T that automatically
386        ///     tracks its own id or has a way for you to resort results.
387        /// ```
388        pub async fn wait_on_batch_async(&mut self, tasks: &TaskBatch) -> TaskResults<R> {
389            let mut results = TaskResults::<R>::default();
390            for task in tasks {
391                results.push(self.wait_on_async(task).await);
392            }
393            results
394        }
395
396        ///
397        /// This method waits until all queued tasks have been processed from the main queue.
398        ///
399        /// We poll the status of the main queue every [DEFAULT_SLEEP_DURATION](DEFAULT_SLEEP_DURATION) ms.
400        ///
401        /// Upon completion,
402        ///
403        /// 1. We collect the results generated (if any).
404        /// 2. We reset the main task and result internal queue states.
405        /// 3. Return the list of results ([TaskResults<R>](TaskResults)).
406        ///
407        /// This operation consumes all the tasks.
408        ///
409        /// ### Note:
410        /// ```text
411        ///     Results returned here are not guaranteed to be in the same order as the order in which
412        ///     the tasks were queued for work. You will need to pass a type as T that automatically
413        ///     tracks its own id or has a way for you to resort results.
414        /// ```
415        pub async fn wait_async(&mut self) -> TaskResults<R> {
416            let task_batch = self
417                .tasks
418                .read()
419                .unwrap()
420                .keys()
421                .cloned()
422                .collect::<Vec<_>>();
423            self.wait_on_batch_async(&task_batch).await
424        }
425
426        ///
427        /// Check if all work has been completed from the task queue.
428        ///
429        /// ## Examples
430        ///
431        /// ### Sync Usage
432        ///
433        ///```
434        /// use rumtk_core::threading::threading_manager::TaskManager;
435        ///
436        /// let manager = TaskManager::<usize>::new(&4).unwrap();
437        ///
438        /// let all_done = manager.is_all_completed();
439        ///
440        /// assert_eq!(all_done, true, "Empty TaskManager reports tasks are not completed!");
441        ///
442        /// ```
443        ///
444        pub fn is_all_completed(&self) -> bool {
445            self._is_all_completed_async()
446        }
447
448        pub async fn is_all_completed_async(&self) -> bool {
449            self._is_all_completed_async()
450        }
451
452        fn _is_all_completed_async(&self) -> bool {
453            for (_, task) in self.tasks.read().unwrap().iter() {
454                if !task.finished {
455                    return false;
456                }
457            }
458
459            true
460        }
461
462        ///
463        /// Check if a task completed
464        ///
465        pub fn is_finished(&self, id: &TaskID) -> bool {
466            match self.tasks.read().unwrap().get(id) {
467                Some(t) => t.finished,
468                None => false,
469            }
470        }
471
472        ///
473        /// Alias for [wait](TaskManager::wait).
474        ///
475        fn gather(&mut self) -> TaskResults<R> {
476            self.wait()
477        }
478
479        pub fn has_job(&self, id: &TaskID) -> bool {
480            match self.tasks.read().unwrap().get(id) {
481                Some(_) => true,
482                None => false,
483            }
484        }
485    }
486}
487
488///
489/// This module contains a few helper.
490///
491/// For example, you can find a function for determining number of threads available in system.
492/// The sleep family of functions are also here.
493///
494pub mod threading_functions {
495    use crate::base::RUMResult;
496    use crate::net::tcp::{AsyncOwnedRwLockReadGuard, AsyncOwnedRwLockWriteGuard, SafeLockReadGuard, SafeLockWriteGuard, SafeTokioRuntime};
497    use crate::threading::thread_primitives::{AsyncRwLock, SafeLock};
498    use num_cpus;
499    use std::future::Future;
500    use std::slice::SliceIndex;
501    use std::sync::{Arc, LazyLock};
502    use std::thread::{available_parallelism, sleep as std_sleep};
503    use std::time::Duration;
504    use tokio::runtime::Runtime;
505    use tokio::task::JoinHandle;
506    use tokio::time::sleep as tokio_sleep;
507    /**************************** Globals **************************************/
508    static mut DEFAULT_RUNTIME: SafeTokioRuntime = SafeTokioRuntime::new();
509    pub static DEFAULT_CPUS: LazyLock<usize> = LazyLock::<usize>::new(|| {
510        let cpus: usize = num_cpus::get();
511        let parallelism = match available_parallelism() {
512            Ok(n) => n.get(),
513            Err(_) => 0,
514        };
515
516        if parallelism >= cpus {
517            parallelism
518        } else {
519            cpus
520        }
521    });
522
523    pub const NANOS_PER_SEC: u64 = 1000000000;
524    pub const MILLIS_PER_SEC: u64 = 1000;
525    pub const MICROS_PER_SEC: u64 = 1000000;
526    pub const DEFAULT_SLEEP_DURATION: f32 = 0.001;
527    /**************************** Helpers **************************************/
528    pub fn init_runtime<'a>(workers: usize) -> &'a Runtime {
529        unsafe {
530            let runtime = DEFAULT_RUNTIME.get_or_init(|| {
531                let mut builder = tokio::runtime::Builder::new_multi_thread();
532                builder.worker_threads(workers);
533                builder.enable_all();
534                match builder.build() {
535                    Ok(handle) => handle,
536                    Err(e) => panic!(
537                        "Unable to initialize threading tokio runtime because {}!",
538                        &e
539                    ),
540                }
541            });
542            runtime
543        }
544    }
545
546    #[inline]
547    pub fn get_default_system_thread_count() -> usize {
548        *DEFAULT_CPUS
549    }
550
551    #[inline]
552    pub fn sleep(s: f32) {
553        let ns = s * NANOS_PER_SEC as f32;
554        let rounded_ns = ns.round() as u64;
555        let duration = Duration::from_nanos(rounded_ns);
556        std_sleep(duration);
557    }
558
559    #[inline]
560    pub async fn async_sleep(s: f32) {
561        let ns = s * NANOS_PER_SEC as f32;
562        let rounded_ns = ns.round() as u64;
563        let duration = Duration::from_nanos(rounded_ns);
564        tokio_sleep(duration).await;
565    }
566
567    ///
568    /// Given a closure task, push it onto the current `tokio` runtime for execution.
569    /// Every [DEFAULT_SLEEP_DURATION] seconds, we check if the task has concluded.
570    /// Once the task has concluded, we call [tokio::block_on](tokio::task::block_in_place) to resolve and extract the task
571    /// result.
572    ///
573    /// Because this helper function can fail, the return value is wrapped inside a [RUMResult].
574    ///
575    /// ## Example
576    ///
577    /// ```
578    /// use rumtk_core::threading::threading_functions::{init_runtime, block_on_task};
579    ///
580    /// const Hello: &str = "World!";
581    ///
582    /// init_runtime(5);
583    ///
584    /// let result = block_on_task(async {
585    ///     Hello
586    /// });
587    ///
588    /// assert_eq!(Hello, result, "Result mismatches expected! {} vs. {}", Hello, result);
589    /// ```
590    ///
591    /// ## Notes
592    /// ```text
593    ///     You need to wrap our call to block_on with a call to tokio::task::block_in_place to force
594    ///     cleanup of async executor and therefore avoid panics from the tokio runtime!
595    ///     Per Tokio's documentation, spawn_blocking would be better since it moves the task to an
596    ///     executor meant for blocking tasks instead of moving tasks out of the current thread and
597    ///     converting the thread into a clocking executor. The reason we don't do that is because
598    ///     the call to this function expects to block the current thread until completion and then
599    ///     return the result. If there's an issue with IO, revisit this function.
600    ///
601    ///     https://docs.rs/tokio/latest/tokio/task/fn.block_in_place.html
602    ///     https://docs.rs/tokio/latest/tokio/runtime/struct.Runtime.html#method.block_on
603    ///     https://docs.rs/tokio/latest/tokio/runtime/struct.Handle.html#method.spawn_blocking
604    ///     https://docs.rs/tokio/latest/tokio/runtime/struct.Handle.html#method.spawn_blocking
605    /// ```
606    ///
607    #[inline]
608    pub fn block_on_task<R, F>(task: F) -> R
609    where
610        F: Future<Output = R> + Send + 'static,
611        F::Output: Send + 'static,
612    {
613        let rt = init_runtime(get_default_system_thread_count());
614        // You need to wrap our call to block_on with a call to tokio::task::block_in_place to force 
615        // cleanup of async executor and therefore avoid panics from the tokio runtime!
616        // Per Tokio's documentation, spawn_blocking would be better since it moves the task to an 
617        // executor meant for blocking tasks instead of moving tasks out of the current thread and 
618        // converting the thread into a clocking executor. The reason we don't do that is because 
619        // the call to this function expects to block the current thread until completion and then 
620        // return the result. If there's an issue with IO, revisit this function.
621        //
622        // https://docs.rs/tokio/latest/tokio/task/fn.block_in_place.html
623        // https://docs.rs/tokio/latest/tokio/runtime/struct.Runtime.html#method.block_on
624        // https://docs.rs/tokio/latest/tokio/runtime/struct.Handle.html#method.spawn_blocking
625        // https://docs.rs/tokio/latest/tokio/runtime/struct.Handle.html#method.spawn_blocking
626        tokio::task::block_in_place(move || {
627            rt.block_on(task)
628        })
629    }
630
631    ///
632    /// This helper should be used for spawning tasks that would normally block the async runtime.
633    /// However, here we use the appropriate `tokio` facilities to signal the runtime on how
634    /// to handle this, potentially blocking, task. For waiting on potentially blocking futures, use
635    /// [block_on_task] instead!
636    ///
637    ///
638    ///
639    /// ## Notes
640    /// ```text
641    ///     You need to wrap our call to block_on with a call to tokio::task::block_in_place to force
642    ///     cleanup of async executor and therefore avoid panics from the tokio runtime!
643    ///     Per Tokio's documentation, spawn_blocking would be better since it moves the task to an
644    ///     executor meant for blocking tasks instead of moving tasks out of the current thread and
645    ///     converting the thread into a clocking executor. The reason we don't do that is because
646    ///     the call to this function expects to block the current thread until completion and then
647    ///     return the result. If there's an issue with IO, revisit this function.
648    ///
649    ///     https://docs.rs/tokio/latest/tokio/task/fn.block_in_place.html
650    ///     https://docs.rs/tokio/latest/tokio/runtime/struct.Runtime.html#method.block_on
651    ///     https://docs.rs/tokio/latest/tokio/runtime/struct.Handle.html#method.spawn_blocking
652    ///     https://docs.rs/tokio/latest/tokio/runtime/struct.Handle.html#method.spawn_blocking
653    /// ```
654    ///
655    pub fn spawn_blocking_sync_task<R, F>(task: F) -> JoinHandle<R>
656    where
657        F: FnOnce() -> R + Send + 'static,
658        R: Send + 'static,
659    {
660        let rt = init_runtime(get_default_system_thread_count());
661        rt.spawn_blocking(task)
662    }
663
664    pub fn new_lock<T>(data: T) -> SafeLock<T> {
665        Arc::new(AsyncRwLock::new(data))
666    }
667
668    ///
669    /// This function gives you read access to underlying structure.
670    ///
671    /// Helper function for executing microtask immediately after locking the spin lock. This function
672    /// should be used in situations in which you want to minimize the risk of Time of Check Time of
673    /// Use security bugs.
674    ///
675    /// ## Example
676    /// ```
677    /// use rumtk_core::base::RUMResult;
678    /// use rumtk_core::threading::thread_primitives::SafeLock;
679    /// use rumtk_core::threading::threading_functions::{new_lock, process_read_critical_section};
680    ///
681    /// let data = 5;
682    /// let lock = new_lock(data.clone());
683    /// let result = process_read_critical_section(lock, |guard| -> RUMResult<i32> {
684    ///     Ok(*guard)
685    /// }).unwrap();
686    ///
687    /// assert_eq!(result, data, "Failed to execute critical section through which we retrieve the locked data!");
688    /// ```
689    ///
690    pub fn process_read_critical_section<T, R, F>(
691        lock: SafeLock<T>,
692        critical_section: F,
693    ) -> R
694    where 
695        F: Fn(SafeLockReadGuard<T>) -> R, T: Send + Sync + 'static,
696    {
697        tokio::task::block_in_place(move || {
698            let read_guard = lock_read(lock);
699            critical_section(read_guard)
700        })
701    }
702
703    ///
704    /// This function gives you write access to underlying structure.
705    ///
706    /// Helper function for executing microtask immediately after locking the spin lock. This function
707    /// should be used in situations in which you want to minimize the risk of Time of Check Time of
708    /// Use security bugs.
709    ///
710    /// ## Example
711    /// ```
712    /// use rumtk_core::base::RUMResult;
713    /// use rumtk_core::threading::thread_primitives::SafeLock;
714    /// use rumtk_core::threading::threading_functions::{new_lock, process_write_critical_section};
715    ///
716    /// let data = 5;
717    /// let lock = new_lock(data.clone());
718    /// let new_data = 10;
719    /// let result = process_write_critical_section(lock, |mut guard| -> RUMResult<i32> {
720    ///     *guard = new_data;
721    ///     Ok(*guard)
722    /// }).unwrap();
723    ///
724    /// assert_eq!(result, new_data, "Failed to execute critical section through which we modify the locked data!");
725    /// ```
726    ///
727    pub fn process_write_critical_section<T, R, F>(
728        lock: SafeLock<T>,
729        critical_section: F,
730    ) -> R
731    where
732        F: Fn(SafeLockWriteGuard<T>) -> R, T: Send + Sync + 'static,
733    {
734        tokio::task::block_in_place(move || {
735            let write_guard = lock_write(lock);
736            critical_section(write_guard)
737        })
738    }
739
740    ///
741    /// Obtain read guard to standard spin lock such that you have a more ergonomic interface to
742    /// locked data.
743    ///
744    /// It is preferable to use [process_read_critical_section] when you must process
745    /// critical logic that is sensitive to time of check time of use security bugs!
746    ///
747    /// ## Example
748    /// ```
749    /// use rumtk_core::threading::thread_primitives::SafeLock;
750    /// use rumtk_core::threading::threading_functions::{new_lock, lock_read};
751    ///
752    /// let data = 5;
753    /// let lock = new_lock(data.clone());
754    /// let result = *lock_read(lock);
755    ///
756    /// assert_eq!(result, data, "Failed to access the locked data!");
757    /// ```
758    ///
759    pub fn lock_read<T: Send + Sync + 'static>(lock: SafeLock<T>) -> AsyncOwnedRwLockReadGuard<T> {
760        block_on_task(async move {
761            lock.read_owned().await
762        })
763    }
764
765    ///
766    /// Obtain write guard to standard spin lock such that you have a more ergonomic interface to
767    /// locked data.
768    ///
769    /// It is preferable to use [process_write_critical_section] when you must process
770    /// critical logic that is sensitive to time of check time of use security bugs!
771    ///
772    /// ## Example
773    /// ```
774    /// use rumtk_core::threading::thread_primitives::SafeLock;
775    /// use rumtk_core::threading::threading_functions::{new_lock, lock_read, lock_write};
776    ///
777    /// let data = 5;
778    /// let lock = new_lock(data.clone());
779    /// let new_data = 10;
780    ///
781    /// *lock_write(lock.clone()) = new_data;
782    ///
783    /// let result = *lock_read(lock);
784    ///
785    /// assert_eq!(result, new_data, "Failed to modify the locked data!");
786    /// ```
787    ///
788    pub fn lock_write<T: Send + Sync + 'static>(lock: SafeLock<T>) -> AsyncOwnedRwLockWriteGuard<T> {
789        block_on_task(async move {
790            lock.write_owned().await
791        })
792    }
793}
794
795///
796/// Main API for interacting with the threading back end. Remember, we use tokio as our executor.
797/// This means that by default, all jobs sent to the thread pool have to be async in nature.
798/// These macros make handling of these jobs at the sync/async boundary more convenient.
799///
800pub mod threading_macros {
801    use crate::threading::thread_primitives;
802    use crate::threading::threading_manager::SafeTaskArgs;
803
804    ///
805    /// First, let's make sure we have *tokio* initialized at least once. The runtime created here
806    /// will be saved to the global context so the next call to this macro will simply grab a
807    /// reference to the previously initialized runtime.
808    ///
809    /// Passing nothing will default to initializing a runtime using the default number of threads
810    /// for this system. This is typically equivalent to number of cores/threads for your CPU.
811    ///
812    /// Passing `threads` number will yield a runtime that allocates that many threads.
813    ///
814    ///
815    /// ## Examples
816    ///
817    /// ```
818    ///     use rumtk_core::{rumtk_init_threads, rumtk_resolve_task, rumtk_create_task_args, rumtk_create_task, rumtk_spawn_task};
819    ///     use rumtk_core::base::RUMResult;
820    ///     use rumtk_core::threading::threading_manager::SafeTaskArgs;
821    ///
822    ///     async fn test(args: &SafeTaskArgs<i32>) -> RUMResult<Vec<i32>> {
823    ///         let mut result = Vec::<i32>::new();
824    ///         for arg in args.read().await.iter() {
825    ///             result.push(*arg);
826    ///         }
827    ///         Ok(result)
828    ///     }
829    ///
830    ///     let args = rumtk_create_task_args!(1);                               // Creates a vector of i32s
831    ///     let task = rumtk_create_task!(test, args);                           // Creates a standard task which consists of a function or closure accepting a Vec<T>
832    ///     let result = rumtk_resolve_task!(task); // Spawn's task and waits for it to conclude.
833    /// ```
834    ///
835    /// ```
836    ///     use rumtk_core::{rumtk_init_threads, rumtk_resolve_task, rumtk_create_task_args, rumtk_create_task, rumtk_spawn_task};
837    ///     use rumtk_core::base::RUMResult;
838    ///     use rumtk_core::threading::threading_manager::SafeTaskArgs;
839    ///
840    ///     async fn test(args: &SafeTaskArgs<i32>) -> RUMResult<Vec<i32>> {
841    ///         let mut result = Vec::<i32>::new();
842    ///         for arg in args.read().await.iter() {
843    ///             result.push(*arg);
844    ///         }
845    ///         Ok(result)
846    ///     }
847    ///
848    ///     let thread_count: usize = 10;
849    ///     let args = rumtk_create_task_args!(1);
850    ///     let task = rumtk_create_task!(test, args);
851    ///     let result = rumtk_resolve_task!(task);
852    /// ```
853    #[macro_export]
854    macro_rules! rumtk_init_threads {
855        ( ) => {{
856            use $crate::threading::threading_functions::{
857                get_default_system_thread_count, init_runtime,
858            };
859            init_runtime(get_default_system_thread_count())
860        }};
861        ( $threads:expr ) => {{
862            use $crate::rumtk_cache_fetch;
863            use $crate::threading::threading_functions::init_runtime;
864            init_runtime($threads)
865        }};
866    }
867
868    ///
869    /// Puts task onto the runtime queue.
870    ///
871    /// The parameters to this macro are a reference to the runtime (`rt`) and a future (`func`).
872    ///
873    /// The return is a [thread_primitives::JoinHandle<T>] instance. If the task was a standard
874    /// framework task, you will get [thread_primitives::AsyncTaskHandle] instead.
875    ///
876    #[macro_export]
877    macro_rules! rumtk_spawn_task {
878        ( $func:expr ) => {{
879            use $crate::rumtk_init_threads;
880            let rt = rumtk_init_threads!();
881            rt.spawn($func)
882        }};
883        ( $rt:expr, $func:expr ) => {{
884            $rt.spawn($func)
885        }};
886    }
887
888    #[macro_export]
889    macro_rules! rumtk_spawn_blocking_task {
890        ( $func:expr ) => {{
891            use $crate::threading::threading_functions::spawn_blocking_sync_task;
892            spawn_blocking_sync_task($func)
893        }}
894    }
895
896    ///
897    /// Using the initialized runtime, wait for the future to resolve in a thread blocking manner!
898    ///
899    /// If you pass a reference to the runtime (`rt`) and an async closure (`func`), we await the
900    /// async closure without passing any arguments.
901    ///
902    /// You can pass a third argument to this macro in the form of any number of arguments (`arg_item`).
903    /// In such a case, we pass those arguments to the call on the async closure and await on results.
904    ///
905    #[macro_export]
906    macro_rules! rumtk_wait_on_task {
907        ( $func:expr ) => {{
908            use $crate::threading::threading_functions::block_on_task;
909            block_on_task(async move {
910                $func().await
911            })
912        }};
913        ( $func:expr, $($arg_items:expr),+ ) => {{
914            use $crate::threading::threading_functions::block_on_task;
915            block_on_task(async move {
916                $func($($arg_items),+).await
917            })
918        }};
919    }
920
921    ///
922    /// This macro awaits a future.
923    ///
924    /// The arguments are a reference to the runtime (`rt) and a future.
925    ///
926    /// If there is a result, you will get the result of the future.
927    ///
928    /// ## Examples
929    ///
930    /// ```
931    ///     use rumtk_core::{rumtk_init_threads, rumtk_resolve_task, rumtk_create_task_args, rumtk_create_task, rumtk_spawn_task};
932    ///     use rumtk_core::base::RUMResult;
933    ///     use rumtk_core::threading::threading_manager::SafeTaskArgs;
934    ///
935    ///     async fn test(args: &SafeTaskArgs<i32>) -> RUMResult<Vec<i32>> {
936    ///         let mut result = Vec::<i32>::new();
937    ///         for arg in args.read().await.iter() {
938    ///             result.push(*arg);
939    ///         }
940    ///         Ok(result)
941    ///     }
942    ///
943    ///     let args = rumtk_create_task_args!(1);
944    ///     let task = rumtk_create_task!(test, args);
945    ///     let result = rumtk_resolve_task!(task);
946    /// ```
947    ///
948    #[macro_export]
949    macro_rules! rumtk_resolve_task {
950        ( $future:expr ) => {{
951            use $crate::threading::threading_functions::block_on_task;
952            // Fun tidbit, the expression rumtk_resolve_task!(&rt, rumtk_spawn_task!(&rt, task)), where
953            // rt is the tokio runtime yields async move { { &rt.spawn(task) } }. However, the whole thing
954            // is technically moved into the async closure and captured so things like mutex guards
955            // technically go out of the outer scope. As a result that expression fails to compile even
956            // though the intent is for rumtk_spawn_task to resolve first and its result get moved
957            // into the async closure. To ensure that happens regardless of given expression, we do
958            // a variable assignment below to force the "future" macro expressions to resolve before
959            // moving into the closure. DO NOT REMOVE OR "SIMPLIFY" THE let future = $future LINE!!!
960            //let future = $future;
961            block_on_task(async move { $future.await })
962        }};
963    }
964
965    ///
966    /// This macro allows to resolve a `sync` closure that was executed in a safe thread.
967    /// You cannot run this macro outside the `async` context.
968    ///
969    #[macro_export]
970    macro_rules! rumtk_resolve_sync_task {
971        ( $closure:expr ) => {{
972            use $crate::threading::threading_functions::spawn_blocking_sync_task;
973            use $crate::strings::rumtk_format;
974            match spawn_blocking_sync_task($closure).await {
975                Ok(result) => result,
976                Err(e) => Err(rumtk_format!("Issue with blocking task => {}", e))
977            }
978        }};
979    }
980
981    ///
982    /// This macro creates an async body that calls the async closure and awaits it.
983    ///
984    /// ## Example
985    ///
986    /// ```
987    /// use std::sync::{Arc, RwLock};
988    /// use tokio::sync::RwLock as AsyncRwLock;
989    /// use rumtk_core::strings::RUMString;
990    /// use rumtk_core::threading::threading_manager::{SafeTaskArgs, TaskItems};
991    ///
992    /// pub type SafeTaskArgs2<T> = Arc<RwLock<TaskItems<T>>>;
993    /// let expected = vec![
994    ///     RUMString::from("Hello"),
995    ///     RUMString::from("World!"),
996    ///     RUMString::from("Overcast"),
997    ///     RUMString::from("and"),
998    ///     RUMString::from("Sad"),
999    ///  ];
1000    /// let locked_args = AsyncRwLock::new(expected.clone());
1001    /// let task_args = SafeTaskArgs::<RUMString>::new(locked_args);
1002    ///
1003    ///
1004    /// ```
1005    ///
1006    #[macro_export]
1007    macro_rules! rumtk_create_task {
1008        ( $func:expr ) => {{
1009            async move {
1010                let f = $func;
1011                f().await
1012            }
1013        }};
1014        ( $func:expr, $args:expr ) => {{
1015            let f = $func;
1016            async move { f(&$args).await }
1017        }};
1018    }
1019
1020    ///
1021    /// Creates an instance of [SafeTaskArgs](SafeTaskArgs) with the arguments passed.
1022    ///
1023    /// ## Note
1024    ///
1025    /// All arguments must be of the same type
1026    ///
1027    #[macro_export]
1028    macro_rules! rumtk_create_task_args {
1029        ( ) => {{
1030            use $crate::threading::threading_manager::{TaskArgs, SafeTaskArgs, TaskItems};
1031            use $crate::threading::thread_primitives::AsyncRwLock;
1032            SafeTaskArgs::new(AsyncRwLock::new(vec![]))
1033        }};
1034        ( $($args:expr),+ ) => {{
1035            use $crate::threading::threading_manager::{SafeTaskArgs};
1036            use $crate::threading::thread_primitives::AsyncRwLock;
1037            SafeTaskArgs::new(AsyncRwLock::new(vec![$($args),+]))
1038        }};
1039    }
1040
1041    ///
1042    /// Convenience macro for packaging the task components and launching the task in one line.
1043    ///
1044    /// One of the advantages is that you can generate a new `tokio` runtime by specifying the
1045    /// number of threads at the end. This is optional. Meaning, we will default to the system's
1046    /// number of threads if that value is not specified.
1047    ///
1048    /// Between the `func` parameter and the optional `threads` parameter, you can specify a
1049    /// variable number of arguments to pass to the task. each argument must be of the same type.
1050    /// If you wish to pass different arguments with different types, please define an abstract type
1051    /// whose underlying structure is a tuple of items and pass that instead.
1052    ///
1053    /// ## Examples
1054    ///
1055    /// ### With Default Thread Count
1056    /// ```
1057    ///     use rumtk_core::{rumtk_exec_task};
1058    ///     use rumtk_core::base::RUMResult;
1059    ///     use rumtk_core::threading::threading_manager::SafeTaskArgs;
1060    ///
1061    ///     async fn test(args: &SafeTaskArgs<i32>) -> RUMResult<Vec<i32>> {
1062    ///         let mut result = Vec::<i32>::new();
1063    ///         for arg in args.read().await.iter() {
1064    ///             result.push(*arg);
1065    ///         }
1066    ///         Ok(result)
1067    ///     }
1068    ///
1069    ///     let result = rumtk_exec_task!(test, vec![5]).unwrap();
1070    ///     assert_eq!(&result.clone(), &vec![5], "Results mismatch");
1071    ///     assert_ne!(&result.clone(), &vec![5, 10], "Results do not mismatch as expected!");
1072    /// ```
1073    ///
1074    /// ### With Custom Thread Count
1075    /// ```
1076    ///     use rumtk_core::{rumtk_exec_task};
1077    ///     use rumtk_core::base::RUMResult;
1078    ///     use rumtk_core::threading::threading_manager::SafeTaskArgs;
1079    ///
1080    ///     async fn test(args: &SafeTaskArgs<i32>) -> RUMResult<Vec<i32>> {
1081    ///         let mut result = Vec::<i32>::new();
1082    ///         for arg in args.read().await.iter() {
1083    ///             result.push(*arg);
1084    ///         }
1085    ///         Ok(result)
1086    ///     }
1087    ///
1088    ///     let result = rumtk_exec_task!(test, vec![5], 5).unwrap();
1089    ///     assert_eq!(&result.clone(), &vec![5], "Results mismatch");
1090    ///     assert_ne!(&result.clone(), &vec![5, 10], "Results do not mismatch as expected!");
1091    /// ```
1092    ///
1093    /// ### With Async Function Body
1094    /// ```
1095    ///     use rumtk_core::{rumtk_exec_task};
1096    ///     use rumtk_core::base::RUMResult;
1097    ///     use rumtk_core::threading::threading_manager::SafeTaskArgs;
1098    ///
1099    ///     let result = rumtk_exec_task!(
1100    ///     async move |args: &SafeTaskArgs<i32>| -> RUMResult<Vec<i32>> {
1101    ///         let mut result = Vec::<i32>::new();
1102    ///         for arg in args.read().await.iter() {
1103    ///             result.push(*arg);
1104    ///         }
1105    ///         Ok(result)
1106    ///     },
1107    ///     vec![5]).unwrap();
1108    ///     assert_eq!(&result.clone(), &vec![5], "Results mismatch");
1109    ///     assert_ne!(&result.clone(), &vec![5, 10], "Results do not mismatch as expected!");
1110    /// ```
1111    ///
1112    /// ### With Async Function Body and No Args
1113    /// ```
1114    ///     use rumtk_core::{rumtk_exec_task};
1115    ///     use rumtk_core::base::RUMResult;
1116    ///     use rumtk_core::threading::threading_manager::SafeTaskArgs;
1117    ///
1118    ///     let result = rumtk_exec_task!(
1119    ///     async || -> RUMResult<Vec<i32>> {
1120    ///         let mut result = Vec::<i32>::new();
1121    ///         Ok(result)
1122    ///     }).unwrap();
1123    ///     let empty = Vec::<i32>::new();
1124    ///     assert_eq!(&result.clone(), &empty, "Results mismatch");
1125    ///     assert_ne!(&result.clone(), &vec![5, 10], "Results do not mismatch as expected!");
1126    /// ```
1127    ///
1128    /// ## Equivalent To
1129    ///
1130    /// ```no_run
1131    ///     use rumtk_core::{rumtk_init_threads, rumtk_resolve_task, rumtk_create_task_args, rumtk_create_task, rumtk_spawn_task};
1132    ///     use rumtk_core::base::RUMResult;
1133    ///     use rumtk_core::threading::threading_manager::SafeTaskArgs;
1134    ///
1135    ///     async fn test(args: &SafeTaskArgs<i32>) -> RUMResult<Vec<i32>> {
1136    ///         let mut result = Vec::<i32>::new();
1137    ///         for arg in args.read().await.iter() {
1138    ///             result.push(*arg);
1139    ///         }
1140    ///         Ok(result)
1141    ///     }
1142    ///
1143    ///     let args = rumtk_create_task_args!(1);
1144    ///     let task = rumtk_create_task!(test, args);
1145    ///     let result = rumtk_resolve_task!(task);
1146    /// ```
1147    ///
1148    #[macro_export]
1149    macro_rules! rumtk_exec_task {
1150        ($func:expr ) => {{
1151            use $crate::{
1152                rumtk_create_task, rumtk_create_task_args, rumtk_init_threads, rumtk_resolve_task,
1153            };
1154            let task = rumtk_create_task!($func);
1155            rumtk_resolve_task!(task)
1156        }};
1157        ($func:expr, $args:expr ) => {{
1158            use $crate::threading::threading_functions::get_default_system_thread_count;
1159            rumtk_exec_task!($func, $args, get_default_system_thread_count())
1160        }};
1161        ($func:expr, $args:expr , $threads:expr ) => {{
1162            use $crate::threading::thread_primitives::AsyncRwLock;
1163            use $crate::{
1164                rumtk_create_task, rumtk_create_task_args, rumtk_init_threads, rumtk_resolve_task,
1165            };
1166            let args = SafeTaskArgs::new(AsyncRwLock::new($args));
1167            let task = rumtk_create_task!($func, args);
1168            rumtk_resolve_task!(task)
1169        }};
1170    }
1171
1172    ///
1173    /// Sleep a duration of time in a sync context, so no await can be call on the result.
1174    ///
1175    /// You can pass any value that can be cast to f32.
1176    ///
1177    /// The precision is up to nanoseconds and it is depicted by the number of decimal places.
1178    ///
1179    /// ## Examples
1180    ///
1181    /// ```
1182    ///     use rumtk_core::rumtk_sleep;
1183    ///     rumtk_sleep!(1);           // Sleeps for 1 second.
1184    ///     rumtk_sleep!(0.001);       // Sleeps for 1 millisecond
1185    ///     rumtk_sleep!(0.000001);    // Sleeps for 1 microsecond
1186    ///     rumtk_sleep!(0.000000001); // Sleeps for 1 nanosecond
1187    /// ```
1188    ///
1189    #[macro_export]
1190    macro_rules! rumtk_sleep {
1191        ( $dur:expr) => {{
1192            use $crate::threading::threading_functions::sleep;
1193            sleep($dur as f32)
1194        }};
1195    }
1196
1197    ///
1198    /// Sleep for some duration of time in an async context. Meaning, we can be awaited.
1199    ///
1200    /// You can pass any value that can be cast to f32.
1201    ///
1202    /// The precision is up to nanoseconds and it is depicted by the number of decimal places.
1203    ///
1204    /// ## Examples
1205    ///
1206    /// ```
1207    ///     use rumtk_core::{rumtk_async_sleep, rumtk_exec_task};
1208    ///     use rumtk_core::base::RUMResult;
1209    ///     rumtk_exec_task!( async || -> RUMResult<()> {
1210    ///             rumtk_async_sleep!(1).await;           // Sleeps for 1 second.
1211    ///             rumtk_async_sleep!(0.001).await;       // Sleeps for 1 millisecond
1212    ///             rumtk_async_sleep!(0.000001).await;    // Sleeps for 1 microsecond
1213    ///             rumtk_async_sleep!(0.000000001).await; // Sleeps for 1 nanosecond
1214    ///             Ok(())
1215    ///         }
1216    ///     );
1217    /// ```
1218    ///
1219    #[macro_export]
1220    macro_rules! rumtk_async_sleep {
1221        ( $dur:expr) => {{
1222            use $crate::threading::threading_functions::async_sleep;
1223            async_sleep($dur as f32)
1224        }};
1225    }
1226
1227    ///
1228    ///
1229    ///
1230    #[macro_export]
1231    macro_rules! rumtk_new_task_queue {
1232        ( $worker_num:expr ) => {{
1233            use $crate::threading::threading_manager::TaskManager;
1234            TaskManager::new($worker_num);
1235        }};
1236    }
1237
1238    ///
1239    /// Creates a new safe lock to guard the given data. This interface was created to cleanup lock
1240    /// management for consumers of framework!
1241    ///
1242    /// ## Example
1243    /// ```
1244    /// use rumtk_core::{rumtk_new_lock};
1245    ///
1246    /// let data = 5;
1247    /// let lock = rumtk_new_lock!(data);
1248    /// ```
1249    ///
1250    #[macro_export]
1251    macro_rules! rumtk_new_lock {
1252        ( $data:expr ) => {{
1253            use $crate::threading::threading_functions::new_lock;
1254            new_lock($data)
1255        }};
1256    }
1257
1258    ///
1259    /// Using a standard spin lock [SafeLock](thread_primitives::SafeLock), lock it and execute the
1260    /// critical section. The critical section itself is a synchronous function or closure. In this case,
1261    /// the critical section simply retrieves a value from a guarded dataset.
1262    ///
1263    /// ## Example
1264    /// ```
1265    /// use rumtk_core::base::RUMResult;
1266    /// use rumtk_core::{rumtk_new_lock, rumtk_critical_section_read};
1267    ///
1268    /// let data = 5;
1269    /// let lock = rumtk_new_lock!(data);
1270    /// let result = rumtk_critical_section_read!(
1271    ///     lock,
1272    ///     |guard| -> RUMResult<i32> {
1273    ///         let result: i32 = *guard;
1274    ///         Ok(result)
1275    ///     }
1276    /// ).expect("No errors locking!");
1277    ///
1278    /// assert_eq!(result, data, "Critical section yielded invalid result!");
1279    /// ```
1280    ///
1281    #[macro_export]
1282    macro_rules! rumtk_critical_section_read {
1283        ( $lock:expr, $function:expr ) => {{
1284            use $crate::threading::threading_functions::process_read_critical_section;
1285            process_read_critical_section($lock, $function)
1286        }};
1287    }
1288
1289    ///
1290    /// Using a standard spin lock [SafeLock](thread_primitives::SafeLock), lock it and execute the
1291    /// critical section. The critical section itself is a synchronous function or closure. In this case,
1292    /// the critical section attempts to modify the internal state of a guarded dataset.
1293    ///
1294    /// ## Example
1295    /// ```
1296    /// use rumtk_core::{rumtk_new_lock, rumtk_critical_section_write};
1297    ///
1298    /// let data = 5;
1299    /// let new_data = 10;
1300    /// let lock = rumtk_new_lock!(data);
1301    /// let result = rumtk_critical_section_write!(
1302    ///     lock,
1303    ///     |mut guard| {
1304    ///         *guard = new_data;
1305    ///     }
1306    /// );
1307    ///
1308    /// assert_eq!(result, (), "Critical section yielded invalid result!");
1309    /// ```
1310    ///
1311    #[macro_export]
1312    macro_rules! rumtk_critical_section_write {
1313        ( $lock:expr, $function:expr ) => {{
1314            use $crate::threading::threading_functions::process_write_critical_section;
1315            process_write_critical_section($lock, $function)
1316        }};
1317    }
1318
1319    ///
1320    /// Framework interface to obtain a `read` guard to the locked data.
1321    /// To access the internal data, you will need to dereference the guard (`*guard`).
1322    ///
1323    /// It is preferred to use [rumtk_critical_section_read] if you need to avoid `time of check time
1324    /// of use` security bugs.
1325    ///
1326    /// ## Example
1327    /// ```
1328    /// use rumtk_core::{rumtk_new_lock, rumtk_lock_read};
1329    ///
1330    /// let data = 5;
1331    /// let lock = rumtk_new_lock!(data.clone());
1332    /// let result = *rumtk_lock_read!(lock);
1333    ///
1334    /// assert_eq!(result, data, "Failed to access locked data.");
1335    /// ```
1336    ///
1337    #[macro_export]
1338    macro_rules! rumtk_lock_read {
1339        ( $lock:expr ) => {{
1340            use $crate::threading::threading_functions::lock_read;
1341            lock_read($lock.clone())
1342        }};
1343    }
1344
1345    ///
1346    /// Framework interface to obtain a `write` guard to the locked data.
1347    /// To access the internal data, you will need to dereference the guard (`*guard`).
1348    ///
1349    /// It is preferred to use [rumtk_critical_section_write] if you need to avoid `time of check time
1350    /// of use` security bugs.
1351    ///
1352    /// ## Example
1353    /// ```
1354    /// use rumtk_core::{rumtk_new_lock, rumtk_lock_read, rumtk_lock_write};
1355    ///
1356    /// let data = 5;
1357    /// let lock = rumtk_new_lock!(data.clone());
1358    /// let new_data = 10;
1359    ///
1360    /// *rumtk_lock_write!(lock) = new_data;
1361    /// let result = *rumtk_lock_read!(lock);
1362    ///
1363    /// assert_eq!(result, new_data, "Failed to modify locked data.");
1364    /// ```
1365    ///
1366    #[macro_export]
1367    macro_rules! rumtk_lock_write {
1368        ( $lock:expr ) => {{
1369            use $crate::threading::threading_functions::lock_write;
1370            lock_write($lock.clone())
1371        }};
1372    }
1373}