#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; }
// Parent 就是之後的 await_resume resume 的handle.
// await_suspend 吃的handle就是await_resume resume的handle.
std::coroutine_handle<> await_suspend(std::coroutine_handle<> parent) {
h_.promise().continuation = parent;
return h_; // jump straight into child; no queue
}
// await_suspend 吃的handle就是await_resume resume的handle.
// co_await 回傳值在此.
// 若 coroutune function return的type 也為 waitable
// 其promise_type::void return_value(T v) 會將return value儲存至
// promise_type, 其final_suspend() return waitable type會resume parent
// 跳至此await_resume().
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;
}
// Minimal single-threaded model of C9-style cancellation.
//
// Key difference from the stop_token version: a cancelled coroutine is NEVER
// resumed. The cancel callback destroys every frame of the logical thread
// (leaf -> entry), running destructors, and then resumes an "on_exit" handle
// so the owner learns the thread is gone. No code checks a token.
//
// Name map to google3/util/c9/internal:
// PromiseLink -> CoroutinePromiseLink (waiter link, co_thread)
// CoThread -> CoThread (cancelled_ AsyncNotification, on_exit)
// SuspensionToken -> SuspensionToken (BecameReady/Cancelled/AlreadyCancelled,
// DestroyCallStack)
// Scheduler -> ThreadState ready queue + Schedule/CancelWaiter
// Yield / CondVar -> leaf awaitables like c9::Event::Awaitable
#include <algorithm>
#include <coroutine>
#include <deque>
#include <exception>
#include <functional>
#include <iostream>
#include <optional>
#include <type_traits>
#include <utility>
struct FrameCounter {
static inline int live = 0;
FrameCounter() { ++live; }
~FrameCounter() { --live; }
};
class CoThread;
// Base of every Co<T> promise. The `waiter` pointers form the intrusive
// "coroutine call stack" that cancellation walks and destroys.
struct PromiseLink {
std::coroutine_handle<> self; // this frame
PromiseLink* waiter = nullptr; // caller frame; nullptr for the entry frame
CoThread* co_thread = nullptr; // logical thread this frame belongs to
};
// One logical thread of execution = one chain of nested co_awaits.
class CoThread {
public:
explicit CoThread(std::coroutine_handle<> on_exit) : on_exit_(on_exit) {}
std::coroutine_handle<> on_exit() const { return on_exit_; }
// Called by the leaf awaitable while suspending. False = already cancelled,
// so the leaf must not park; it must destroy the stack instead.
bool Register(std::function<void()> on_cancel) {
if (cancelled_) return false;
on_cancel_ = std::move(on_cancel);
return true;
}
void Unregister() { on_cancel_ = nullptr; }
// Request cancellation. If a leaf is parked, its callback runs now and
// destroys the whole chain. If the chain is currently running, nothing
// happens here; it dies at its next suspension point (Register -> false).
void Cancel() {
cancelled_ = true;
if (on_cancel_) std::exchange(on_cancel_, nullptr)();
}
private:
std::coroutine_handle<> on_exit_;
bool cancelled_ = false;
std::function<void()> on_cancel_;
};
// Handed to a leaf awaitable at suspension. Exactly one of BecameReady,
// Cancelled, AlreadyCancelled is called for each suspension.
class SuspensionToken {
public:
explicit SuspensionToken(PromiseLink* leaf) : leaf_(leaf) {}
bool operator==(const SuspensionToken&) const = default;
bool RegisterForCancellation(std::function<void()> on_cancel) {
return co_thread().Register(std::move(on_cancel));
}
// Normal wakeup: stop listening for cancel; caller resumes the leaf.
std::coroutine_handle<> BecameReady() {
co_thread().Unregister();
return leaf_->self;
}
// Cancel callback won the race: destroy the chain, tell the owner. The leaf
// coroutine is never resumed.
void Cancelled() {
std::coroutine_handle<> on_exit = co_thread().on_exit(); // read first
DestroyCallStack(); // frees leaf_
on_exit.resume();
}
// Register failed (cancelled before we could park): destroy the chain and
// return on_exit for the caller to symmetric-transfer into.
std::coroutine_handle<> AlreadyCancelled() {
std::coroutine_handle<> on_exit = co_thread().on_exit();
DestroyCallStack();
return on_exit;
}
private:
CoThread& co_thread() { return *leaf_->co_thread; }
// Callee before caller, iteratively, so OS stack depth stays O(1).
void DestroyCallStack() {
for (PromiseLink* p = leaf_; p != nullptr;) {
PromiseLink* next = p->waiter;
p->self.destroy(); // runs this frame's destructors (RAII)
p = next;
}
}
PromiseLink* leaf_;
};
// Ready queue of wakeups. Holds tokens, not raw handles, so BecameReady runs
// (unregistering the cancel callback) right before resume.
class Scheduler {
public:
void Schedule(SuspensionToken t) { ready_.push_back(t); }
// Pull a scheduled wakeup back out so the frame can be destroyed instead.
// False = already dequeued and running; cancellation then takes effect at
// the coroutine's next suspension point.
bool CancelWaiter(const SuspensionToken& t) {
return std::erase(ready_, t) > 0;
}
void Run() {
while (!ready_.empty()) {
SuspensionToken t = ready_.front();
ready_.pop_front();
t.BecameReady().resume();
}
}
private:
std::deque<SuspensionToken> ready_;
};
// Awaitable child coroutine. Lazy start, symmetric transfer, frame destroys
// itself at final_suspend (result is written into the parent's frame first).
template <typename T>
class Co {
static constexpr bool kVoid = std::is_void_v<T>;
struct Empty {};
using Slot = std::conditional_t<kVoid, Empty, std::optional<T>>;
template <typename U>
struct ReturnValue {
std::optional<U>* result = nullptr; // lives in the awaiting parent
void return_value(U v) { result->emplace(std::move(v)); }
};
struct ReturnVoid {
void return_void() {}
};
public:
struct promise_type
: PromiseLink,
FrameCounter,
std::conditional_t<kVoid, ReturnVoid, ReturnValue<T>> {
Co get_return_object() {
auto h = std::coroutine_handle<promise_type>::from_promise(*this);
this->self = h;
return Co(h);
}
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 {
PromiseLink* waiter = h.promise().waiter;
std::coroutine_handle<> on_exit = h.promise().co_thread->on_exit();
h.destroy(); // frame gone; only locals used from here on
return waiter != nullptr ? waiter->self : on_exit;
}
void await_resume() noexcept {}
};
return Final{};
}
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(); // only if never started
}
bool await_ready() const noexcept { return false; }
template <typename P>
std::coroutine_handle<> await_suspend(std::coroutine_handle<P> parent) {
promise_type& p = h_.promise();
p.waiter = &parent.promise();
p.co_thread = parent.promise().co_thread;
if constexpr (!kVoid) p.result = &result_;
return std::exchange(h_, {}); // frame now owns itself
}
T await_resume() {
if constexpr (!kVoid) return std::move(*result_);
}
// For bridges only: hand the entry frame to a new CoThread.
std::coroutine_handle<promise_type> Release() {
return std::exchange(h_, {});
}
private:
std::coroutine_handle<promise_type> h_;
Slot result_;
};
// Tiny coroutine used as a CoThread's on_exit handle.
struct OnExit {
struct promise_type : FrameCounter {
OnExit get_return_object() {
return {std::coroutine_handle<promise_type>::from_promise(*this)};
}
std::suspend_always initial_suspend() noexcept { return {}; }
std::suspend_never final_suspend() noexcept { return {}; }
void return_void() {}
void unhandled_exception() { std::terminate(); }
};
std::coroutine_handle<> handle;
};
OnExit ReportExit(const char* name) {
std::cout << name << ": thread exited\n";
co_return;
}
// Bridge: start `co` as the entry frame of a fresh CoThread (what C9 bridges,
// RunConcurrently and WithDeadline do for their children).
void Spawn(Scheduler& sched, CoThread& ct, Co<void> co) {
auto h = co.Release();
h.promise().co_thread = &ct; // waiter stays nullptr: entry frame
sched.Schedule(SuspensionToken(&h.promise()));
}
// ---- Leaf awaitables (the only place cancellation is "implemented") ----
// Re-queue the current coroutine (simulates work). Cancellable.
class Yield {
public:
explicit Yield(Scheduler& s) : sched_(s) {}
bool await_ready() const noexcept { return false; }
template <typename P>
std::coroutine_handle<> await_suspend(std::coroutine_handle<P> h) {
token_.emplace(&h.promise());
if (!token_->RegisterForCancellation([this] { Cancel(); })) {
return token_->AlreadyCancelled();
}
sched_.Schedule(*token_);
return std::noop_coroutine();
}
void await_resume() const noexcept {}
private:
void Cancel() {
// Copy out: Cancelled() destroys the frame that holds *this.
Scheduler& sched = sched_;
SuspensionToken token = *token_;
if (sched.CancelWaiter(token)) token.Cancelled();
}
Scheduler& sched_;
std::optional<SuspensionToken> token_;
};
// Condition variable. Note Wait() takes no token: cancellation is implicit.
class CondVar {
public:
explicit CondVar(Scheduler& sched) : sched_(sched) {}
class Awaiter {
public:
explicit Awaiter(CondVar& cv) : cv_(cv) {}
bool await_ready() const noexcept { return false; }
template <typename P>
std::coroutine_handle<> await_suspend(std::coroutine_handle<P> h) {
token_.emplace(&h.promise());
if (!token_->RegisterForCancellation([this] { Cancel(); })) {
return token_->AlreadyCancelled();
}
cv_.waiters_.push_back(*token_);
return std::noop_coroutine();
}
void await_resume() const noexcept {}
private:
void Cancel() {
CondVar& cv = cv_;
SuspensionToken token = *token_;
// Either still parked here, or NotifyOne already moved it to the ready
// queue. Pull it back from wherever it is, then destroy.
if (std::erase(cv.waiters_, token) > 0 || cv.sched_.CancelWaiter(token)) {
token.Cancelled();
}
}
CondVar& cv_;
std::optional<SuspensionToken> token_;
};
Awaiter Wait() { return Awaiter(*this); }
void NotifyOne() {
if (waiters_.empty()) return;
sched_.Schedule(waiters_.front());
waiters_.pop_front();
}
private:
Scheduler& sched_;
std::deque<SuspensionToken> waiters_;
};
// ---- User code: no token anywhere ----
class Channel {
public:
explicit Channel(Scheduler& sched) : cv_(sched) {}
void Push(int v) {
items_.push_back(v);
cv_.NotifyOne();
}
Co<int> Pop() {
while (items_.empty()) co_await cv_.Wait();
int v = items_.front();
items_.pop_front();
co_return v;
}
private:
std::deque<int> items_;
CondVar cv_;
};
Co<void> Consumer(Scheduler& sched, Channel& ch, int id) {
struct Cleanup {
int id;
~Cleanup() { std::cout << "consumer " << id << " cleanup (RAII)\n"; }
} cleanup{id};
for (;;) { // no exit path: cancellation destroys this frame instead
int item = co_await ch.Pop();
std::cout << "consumer " << id << " processes " << item << "\n";
co_await Yield(sched);
}
}
Co<void> Producer(Scheduler& sched, Channel& ch) {
for (int i = 0;; ++i) {
std::cout << "producer pushes " << i << "\n";
ch.Push(i);
co_await Yield(sched);
}
}
Co<void> StopAfter(Scheduler& sched, std::deque<CoThread*> threads, int ticks) {
for (int i = 0; i < ticks; ++i) co_await Yield(sched);
std::cout << "=== cancel requested ===\n";
for (CoThread* t : threads) t->Cancel(); // callbacks destroy the stacks
}
int main() {
Scheduler sched;
Channel ch(sched);
CoThread c1(ReportExit("consumer 1").handle);
CoThread c2(ReportExit("consumer 2").handle);
CoThread prod(ReportExit("producer").handle);
CoThread stopper(ReportExit("stopper").handle);
Spawn(sched, c1, Consumer(sched, ch, 1));
Spawn(sched, c2, Consumer(sched, ch, 2));
Spawn(sched, prod, Producer(sched, ch));
Spawn(sched, stopper, StopAfter(sched, {&c1, &c2, &prod}, 3));
sched.Run();
std::cout << "live coroutine frames after Run(): " << FrameCounter::live
<< "\n";
return FrameCounter::live == 0 ? 0 : 1;
}
