| 833 | |
| 834 | #[test] |
| 835 | fn test_builder_thread_pool() { |
| 836 | let pool = BuilderThreadPool::new(); |
| 837 | |
| 838 | // Test we can run 10 tasks in parallel |
| 839 | let active_tasks = AtomicUsize::new(0); |
| 840 | let mut tasks = Vec::new(); |
| 841 | for _i in 0..10 { |
| 842 | tasks.push(|| { |
| 843 | active_tasks.fetch_add(1, Ordering::Relaxed); |
| 844 | thread::sleep(Duration::from_millis(100)); |
| 845 | Result::<(), ()>::Ok(()) |
| 846 | }); |
| 847 | } |
| 848 | pool.run( |
| 849 | tasks, |
| 850 | NonZeroUsize::new(3).unwrap(), |
| 851 | &AtomicAbortStatus::new(AbortStatus::None), |
| 852 | &|| AbortStatus::None, |
| 853 | ) |
| 854 | .unwrap(); |
| 855 | |
| 856 | assert_eq!(active_tasks.load(Ordering::Relaxed), 10); |
| 857 | |
| 858 | // Test max parallelism |
| 859 | let active_tasks = AtomicUsize::new(0); |
| 860 | let max_active_task = AtomicUsize::new(0); |
| 861 | let mut tasks = Vec::new(); |
| 862 | for _i in 0..10 { |
| 863 | tasks.push(|| { |
| 864 | let nb_tasks = active_tasks.fetch_add(1, Ordering::Relaxed); |
| 865 | max_active_task.fetch_max(nb_tasks + 1, Ordering::Relaxed); |
| 866 | thread::sleep(Duration::from_millis(1000)); |
| 867 | active_tasks.fetch_sub(1, Ordering::Relaxed); |
| 868 | Result::<(), ()>::Ok(()) |
| 869 | }); |
| 870 | } |
| 871 | pool.run( |
| 872 | tasks, |
| 873 | NonZeroUsize::new(3).unwrap(), |
| 874 | &AtomicAbortStatus::new(AbortStatus::None), |
| 875 | &|| AbortStatus::None, |
| 876 | ) |
| 877 | .unwrap(); |
| 878 | |
| 879 | assert_eq!(active_tasks.load(Ordering::Relaxed), 0); |
| 880 | assert_eq!(max_active_task.load(Ordering::Relaxed), 3); |
| 881 | |
| 882 | // Test we get our error, and that we try to stop all tasks on first error |
| 883 | let mut tasks = Vec::new(); |
| 884 | let active_taks = Arc::new(AtomicUsize::new(0)); |
| 885 | let should_abort_flag = AtomicAbortStatus::new(AbortStatus::None); |
| 886 | for i in 0..10 { |
| 887 | tasks.push({ |
| 888 | let active_tasks = active_taks.clone(); |
| 889 | move || { |
| 890 | active_tasks.fetch_add(1, Ordering::Relaxed); |
| 891 | thread::sleep(Duration::from_millis(100)); |
| 892 | if i == 5 { |