From 6ec864ad48a154030a504b9908090e9d6e6ff760 Mon Sep 17 00:00:00 2001 From: Coldwings Date: Thu, 13 Aug 2026 17:38:17 +0800 Subject: [PATCH 1/2] fix(coro): protect selected completion wakes --- CHANGELOG.md | 5 + .../elio/coro/detail/completion_waiter.hpp | 141 ++++++++++++- include/elio/coro/task.hpp | 3 +- include/elio/coro/task_group.hpp | 24 ++- include/elio/coro/task_handle.hpp | 6 +- include/elio/rdma_ibverbs/endpoint.hpp | 3 +- include/elio/sync/object_cache.hpp | 3 +- tests/unit/test_object_cache.cpp | 88 ++++++++ tests/unit/test_task.cpp | 193 ++++++++++++++++++ tests/unit/test_task_group.cpp | 56 +++++ wiki/API-Contracts.md | 2 +- 11 files changed, 503 insertions(+), 21 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 178a628d..6226909e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -17,6 +17,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- **Completion waiter wake lifetime**: Join handles, task handles, task groups, + object-cache release waits, and RDMA pump-exit waits now retain a slot-owned + selected-wake lease between dequeue and scheduling. Destroying an awaiting + coroutine before scheduling ownership is claimed abandons that wake instead + of leaving completion with a stale raw coroutine handle (#1049). - **Autoscaler example task lifetime**: Each load phase now retains join handles and drains every submitted task before releasing its phase-local completion counter or beginning the low-load phase. The shorter workload and CI smoke diff --git a/include/elio/coro/detail/completion_waiter.hpp b/include/elio/coro/detail/completion_waiter.hpp index 06b6b34c..382c9a6b 100644 --- a/include/elio/coro/detail/completion_waiter.hpp +++ b/include/elio/coro/detail/completion_waiter.hpp @@ -1,6 +1,8 @@ #pragma once #include +#include +#include #include #include #include @@ -9,6 +11,50 @@ namespace elio::coro::detail { class completion_waiter_slot; +#ifdef ELIO_RUNTIME_TEST_HOOKS +inline std::atomic pause_before_completion_wake_claim_for_test{false}; +inline std::atomic completion_wake_claim_paused_for_test{false}; +#endif + +/// Move-only ownership of a completion wake selected by a producer. +/// +/// The slot, rather than the awaiting coroutine frame, retains the selected +/// handle until claim() transfers scheduling ownership. This lets waiter +/// destruction abandon a wake after dequeue without leaving the producer with +/// a stale raw coroutine handle. The slot owner must outlive every unclaimed +/// lease returned by that slot. +class completion_wake_lease { +public: + completion_wake_lease() noexcept = default; + ~completion_wake_lease(); + + completion_wake_lease(const completion_wake_lease&) = delete; + completion_wake_lease& operator=(const completion_wake_lease&) = delete; + + completion_wake_lease(completion_wake_lease&& other) noexcept; + completion_wake_lease& operator=(completion_wake_lease&& other) noexcept; + + [[nodiscard]] explicit operator bool() const noexcept { + return slot_ != nullptr; + } + + /// Claim scheduling ownership of the selected handle. Returns an empty + /// handle if waiter destruction abandoned this generation first. + [[nodiscard]] std::coroutine_handle<> claim() noexcept; + +private: + completion_wake_lease(completion_waiter_slot& slot, + size_t generation) noexcept + : slot_(&slot), generation_(generation) {} + + void reset() noexcept; + + completion_waiter_slot* slot_ = nullptr; + size_t generation_ = 0; + + friend class completion_waiter_slot; +}; + /// Awaiter-owned registration for a single completion waiter. class completion_waiter { public: @@ -42,7 +88,7 @@ class completion_waiter_slot { ~completion_waiter_slot() { std::lock_guard lock(mutex_); - assert(waiter_ == nullptr && + assert(waiter_ == nullptr && !selected() && "completion waiter slot destroyed with a pending waiter"); } @@ -68,22 +114,53 @@ class completion_waiter_slot { return true; } - std::coroutine_handle<> take() noexcept { + completion_wake_lease take() noexcept { std::lock_guard lock(mutex_); - if (!waiter_) { + if (!waiter_ || selected()) { return {}; } - auto handle = waiter_->handle_; + selected_handle_ = waiter_->handle_; + ++generation_; waiter_->handle_ = {}; + return completion_wake_lease(*this, generation_); + } + +private: + [[nodiscard]] bool selected() const noexcept { + return (generation_ & size_t{1}) != 0; + } + + std::coroutine_handle<> claim(size_t generation) noexcept { + std::lock_guard lock(mutex_); + if (!selected() || generation_ != generation) { + return {}; + } + + auto handle = selected_handle_; + ++generation_; waiter_ = nullptr; + selected_handle_ = {}; return handle; } -private: + void abandon(size_t generation) noexcept { + std::lock_guard lock(mutex_); + if (!selected() || generation_ != generation) { + return; + } + ++generation_; + waiter_ = nullptr; + selected_handle_ = {}; + } + void remove(completion_waiter& waiter) noexcept { std::lock_guard lock(mutex_); if (waiter_ == &waiter) { + if (selected()) { + ++generation_; + selected_handle_ = {}; + } waiter_ = nullptr; } waiter.handle_ = {}; @@ -104,10 +181,64 @@ class completion_waiter_slot { std::mutex mutex_; completion_waiter* waiter_ = nullptr; + std::coroutine_handle<> selected_handle_{}; + // Odd generations hold a selected, cancelable wake. Claim or abandonment + // advances to the next even generation before the slot can be reused. + size_t generation_ = 0; friend class completion_waiter; + friend class completion_wake_lease; }; +inline completion_wake_lease::~completion_wake_lease() { + reset(); +} + +inline completion_wake_lease::completion_wake_lease( + completion_wake_lease&& other) noexcept + : slot_(std::exchange(other.slot_, nullptr)) + , generation_(std::exchange(other.generation_, 0)) {} + +inline completion_wake_lease& completion_wake_lease::operator=( + completion_wake_lease&& other) noexcept { + if (this != &other) { + reset(); + slot_ = std::exchange(other.slot_, nullptr); + generation_ = std::exchange(other.generation_, 0); + } + return *this; +} + +inline std::coroutine_handle<> completion_wake_lease::claim() noexcept { +#ifdef ELIO_RUNTIME_TEST_HOOKS + if (pause_before_completion_wake_claim_for_test.load( + std::memory_order_acquire)) { + completion_wake_claim_paused_for_test.store( + true, std::memory_order_release); + completion_wake_claim_paused_for_test.notify_all(); + while (pause_before_completion_wake_claim_for_test.load( + std::memory_order_acquire)) { + pause_before_completion_wake_claim_for_test.wait( + true, std::memory_order_acquire); + } + completion_wake_claim_paused_for_test.store( + false, std::memory_order_release); + completion_wake_claim_paused_for_test.notify_all(); + } +#endif + auto* slot = std::exchange(slot_, nullptr); + const auto generation = std::exchange(generation_, 0); + return slot ? slot->claim(generation) : std::coroutine_handle<>{}; +} + +inline void completion_wake_lease::reset() noexcept { + auto* slot = std::exchange(slot_, nullptr); + const auto generation = std::exchange(generation_, 0); + if (slot) { + slot->abandon(generation); + } +} + inline completion_waiter::~completion_waiter() { if (slot_) { slot_->remove(*this); diff --git a/include/elio/coro/task.hpp b/include/elio/coro/task.hpp index 25301095..f645150f 100644 --- a/include/elio/coro/task.hpp +++ b/include/elio/coro/task.hpp @@ -143,7 +143,8 @@ struct join_state_base { void complete() { completed_.store(true, std::memory_order_release); - auto waiter = waiter_.take(); + auto wake = waiter_.take(); + auto waiter = wake.claim(); if (waiter) { runtime::schedule_handle(waiter); } diff --git a/include/elio/coro/task_group.hpp b/include/elio/coro/task_group.hpp index 6fe8a04a..3354d2d8 100644 --- a/include/elio/coro/task_group.hpp +++ b/include/elio/coro/task_group.hpp @@ -76,19 +76,23 @@ inline std::atomic task_group_join_observed_pending_for_test{false}; struct task_group_child_completion final { runtime::scheduler* scheduler = nullptr; - std::coroutine_handle<> waiter{}; + completion_wake_lease wake{}; }; inline void resume_task_group_join_waiter( task_group_child_completion completion) noexcept { - if (completion.scheduler->try_schedule(completion.waiter)) { + auto waiter = completion.wake.claim(); + if (!waiter) { + return; + } + if (completion.scheduler->try_schedule(waiter)) { return; } if (is_scheduler_worker(completion.scheduler)) { - auto* promise = get_promise_base(completion.waiter.address()); + auto* promise = get_promise_base(waiter.address()); detail::frame_context_scope frame_scope(promise); - completion.waiter.resume(); + waiter.resume(); return; } @@ -97,7 +101,7 @@ inline void resume_task_group_join_waiter( // lifetime contract that the scheduler runs until the group drains. while (completion.scheduler->is_running()) { std::this_thread::yield(); - if (completion.scheduler->try_schedule(completion.waiter)) { + if (completion.scheduler->try_schedule(waiter)) { return; } } @@ -162,9 +166,9 @@ class task_group_completion_state final { } if (completed) { - auto waiter = all_done_waiter_.take(); - if (waiter) { - return {scheduler_, waiter}; + auto wake = all_done_waiter_.take(); + if (wake) { + return {scheduler_, std::move(wake)}; } } return {}; @@ -426,10 +430,10 @@ class task_group_child_registration final { } #endif auto completion = completion_state->child_finished(); - if (completion.waiter) { + if (completion.wake) { // Failure state ownership is already released before either // readiness or a registered join waiter can observe completion. - resume_task_group_join_waiter(completion); + resume_task_group_join_waiter(std::move(completion)); } } } diff --git a/include/elio/coro/task_handle.hpp b/include/elio/coro/task_handle.hpp index 13bf3106..b5b0bcaa 100644 --- a/include/elio/coro/task_handle.hpp +++ b/include/elio/coro/task_handle.hpp @@ -122,7 +122,8 @@ struct task_state { } void notify_waiter() { - auto waiter = waiter_.take(); + auto wake = waiter_.take(); + auto waiter = wake.claim(); if (waiter) { runtime::schedule_handle(waiter); } @@ -204,7 +205,8 @@ struct task_state { } void notify_waiter() { - auto waiter = waiter_.take(); + auto wake = waiter_.take(); + auto waiter = wake.claim(); if (waiter) { runtime::schedule_handle(waiter); } diff --git a/include/elio/rdma_ibverbs/endpoint.hpp b/include/elio/rdma_ibverbs/endpoint.hpp index 58cc61a3..51b3555c 100644 --- a/include/elio/rdma_ibverbs/endpoint.hpp +++ b/include/elio/rdma_ibverbs/endpoint.hpp @@ -151,7 +151,8 @@ class pump_exit_state { void mark_exited() noexcept { exited_.store(true, std::memory_order_release); - if (auto waiter = waiter_.take()) { + auto wake = waiter_.take(); + if (auto waiter = wake.claim()) { elio::runtime::schedule_handle(waiter); } } diff --git a/include/elio/sync/object_cache.hpp b/include/elio/sync/object_cache.hpp index e72d3461..33bfbea2 100644 --- a/include/elio/sync/object_cache.hpp +++ b/include/elio/sync/object_cache.hpp @@ -414,7 +414,8 @@ class object_cache { detail_oc::release_waiter_probes_for_test.fetch_add( 1, std::memory_order_relaxed); #endif - auto waiter = entry_->release_waiter_.take(); + auto wake = entry_->release_waiter_.take(); + auto waiter = wake.claim(); if (waiter) { runtime::schedule_handle(waiter); } diff --git a/tests/unit/test_object_cache.cpp b/tests/unit/test_object_cache.cpp index 94adef76..26a5c1b8 100644 --- a/tests/unit/test_object_cache.cpp +++ b/tests/unit/test_object_cache.cpp @@ -6,6 +6,7 @@ #include #include +#include #include #include #include @@ -642,6 +643,93 @@ TEST_CASE("object_cache release unregisters waiter when suspended release is des REQUIRE(release_continuations.load(std::memory_order_acquire) == 0); } +TEST_CASE("object_cache release skips a waiter destroyed after selection", + "[object_cache][release][cancellation][lifetime][regression]") { + using cache_type = object_cache; + using namespace elio::coro::detail; + + cache_type cache({.num_shards = 4}); + std::optional other; + std::atomic published{0}; + std::atomic release_continuations{0}; + + struct hook_context { + std::atomic* published; + } ctx{&published}; + + struct hook_guard { + hook_guard(void (*cb)(void*), void* ctx) { + detail_oc::release_waiter_published_hook.context.store( + ctx, std::memory_order_release); + detail_oc::release_waiter_published_hook.callback.store( + cb, std::memory_order_release); + } + + ~hook_guard() { + detail_oc::release_waiter_published_hook.callback.store( + nullptr, std::memory_order_release); + detail_oc::release_waiter_published_hook.context.store( + nullptr, std::memory_order_release); + } + } guard{ + [](void* raw) noexcept { + auto* c = static_cast(raw); + c->published->fetch_add(1, std::memory_order_acq_rel); + }, + &ctx}; + + auto waiter_task = [&]() -> task { + auto releaser = co_await cache.get("selected_key", []() -> task { + co_return 17; + }); + other.emplace(co_await cache.get( + "selected_key", []() -> task { co_return 19; })); + + auto owned = co_await releaser.release(); + release_continuations.fetch_add(1, std::memory_order_relaxed); + REQUIRE(owned != nullptr); + }; + + auto waiter = waiter_task(); + auto handle = task_access::release(std::move(waiter)); + handle.resume(); + REQUIRE(published.load(std::memory_order_acquire) == 1); + REQUIRE(other.has_value()); + REQUIRE(release_continuations.load(std::memory_order_acquire) == 0); + + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + pause_before_completion_wake_claim_for_test.store( + true, std::memory_order_release); + std::thread producer([&other] { other.reset(); }); + + const auto deadline = std::chrono::steady_clock::now() + + std::chrono::seconds(5); + while (!completion_wake_claim_paused_for_test.load( + std::memory_order_acquire) && + std::chrono::steady_clock::now() < deadline) { + std::this_thread::yield(); + } + const bool claim_paused = completion_wake_claim_paused_for_test.load( + std::memory_order_acquire); + if (claim_paused) { + handle.destroy(); + } + pause_before_completion_wake_claim_for_test.store( + false, std::memory_order_release); + pause_before_completion_wake_claim_for_test.notify_all(); + producer.join(); + + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + if (!claim_paused) { + handle.destroy(); + } + REQUIRE(claim_paused); + REQUIRE_FALSE(other.has_value()); + REQUIRE(release_continuations.load(std::memory_order_acquire) == 0); +} + TEST_CASE("object_cache TTL expiry", "[object_cache]") { scheduler sched(1); sched.start(); diff --git a/tests/unit/test_task.cpp b/tests/unit/test_task.cpp index 3ef38bf0..0ec81a26 100644 --- a/tests/unit/test_task.cpp +++ b/tests/unit/test_task.cpp @@ -14,7 +14,9 @@ #include #include #include +#include #include +#include #include #include #include "../test_main.cpp" // For scaled timeouts @@ -609,6 +611,197 @@ TEST_CASE("destroyed join_handle waiter is unregistered", state->set_value(); } +TEST_CASE("completion wake lease preserves slot ownership transitions", + "[task][completion_waiter][lifetime][regression]") { + using elio::coro::detail::completion_waiter; + using elio::coro::detail::completion_waiter_slot; + using elio::coro::detail::completion_wake_lease; + + STATIC_REQUIRE(std::is_nothrow_move_constructible_v); + STATIC_REQUIRE(std::is_nothrow_move_assignable_v); + STATIC_REQUIRE(std::is_nothrow_move_constructible_v); + STATIC_REQUIRE(std::is_nothrow_move_assignable_v); + STATIC_REQUIRE(sizeof(completion_waiter) == sizeof(completion_wake_lease)); + STATIC_REQUIRE(sizeof(completion_waiter_slot) <= + sizeof(std::mutex) + 4 * sizeof(void*)); + INFO("completion_waiter_slot=" << sizeof(completion_waiter_slot) + << ", completion_waiter=" << sizeof(completion_waiter) + << ", completion_wake_lease=" << sizeof(completion_wake_lease) + << ", join_state_base=" + << sizeof(elio::coro::detail::join_state_base) + << ", task_state=" + << sizeof(elio::coro::detail::task_state) + << ", task_state=" + << sizeof(elio::coro::detail::task_state)); + + SECTION("lease destruction abandons the selected wake") { + completion_waiter_slot slot; + completion_waiter waiter(slot); + const auto handle = std::noop_coroutine(); + + REQUIRE(slot.register_waiter(waiter, handle, [] { return false; })); + { + auto wake = slot.take(); + REQUIRE(wake); + REQUIRE_FALSE(slot.take()); + STATIC_REQUIRE(noexcept(wake.claim())); + } + + REQUIRE(slot.register_waiter(waiter, handle, [] { return false; })); + auto wake = slot.take(); + REQUIRE(wake.claim() == handle); + } + + SECTION("ready registration remains side-effect free") { + completion_waiter_slot slot; + completion_waiter waiter(slot); + const auto handle = std::noop_coroutine(); + + REQUIRE_FALSE(slot.register_waiter( + waiter, handle, [] { return true; })); + REQUIRE_FALSE(slot.take()); + } + + SECTION("an abandoned generation cannot claim a reused slot") { + completion_waiter_slot slot; + std::optional first(std::in_place, slot); + const auto handle = std::noop_coroutine(); + + REQUIRE(slot.register_waiter(*first, handle, [] { return false; })); + auto stale_wake = slot.take(); + first.reset(); + + completion_waiter second(slot); + REQUIRE(slot.register_waiter(second, handle, [] { return false; })); + auto current_wake = slot.take(); + REQUIRE_FALSE(stale_wake.claim()); + REQUIRE(current_wake.claim() == handle); + } + + SECTION("moving a registered waiter preserves its registration") { + completion_waiter_slot slot; + completion_waiter first(slot); + const auto handle = std::noop_coroutine(); + + REQUIRE(slot.register_waiter(first, handle, [] { return false; })); + completion_waiter second(std::move(first)); + auto wake = slot.take(); + REQUIRE(wake.claim() == handle); + } + + SECTION("moving then destroying a selected waiter abandons its wake") { + completion_waiter_slot slot; + completion_waiter first(slot); + const auto handle = std::noop_coroutine(); + + REQUIRE(slot.register_waiter(first, handle, [] { return false; })); + auto wake = slot.take(); + { + completion_waiter second(std::move(first)); + } + REQUIRE_FALSE(wake.claim()); + } +} + +TEST_CASE("join completion skips a waiter destroyed after selection", + "[task][join_handle][cancellation][lifetime][regression]") { + auto state = std::make_shared>(); + std::atomic resumed{false}; + + auto waiter_task = [state, &resumed]() -> task { + join_handle handle(state); + co_await handle; + resumed.store(true, std::memory_order_release); + }; + + auto waiter = waiter_task(); + auto h = elio::coro::detail::task_access::release(std::move(waiter)); + h.resume(); + REQUIRE_FALSE(h.done()); + + using namespace elio::coro::detail; + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + pause_before_completion_wake_claim_for_test.store( + true, std::memory_order_release); + + std::thread producer([state] { state->set_value(); }); + const auto deadline = std::chrono::steady_clock::now() + scaled_sec(5); + while (!completion_wake_claim_paused_for_test.load( + std::memory_order_acquire) && + std::chrono::steady_clock::now() < deadline) { + std::this_thread::yield(); + } + const bool selection_paused = + completion_wake_claim_paused_for_test.load(std::memory_order_acquire); + if (selection_paused) { + h.destroy(); + } + + pause_before_completion_wake_claim_for_test.store( + false, std::memory_order_release); + pause_before_completion_wake_claim_for_test.notify_all(); + producer.join(); + + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + if (!selection_paused) { + h.destroy(); + } + REQUIRE(selection_paused); + REQUIRE_FALSE(resumed.load(std::memory_order_acquire)); +} + +TEST_CASE("task handle completion skips a waiter destroyed after selection", + "[task][task_handle][cancellation][lifetime][regression]") { + auto state = std::make_shared>(); + task_handle handle(state); + std::atomic resumed{false}; + + auto waiter_task = [&handle, &resumed]() -> task { + auto result = co_await handle; + (void)result; + resumed.store(true, std::memory_order_release); + }; + + auto waiter = waiter_task(); + auto h = elio::coro::detail::task_access::release(std::move(waiter)); + h.resume(); + REQUIRE_FALSE(h.done()); + + using namespace elio::coro::detail; + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + pause_before_completion_wake_claim_for_test.store( + true, std::memory_order_release); + + std::thread producer([state] { state->set_value(); }); + const auto deadline = std::chrono::steady_clock::now() + scaled_sec(5); + while (!completion_wake_claim_paused_for_test.load( + std::memory_order_acquire) && + std::chrono::steady_clock::now() < deadline) { + std::this_thread::yield(); + } + const bool selection_paused = + completion_wake_claim_paused_for_test.load(std::memory_order_acquire); + if (selection_paused) { + h.destroy(); + } + + pause_before_completion_wake_claim_for_test.store( + false, std::memory_order_release); + pause_before_completion_wake_claim_for_test.notify_all(); + producer.join(); + + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + if (!selection_paused) { + h.destroy(); + } + REQUIRE(selection_paused); + REQUIRE_FALSE(resumed.load(std::memory_order_acquire)); +} + TEST_CASE("task co_return value", "[task]") { auto t = simple_return_value(); diff --git a/tests/unit/test_task_group.cpp b/tests/unit/test_task_group.cpp index 6ef6146c..fce772e9 100644 --- a/tests/unit/test_task_group.cpp +++ b/tests/unit/test_task_group.cpp @@ -4,6 +4,7 @@ #include #include #include +#include #include #include #include @@ -776,6 +777,61 @@ TEST_CASE("task_scope retains its body callable until children drain", sched.shutdown(); } +TEST_CASE("task_group completion skips a join waiter destroyed after selection", + "[task_group][cancellation][lifetime][regression]") { + using namespace elio::coro::detail; + + scheduler sched(1); + sched.start(); + auto state = std::make_shared(sched); + state->register_child(); + std::atomic resumed{false}; + std::atomic selected{false}; + + auto waiter_task = [state, &resumed]() -> task { + co_await task_group_completion_state::all_done_awaitable(state); + resumed.store(true, std::memory_order_release); + }; + auto waiter = waiter_task(); + auto handle = task_access::release(std::move(waiter)); + handle.resume(); + REQUIRE_FALSE(handle.done()); + + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + pause_before_completion_wake_claim_for_test.store( + true, std::memory_order_release); + std::thread producer([state, &selected] { + auto completion = state->child_finished(); + selected.store(static_cast(completion.wake), + std::memory_order_release); + resume_task_group_join_waiter(std::move(completion)); + }); + + const bool claim_paused = wait_for_flag( + completion_wake_claim_paused_for_test); + if (claim_paused) { + handle.destroy(); + } + pause_before_completion_wake_claim_for_test.store( + false, std::memory_order_release); + pause_before_completion_wake_claim_for_test.notify_all(); + producer.join(); + + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + const bool drained = sched.wait_for_idle(std::chrono::seconds(5)); + sched.shutdown(); + if (!claim_paused) { + handle.destroy(); + } + REQUIRE(claim_paused); + REQUIRE(selected.load(std::memory_order_acquire)); + REQUIRE_FALSE(resumed.load(std::memory_order_acquire)); + REQUIRE(state->outstanding_children() == 0); + REQUIRE(drained); +} + TEST_CASE("task_scope survives rejected handoff after external body wakeup", "[task_group][task_scope][structured][scheduler_domain][failure]") { scheduler sched(1); diff --git a/wiki/API-Contracts.md b/wiki/API-Contracts.md index b29247cf..0682b46f 100644 --- a/wiki/API-Contracts.md +++ b/wiki/API-Contracts.md @@ -33,7 +33,7 @@ by a broad module heading without checking the feature page or header comment. | Area | Elio guarantees | Caller must guarantee | |------|-----------------|-----------------------| -| Object lifetime | Public APIs preserve memory safety when documented lifetimes are met. Awaiters clean up registered waiters when their owning coroutine frame is destroyed according to the API contract. | Keep borrowed buffers, spans, callbacks, streams, schedulers, TLS contexts, RDMA objects, and other externally owned resources alive for the documented duration. | +| Object lifetime | Public APIs preserve memory safety when documented lifetimes are met. Awaiters clean up registered waiters when their owning coroutine frame is destroyed according to the API contract. Completion waits also abandon a wake selected by the producer when the frame is destroyed before scheduling ownership is claimed. | Keep borrowed buffers, spans, callbacks, streams, schedulers, TLS contexts, RDMA objects, and other externally owned resources alive for the documented duration. | | Thread safety | APIs documented as concurrent-safe protect their own internal state. Protocol send helpers that document frame serialization prevent interleaved wire frames. | Unless documented otherwise, serialize access to mutable objects from multiple coroutines or threads. Do not issue overlapping operations on non-concurrent streams or clients. | | Cancellation and timeouts | Cancellation requests are cooperative and best-effort. APIs with `cancel_token` or timeout parameters preserve one terminal operation state and report cancellation through their documented result channel. A request does not roll back I/O side effects, and an actual completion may win a cancellation race. | Pass the relevant token into operations that support cancellation and handle the documented result. A runtime or wrapper request does not forcibly destroy a child operation or wake a wait that does not observe that token. | | Results and terminal states | Operations report errors through their documented result type, optional return, `errno`, or exception path. Protocol objects transition to closed/error states as documented. | Check operation results before using returned values. Do not continue using an object after a documented terminal state unless the API documents reuse. | From ffe253a26d375f6a7009a354dcba48e3540d8157 Mon Sep 17 00:00:00 2001 From: Coldwings Date: Thu, 13 Aug 2026 17:49:38 +0800 Subject: [PATCH 2/2] test(rdma): cover selected pump exit waiter destruction --- tests/unit/test_rdma_ibverbs_endpoint.cpp | 51 +++++++++++++++++++++++ 1 file changed, 51 insertions(+) diff --git a/tests/unit/test_rdma_ibverbs_endpoint.cpp b/tests/unit/test_rdma_ibverbs_endpoint.cpp index 52e7a6bb..0d559382 100644 --- a/tests/unit/test_rdma_ibverbs_endpoint.cpp +++ b/tests/unit/test_rdma_ibverbs_endpoint.cpp @@ -23,6 +23,8 @@ #include #include +#include +#include #include #include @@ -151,6 +153,55 @@ TEST_CASE("endpoint pump exit signal wakes and unregisters coroutine waiters", REQUIRE_NOTHROW(state.mark_exited()); REQUIRE(state.is_exited()); } + + SECTION("destroyed shutdown waiter abandons a selected pump exit wake") { + using namespace elio::coro::detail; + + pump_exit_state state; + std::atomic resumed{false}; + auto wait_for_abandoned_exit = [&]() -> elio::coro::task { + co_await state.wait(); + resumed.store(true, std::memory_order_release); + }; + auto waiting = wait_for_abandoned_exit(); + auto handle = task_access::release(std::move(waiting)); + handle.resume(); + REQUIRE_FALSE(handle.done()); + + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + pause_before_completion_wake_claim_for_test.store( + true, std::memory_order_release); + std::thread producer([&state] { state.mark_exited(); }); + + const auto deadline = std::chrono::steady_clock::now() + + std::chrono::seconds(5); + while (!completion_wake_claim_paused_for_test.load( + std::memory_order_acquire) + && std::chrono::steady_clock::now() < deadline) { + std::this_thread::yield(); + } + const bool claim_paused = + completion_wake_claim_paused_for_test.load( + std::memory_order_acquire); + if (claim_paused) { + handle.destroy(); + } + + pause_before_completion_wake_claim_for_test.store( + false, std::memory_order_release); + pause_before_completion_wake_claim_for_test.notify_all(); + producer.join(); + completion_wake_claim_paused_for_test.store(false, + std::memory_order_release); + if (!claim_paused) { + handle.destroy(); + } + + REQUIRE(claim_paused); + REQUIRE(state.is_exited()); + REQUIRE_FALSE(resumed.load(std::memory_order_acquire)); + } } TEST_CASE("endpoint resource teardown retains failures for retry",