diff --git a/CHANGELOG.md b/CHANGELOG.md index 1bf5cc85..178a628d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -56,6 +56,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Changed +- **Allocation-free ready semaphore acquires**: Non-cancellable acquires now + defer shared wake-state allocation until the ready permit check fails. + Entering `await_suspend` allocates before taking the semaphore queue lock, + even if its locked recheck consumes a released permit. This preserves permit + accounting and dequeue-versus-frame-destruction safety. Token-aware acquires + retain eager cancellation state (#1039). - **Allocation-free uncontended mutex locks**: Non-cancellable locks now defer shared wake-state allocation until both the initial acquisition and the suspension-entry recheck observe contention. Truly parked waiters retain diff --git a/examples/microbench.cpp b/examples/microbench.cpp index 74362320..04256e1f 100644 --- a/examples/microbench.cpp +++ b/examples/microbench.cpp @@ -3,6 +3,7 @@ #include #include #include +#include #include #include #include @@ -47,6 +48,22 @@ coro::task mutex_handoffs(sync::mutex& mutex, size_t iterations) { } } +coro::task ready_semaphore_acquires(sync::semaphore& semaphore, + size_t iterations) { + for (size_t i = 0; i < iterations; ++i) { + co_await semaphore.acquire(); + semaphore.release(); + } +} + +coro::task semaphore_handoffs(sync::semaphore& semaphore, + size_t iterations) { + for (size_t i = 0; i < iterations; ++i) { + co_await semaphore.acquire(); + semaphore.release(); + } +} + int main() { log::logger::instance().set_level(log::level::error); @@ -217,7 +234,105 @@ int main() { << " ns/handoff" << std::endl; } - // 7. Measure MPSC push only (no scheduler overhead) + // 7. Measure ready semaphore acquire/release. No waiter is published, so + // release() does not reserve storage for a wake-state vector. One + // long-lived frame keeps coroutine construction outside the timed loop. + { + constexpr size_t acquire_iterations = 1000000; + sync::semaphore semaphore(1); + auto acquires = + ready_semaphore_acquires(semaphore, acquire_iterations); + auto handle = coro::detail::task_access::handle(acquires); + + auto start = high_resolution_clock::now(); + { + coro::detail::frame_context_scope frame_scope( + std::addressof(handle.promise())); + handle.resume(); + } + auto end = high_resolution_clock::now(); + if (!handle.done() || semaphore.count() != 1) { + std::abort(); + } + auto ns = duration_cast(end - start).count(); + + std::cout << "Ready semaphore acquire/release: " + << (static_cast(ns) / acquire_iterations) + << " ns/iteration" << std::endl; + } + + // 8. Isolate construction, publication, and unlinking of a parked + // semaphore waiter without release()'s existing wake-vector allocation. + // This is a mechanism diagnostic rather than an application throughput + // benchmark. + { + constexpr size_t park_iterations = 200000; + sync::semaphore semaphore(0); + + auto start = high_resolution_clock::now(); + for (size_t i = 0; i < park_iterations; ++i) { + auto waiter = semaphore.acquire(); + if (waiter.await_ready() || + !waiter.await_suspend(std::noop_coroutine())) { + std::abort(); + } + } + auto end = high_resolution_clock::now(); + if (semaphore.count() != 0) { + std::abort(); + } + auto ns = duration_cast(end - start).count(); + + std::cout << "Semaphore park/unlink diagnostic: " + << (static_cast(ns) / park_iterations) + << " ns/iteration" << std::endl; + } + + // 9. Measure forced permit handoff between two long-lived frames. Each + // release of an already-published waiter also exercises release()'s + // existing one-element wake-vector allocation, so the park/unlink result + // above remains the allocation-codegen control. + { + constexpr size_t handoff_iterations_per_task = 100000; + constexpr size_t total_handoffs = handoff_iterations_per_task * 2; + sync::semaphore semaphore(0); + + auto first = + semaphore_handoffs(semaphore, handoff_iterations_per_task); + auto second = + semaphore_handoffs(semaphore, handoff_iterations_per_task); + auto first_handle = coro::detail::task_access::handle(first); + auto second_handle = coro::detail::task_access::handle(second); + { + coro::detail::frame_context_scope frame_scope( + std::addressof(first_handle.promise())); + first_handle.resume(); + } + { + coro::detail::frame_context_scope frame_scope( + std::addressof(second_handle.promise())); + second_handle.resume(); + } + if (first_handle.done() || second_handle.done()) { + std::abort(); + } + + auto start = high_resolution_clock::now(); + semaphore.release(); + auto end = high_resolution_clock::now(); + if (!first_handle.done() || !second_handle.done() || + semaphore.count() != 1 || !semaphore.try_acquire() || + semaphore.count() != 0) { + std::abort(); + } + auto ns = duration_cast(end - start).count(); + + std::cout << "Forced two-task semaphore handoff: " + << (static_cast(ns) / total_handoffs) + << " ns/handoff" << std::endl; + } + + // 10. Measure MPSC push only (no scheduler overhead) { runtime::mpsc_queue queue; @@ -234,7 +349,7 @@ int main() { while (queue.pop()) {} } - // 8. Measure Chase-Lev push only + // 11. Measure Chase-Lev push only { runtime::chase_lev_deque queue; @@ -251,7 +366,7 @@ int main() { while (queue.pop()) {} } - // 9. Compare atomic RMW with single-writer snapshot publication + // 12. Compare atomic RMW with single-writer snapshot publication { std::atomic published{0}; @@ -283,7 +398,7 @@ int main() { << " ns/update" << std::endl; } - // 10. Compare exact timestamps with the disabled diagnostic fast path + // 13. Compare exact timestamps with the disabled diagnostic fast path { std::atomic last_task_time{ steady_clock::now()}; @@ -318,7 +433,7 @@ int main() { << " ns/update" << std::endl; } - // 11. Measure atomic fence alone + // 14. Measure atomic fence alone { auto start = high_resolution_clock::now(); for (int i = 0; i < N; ++i) { @@ -330,7 +445,7 @@ int main() { std::cout << "Atomic release fence: " << (ns / N) << " ns" << std::endl; } - // 12. Measure eventfd write + // 15. Measure eventfd write { int fd = eventfd(0, EFD_NONBLOCK); uint64_t val = 1; @@ -346,7 +461,7 @@ int main() { close(fd); } - // 13. Full spawn path (with running scheduler) - includes alloc + spawn + // 16. Full spawn path (with running scheduler) - includes alloc + spawn { runtime::scheduler sched(4); sched.start(); @@ -368,7 +483,7 @@ int main() { sched.shutdown(); } - // 14. Measure warmed-up worker overhead + // 17. Measure warmed-up worker overhead { runtime::scheduler sched(4); sched.start(); diff --git a/include/elio/sync/semaphore.hpp b/include/elio/sync/semaphore.hpp index 629baeb3..a5bb7de9 100644 --- a/include/elio/sync/semaphore.hpp +++ b/include/elio/sync/semaphore.hpp @@ -4,7 +4,10 @@ #include #include #include +#include +#include #include +#include #include #include #include @@ -47,18 +50,16 @@ class semaphore { class acquire_waiter : public elio::detail::intrusive_list_node { public: explicit acquire_waiter(semaphore& s) - : sem_(s) - , wake_state_(detail::make_wake_state()) {} + : sem_(s) {} acquire_waiter(semaphore& s, bool cancellable) : sem_(s) - , wake_state_(detail::make_wake_state()) - , cancellable_(cancellable) {} + , waiter_state_(cancellable) {} ~acquire_waiter() { // Fast path: if we never suspended, we were never enqueued, // so no wake function could hold a reference to us. - if (!suspended_) return; + if (!waiter_state_.suspended) return; detail::wake_state_ptr to_schedule; // Slow path: acquire mutex to prevent race with release() @@ -66,13 +67,14 @@ class semaphore { std::lock_guard guard(sem_.mutex_); if (this->is_linked()) { sem_.waiters_.remove(this); - detail::cancel_wake_state(wake_state_); - } else if (grant_pending_ && !resumed_) { - detail::cancel_wake_state(wake_state_); - grant_pending_ = false; + detail::cancel_wake_state(waiter_state_.wake()); + } else if (waiter_state_.grant_pending && + !waiter_state_.resumed) { + detail::cancel_wake_state(waiter_state_.wake()); + waiter_state_.grant_pending = false; to_schedule = sem_.recover_cancelled_handoff_locked(); } else { - detail::cancel_wake_state(wake_state_); + detail::cancel_wake_state(waiter_state_.wake()); } } @@ -82,17 +84,17 @@ class semaphore { } bool await_ready_impl() const { - if (!cancellable_) { + if (!waiter_state_.cancellable) { return sem_.try_acquire(); } - if (wake_state_->was_cancelled()) { + if (waiter_state_->was_cancelled()) { return true; } if (!sem_.try_acquire()) { return false; } - if (detail::claim_wake_state(wake_state_) != + if (detail::claim_wake_state(waiter_state_.wake()) != detail::wake_action::rejected) { return true; } @@ -102,35 +104,53 @@ class semaphore { return true; } - bool await_suspend_impl(std::coroutine_handle<> awaiter) noexcept { - if (!cancellable_) { - std::lock_guard guard(sem_.mutex_); - if (sem_.count_ > 0) { - --sem_.count_; - return false; - } + bool await_suspend_impl(std::coroutine_handle<> awaiter) { + if (waiter_state_.cancellable) { + return await_suspend_cancellable_impl(awaiter); + } + return await_suspend_non_cancellable_impl(awaiter); + } + + bool await_suspend_non_cancellable_impl( + std::coroutine_handle<> awaiter) { + assert(!waiter_state_.cancellable); - wake_state_->set_handle(awaiter); - sem_.waiters_.push_back(this); - suspended_ = true; + // A release may outlive the coroutine frame after dequeuing this + // waiter, so an acquire that reaches suspension still needs + // independent shared ownership. Allocate before taking the queue + // lock so failure leaves both the permit count and queue unchanged. + waiter_state_.emplace(); + + std::lock_guard guard(sem_.mutex_); + if (sem_.count_ > 0) { + --sem_.count_; + return false; + } + + waiter_state_->set_handle(awaiter); + sem_.waiters_.push_back(this); + waiter_state_.suspended = true; #ifdef ELIO_RUNTIME_TEST_HOOKS - detail::semaphore_waiter_publications_for_test.fetch_add( - 1, std::memory_order_release); + detail::semaphore_waiter_publications_for_test.fetch_add( + 1, std::memory_order_release); #endif - return true; - } + return true; + } + bool await_suspend_cancellable_impl( + std::coroutine_handle<> awaiter) noexcept { + assert(waiter_state_.cancellable); detail::wake_state_ptr to_schedule; { std::lock_guard guard(sem_.mutex_); - if (wake_state_->was_cancelled()) { + if (waiter_state_->was_cancelled()) { return false; } if (sem_.count_ > 0) { --sem_.count_; - if (detail::claim_wake_state(wake_state_) != + if (detail::claim_wake_state(waiter_state_.wake()) != detail::wake_action::rejected) { return false; } @@ -139,22 +159,22 @@ class semaphore { // another live waiter, or restore it to the count. to_schedule = sem_.recover_cancelled_handoff_locked(); } else { - if (!wake_state_->set_handle_blocked(awaiter)) { + if (!waiter_state_->set_handle_blocked(awaiter)) { return false; } sem_.waiters_.push_back(this); - suspended_ = true; + waiter_state_.suspended = true; #ifdef ELIO_RUNTIME_TEST_HOOKS detail::semaphore_waiter_publications_for_test.fetch_add( 1, std::memory_order_release); #endif - if (wake_state_->unblock_after_publish()) { + if (waiter_state_->unblock_after_publish()) { return true; } sem_.waiters_.remove(this); - suspended_ = false; + waiter_state_.suspended = false; } } @@ -165,42 +185,109 @@ class semaphore { } coro::cancel_result await_resume_impl() noexcept { - if (!cancellable_) { - resumed_ = true; - grant_pending_ = false; - suspended_ = false; + if (!waiter_state_.cancellable) { + waiter_state_.resumed = true; + waiter_state_.grant_pending = false; + waiter_state_.suspended = false; return coro::cancel_result::completed; } - if (wake_state_->was_cancelled()) { - if (suspended_) { + if (waiter_state_->was_cancelled()) { + if (waiter_state_.suspended) { std::lock_guard guard(sem_.mutex_); if (this->is_linked()) { sem_.waiters_.remove(this); } - suspended_ = false; + waiter_state_.suspended = false; } return coro::cancel_result::cancelled; } - resumed_ = true; - grant_pending_ = false; - suspended_ = false; + waiter_state_.resumed = true; + waiter_state_.grant_pending = false; + waiter_state_.suspended = false; return coro::cancel_result::completed; } protected: const detail::wake_state_ptr& cancellation_wake_state() const noexcept { - return wake_state_; + return waiter_state_.wake(); } private: + class waiter_state { + public: + waiter_state() noexcept = default; + + explicit waiter_state(bool is_cancellable) + : cancellable(is_cancellable) { + if (!is_cancellable) return; + + ::new (static_cast(storage_)) + detail::wake_state_ptr(detail::make_wake_state()); + // Publish engagement only after placement construction + // succeeds; a throwing allocation leaves no active object. + engaged = true; + } + + ~waiter_state() { + if (engaged) { + std::destroy_at(std::addressof(storage_ref())); + } + } + + waiter_state(const waiter_state&) = delete; + waiter_state& operator=(const waiter_state&) = delete; + waiter_state(waiter_state&&) = delete; + waiter_state& operator=(waiter_state&&) = delete; + + void emplace() { + assert(!engaged); + ::new (static_cast(storage_)) + detail::wake_state_ptr(detail::make_wake_state()); + // A failed allocation leaves the slot disengaged. + engaged = true; + } + + [[nodiscard]] const detail::wake_state_ptr& wake() const noexcept { + assert(engaged); + return storage_ref(); + } + + [[nodiscard]] detail::wake_state* operator->() const noexcept { + return wake().get(); + } + + bool engaged = false; + bool cancellable = false; + bool suspended = false; + bool resumed = false; + bool grant_pending = false; + + private: + detail::wake_state_ptr& storage_ref() noexcept { + return *std::launder(storage_ptr()); + } + + const detail::wake_state_ptr& storage_ref() const noexcept { + return *std::launder(storage_ptr()); + } + + detail::wake_state_ptr* storage_ptr() noexcept { + return reinterpret_cast(storage_); + } + + const detail::wake_state_ptr* storage_ptr() const noexcept { + return reinterpret_cast( + storage_); + } + + alignas(detail::wake_state_ptr) + std::byte storage_[sizeof(detail::wake_state_ptr)]; + }; + semaphore& sem_; - detail::wake_state_ptr wake_state_; - bool cancellable_ = false; - bool suspended_ = false; // True if enqueued in waiters_ - bool resumed_ = false; // True after a popped waiter resumes normally - bool grant_pending_ = false; // True after release() transfers a permit + waiter_state waiter_state_; friend class semaphore; }; @@ -219,6 +306,11 @@ class semaphore { cancel_registration_.unregister(); } + bool await_suspend_impl( + std::coroutine_handle<> awaiter) noexcept { + return await_suspend_cancellable_impl(awaiter); + } + coro::cancel_result await_resume_cancellable() noexcept { cancel_registration_.unregister(); return await_resume_impl(); @@ -233,8 +325,8 @@ class semaphore { explicit acquire_awaitable(semaphore& s) : waiter_(s) {} bool await_ready() const noexcept { return waiter_.await_ready_impl(); } - bool await_suspend(std::coroutine_handle<> awaiter) noexcept { - return waiter_.await_suspend_impl(awaiter); + bool await_suspend(std::coroutine_handle<> awaiter) { + return waiter_.await_suspend_non_cancellable_impl(awaiter); } void await_resume() noexcept { (void)waiter_.await_resume_impl(); @@ -262,7 +354,9 @@ class semaphore { cancellable_acquire_waiter waiter_; }; - /// Acquire (decrement) the semaphore. + /// Acquire (decrement) the semaphore. A ready no-token acquire completes + /// without allocation. Suspension allocates shared wake state and can + /// propagate std::bad_alloc from await_suspend(). auto acquire() { return acquire_awaitable(*this); } @@ -325,14 +419,14 @@ class semaphore { detail::wake_state_ptr claim_waiter_locked() noexcept { while (!waiters_.empty()) { auto* waiter = waiters_.pop_front(); - if (waiter->cancellable_) { - if (detail::claim_wake_state(waiter->wake_state_) == + if (waiter->waiter_state_.cancellable) { + if (detail::claim_wake_state(waiter->waiter_state_.wake()) == detail::wake_action::rejected) { continue; } } - waiter->grant_pending_ = true; - return waiter->wake_state_; + waiter->waiter_state_.grant_pending = true; + return waiter->waiter_state_.wake(); } return nullptr; diff --git a/tests/unit/test_sync_cancellation.cpp b/tests/unit/test_sync_cancellation.cpp index 6a44b3f5..a263dca4 100644 --- a/tests/unit/test_sync_cancellation.cpp +++ b/tests/unit/test_sync_cancellation.cpp @@ -417,6 +417,240 @@ TEST_CASE("mutex fast paths defer wake-state allocation", } } +TEST_CASE("semaphore ready path defers wake-state allocation", + "[sync][semaphore][allocation]") { + auto& allocations = + elio::sync::detail::wake_state_allocations_for_test; + auto& publications = + elio::sync::detail::semaphore_waiter_publications_for_test; + auto& fail_next = + elio::sync::detail::fail_next_wake_state_allocation_for_test; + fail_next.store(false, std::memory_order_relaxed); + allocations.store(0, std::memory_order_relaxed); + publications.store(0, std::memory_order_relaxed); + + SECTION("ready acquire") { + semaphore sem(1); + auto waiter = sem.acquire(); + static_assert(!noexcept( + waiter.await_suspend(std::noop_coroutine()))); + + REQUIRE(waiter.await_ready()); + waiter.await_resume(); + + REQUIRE(allocations.load(std::memory_order_relaxed) == 0); + REQUIRE(publications.load(std::memory_order_relaxed) == 0); + REQUIRE(sem.count() == 0); + } + + SECTION("parked acquire") { + semaphore sem(0); + auto waiter = sem.acquire(); + + REQUIRE_FALSE(waiter.await_ready()); + REQUIRE(waiter.await_suspend(std::noop_coroutine())); + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE(publications.load(std::memory_order_relaxed) == 1); + + sem.release(); + REQUIRE(sem.count() == 0); + waiter.await_resume(); + REQUIRE(sem.count() == 0); + } + + SECTION("explicit non-cancellable waiter defers its state") { + semaphore ready_sem(1); + { + semaphore::acquire_waiter waiter(ready_sem, false); + + REQUIRE(allocations.load(std::memory_order_relaxed) == 0); + REQUIRE(waiter.await_ready_impl()); + REQUIRE(waiter.await_resume_impl() == cancel_result::completed); + REQUIRE(allocations.load(std::memory_order_relaxed) == 0); + REQUIRE(publications.load(std::memory_order_relaxed) == 0); + } + REQUIRE(ready_sem.count() == 0); + + semaphore parked_sem(0); + { + semaphore::acquire_waiter waiter(parked_sem, false); + + REQUIRE(allocations.load(std::memory_order_relaxed) == 0); + REQUIRE_FALSE(waiter.await_ready_impl()); + REQUIRE(waiter.await_suspend_impl(std::noop_coroutine())); + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE(publications.load(std::memory_order_relaxed) == 1); + + parked_sem.release(); + REQUIRE(waiter.await_resume_impl() == cancel_result::completed); + } + REQUIRE(parked_sem.count() == 0); + } + + SECTION("explicit cancellable base waiter keeps suspend compatibility") { + semaphore sem(0); + semaphore::acquire_waiter waiter(sem, true); + static_assert(!noexcept( + waiter.await_suspend_impl(std::noop_coroutine()))); + + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE_FALSE(waiter.await_ready_impl()); + REQUIRE(waiter.await_suspend_impl(std::noop_coroutine())); + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE(publications.load(std::memory_order_relaxed) == 1); + + sem.release(); + REQUIRE(waiter.await_resume_impl() == cancel_result::completed); + REQUIRE(sem.count() == 0); + } + + SECTION("direct cancellable waiter keeps noexcept suspend compatibility") { + semaphore sem(0); + cancel_source source; + semaphore::cancellable_acquire_waiter waiter( + sem, source.get_token()); + static_assert(noexcept( + waiter.await_suspend_impl(std::noop_coroutine()))); + + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE_FALSE(waiter.await_ready_impl()); + REQUIRE(waiter.await_suspend_impl(std::noop_coroutine())); + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE(publications.load(std::memory_order_relaxed) == 1); + + sem.release(); + REQUIRE(waiter.await_resume_cancellable() == + cancel_result::completed); + REQUIRE(sem.count() == 0); + } + + SECTION("release between ready and suspend") { + semaphore sem(0); + auto waiter = sem.acquire(); + + REQUIRE_FALSE(waiter.await_ready()); + sem.release(); + REQUIRE_FALSE(waiter.await_suspend(std::noop_coroutine())); + waiter.await_resume(); + + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE(publications.load(std::memory_order_relaxed) == 0); + REQUIRE(sem.count() == 0); + } + + SECTION("another acquire may consume an unpublished permit") { + semaphore sem(0); + auto first = sem.acquire(); + REQUIRE_FALSE(first.await_ready()); + + sem.release(); + auto next = sem.acquire(); + REQUIRE(next.await_ready()); + next.await_resume(); + REQUIRE(sem.count() == 0); + + REQUIRE(first.await_suspend(std::noop_coroutine())); + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE(publications.load(std::memory_order_relaxed) == 1); + + sem.release(); + first.await_resume(); + REQUIRE(sem.count() == 0); + } + + SECTION("allocation failure leaves an empty semaphore unchanged") { + semaphore sem(0); + { + auto waiter = sem.acquire(); + REQUIRE_FALSE(waiter.await_ready()); + fail_next.store(true, std::memory_order_release); + REQUIRE_THROWS_AS( + waiter.await_suspend(std::noop_coroutine()), std::bad_alloc); + } + + REQUIRE_FALSE(fail_next.load(std::memory_order_acquire)); + REQUIRE(allocations.load(std::memory_order_relaxed) == 0); + REQUIRE(publications.load(std::memory_order_relaxed) == 0); + REQUIRE(sem.count() == 0); + + sem.release(); + auto next = sem.acquire(); + REQUIRE(next.await_ready()); + next.await_resume(); + REQUIRE(sem.count() == 0); + } + + SECTION("allocation failure preserves a concurrently released permit") { + semaphore sem(0); + { + auto waiter = sem.acquire(); + REQUIRE_FALSE(waiter.await_ready()); + sem.release(); + fail_next.store(true, std::memory_order_release); + REQUIRE_THROWS_AS( + waiter.await_suspend(std::noop_coroutine()), std::bad_alloc); + } + + REQUIRE_FALSE(fail_next.load(std::memory_order_acquire)); + REQUIRE(allocations.load(std::memory_order_relaxed) == 0); + REQUIRE(publications.load(std::memory_order_relaxed) == 0); + REQUIRE(sem.count() == 1); + + auto next = sem.acquire(); + REQUIRE(next.await_ready()); + next.await_resume(); + REQUIRE(sem.count() == 0); + } + + SECTION("cancellable ready acquire keeps eager arbitration state") { + semaphore sem(1); + cancel_source source; + auto waiter = sem.acquire(source.get_token()); + static_assert(noexcept( + waiter.await_suspend(std::noop_coroutine()))); + + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE(waiter.await_ready()); + REQUIRE(waiter.await_resume() == cancel_result::completed); + REQUIRE(publications.load(std::memory_order_relaxed) == 0); + REQUIRE(sem.count() == 0); + } + + SECTION("cancellable parked acquire keeps eager arbitration state") { + semaphore sem(0); + cancel_source source; + auto waiter = sem.acquire(source.get_token()); + + REQUIRE(allocations.load(std::memory_order_relaxed) == 1); + REQUIRE_FALSE(waiter.await_ready()); + REQUIRE(waiter.await_suspend(std::noop_coroutine())); + REQUIRE(publications.load(std::memory_order_relaxed) == 1); + + sem.release(); + REQUIRE(waiter.await_resume() == cancel_result::completed); + REQUIRE(sem.count() == 0); + } + + SECTION("cancellable construction failure leaves semaphore unchanged") { + semaphore sem(1); + cancel_source source; + fail_next.store(true, std::memory_order_release); + + REQUIRE_THROWS_AS( + (void)sem.acquire(source.get_token()), std::bad_alloc); + + REQUIRE_FALSE(fail_next.load(std::memory_order_acquire)); + REQUIRE(allocations.load(std::memory_order_relaxed) == 0); + REQUIRE(publications.load(std::memory_order_relaxed) == 0); + REQUIRE(sem.count() == 1); + + auto next = sem.acquire(); + REQUIRE(next.await_ready()); + next.await_resume(); + REQUIRE(sem.count() == 0); + } +} + TEST_CASE("runtime cancellation wakes basic sync waits", "[sync][cancellation][cancel_token][runtime]") { scheduler sched(2); @@ -1720,6 +1954,7 @@ TEST_CASE("semaphore release does not schedule a waiter destroyed after dequeue" sem.release(2); destroyer.join(); + REQUIRE(sem.count() == 1); } TEST_CASE("shared_mutex reader wake does not schedule a waiter destroyed after dequeue", diff --git a/wiki/API-Reference.md b/wiki/API-Reference.md index c7f8fd62..820fe7df 100644 --- a/wiki/API-Reference.md +++ b/wiki/API-Reference.md @@ -3537,6 +3537,17 @@ public: }; ``` +The no-token `acquire()` completes without wake-state allocation when a permit +is immediately available. If the ready check observes no permit, suspension +creates independently owned wake state before entering the semaphore queue +critical section. `await_suspend()` may therefore propagate `std::bad_alloc`, +but allocation failure does not consume a permit or publish a waiter. If a +permit released between `await_ready()` and `await_suspend()` remains +available, the locked recheck consumes it; that race can perform one unused +allocation without losing the permit. Token-aware acquires create arbitration +state eagerly so cancellation can race permit handoff with exactly one terminal +result. + --- ## Timers (`elio::time`) diff --git a/wiki/Performance-Tuning.md b/wiki/Performance-Tuning.md index 386da46b..1b9f927d 100644 --- a/wiki/Performance-Tuning.md +++ b/wiki/Performance-Tuning.md @@ -468,6 +468,30 @@ if (mtx.try_lock()) { } ``` +### Semaphore Performance + +A non-cancellable `semaphore::acquire()` consumes an immediately available +permit without allocating shared wake state. If no permit is ready, the +awaiter creates independently owned state before taking the queue lock and +rechecks the count under that lock before parking. This preserves permit +accounting when `release()` races the transition from `await_ready()` to +`await_suspend()`, although that race can create one state that is not used for +suspension. Token-aware acquires retain eager allocation for cancellation and +permit-handoff arbitration. + +For steady-state permit guards, pair each successful acquire with the matching +release: + +```cpp +sync::semaphore permits(1); + +coro::task use_permit() { + co_await permits.acquire(); + // ... bounded work ... + permits.release(); +} +``` + ### Reader-Writer Lock For read-heavy workloads: