#include <coroutine>
#include <deque>
#include <exception>
#include <iostream>
#include <utility>
// Single-threaded scheduler: a queue of coroutines ready to resume.
class Scheduler {
public:
void Schedule(std::coroutine_handle<> h) {
std::cout << " [queue push_back " << h.address() << "]\n";
ready_.push_back(h);
}
void Run() {
while (!ready_.empty()) {
std::coroutine_handle<> h = ready_.front();
ready_.pop_front();
std::cout << " [queue pop_front -> resume " << h.address() << "]\n";
h.resume();
}
}
private:
std::deque<std::coroutine_handle<>> ready_;
};
// Coroutine condition variable. Wait() parks the handle;
// Notify*() moves parked handles to the scheduler's ready queue.
class CondVar {
public:
explicit CondVar(Scheduler& sched) : sched_(sched) {}
auto Wait() {
struct Awaiter {
CondVar& cv;
bool await_ready() const noexcept { return false; }
void await_suspend(std::coroutine_handle<> h) {
std::cout << " [cv parks " << h.address() << "]\n";
cv.waiters_.push_back(h);
}
void await_resume() const noexcept {}
};
return Awaiter{*this};
}
void NotifyOne() {
if (waiters_.empty()) return;
sched_.Schedule(waiters_.front());
waiters_.pop_front();
}
void NotifyAll() {
while (!waiters_.empty()) NotifyOne();
}
private:
Scheduler& sched_;
std::deque<std::coroutine_handle<>> waiters_;
};
// Fire-and-forget top-level coroutine (unchanged from the original).
struct Task {
struct promise_type {
Task get_return_object() { return {}; }
std::suspend_never initial_suspend() noexcept { return {}; }
std::suspend_never final_suspend() noexcept { return {}; }
void return_void() {}
void unhandled_exception() { std::terminate(); }
};
};
// Awaitable child coroutine: lazy start + symmetric transfer both ways.
struct Co {
struct promise_type {
std::coroutine_handle<> continuation; // parent to resume when done
Co get_return_object() {
return Co{std::coroutine_handle<promise_type>::from_promise(*this)};
}
// Don't run until the parent co_awaits us (so continuation is set).
std::suspend_always initial_suspend() noexcept { return {}; }
auto final_suspend() noexcept {
struct Final {
bool await_ready() noexcept { return false; }
std::coroutine_handle<> await_suspend(
std::coroutine_handle<promise_type> h) noexcept {
return h.promise().continuation; // jump straight back to parent
}
void await_resume() noexcept {}
};
return Final{};
}
void return_void() {}
void unhandled_exception() { std::terminate(); }
};
explicit Co(std::coroutine_handle<promise_type> h) : h_(h) {}
Co(Co&& o) noexcept : h_(std::exchange(o.h_, {})) {}
~Co() {
if (h_) h_.destroy(); // frame lives until the co_await expression ends
}
bool await_ready() const noexcept { return false; }
std::coroutine_handle<> await_suspend(std::coroutine_handle<> parent) {
h_.promise().continuation = parent;
return h_; // jump straight into child; no queue
}
void await_resume() const noexcept {}
std::coroutine_handle<promise_type> h_;
};
bool data_ready = false;
Co WaitForData(CondVar& cv, int id) {
while (!data_ready) { // Predicate loop, same as std::condition_variable.
std::cout << "child " << id << " waits\n";
co_await cv.Wait();
}
std::cout << "child " << id << " sees data\n";
}
Task Consumer(CondVar& cv, int id) {
std::cout << "consumer " << id << " calls child\n";
co_await WaitForData(cv, id); // nested co_await
std::cout << "consumer " << id << " woke up\n";
}
Task Producer(CondVar& cv) {
data_ready = true;
std::cout << "producer notifies\n";
cv.NotifyAll();
co_return;
}
int main() {
Scheduler sched;
CondVar cv(sched);
Consumer(cv, 1);
Consumer(cv, 2);
Producer(cv);
std::cout << "--- sched.Run() ---\n";
sched.Run();
}
In a C++20 coroutine, a suspension point is any location where the coroutine can pause execution, save its local state to the heap-allocated coroutine frame, and return control to its caller or resumer.
There are four types of suspension points, split between those you write explicitly in the function body and those the compiler injects automatically around your code.
Where They Live: Compiler Transformation
When write a coroutine, the C++20 compiler rewrites its body into a state machine wrapped in boilerplate that invokes your promise_type. Every coroutine contains at least the two implicit suspension points, plus any explicit ones we write:
// What the compiler generates behind the scenes for your coroutine:
{
promise_type promise;
auto return_object = promise.get_return_object();
// --- 1. IMPLICIT: Initial Suspension Point ---
co_await promise.initial_suspend();
try {
// --- YOUR COROUTINE BODY STARTS HERE ---
// --- 2. EXPLICIT: co_await Suspension Point ---
auto data = co_await async_read();
// --- 3. EXPLICIT: co_yield Suspension Point ---
co_yield data; // Rewritten as: co_await promise.yield_value(data);
// co_return is NOT a suspension point; it sets the result and jumps to final_suspend
co_return; // Calls promise.return_void() and goto final_suspend;
// --- YOUR COROUTINE BODY ENDS HERE ---
} catch (...) {
promise.unhandled_exception();
}
final_suspend:
// --- 4. IMPLICIT: Final Suspension Point ---
co_await promise.final_suspend();
} // Coroutine frame is automatically destroyed here ONLY if final_suspend does not suspendDoes a Suspension Point Always Suspend?
Hitting a suspension point means the coroutine may suspend, not that it must suspend. At every co_await (including those generated by co_yield, initial_suspend, and final_suspend), the compiler queries an Awaiter object using a 3-step protocol:
awaiter.await_ready(): Checked first as a fast-path optimization.If it returns
true, the result is already available. The coroutine does not suspend and immediately callsawait_resume().If it returns
false, the coroutine prepares to suspend by saving its instruction pointer and local registers into the coroutine frame.
awaiter.await_suspend(handle): Called right after the coroutine state is saved. Its return type controls what happens next:void: Truly suspends the coroutine and returns control to the caller/resumer.bool: Returningtruesuspends and returns to the caller; returningfalseaborts the suspension and immediately resumes the coroutine on the current thread.std::coroutine_handle: Suspends the current coroutine and immediately resumes the returned handle via symmetric transfer (without growing the call stack).
awaiter.await_resume(): Called once the coroutine is resumed (or immediately ifawait_ready()returnedtrue). Its return value becomes the result of theco_awaitexpression.
#include <coroutine>
#include <deque>
#include <exception>
#include <iostream>
#include <optional>
#include <stop_token>
#include <utility>
// Counts live coroutine frames so main() can prove nothing leaked.
struct FrameCounter {
static inline int live = 0;
FrameCounter() { ++live; }
~FrameCounter() { --live; }
};<
// Single-threaded scheduler: a queue of coroutines ready to resume.
class Scheduler {
public:
void Schedule(std::coroutine_handle<> h) { ready_.push_back(h); }
// Returns when no coroutine is runnable. After a stop request, that means
// every task has exited (graceful shutdown), unless one is stuck on
// something that ignores the stop token.
void Run() {
while (!ready_.empty()) {
std::coroutine_handle<> h = ready_.front();
ready_.pop_front();
h.resume();
}
}
private:
std::deque<std::coroutine_handle<>> ready_;
};
// Re-queues the current coroutine so others get a turn (simulates work).
auto Yield(Scheduler& sched) {
struct Awaiter {
Scheduler& sched;
bool await_ready() const noexcept { return false; }
void await_suspend(std::coroutine_handle<> h) { sched.Schedule(h); }
void await_resume() const noexcept {}
};
return Awaiter{sched};
}
// Coroutine condition variable with cancellation.
// Wait(token) parks the handle. It is woken by NotifyOne/NotifyAll OR by
// token's stop request, whichever comes first. Callers must re-check their
// predicate and the token after waking (same as std::condition_variable_any).
class CondVar {
public:
explicit CondVar(Scheduler& sched) : sched_(sched) {}
class WaitAwaiter {
public:
WaitAwaiter(CondVar& cv, std::stop_token token)
: cv_(cv), token_(std::move(token)) {}
// Already stopped: don't suspend at all.
bool await_ready() const noexcept { return token_.stop_requested(); }
void await_suspend(std::coroutine_handle<> h) {
h_ = h;
cv_.waiters_.push_back(h);
// Runs OnStop inline if stop was already requested. The callback is
// deregistered when this awaiter is destroyed (end of co_await).
on_stop_.emplace(token_, OnStop{this});
}
void await_resume() const noexcept {}
private:
struct OnStop {
WaitAwaiter* self;
void operator()() const noexcept { self->cv_.CancelWait(self->h_); }
};
CondVar& cv_;
std::stop_token token_;
std::coroutine_handle<> h_;
std::optional<std::stop_callback<OnStop>> on_stop_;
};
WaitAwaiter Wait(std::stop_token token) {
return WaitAwaiter(*this, std::move(token));
}
void NotifyOne() {
if (waiters_.empty()) return;
sched_.Schedule(waiters_.front());
waiters_.pop_front();
}
void NotifyAll() {
while (!waiters_.empty()) NotifyOne();
}
private:
// Wakes h only if it is still parked here. If NotifyOne already moved it
// to the ready queue, do nothing; scheduling it twice would resume a
// coroutine that is not suspended (UB).
void CancelWait(std::coroutine_handle<> h) {
for (auto it = waiters_.begin(); it != waiters_.end(); ++it) {
if (*it == h) {
waiters_.erase(it);
sched_.Schedule(h);
return;
}
}
}
Scheduler& sched_;
std::deque<std::coroutine_handle<>> waiters_;
};
// Fire-and-forget top-level coroutine. Frame self-destroys at the end.
struct Task {
struct promise_type : FrameCounter {
Task get_return_object() { return {}; }
std::suspend_never initial_suspend() noexcept { return {}; }
std::suspend_never final_suspend() noexcept { return {}; }
void return_void() {}
void unhandled_exception() { std::terminate(); }
};
};
// Awaitable child coroutine returning T: lazy start + symmetric transfer.
template <typename T>
class Co {
public:
struct promise_type : FrameCounter {
std::coroutine_handle<> continuation; // parent to resume when done
std::optional<T> result;
Co get_return_object() {
return Co(std::coroutine_handle<promise_type>::from_promise(*this));
}
std::suspend_always initial_suspend() noexcept { return {}; }
auto final_suspend() noexcept {
struct Final {
bool await_ready() noexcept { return false; }
std::coroutine_handle<> await_suspend(
std::coroutine_handle<promise_type> h) noexcept {
return h.promise().continuation; // jump straight back to parent
}
void await_resume() noexcept {}
};
return Final{};
}
void return_value(T v) { result.emplace(std::move(v)); }
void unhandled_exception() { std::terminate(); }
};
explicit Co(std::coroutine_handle<promise_type> h) : h_(h) {}
Co(Co&& o) noexcept : h_(std::exchange(o.h_, {})) {}
~Co() {
if (h_) h_.destroy();
}
bool await_ready() const noexcept { return false; }
std::coroutine_handle<> await_suspend(std::coroutine_handle<> parent) {
h_.promise().continuation = parent;
return h_; // jump straight into child; no queue
}
T await_resume() { return std::move(*h_.promise().result); }
private:
std::coroutine_handle<promise_type> h_;
};
// Unbounded work queue.
class Channel {
public:
explicit Channel(Scheduler& sched) : cv_(sched) {}
void Push(int v) {
items_.push_back(v);
cv_.NotifyOne();
}
// Graceful semantics: keeps handing out queued items even after stop, and
// returns nullopt only when the queue is empty AND stop was requested.
Co<std::optional<int>> Pop(std::stop_token token) {
while (items_.empty()) {
if (token.stop_requested()) co_return std::nullopt;
co_await cv_.Wait(token); // nested co_await; parks THIS child frame
}
int v = items_.front();
items_.pop_front();
co_return v;
}
private:
std::deque<int> items_;
CondVar cv_;
};
Task Consumer(Scheduler& sched, Channel& ch, int id, std::stop_token token) {
struct Cleanup { // proves destructors run on shutdown
int id;
~Cleanup() { std::cout << "consumer " << id << " cleanup (RAII)\n"; }
} cleanup{id};
while (std::optional<int> item = co_await ch.Pop(token)) {
std::cout << "consumer " << id << " processes " << *item << "\n";
co_await Yield(sched); // simulate work
}
std::cout << "consumer " << id << " exits\n";
}
Task Producer(Scheduler& sched, Channel& ch, std::stop_token token) {
for (int i = 0; !token.stop_requested(); ++i) {
std::cout << "producer pushes " << i << "\n";
ch.Push(i);
co_await Yield(sched);
}
std::cout << "producer exits (stop requested)\n";
}
// Stands in for a SIGTERM handler / deadline: requests stop after `ticks`.
Task StopAfter(Scheduler& sched, std::stop_source& stop, int ticks) {
for (int i = 0; i < ticks; ++i) co_await Yield(sched);
std::cout << "=== shutdown requested ===\n";
stop.request_stop(); // wakes parked waiters via their stop_callbacks
}
int main() {
Scheduler sched;
Channel ch(sched);
std::stop_source stop;
Consumer(sched, ch, 1, stop.get_token());
Consumer(sched, ch, 2, stop.get_token());
Producer(sched, ch, stop.get_token());
StopAfter(sched, stop, 3);
sched.Run(); // returns once every task has exited
std::cout << "live coroutine frames after Run(): " << FrameCounter::live
<< "\n";
return FrameCounter::live == 0 ? 0 : 1;
}
No comments:
Post a Comment
Note: Only a member of this blog may post a comment.