1use std::{
4 any::Any,
5 ops::ControlFlow,
6 panic::{self, AssertUnwindSafe, catch_unwind},
7 sync::{
8 Arc, OnceLock, SyncView,
9 atomic::{AtomicBool, AtomicUsize, Ordering},
10 mpmc,
11 },
12 time::Duration,
13};
14
15use fixedbitset::FixedBitSet;
16use parking_lot::{Condvar, Mutex, RwLock};
17use tokio::{runtime::Handle, task::AbortHandle};
18use tracing::{Span, info_span};
19
20use crate::{TurboTasksApi, manager::try_turbo_tasks, turbo_tasks_scope};
21
22const WORKER_IDLE_TIMEOUT: Duration = Duration::from_micros(100);
27
28pub fn scope_unbounded<'env, T, F>(initial: impl IntoIterator<Item = T>, run: F)
41where
42 T: Send + 'static,
43 F: Fn(&Scope<'_, T, ()>, T) -> ControlFlow<()> + Send + Sync + 'env,
44{
45 scope_unbounded_with(
46 initial,
47 || (),
48 |spawner, item, ()| run(spawner, item),
49 |(), ()| (),
50 )
51}
52
53pub fn scope_unbounded_with<'env, T, R, F, Init, Merge>(
77 initial: impl IntoIterator<Item = T>,
78 init: Init,
79 run: F,
80 merge: Merge,
81) -> R
82where
83 T: Send + 'static,
84 R: Send + 'env,
85 F: Fn(&Scope<'_, T, R>, T, &mut R) -> ControlFlow<()> + Send + Sync + 'env,
86 Init: Fn() -> R + Send + Sync + 'env,
87 Merge: Fn(R, R) -> R + Send + Sync + 'env,
88{
89 let handle = Handle::current();
90 let max_workers = handle.metrics().num_workers().saturating_sub(1);
92 let span = Span::current();
93
94 let init_ref: &(dyn Fn() -> R + Send + Sync + '_) = &init;
98 let merge_ref: &(dyn Fn(R, R) -> R + Send + Sync + '_) = &merge;
99
100 let (sender, receiver) = mpmc::channel();
101 let mut inner = ScopeInner {
102 remaining_tasks: AtomicUsize::new(0),
103 panic: OnceLock::new(),
104 work_queue: receiver,
105 work_queue_sender: RwLock::new(Some(sender)),
106 aborted: AtomicBool::new(false),
107 available_slots: AtomicUsize::new(max_workers),
108 workers: Mutex::new(WorkerSlots::new(max_workers)),
109 handle: handle.clone(),
110 span: span.clone(),
111 workers_idle: Condvar::new(),
112 turbo_tasks: try_turbo_tasks(),
113 run: &run,
114 results: Mutex::new(None),
115 init: init_ref,
116 merge: merge_ref,
117 };
118
119 let joiner = Joiner { inner: &inner };
121
122 inner.remaining_tasks.fetch_add(1, Ordering::Relaxed);
124 for item in initial {
125 enqueue(&inner, item);
126 }
127
128 drop(joiner);
131
132 if let Some(err) = inner.panic.take() {
133 panic::resume_unwind(err.into_inner());
134 }
135
136 inner.results.lock().take().unwrap_or_else(init)
137}
138
139pub struct Scope<'scope, T: Send + 'static, R = ()> {
142 inner: &'scope ScopeInner<'scope, T, R>,
143}
144
145impl<T: Send + 'static, R: Send> Scope<'_, T, R> {
146 pub fn spawn(&self, item: T) {
151 enqueue(self.inner, item);
152 }
153}
154
155type RunFn<'run, T, R> =
159 &'run (dyn Fn(&Scope<'_, T, R>, T, &mut R) -> ControlFlow<()> + Send + Sync + 'run);
160
161struct ScopeInner<'run, T: Send + 'static, R> {
168 remaining_tasks: AtomicUsize,
171 panic: OnceLock<SyncView<Box<dyn Any + Send + 'static>>>,
173 work_queue: mpmc::Receiver<T>,
175 work_queue_sender: RwLock<Option<mpmc::Sender<T>>>,
177 aborted: AtomicBool,
178 available_slots: AtomicUsize,
180 workers: Mutex<WorkerSlots>,
181 workers_idle: Condvar,
183 handle: Handle,
184 span: Span,
185 turbo_tasks: Option<Arc<dyn TurboTasksApi>>,
186 run: RunFn<'run, T, R>,
188 init: &'run (dyn Fn() -> R + Send + Sync + 'run),
189 merge: &'run (dyn Fn(R, R) -> R + Send + Sync + 'run),
190 results: Mutex<Option<R>>,
193}
194
195impl<T: Send + 'static, R> ScopeInner<'_, T, R> {
196 fn close(&self) {
200 drop(self.work_queue_sender.write().take());
201 }
202
203 fn abort(&self) {
205 self.aborted.store(true, Ordering::Release);
206 self.close();
207 }
208
209 fn on_item_finished(&self) {
210 if self.remaining_tasks.fetch_sub(1, Ordering::Release) == 1 {
211 self.close();
212 }
213 }
214
215 fn record_panic(&self, err: Box<dyn Any + Send + 'static>) {
217 self.abort();
218 let _ = self.panic.set(SyncView::new(err));
219 }
220
221 fn drain(&self, is_worker: bool) {
227 if is_worker && let Some(turbo_tasks) = &self.turbo_tasks {
228 turbo_tasks_scope(turbo_tasks.clone(), || self.drain_loop(is_worker))
229 } else {
230 self.drain_loop(is_worker)
231 }
232 }
233
234 fn drain_loop(&self, is_worker: bool) {
235 let mut acc: Option<R> = None;
236 while let Some(item) = if is_worker {
237 self.work_queue.recv_timeout(WORKER_IDLE_TIMEOUT).ok()
238 } else {
239 self.work_queue.recv().ok()
240 } {
241 if self.aborted.load(Ordering::Acquire) {
243 self.on_item_finished();
244 continue;
245 }
246 let spawner = Scope { inner: self };
247 let result = catch_unwind(AssertUnwindSafe(|| {
248 let acc = acc.get_or_insert_with(self.init);
251 (self.run)(&spawner, item, acc)
252 }));
253
254 match result {
255 Ok(ControlFlow::Continue(())) => {}
256 Ok(ControlFlow::Break(())) => {
257 self.abort();
258 }
259 Err(panic) => {
261 self.record_panic(panic);
262 }
263 };
264 self.on_item_finished();
265 }
266
267 if let Some(acc) = acc {
269 let merged = catch_unwind(AssertUnwindSafe(|| {
270 let mut results = self.results.lock();
271 *results = Some(match results.take() {
272 Some(existing) => (self.merge)(existing, acc),
273 None => acc,
274 });
275 }));
276 if let Err(panic) = merged {
277 self.record_panic(panic);
278 }
279 }
280 }
281}
282
283fn enqueue<T: Send + 'static, R: Send>(inner: &ScopeInner<'_, T, R>, item: T) {
287 if inner.aborted.load(Ordering::Acquire) {
288 return;
289 }
290 let num_tasks = inner.remaining_tasks.fetch_add(1, Ordering::Relaxed) + 1;
291 let sent = {
292 let sender = inner.work_queue_sender.read();
293 match sender.as_ref() {
294 Some(sender) => sender.send(item).is_ok(),
295 None => false,
297 }
298 };
299 if !sent {
300 inner.on_item_finished(); return;
302 }
303 spawn_worker_if_needed(inner, num_tasks);
304}
305
306fn spawn_worker_if_needed<T: Send + 'static, R: Send>(
308 inner: &ScopeInner<'_, T, R>,
309 num_enqueued_tasks: usize,
310) {
311 if num_enqueued_tasks <= 1
312 || inner.available_slots.load(Ordering::Relaxed) == 0
313 || inner.aborted.load(Ordering::Acquire)
314 {
315 return;
316 }
317
318 let erased: &(dyn Drainable + Send + Sync + '_) = inner;
322 let erased: &'static (dyn Drainable + Send + Sync + 'static) = unsafe {
323 std::mem::transmute::<
324 &(dyn Drainable + Send + Sync + '_),
325 &'static (dyn Drainable + Send + Sync + 'static),
326 >(erased)
327 };
328
329 let mut slots = inner.workers.lock();
330 let Some(slot) = slots.free_slot() else {
331 return;
332 };
333 let span = inner.span.clone();
334 let guard = erased.claim_worker_slot(slot);
336 let handle = inner
337 .handle
338 .spawn(async move {
339 let _span = span.entered();
340 let _guard = guard;
341 erased.drain(true);
342 })
343 .abort_handle();
344 slots.occupy(slot, handle, &inner.available_slots);
345}
346
347trait Drainable {
354 fn drain(&self, is_worker: bool);
355 fn claim_worker_slot(&self, slot: usize) -> WorkerGuard<'_>;
358}
359
360impl<T: Send + 'static, R> Drainable for ScopeInner<'_, T, R> {
361 fn drain(&self, is_worker: bool) {
362 ScopeInner::drain(self, is_worker)
363 }
364
365 fn claim_worker_slot(&self, slot: usize) -> WorkerGuard<'_> {
366 WorkerGuard {
367 slots: &self.workers,
368 workers_idle: &self.workers_idle,
369 available_slots: &self.available_slots,
370 slot,
371 }
372 }
373}
374
375struct WorkerSlots {
380 handles: Vec<Option<AbortHandle>>,
381 occupied: FixedBitSet,
384}
385
386impl WorkerSlots {
387 fn new(max_workers: usize) -> Self {
388 Self {
389 handles: vec![None; max_workers],
390 occupied: FixedBitSet::with_capacity(max_workers),
391 }
392 }
393
394 fn free_slot(&self) -> Option<usize> {
395 self.occupied.zeroes().next()
396 }
397
398 fn occupy(&mut self, slot: usize, handle: AbortHandle, available_slots: &AtomicUsize) {
400 debug_assert!(
401 !self.occupied.contains(slot),
402 "slot {slot} already occupied"
403 );
404 available_slots.fetch_sub(1, Ordering::Relaxed);
405 self.occupied.insert(slot);
406 let previous = self.handles[slot].replace(handle);
407 debug_assert!(previous.is_none(), "slot {slot} held a live handle");
408 }
409
410 fn release(&mut self, slot: usize, available_slots: &AtomicUsize) {
412 available_slots.fetch_add(1, Ordering::Relaxed);
413 self.occupied.remove(slot);
414 self.handles[slot] = None;
415 }
416
417 fn is_idle(&self) -> bool {
420 self.occupied.is_clear()
421 }
422}
423
424struct WorkerGuard<'a> {
426 slots: &'a Mutex<WorkerSlots>,
427 workers_idle: &'a Condvar,
428 available_slots: &'a AtomicUsize,
429 slot: usize,
430}
431
432impl Drop for WorkerGuard<'_> {
433 fn drop(&mut self) {
434 let mut slots_guard = self.slots.lock();
435 slots_guard.release(self.slot, self.available_slots);
436 if slots_guard.is_idle() {
437 self.workers_idle.notify_all();
440 }
441 }
442}
443
444struct Joiner<'a, 'run, T: Send + 'static, R> {
446 inner: &'a ScopeInner<'run, T, R>,
447}
448
449impl<T: Send + 'static, R> Drop for Joiner<'_, '_, T, R> {
450 fn drop(&mut self) {
451 self.inner.on_item_finished();
453 self.inner.drain(false);
455 let _span = info_span!("blocking: waiting for scope to end").entered();
458 let handles: Vec<_> = {
459 let mut slots = self.inner.workers.lock();
460 slots.handles.iter_mut().filter_map(Option::take).collect()
461 };
462
463 for handle in handles {
466 handle.abort();
467 }
468
469 let mut slots = self.inner.workers.lock();
473 while !slots.is_idle() {
474 self.inner.workers_idle.wait(&mut slots);
475 }
476 }
477}
478#[cfg(test)]
479mod tests {
480 use std::{
481 sync::{Arc, atomic::AtomicUsize},
482 thread,
483 time::Duration,
484 };
485
486 use super::*;
487
488 fn with_runtime<F, T>(worker_threads: usize, body: F) -> T
497 where
498 F: FnOnce() -> T + Send + 'static,
499 T: Send + 'static,
500 {
501 let runtime = tokio::runtime::Builder::new_multi_thread()
502 .worker_threads(worker_threads)
503 .enable_time()
504 .build()
505 .unwrap();
506 runtime.block_on(async { tokio::task::spawn_blocking(body).await.unwrap() })
507 }
508
509 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
511 async fn test_unbounded_wide_burst_of_leaves() {
512 const CHILDREN: usize = 1000;
513 let processed = Arc::new(AtomicUsize::new(0));
514 let processed_clone = processed.clone();
515 tokio::task::spawn_blocking(move || {
516 scope_unbounded(std::iter::once(0usize), move |spawner, item| {
517 processed_clone.fetch_add(1, Ordering::SeqCst);
518 if item == 0 {
519 for i in 0..CHILDREN {
521 spawner.spawn(1 + i);
522 }
523 }
524 ControlFlow::Continue(())
525 });
526 })
527 .await
528 .unwrap();
529 assert_eq!(processed.load(Ordering::SeqCst), 1 + CHILDREN);
531 }
532
533 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
535 async fn test_unbounded_slow_seeding_iterator_completes() {
536 const SEEDS: usize = 16;
537 let processed = Arc::new(AtomicUsize::new(0));
538 let processed_clone = processed.clone();
539 tokio::task::spawn_blocking(move || {
540 let slow_seeds = std::iter::from_fn({
543 let mut next = 0;
544 move || {
545 if next == SEEDS {
546 return None;
547 }
548 thread::sleep(Duration::from_millis(2));
549 next += 1;
550 Some(next - 1)
551 }
552 });
553 scope_unbounded(slow_seeds, move |_spawner, _item| {
554 processed_clone.fetch_add(1, Ordering::SeqCst);
555 ControlFlow::Continue(())
556 });
557 })
558 .await
559 .unwrap();
560 assert_eq!(
561 processed.load(Ordering::SeqCst),
562 SEEDS,
563 "seeds produced after the queue briefly drained must still be processed"
564 );
565 }
566
567 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
569 async fn test_unbounded_abort_during_cascade() {
570 const MAX_ID: usize = 1 << 14;
573 let processed = Arc::new(AtomicUsize::new(0));
574 let processed_clone = processed.clone();
575 tokio::task::spawn_blocking(move || {
576 scope_unbounded(std::iter::once(1usize), move |spawner, id| {
577 let n = processed_clone.fetch_add(1, Ordering::SeqCst);
578 if n == 100 {
579 return ControlFlow::Break(());
580 }
581 let (left, right) = (id * 2, id * 2 + 1);
582 if left <= MAX_ID {
583 spawner.spawn(left);
584 }
585 if right <= MAX_ID {
586 spawner.spawn(right);
587 }
588 ControlFlow::Continue(())
589 });
590 })
591 .await
592 .unwrap();
593 let count = processed.load(Ordering::SeqCst);
594 assert!(
595 count < MAX_ID,
596 "abort must cut the cascade short, but {count} items ran"
597 );
598 }
599
600 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
605 async fn test_unbounded_spawn_after_abort_is_dropped() {
606 let processed = Arc::new(AtomicUsize::new(0));
607 let processed_clone = processed.clone();
608 const SEEDS: usize = 64;
609 tokio::task::spawn_blocking(move || {
610 scope_unbounded(0..SEEDS, move |spawner, item| {
611 processed_clone.fetch_add(1, Ordering::SeqCst);
612 for i in 0..1000 {
615 spawner.spawn(SEEDS + item * 1000 + i);
616 }
617 ControlFlow::Break(())
618 });
619 })
620 .await
621 .unwrap();
622 let count = processed.load(Ordering::SeqCst);
625 assert!(
626 count <= SEEDS,
627 "post-abort spawns must be dropped, but {count} items ran"
628 );
629 }
630
631 #[tokio::test(flavor = "current_thread")]
633 async fn test_unbounded_abort_current_thread_runtime() {
634 let processed = Arc::new(AtomicUsize::new(0));
635 let processed_clone = processed.clone();
636 tokio::task::spawn_blocking(move || {
637 scope_unbounded(0..1000usize, move |spawner, _item| {
638 processed_clone.fetch_add(1, Ordering::SeqCst);
639 spawner.spawn(9999);
640 ControlFlow::Break(())
641 });
642 })
643 .await
644 .unwrap();
645 assert_eq!(processed.load(Ordering::SeqCst), 1);
647 }
648
649 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
653 async fn test_unbounded_abort_then_panic() {
654 let result = catch_unwind(AssertUnwindSafe(|| {
655 scope_unbounded(0..1000usize, |_spawner, item| {
656 if item == 0 {
657 panic!("Intentional panic");
658 }
659 ControlFlow::Break(())
660 });
661 unreachable!();
662 }));
663 let err = result.expect_err("the panic must propagate even though the scope aborted");
664 assert_eq!(err.downcast_ref::<&str>(), Some(&"Intentional panic"));
665 }
666
667 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
674 async fn test_unbounded_panic_propagates_and_abandons_queue() {
675 const MAX_ID: usize = 1 << 14;
676 let processed = Arc::new(AtomicUsize::new(0));
677 let processed_clone = processed.clone();
678 let result = catch_unwind(AssertUnwindSafe(|| {
679 scope_unbounded(std::iter::once(1usize), move |spawner, id| {
680 let n = processed_clone.fetch_add(1, Ordering::SeqCst);
681 if n == 100 {
682 panic!("Intentional panic");
683 }
684 let (left, right) = (id * 2, id * 2 + 1);
685 if left <= MAX_ID {
686 spawner.spawn(left);
687 }
688 if right <= MAX_ID {
689 spawner.spawn(right);
690 }
691 ControlFlow::Continue(())
692 });
693 unreachable!();
694 }));
695 let err = result.expect_err("the panic must propagate");
696 assert_eq!(err.downcast_ref::<&str>(), Some(&"Intentional panic"));
697 let count = processed.load(Ordering::SeqCst);
698 assert!(
699 count < MAX_ID,
700 "a panic must cut the cascade short, but {count} items ran"
701 );
702 }
703
704 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
711 async fn test_unbounded_with_collects_all_values() {
712 const ITEMS: usize = 500;
713 let mut collected = tokio::task::spawn_blocking(|| {
714 scope_unbounded_with(
715 0..ITEMS,
716 Vec::new,
717 |_spawner, item: usize, acc: &mut Vec<usize>| {
718 acc.push(item);
719 ControlFlow::Continue(())
720 },
721 |mut a: Vec<usize>, b| {
722 a.extend(b);
723 a
724 },
725 )
726 })
727 .await
728 .unwrap();
729 collected.sort_unstable();
730 assert_eq!(collected, (0..ITEMS).collect::<Vec<_>>());
731 }
732
733 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
736 async fn test_unbounded_with_empty_returns_init() {
737 let total = tokio::task::spawn_blocking(|| {
738 scope_unbounded_with(
739 std::iter::empty::<usize>(),
740 || 42usize,
741 |_spawner, _item, _acc| ControlFlow::Continue(()),
742 |a, b| a + b,
743 )
744 })
745 .await
746 .unwrap();
747 assert_eq!(total, 42, "expected exactly one init(), got {total}");
748 }
749
750 #[tokio::test(flavor = "current_thread")]
753 async fn test_unbounded_current_thread_direct_call_completes() {
754 let processed = Arc::new(AtomicUsize::new(0));
755 let processed_clone = processed.clone();
756 scope_unbounded(0..8usize, move |spawner, item| {
757 processed_clone.fetch_add(1, Ordering::SeqCst);
758 if item < 3 {
759 spawner.spawn(100 + item);
760 }
761 ControlFlow::Continue(())
762 });
763 assert_eq!(processed.load(Ordering::SeqCst), 11);
764 }
765
766 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
769 async fn test_unbounded_with_abort_returns_partial_results() {
770 let processed = tokio::task::spawn_blocking(|| {
771 scope_unbounded_with(
772 0..1000usize,
773 || 0usize,
774 |_spawner, item, acc| {
775 *acc += 1;
776 if item == 0 {
777 return ControlFlow::Break(());
778 }
779 ControlFlow::Continue(())
780 },
781 |a, b| a + b,
782 )
783 })
784 .await
785 .unwrap();
786 assert!(processed >= 1, "expected the aborting item to be counted");
788 assert!(
789 processed < 1000,
790 "abort should abandon queued items, but all {processed} ran"
791 );
792 }
793
794 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
797 async fn test_unbounded_with_panic_propagates() {
798 let result = catch_unwind(AssertUnwindSafe(|| {
799 scope_unbounded_with(
800 0..100usize,
801 || 0usize,
802 |_spawner, item, acc| {
803 if item == 50 {
804 panic!("Intentional panic");
805 }
806 *acc += 1;
807 ControlFlow::Continue(())
808 },
809 |a, b| a + b,
810 );
811 unreachable!();
812 }));
813 let err = result.expect_err("the panic must propagate out of the fold API");
814 assert_eq!(err.downcast_ref::<&str>(), Some(&"Intentional panic"));
815 }
816
817 #[test]
819 fn test_unbounded_with_borrowed_accumulator() {
820 let label = String::from("item");
821 let count = with_runtime(4, move || {
822 let label = &label;
823 scope_unbounded_with(
824 0..32usize,
825 Vec::new,
826 |_spawner, item: usize, acc: &mut Vec<String>| {
827 acc.push(format!("{label}-{item}"));
828 ControlFlow::Continue(())
829 },
830 |mut a: Vec<String>, b| {
831 a.extend(b);
832 a
833 },
834 )
835 .len()
836 });
837 assert_eq!(count, 32);
838 }
839
840 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
855 async fn test_unbounded_worker_respawns_after_going_idle() {
856 let inits = Arc::new(AtomicUsize::new(0));
857 let counted = inits.clone();
858 let processed = Arc::new(AtomicUsize::new(0));
859 let ran = processed.clone();
860 tokio::task::spawn_blocking(move || {
861 scope_unbounded_with(
862 std::iter::once(0usize),
863 move || {
864 counted.fetch_add(1, Ordering::SeqCst);
865 },
866 move |spawner, item, ()| {
867 ran.fetch_add(1, Ordering::SeqCst);
868 if item == 0 {
869 thread::sleep(Duration::from_millis(100));
873 spawner.spawn(1);
874 }
875 ControlFlow::Continue(())
876 },
877 |(), ()| (),
878 )
879 })
880 .await
881 .unwrap();
882 assert_eq!(
883 processed.load(Ordering::SeqCst),
884 2,
885 "work spawned after the pool went idle must still run"
886 );
887 let count = inits.load(Ordering::SeqCst);
890 assert!(
891 (1..=2).contains(&count),
892 "expected 1-2 drainer lifetimes, got {count}"
893 );
894 }
895
896 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
898 async fn test_unbounded_empty_spawns_no_workers() {
899 let inits = Arc::new(AtomicUsize::new(0));
900 let counted = inits.clone();
901 tokio::task::spawn_blocking(move || {
902 scope_unbounded_with(
903 std::iter::empty::<usize>(),
904 move || counted.fetch_add(1, Ordering::SeqCst),
905 |_spawner, _item, _acc| ControlFlow::Continue(()),
906 |a, b| a + b,
907 )
908 })
909 .await
910 .unwrap();
911 assert_eq!(
912 inits.load(Ordering::SeqCst),
913 1,
914 "expected only the terminal identity init(), not a per-drainer one"
915 );
916 }
917
918 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
920 async fn test_unbounded_busy_queue_does_not_churn_workers() {
921 const ITEMS: usize = 20_000;
922 let inits = Arc::new(AtomicUsize::new(0));
923 let counted = inits.clone();
924 let processed = tokio::task::spawn_blocking(move || {
925 scope_unbounded_with(
926 0..ITEMS,
927 move || {
928 counted.fetch_add(1, Ordering::SeqCst);
929 0usize
930 },
931 |_spawner, _item, acc| {
932 *acc += 1;
933 ControlFlow::Continue(())
934 },
935 |a, b| a + b,
936 )
937 })
938 .await
939 .unwrap();
940 assert_eq!(processed, ITEMS, "every item must run");
941 let count = inits.load(Ordering::SeqCst);
946 assert!(
947 count <= 16,
948 "a saturated queue should not churn drainers, got {count} lifetimes for {ITEMS} items"
949 );
950 }
951
952 #[test]
955 fn test_unbounded_join_under_thread_starvation() {
956 let runtime = tokio::runtime::Builder::new_multi_thread()
957 .worker_threads(2)
958 .enable_time()
959 .build()
960 .unwrap();
961 runtime.block_on(async {
962 let mut scopes = Vec::new();
963 for _ in 0..8 {
964 scopes.push(tokio::task::spawn_blocking(|| {
965 let processed = Arc::new(AtomicUsize::new(0));
966 let counted = processed.clone();
967 scope_unbounded(0..200usize, move |spawner, item| {
968 counted.fetch_add(1, Ordering::SeqCst);
969 if item < 200 {
971 spawner.spawn(1000 + item);
972 }
973 thread::yield_now();
974 ControlFlow::Continue(())
975 });
976 processed.load(Ordering::SeqCst)
977 }));
978 }
979 for scope in scopes {
980 assert_eq!(scope.await.unwrap(), 400, "every scope must drain fully");
981 }
982 });
983 }
984}