Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
141 changes: 136 additions & 5 deletions include/elio/coro/detail/completion_waiter.hpp
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
#pragma once

#include <cassert>
#include <atomic>
#include <cstddef>
#include <coroutine>
#include <mutex>
#include <utility>
Expand All @@ -9,6 +11,50 @@ namespace elio::coro::detail {

class completion_waiter_slot;

#ifdef ELIO_RUNTIME_TEST_HOOKS
inline std::atomic<bool> pause_before_completion_wake_claim_for_test{false};
inline std::atomic<bool> 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:
Expand Down Expand Up @@ -42,7 +88,7 @@ class completion_waiter_slot {

~completion_waiter_slot() {
std::lock_guard<std::mutex> lock(mutex_);
assert(waiter_ == nullptr &&
assert(waiter_ == nullptr && !selected() &&
"completion waiter slot destroyed with a pending waiter");
}

Expand All @@ -68,22 +114,53 @@ class completion_waiter_slot {
return true;
}

std::coroutine_handle<> take() noexcept {
completion_wake_lease take() noexcept {
std::lock_guard<std::mutex> 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<std::mutex> 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<std::mutex> lock(mutex_);
if (!selected() || generation_ != generation) {
return;
}
++generation_;
waiter_ = nullptr;
selected_handle_ = {};
}

void remove(completion_waiter& waiter) noexcept {
std::lock_guard<std::mutex> lock(mutex_);
if (waiter_ == &waiter) {
if (selected()) {
++generation_;
selected_handle_ = {};
}
waiter_ = nullptr;
}
waiter.handle_ = {};
Expand All @@ -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);
Expand Down
3 changes: 2 additions & 1 deletion include/elio/coro/task.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
24 changes: 14 additions & 10 deletions include/elio/coro/task_group.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -76,19 +76,23 @@ inline std::atomic<bool> 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;
}

Expand All @@ -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;
}
}
Expand Down Expand Up @@ -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 {};
Expand Down Expand Up @@ -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));
}
}
}
Expand Down
6 changes: 4 additions & 2 deletions include/elio/coro/task_handle.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down Expand Up @@ -204,7 +205,8 @@ struct task_state<void> {
}

void notify_waiter() {
auto waiter = waiter_.take();
auto wake = waiter_.take();
auto waiter = wake.claim();
if (waiter) {
runtime::schedule_handle(waiter);
}
Expand Down
3 changes: 2 additions & 1 deletion include/elio/rdma_ibverbs/endpoint.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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()) {
Comment thread
Coldwings marked this conversation as resolved.
elio::runtime::schedule_handle(waiter);
}
}
Expand Down
3 changes: 2 additions & 1 deletion include/elio/sync/object_cache.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
Loading