@@ -104,6 +104,28 @@ namespace threading {
104104 }
105105 return true ;
106106 }
107+
108+ // / Returns true when every concurrency slot can be acquired via the fast path.
109+ // / Held locks are released when the function returns.
110+ bool has_full_capacity (Group& group, const int concurrency) {
111+ std::vector<std::unique_ptr<Lock>> held;
112+ held.reserve (static_cast <std::size_t >(concurrency));
113+ for (int i = 0 ; i < concurrency; ++i) {
114+ auto lock = group.try_acquire_running_lock ();
115+ if (lock == nullptr ) {
116+ return false ;
117+ }
118+ held.push_back (std::move (lock));
119+ }
120+ return true ;
121+ }
122+
123+ // / Spin until all concurrency slots are acquirable or `timeout` elapses.
124+ bool wait_for_full_capacity (Group& group,
125+ const int concurrency,
126+ const std::chrono::milliseconds timeout) {
127+ return wait_for ([&] { return has_full_capacity (group, concurrency); }, timeout);
128+ }
107129 } // namespace
108130
109131 SCENARIO (" When there are no tokens available the lock should be false" ) {
@@ -471,35 +493,6 @@ namespace threading {
471493 }
472494 }
473495
474- SCENARIO (" Opportunistic drain during park publish must not leak group tokens" ) {
475- GIVEN (" A group with one token and a slow-path holder" ) {
476- auto group = make_group (1 );
477- NUClear::id_t task_id_source = 0 ;
478-
479- Group::TestAccess::set_capture_drains (*group, true );
480-
481- std::unique_ptr<Lock> slow_lock = group->lock (++task_id_source, 1 , [] {});
482- CHECK (slow_lock->lock () == true );
483-
484- WHEN (" A fast waiter publishes, the slow lock releases, then the waiter reconciles" ) {
485- auto slot = Group::TestAccess::park_publish (*group, make_test_task (), nullptr , false );
486-
487- slow_lock.reset ();
488-
489- Group::TestAccess::park_reconcile (*group, slot);
490-
491- THEN (" All tokens are restored after quiescing and the group is not deadlocked" ) {
492- auto captured = Group::TestAccess::take_captured_drains (*group);
493- REQUIRE (captured.size () == 1 );
494- captured.front ().lock .reset ();
495-
496- CHECK (Group::TestAccess::tokens (*group) == group->descriptor ->concurrency );
497- CHECK (Group::TestAccess::try_acquire_running_lock (*group) != nullptr );
498- }
499- }
500- }
501- }
502-
503496 SCENARIO (" Concurrent fast and slow path traffic never leaks group tokens or deadlocks" ) {
504497 const int concurrency = GENERATE (1 , 2 , 3 );
505498 CAPTURE (concurrency);
@@ -598,12 +591,8 @@ namespace threading {
598591 // next acquire. Without this the next round's lock() re-raises slow_pending and
599592 // legitimately defers not-yet-drained fast waiters (slow path has priority);
600593 // that is expected scheduler behaviour, not a leak.
601- const bool quiesced = wait_for (
602- [&] {
603- return Group::TestAccess::tokens (*groups[0 ]) == concurrency
604- && Group::TestAccess::tokens (*groups[1 ]) == concurrency;
605- },
606- std::chrono::seconds (10 ));
594+ const bool quiesced = wait_for_full_capacity (*groups[0 ], concurrency, std::chrono::seconds (10 ))
595+ && wait_for_full_capacity (*groups[1 ], concurrency, std::chrono::seconds (10 ));
607596 REQUIRE (quiesced);
608597 }
609598
@@ -627,11 +616,11 @@ namespace threading {
627616
628617 // (b) No leaked/duplicated tokens, and the group is still usable.
629618 for (auto & g : groups) {
630- CHECK (Group::TestAccess::tokens (*g) == concurrency);
631- auto fresh = Group::TestAccess:: try_acquire_running_lock (*g );
619+ CHECK (has_full_capacity (*g, concurrency) );
620+ auto fresh = g-> try_acquire_running_lock ();
632621 CHECK (fresh != nullptr );
633622 fresh.reset ();
634- CHECK (Group::TestAccess::tokens (*g) == concurrency);
623+ CHECK (has_full_capacity (*g, concurrency) );
635624 }
636625 }
637626
@@ -658,7 +647,7 @@ namespace threading {
658647
659648 WHEN (" The fast path tries to acquire a running lock" ) {
660649 THEN (" No token is handed out until the slow lock releases" ) {
661- CHECK (Group::TestAccess:: try_acquire_running_lock (*group ) == nullptr );
650+ CHECK (group-> try_acquire_running_lock () == nullptr );
662651 }
663652 }
664653 }
@@ -690,9 +679,12 @@ namespace threading {
690679 auto pool = std::make_unique<Pool>(*scheduler, pool_desc);
691680 auto group = make_group (1 );
692681
693- Group::TestAccess::park_publish (*group, make_test_task (), pool.get (), false );
682+ // A slow-path waiter blocks the fast path without holding a token.
683+ std::unique_ptr<Lock> slow_lock = group->lock (1 , 1 , [] {});
684+ CHECK_FALSE (group->try_submit (make_test_task (), pool.get (), false ));
694685
695686 WHEN (" The group is destroyed without draining the parked waiter" ) {
687+ slow_lock.reset ();
696688 group.reset ();
697689
698690 THEN (" The pool can still shut down cleanly because external waiters were balanced" ) {
@@ -703,53 +695,6 @@ namespace threading {
703695 }
704696 }
705697
706- SCENARIO (" Releasing a locked slow-path lock drains a committed fast waiter when tokens are negative" ) {
707- GIVEN (" A group with one token, a locked slow-path holder, and a parked fast waiter" ) {
708- auto group = make_group (1 );
709-
710- Group::TestAccess::set_capture_drains (*group, true );
711-
712- std::unique_ptr<Lock> slow_lock = group->lock (1 , 1 , [] {});
713- CHECK (slow_lock->lock () == true );
714-
715- auto slot = Group::TestAccess::park_publish (*group, make_test_task (), nullptr , false );
716- Group::TestAccess::park_reconcile (*group, slot);
717-
718- WHEN (" The slow lock releases while a fast waiter has already reserved a slot" ) {
719- slow_lock.reset ();
720-
721- THEN (" The committed waiter is drained and tokens return to concurrency" ) {
722- auto captured = Group::TestAccess::take_captured_drains (*group);
723- REQUIRE (captured.size () == 1 );
724- captured.front ().lock .reset ();
725-
726- CHECK (Group::TestAccess::tokens (*group) == group->descriptor ->concurrency );
727- }
728- }
729- }
730- }
731-
732- SCENARIO (" Park reconcile with a free token drains an earlier uncounted waiter" ) {
733- GIVEN (" A group with spare tokens and two parked fast waiters" ) {
734- auto group = make_group (2 );
735-
736- Group::TestAccess::set_capture_drains (*group, true );
737-
738- Group::TestAccess::park_publish (*group, make_test_task (), nullptr , false );
739- auto slot2 = Group::TestAccess::park_publish (*group, make_test_task (), nullptr , false );
740-
741- WHEN (" The second waiter reconciles while the first is still uncounted" ) {
742- Group::TestAccess::park_reconcile (*group, slot2);
743-
744- THEN (" The first waiter is opportunistically drained" ) {
745- auto captured = Group::TestAccess::take_captured_drains (*group);
746- REQUIRE (captured.size () == 1 );
747- captured.front ().lock .reset ();
748- }
749- }
750- }
751- }
752-
753698 SCENARIO (" try_submit parks while slow-path waiters hold priority" ) {
754699 GIVEN (" A group whose sole token is held by a slow-path lock" ) {
755700 auto scheduler = std::make_unique<Scheduler>(1 );
@@ -777,13 +722,13 @@ namespace threading {
777722
778723 SCENARIO (" try_acquire_running_lock returns nullptr when every token is in use" ) {
779724 GIVEN (" A group with one token acquired via the fast path" ) {
780- auto group = make_group (1 );
781- auto running = Group::TestAccess:: try_acquire_running_lock (*group );
725+ auto group = make_group (1 );
726+ auto running = group-> try_acquire_running_lock ();
782727 REQUIRE (running != nullptr );
783728
784729 WHEN (" Another fast-path acquisition is attempted" ) {
785730 THEN (" No second token is available" ) {
786- CHECK (Group::TestAccess:: try_acquire_running_lock (*group ) == nullptr );
731+ CHECK (group-> try_acquire_running_lock () == nullptr );
787732 }
788733 }
789734 }
0 commit comments