1use std::collections::VecDeque;
2use std::fmt;
3use std::sync::Arc;
4use std::sync::Mutex;
5
6use tokio::sync::mpsc;
7use tokio::sync::oneshot;
8use tokio::sync::watch;
9use tokio::task::JoinHandle;
10
11use super::channels::JobCancelReceiver;
12use super::channels::JobCancelSender;
13use super::channels::JobResultSender;
14use super::errors::AddJobError;
15use super::errors::StopQueueError;
16use super::job_completion::JobCompletion;
17use super::job_handle::JobHandle;
18use super::job_id::JobId;
19use super::traits::Job;
20
21#[derive(Debug)]
23pub struct JobQueue<P: Ord + Send + Sync + 'static> {
24 shared_queue: Arc<Mutex<SharedQueue<P>>>,
26
27 tx_job_added: mpsc::UnboundedSender<()>,
29
30 tx_stop: tokio::sync::watch::Sender<()>,
32
33 process_jobs_task_handle: Option<JoinHandle<()>>,
35}
36
37impl<P: Ord + Send + Sync + 'static> Drop for JobQueue<P> {
39 fn drop(&mut self) {
40 tracing::debug!("in JobQueue::drop()");
41
42 if !self.tx_stop.is_closed() {
43 if let Err(e) = self.tx_stop.send(()) {
44 tracing::error!("{}", e);
45 }
46 }
47 }
48}
49
50impl<P: Ord + Send + Sync + 'static> JobQueue<P> {
51 pub fn start() -> Self {
55 let shared_queue = SharedQueue {
57 jobs: VecDeque::new(),
58 current_job: None,
59 };
60 let shared_queue: Arc<Mutex<SharedQueue<P>>> = Arc::new(Mutex::new(shared_queue));
61
62 let (tx_job_added, rx_job_added) = mpsc::unbounded_channel();
64 let (tx_stop, rx_stop) = watch::channel(());
65
66 let process_jobs_task_handle =
68 tokio::spawn(process_jobs(shared_queue.clone(), rx_stop, rx_job_added));
69
70 tracing::debug!("JobQueue: started new queue.");
71
72 Self {
74 tx_job_added,
75 tx_stop,
76 shared_queue,
77 process_jobs_task_handle: Some(process_jobs_task_handle),
78 }
79 }
80
81 pub async fn stop(mut self) -> Result<(), StopQueueError> {
92 tracing::info!("JobQueue: stopping.");
93
94 self.tx_stop.send(())?;
96
97 if let Some(jh) = self.process_jobs_task_handle.take() {
99 jh.await?;
100 }
101
102 Ok(())
103 }
104
105 pub fn add_job(
112 &self,
113 job: impl Into<Box<dyn Job>>,
114 priority: P,
115 ) -> Result<JobHandle, AddJobError> {
116 let (result_tx, result_rx) = oneshot::channel();
117 let (cancel_tx, cancel_rx) = watch::channel::<()>(());
118
119 let job_id = JobId::random();
121
122 let m = QueuedJob {
124 job: job.into(),
125 job_id,
126 result_tx,
127 cancel_tx: cancel_tx.clone(),
128 cancel_rx,
129 priority,
130 };
131
132 let (num_jobs, job_running) = {
134 let mut guard = self.shared_queue.lock().unwrap();
136
137 guard.jobs.push_back(m);
139
140 let job_running = match &guard.current_job {
141 Some(j) => format!("#{} - {}", j.job_num, j.job_id),
142 None => "none".to_string(),
143 };
144 (guard.jobs.len(), job_running)
145 }; self.tx_job_added.send(())?;
149
150 tracing::debug!(
152 "JobQueue: job added - {} {} queued job(s). job running: {}",
153 job_id,
154 num_jobs,
155 job_running
156 );
157
158 Ok(JobHandle::new(job_id, result_rx, cancel_tx))
160 }
161
162 pub fn add_job_mut(
181 &mut self,
182 job: impl Into<Box<dyn Job>>,
183 priority: P,
184 ) -> Result<JobHandle, AddJobError> {
185 self.add_job(job, priority)
186 }
187
188 pub fn num_jobs(&self) -> usize {
190 let guard = self.shared_queue.lock().unwrap();
191 guard.jobs.len() + guard.current_job.as_ref().map(|_| 1).unwrap_or(0)
192 }
193
194 pub fn num_queued_jobs(&self) -> usize {
196 self.shared_queue.lock().unwrap().jobs.len()
197 }
198}
199
200async fn process_jobs<P: Ord + Send + Sync + 'static>(
220 shared_queue: Arc<Mutex<SharedQueue<P>>>,
221 mut rx_stop: watch::Receiver<()>,
222 mut rx_job_added: mpsc::UnboundedReceiver<()>,
223) {
224 let mut job_num: usize = 1;
228
229 while rx_job_added.recv().await.is_some() {
235 tracing::debug!("task process_jobs received JobAdded message.");
237 let (next_job, num_pending) = {
238 let mut guard = shared_queue.lock().unwrap();
240
241 guard
243 .jobs
244 .make_contiguous()
245 .sort_by(|a, b| b.priority.cmp(&a.priority));
246 let job = guard.jobs.pop_front().unwrap();
247
248 guard.current_job = Some(CurrentJob {
250 job_num,
251 job_id: job.job_id,
252 cancel_tx: job.cancel_tx.clone(),
253 });
254
255 (job, guard.jobs.len())
256 }; tracing::debug!(
260 " *** JobQueue: begin job #{} - {} - {} queued job(s) ***",
261 job_num,
262 next_job.job_id,
263 num_pending
264 );
265
266 let timer = tokio::time::Instant::now();
268
269 let job_task_handle = if next_job.job.is_async() {
271 tokio::spawn(
272 async move { next_job.job.run_async_cancellable(next_job.cancel_rx).await },
273 )
274 } else {
275 tokio::task::spawn_blocking(move || next_job.job.run(next_job.cancel_rx))
276 };
277
278 let job_task_result = tokio::select! {
280 job_task_result = job_task_handle => job_task_result,
282
283 _ = rx_stop.changed() => {
285
286 handle_stop_signal(&shared_queue).await;
287
288 break;
290 },
291 };
292
293 let job_completion = match job_task_result {
295 Ok(jc) => jc,
296 Err(e) => {
297 if e.is_panic() {
298 JobCompletion::Panicked(e.into_panic())
299 } else if e.is_cancelled() {
300 JobCompletion::Cancelled
301 } else {
302 unreachable!()
303 }
304 }
305 };
306
307 tracing::debug!(
309 " *** JobQueue: ended job #{} - {} - Completion: {} - {} secs ***",
310 job_num,
311 next_job.job_id,
312 job_completion,
313 timer.elapsed().as_secs_f32()
314 );
315 job_num += 1;
316
317 shared_queue.lock().unwrap().current_job = None;
319
320 if let Err(e) = next_job.result_tx.send(job_completion) {
322 tracing::warn!("job-handle dropped? {}", e);
323 }
324 }
325 tracing::debug!("task process_jobs exiting");
326}
327
328async fn handle_stop_signal<P: Ord + Send + Sync + 'static>(
330 shared_queue: &Arc<Mutex<SharedQueue<P>>>,
331) {
332 tracing::debug!("task process_jobs received Stop message.");
333
334 let maybe_info = shared_queue
336 .lock()
337 .unwrap()
338 .current_job
339 .as_ref()
340 .map(|cj| (cj.job_id, cj.cancel_tx.clone()));
341
342 if let Some((job_id, cancel_tx)) = maybe_info {
345 match cancel_tx.send(()) {
346 Ok(()) => {
347 tracing::debug!(
349 "JobQueue: notified current job {} to cancel. waiting...",
350 job_id
351 );
352 cancel_tx.closed().await;
353 tracing::debug!("JobQueue: current job {} has cancelled.", job_id);
354 }
355 Err(e) => {
356 tracing::warn!(
357 "could not send cancellation msg to current job {}. {}",
358 job_id,
359 e
360 )
361 }
362 }
363 }
364}
365
366pub(super) struct QueuedJob<P> {
368 job: Box<dyn Job>,
369 job_id: JobId,
370 result_tx: JobResultSender,
371 cancel_tx: JobCancelSender,
372 cancel_rx: JobCancelReceiver,
373 priority: P,
374}
375
376impl<P: fmt::Debug> fmt::Debug for QueuedJob<P> {
377 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
378 f.debug_struct("QueuedJob")
379 .field("job", &"Box<dyn Job>")
380 .field("job_id", &self.job_id)
381 .field("result_tx", &"JobResultSender")
382 .field("cancel_tx", &"JobCancelSender")
383 .field("cancel_rx", &"JobCancelReceiver")
384 .field("priority", &self.priority)
385 .finish()
386 }
387}
388
389#[derive(Debug)]
391pub(super) struct CurrentJob {
392 job_num: usize,
393 job_id: JobId,
394 cancel_tx: JobCancelSender,
395}
396
397#[derive(Debug)]
399pub(super) struct SharedQueue<P: Ord> {
400 jobs: VecDeque<QueuedJob<P>>,
401 current_job: Option<CurrentJob>,
402}
403
404#[cfg(test)]
405#[cfg_attr(coverage_nightly, coverage(off))]
406mod tests {
407 use std::time::Instant;
408
409 use tracing_test::traced_test;
410
411 use super::*;
412
413 #[tokio::test(flavor = "multi_thread")]
414 #[traced_test]
415 async fn run_sync_jobs_by_priority() -> anyhow::Result<()> {
416 workers::run_jobs_by_priority(false).await
417 }
418
419 #[tokio::test(flavor = "multi_thread")]
420 #[traced_test]
421 async fn run_async_jobs_by_priority() -> anyhow::Result<()> {
422 workers::run_jobs_by_priority(true).await
423 }
424
425 #[tokio::test(flavor = "multi_thread")]
426 #[traced_test]
427 async fn get_sync_job_result() -> anyhow::Result<()> {
428 workers::get_job_result(false).await
429 }
430
431 #[tokio::test(flavor = "multi_thread")]
432 #[traced_test]
433 async fn get_async_job_result() -> anyhow::Result<()> {
434 workers::get_job_result(true).await
435 }
436
437 #[tokio::test(flavor = "multi_thread")]
438 #[traced_test]
439 async fn cancel_sync_job() -> anyhow::Result<()> {
440 workers::cancel_job(false).await
441 }
442
443 #[tokio::test(flavor = "multi_thread")]
444 #[traced_test]
445 async fn cancel_async_job() -> anyhow::Result<()> {
446 workers::cancel_job(true).await
447 }
448
449 #[tokio::test(flavor = "multi_thread")]
450 #[traced_test]
451 async fn cancel_sync_job_in_select() -> anyhow::Result<()> {
452 workers::cancel_job_in_select(false).await
453 }
454
455 #[tokio::test(flavor = "multi_thread")]
456 #[traced_test]
457 async fn cancel_async_job_in_select() -> anyhow::Result<()> {
458 workers::cancel_job_in_select(true).await
459 }
460
461 #[test]
462 #[traced_test]
463 fn runtime_shutdown_timeout_force_cancels_sync_job() -> anyhow::Result<()> {
464 workers::runtime_shutdown_timeout_force_cancels_job(false)
465 }
466
467 #[test]
468 #[traced_test]
469 fn runtime_shutdown_timeout_force_cancels_async_job() -> anyhow::Result<()> {
470 workers::runtime_shutdown_timeout_force_cancels_job(true)
471 }
472
473 #[test]
474 #[traced_test]
475 fn runtime_shutdown_cancels_sync_job() {
476 let _ = workers::runtime_shutdown_cancels_job(false);
477 }
478
479 #[test]
480 #[traced_test]
481 fn runtime_shutdown_cancels_async_job() -> anyhow::Result<()> {
482 workers::runtime_shutdown_cancels_job(true)
483 }
484
485 #[test]
486 #[traced_test]
487 fn spawned_tasks_live_as_long_as_jobqueue() -> anyhow::Result<()> {
488 workers::spawned_tasks_live_as_long_as_jobqueue(true)
489 }
490
491 #[tokio::test(flavor = "multi_thread")]
492 #[traced_test]
493 async fn panic_in_async_job_ends_job_cleanly() -> anyhow::Result<()> {
494 workers::panics::panic_in_job_ends_job_cleanly(true).await
495 }
496
497 #[tokio::test(flavor = "multi_thread")]
498 #[traced_test]
499 async fn panic_in_blocking_job_ends_job_cleanly() -> anyhow::Result<()> {
500 workers::panics::panic_in_job_ends_job_cleanly(false).await
501 }
502
503 #[tokio::test(flavor = "multi_thread")]
504 #[traced_test]
505 async fn stop_queue() -> anyhow::Result<()> {
506 workers::stop_queue().await
507 }
508
509 #[tokio::test(flavor = "multi_thread")]
510 #[traced_test]
511 async fn job_result_wrapper() -> anyhow::Result<()> {
512 workers::job_result_wrapper().await
513 }
514
515 mod workers {
516 use super::*;
517 use crate::errors::JobHandleError;
518 use crate::traits::JobResult;
519 use crate::JobResultWrapper;
520
521 #[derive(Debug, PartialEq, Eq, PartialOrd, Ord)]
522 pub enum DoubleJobPriority {
523 Low = 1,
524 Medium = 2,
525 High = 3,
526 }
527
528 type DoubleJobResult = JobResultWrapper<(u64, u64, Instant)>;
529
530 #[derive(Debug)]
532 struct DoubleJob {
533 data: u64,
534 duration: std::time::Duration,
535 is_async: bool,
536 }
537
538 #[async_trait::async_trait]
539 impl Job for DoubleJob {
540 fn is_async(&self) -> bool {
541 self.is_async
542 }
543
544 fn run(&self, cancel_rx: JobCancelReceiver) -> JobCompletion {
545 let start = Instant::now();
546 let sleep_time =
547 std::cmp::min(std::time::Duration::from_micros(100), self.duration);
548
549 let r = loop {
550 if start.elapsed() < self.duration {
551 match cancel_rx.has_changed() {
552 Ok(changed) if changed => break JobCompletion::Cancelled,
553 Err(_) => break JobCompletion::Cancelled,
554 _ => {}
555 }
556
557 std::thread::sleep(sleep_time);
558 } else {
559 break JobCompletion::Finished(
560 DoubleJobResult::new((self.data, self.data * 2, Instant::now())).into(),
561 );
562 }
563 };
564
565 tracing::info!("results: {:?}", r);
566 r
567 }
568
569 async fn run_async(&self) -> Box<dyn JobResult> {
570 tokio::time::sleep(self.duration).await;
571 let r = DoubleJobResult::new((self.data, self.data * 2, Instant::now()));
572 tracing::info!("results: {:?}", r);
573 r.into()
574 }
575 }
576
577 pub(super) async fn run_jobs_by_priority(is_async: bool) -> anyhow::Result<()> {
581 let start_of_test = Instant::now();
582
583 let mut job_queue = JobQueue::start();
585
586 let mut handles = vec![];
587 let duration = std::time::Duration::from_millis(20);
588
589 for i in (1..10).rev() {
591 let job1 = DoubleJob {
592 data: i,
593 duration,
594 is_async,
595 };
596 let job2 = DoubleJob {
597 data: i * 100,
598 duration,
599 is_async,
600 };
601 let job3 = DoubleJob {
602 data: i * 1000,
603 duration,
604 is_async,
605 };
606
607 handles.push(job_queue.add_job_mut(job1, DoubleJobPriority::Low)?);
609 handles.push(job_queue.add_job_mut(job2, DoubleJobPriority::Medium)?);
610 handles.push(job_queue.add_job_mut(job3, DoubleJobPriority::High)?);
611 }
612
613 assert!(job_queue.num_jobs() > 0);
615 assert!(job_queue.num_queued_jobs() > 0);
616
617 let mut results = futures::future::join_all(handles).await;
619
620 assert_eq!(0, job_queue.num_jobs());
621 assert_eq!(0, job_queue.num_queued_jobs());
622
623 results.sort_by(|a_completion, b_completion| {
626 let a = <&DoubleJobResult>::try_from(a_completion.as_ref().unwrap())
627 .unwrap()
628 .2;
629 let b = <&DoubleJobResult>::try_from(b_completion.as_ref().unwrap())
630 .unwrap()
631 .2;
632
633 a.cmp(&b)
634 });
635
636 let mut prev = DoubleJobResult::new((9999, 0, start_of_test));
641 for (i, c) in results.into_iter().enumerate() {
642 let job_result = DoubleJobResult::try_from(c?)?;
643
644 assert!(job_result.2 > prev.2);
645
646 if i != 1 {
650 assert!(job_result.0 < prev.0);
651 }
652
653 prev = job_result;
654 }
655
656 Ok(())
657 }
658
659 pub(super) async fn get_job_result(is_async: bool) -> anyhow::Result<()> {
662 let mut job_queue = JobQueue::start();
664 let duration = std::time::Duration::from_millis(20);
665
666 for i in 0..10 {
668 let job = DoubleJob {
669 data: i,
670 duration,
671 is_async,
672 };
673
674 let completion = job_queue.add_job_mut(job, DoubleJobPriority::Low)?.await?;
675
676 let job_result = DoubleJobResult::try_from(completion)?;
677
678 assert_eq!(i, job_result.0);
679 assert_eq!(i * 2, job_result.1);
680 }
681
682 Ok(())
683 }
684
685 pub(super) async fn stop_queue() -> anyhow::Result<()> {
688 let mut job_queue = JobQueue::start();
690 let duration = std::time::Duration::from_secs(3600); let job = DoubleJob {
694 data: 10,
695 duration,
696 is_async: true,
697 };
698 let job2 = DoubleJob {
699 data: 10,
700 duration,
701 is_async: true,
702 };
703 let job_handle = job_queue.add_job_mut(job, DoubleJobPriority::Low)?;
704 let job2_handle = job_queue.add_job_mut(job2, DoubleJobPriority::Low)?;
705
706 println!("job-queue: {:?}", job_queue);
708
709 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
710
711 job_queue.stop().await?;
712
713 assert!(job_handle.is_finished());
714 assert!(job2_handle.is_finished());
715
716 assert!(matches!(
717 job_handle.await,
718 Err(JobHandleError::JobResultError(_))
719 ));
720 assert!(matches!(
721 job2_handle.await,
722 Err(JobHandleError::JobResultError(_))
723 ));
724
725 Ok(())
726 }
727
728 pub(super) async fn cancel_job(is_async: bool) -> anyhow::Result<()> {
730 let mut job_queue = JobQueue::start();
732 let duration = std::time::Duration::from_secs(3600); let job = DoubleJob {
736 data: 10,
737 duration,
738 is_async,
739 };
740 let job_handle = job_queue.add_job_mut(job, DoubleJobPriority::Low)?;
741
742 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
743
744 job_handle.cancel().unwrap();
745 let completion = job_handle.await.unwrap();
746 assert!(matches!(completion, JobCompletion::Cancelled));
747
748 Ok(())
749 }
750
751 pub async fn cancel_job_in_select(is_async: bool) -> anyhow::Result<()> {
760 async fn do_some_work(
761 is_async: bool,
762 cancel_work_rx: tokio::sync::oneshot::Receiver<()>,
763 ) -> Result<DoubleJobResult, JobHandleError> {
764 let mut job_queue = JobQueue::start();
766
767 let duration = std::time::Duration::from_secs(3600); let job = DoubleJob {
771 data: 10,
772 duration,
773 is_async,
774 };
775
776 let job_handle = job_queue.add_job_mut(job, DoubleJobPriority::Low).unwrap();
778
779 tokio::pin!(job_handle);
782
783 let completion = tokio::select! {
785 completion = &mut job_handle => completion,
787
788 _ = cancel_work_rx => {
790 job_handle.cancel()?;
791 job_handle.await
792 }
793 };
794
795 println!("job_completion: {:#?}", completion);
796
797 let result = DoubleJobResult::try_from(completion?)?;
799
800 println!("job_result: {:#?}", result);
801
802 Ok(result)
803 }
804
805 let (cancel_tx, cancel_rx) = tokio::sync::oneshot::channel::<()>();
807
808 let worker_task = async move { do_some_work(is_async, cancel_rx).await };
810
811 let jh = tokio::task::spawn(worker_task);
813
814 cancel_tx.send(()).unwrap();
816
817 let job_handle_error = jh.await?.unwrap_err();
819
820 assert!(matches!(job_handle_error, JobHandleError::JobCancelled));
822
823 Ok(())
824 }
825
826 pub(super) fn runtime_shutdown_timeout_force_cancels_job(
847 is_async: bool,
848 ) -> anyhow::Result<()> {
849 let rt = tokio::runtime::Runtime::new()?;
850 let result = rt.block_on(async {
851 let mut job_queue = JobQueue::start();
853 let duration = std::time::Duration::from_secs(3600); let job = DoubleJob {
857 data: 10,
858 duration,
859 is_async,
860 };
861 let _rx = job_queue.add_job_mut(job, DoubleJobPriority::Low)?;
862
863 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
864 println!("finished scope");
865
866 Ok(())
867 });
868
869 let start = std::time::Instant::now();
870
871 println!("waiting 1 second for job before shutdown runtime");
872 rt.shutdown_timeout(tokio::time::Duration::from_secs(1));
873
874 assert!(start.elapsed() < std::time::Duration::from_secs(2));
875
876 result
877 }
878
879 pub(super) fn runtime_shutdown_cancels_job(is_async: bool) -> anyhow::Result<()> {
895 let rt = tokio::runtime::Runtime::new()?;
896 let start = tokio::time::Instant::now();
897
898 let result = rt.block_on(async {
899 let mut job_queue = JobQueue::start();
901
902 let duration = std::time::Duration::from_secs(5);
904
905 let job = DoubleJob {
906 data: 10,
907 duration,
908 is_async,
909 };
910
911 let rx_handle = job_queue.add_job_mut(job, DoubleJobPriority::Low)?;
912 drop(rx_handle);
913
914 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
916
917 Ok(())
918 });
919
920 drop(rt);
924
925 assert!(start.elapsed() < std::time::Duration::from_secs(5));
932
933 result
934 }
935
936 pub(super) fn spawned_tasks_live_as_long_as_jobqueue(is_async: bool) -> anyhow::Result<()> {
950 let rt = tokio::runtime::Runtime::new()?;
951
952 let result_ok: Arc<Mutex<bool>> = Arc::new(Mutex::new(true));
953
954 let result_ok_clone = result_ok.clone();
955 rt.block_on(async {
956 let job_queue = Arc::new(JobQueue::start());
958
959 let job_queue_cloned = job_queue.clone();
961 let jh = tokio::spawn(async move {
962 std::thread::sleep(std::time::Duration::from_millis(200));
967
968 let job = DoubleJob {
969 data: 10,
970 duration: std::time::Duration::from_secs(1),
971 is_async,
972 };
973
974 let result = job_queue_cloned.add_job(job, DoubleJobPriority::Low);
976
977 *result_ok_clone.lock().unwrap() = result.is_ok();
982 });
983
984 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
986
987 jh.abort();
990 let _ = jh.await;
991 });
992
993 drop(rt);
995
996 assert!(*result_ok.lock().unwrap());
997
998 Ok(())
999 }
1000
1001 pub mod panics {
1002 use super::*;
1003
1004 const PANIC_STR: &str = "job panics unexpectedly";
1005
1006 struct PanicJob {
1007 is_async: bool,
1008 }
1009
1010 #[async_trait::async_trait]
1011 impl Job for PanicJob {
1012 fn is_async(&self) -> bool {
1013 self.is_async
1014 }
1015
1016 fn run(&self, _cancel_rx: JobCancelReceiver) -> JobCompletion {
1017 panic!("{}", PANIC_STR);
1018 }
1019
1020 async fn run_async_cancellable(
1021 &self,
1022 _cancel_rx: JobCancelReceiver,
1023 ) -> JobCompletion {
1024 panic!("{}", PANIC_STR);
1025 }
1026 }
1027
1028 pub async fn panic_in_job_ends_job_cleanly(async_job: bool) -> anyhow::Result<()> {
1039 let mut job_queue = JobQueue::start();
1041
1042 let job = PanicJob {
1043 is_async: async_job,
1044 };
1045 let job_handle = job_queue.add_job_mut(job, DoubleJobPriority::Low)?;
1046
1047 let job_result = job_handle.await?.result();
1048
1049 assert!(!job_queue.tx_job_added.is_closed());
1051 assert!(!job_queue.tx_stop.is_closed());
1052
1053 assert!(matches!(
1055 job_result,
1056 Err(e) if e.panic_message() == Some((*PANIC_STR).to_string())
1057 ));
1058
1059 let newjob = DoubleJob {
1061 data: 10,
1062 duration: std::time::Duration::from_millis(50),
1063 is_async: false,
1064 };
1065
1066 let new_job_handle = job_queue.add_job_mut(newjob, DoubleJobPriority::Low)?;
1068
1069 assert!(new_job_handle.await?.result().is_ok());
1071
1072 Ok(())
1073 }
1074 }
1075
1076 pub(super) async fn job_result_wrapper() -> anyhow::Result<()> {
1078 type MyJobResult = JobResultWrapper<(u64, u64, Instant)>;
1079
1080 #[derive(Debug)]
1082 struct MyJob {
1083 data: u64,
1084 duration: std::time::Duration,
1085 }
1086
1087 #[async_trait::async_trait]
1088 impl Job for MyJob {
1089 fn is_async(&self) -> bool {
1090 true
1091 }
1092
1093 async fn run_async(&self) -> Box<dyn JobResult> {
1094 tokio::time::sleep(self.duration).await;
1095 MyJobResult::new((self.data, self.data * 2, Instant::now())).into()
1096 }
1097 }
1098
1099 let mut job_queue = JobQueue::start();
1100 let job = MyJob {
1101 data: 15,
1102 duration: std::time::Duration::from_secs(5),
1103 };
1104 let job_handle = job_queue.add_job_mut(job, 10usize)?;
1105 let completion = job_handle.await?;
1106 let job_result = MyJobResult::try_from(completion)?;
1107 let answer = job_result.into_inner();
1108
1109 assert_eq!(answer.0 * 2, answer.1);
1110
1111 Ok(())
1112 }
1113 }
1114}