Skip to content

Commit e1283aa

Browse files
committed
io_stream implementation returns coroutine handles
1 parent b44611f commit e1283aa

8 files changed

Lines changed: 246 additions & 44 deletions

File tree

‎doc/scheduler.md‎

Lines changed: 198 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,198 @@
1+
# Scheduler Architecture
2+
3+
This document describes the architectural differences between Boost.Asio and Corosio's coroutine scheduling mechanisms, and outlines the implementation approach for Corosio.
4+
5+
## Overview
6+
7+
Asio and Corosio take fundamentally different approaches to coroutine management:
8+
9+
| Aspect | Asio | Corosio |
10+
|--------|------|---------|
11+
| Symmetric transfer | Simulated via `pump()` loop | Native language mechanism |
12+
| Type system | Closed (`asio::awaitable` only) | Open (`IoAwaitable` concept) |
13+
| `await_suspend` return | `void` | `coroutine_handle` |
14+
| I/O initiation timing | `after_suspend_fn_` callback | TBD |
15+
16+
## Symmetric Transfer
17+
18+
### The Problem
19+
20+
When coroutine A awaits coroutine B, naive implementations create stack growth:
21+
22+
```
23+
A.resume()
24+
└─> B.resume()
25+
└─> C.resume()
26+
└─> ... // unbounded stack growth
27+
```
28+
29+
C++20 symmetric transfer solves this by allowing `await_suspend` to return a `coroutine_handle`, which the compiler tail-calls instead of returning to the caller.
30+
31+
### Asio's Approach: Manual Simulation
32+
33+
Asio's `await_suspend` returns `void`, forfeiting language-based symmetric transfer:
34+
35+
```cpp
36+
// asio::awaitable - await_suspend returns void
37+
template <class U>
38+
void await_suspend(
39+
detail::coroutine_handle<detail::awaitable_frame<U, Executor>> h)
40+
{
41+
frame_->push_frame(&h.promise()); // builds linked list
42+
}
43+
```
44+
45+
Instead, Asio maintains a manual stack of frames and uses `pump()` to simulate symmetric transfer:
46+
47+
```cpp
48+
void pump()
49+
{
50+
do
51+
bottom_of_stack_.frame_->top_of_stack_->resume();
52+
while (bottom_of_stack_.frame_ && bottom_of_stack_.frame_->top_of_stack_);
53+
// ...
54+
}
55+
```
56+
57+
The `pump()` loop repeatedly calls `resume()` on the top frame until the stack empties or transfers to another thread. This achieves the same bounded-stack behavior but through explicit frame management.
58+
59+
### Corosio's Approach: Native Symmetric Transfer
60+
61+
Corosio's `task<T>` returns `coroutine_handle` from `await_suspend`, enabling compiler-optimized tail calls:
62+
63+
```cpp
64+
// task<T>::await_suspend - returns coroutine_handle
65+
coro await_suspend(coro cont, executor_ref caller_ex, std::stop_token token)
66+
{
67+
h_.promise().set_continuation(cont, caller_ex);
68+
h_.promise().set_executor(caller_ex);
69+
h_.promise().set_stop_token(token);
70+
return h_; // compiler tail-calls this handle
71+
}
72+
```
73+
74+
Similarly, `final_suspend` returns the continuation handle:
75+
76+
```cpp
77+
auto final_suspend() noexcept
78+
{
79+
struct awaiter
80+
{
81+
coro await_suspend(coro) const noexcept
82+
{
83+
return p_->complete(); // returns continuation
84+
}
85+
// ...
86+
};
87+
return awaiter{this};
88+
}
89+
```
90+
91+
This means task-to-task awaits have zero overhead beyond what the language provides. No pump loop needed for non-I/O transitions.
92+
93+
## Type System
94+
95+
### Asio's Closed System
96+
97+
Asio's `await_suspend` only accepts handles to `awaitable_frame`:
98+
99+
```cpp
100+
void await_suspend(
101+
detail::coroutine_handle<detail::awaitable_frame<U, Executor>> h)
102+
```
103+
104+
This creates a closed type system where only `asio::awaitable<T>` coroutines can participate in I/O chains. User-defined coroutine types with different promise types cannot `co_await` an `asio::awaitable`.
105+
106+
### Corosio's Open System
107+
108+
Corosio uses the `IoAwaitable` concept, allowing any conforming type to participate:
109+
110+
```cpp
111+
template<class Awaitable>
112+
auto transform_awaitable(Awaitable&& a)
113+
{
114+
using A = std::decay_t<Awaitable>;
115+
if constexpr (IoAwaitable<A>)
116+
{
117+
return transform_awaiter<Awaitable>{
118+
std::forward<Awaitable>(a), this};
119+
}
120+
// ...
121+
}
122+
```
123+
124+
The `await_suspend` signature accepts additional context parameters:
125+
126+
```cpp
127+
coro await_suspend(coro cont, executor_ref caller_ex, std::stop_token token)
128+
```
129+
130+
This design allows third-party awaitable types to integrate with Corosio's I/O system by satisfying the `IoAwaitable` concept.
131+
132+
## I/O Initiation Timing
133+
134+
### The Suspension Race Problem
135+
136+
A critical issue in coroutine-based I/O is ensuring the I/O operation isn't initiated until the coroutine is fully suspended. If the completion handler fires before suspension completes, the coroutine may be resumed while still in the middle of suspending—undefined behavior.
137+
138+
### Asio's Solution: `after_suspend_fn_`
139+
140+
Asio solves this with a deferred callback mechanism:
141+
142+
```cpp
143+
struct resume_context
144+
{
145+
void (*after_suspend_fn_)(void*) = nullptr;
146+
void *after_suspend_arg_ = nullptr;
147+
};
148+
149+
void resume()
150+
{
151+
resume_context context;
152+
resume_context_ = &context;
153+
coro_.resume(); // coroutine runs until it suspends
154+
if (context.after_suspend_fn_)
155+
context.after_suspend_fn_(context.after_suspend_arg_); // NOW safe to initiate I/O
156+
}
157+
```
158+
159+
Within `await_suspend`, true I/O operations register their initiation function:
160+
161+
```cpp
162+
// awaitable_async_op::await_suspend
163+
void await_suspend(coroutine_handle<void>)
164+
{
165+
frame_->after_suspend(
166+
[](void* arg)
167+
{
168+
awaitable_async_op* self = static_cast<awaitable_async_op*>(arg);
169+
// Actually initiate the I/O operation here
170+
std::forward<Op&&>(self->op_)(
171+
handler_type(self->frame_->detach_thread(), self->result_));
172+
}, this);
173+
}
174+
```
175+
176+
Key distinction:
177+
- **Task-to-task awaits**: Use `push_frame()`, no `after_suspend_fn_` set
178+
- **True I/O awaits**: Set `after_suspend_fn_` to defer initiation
179+
180+
### Corosio's Approach
181+
182+
TBD - Document Corosio's mechanism for safe I/O initiation timing.
183+
184+
## Scheduler Implementation
185+
186+
TBD - Document:
187+
- Event loop design
188+
- Platform-specific backends (epoll, IOCP)
189+
- Threading model
190+
- Work stealing / distribution
191+
192+
## Implementation Plan
193+
194+
TBD - To be developed after gathering additional facts about:
195+
- Corosio's I/O initiation mechanism
196+
- Scheduler event loop design
197+
- Platform-specific details
198+
- Threading model

‎include/boost/corosio/io_stream.hpp‎

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -232,8 +232,7 @@ class BOOST_COROSIO_DECL io_stream : public io_object
232232
std::stop_token token) -> std::coroutine_handle<>
233233
{
234234
token_ = std::move(token);
235-
ios_.get().read_some(h, ex, buffers_, token_, &ec_, &bytes_transferred_);
236-
return std::noop_coroutine();
235+
return ios_.get().read_some(h, ex, buffers_, token_, &ec_, &bytes_transferred_);
237236
}
238237
};
239238

@@ -273,8 +272,7 @@ class BOOST_COROSIO_DECL io_stream : public io_object
273272
std::stop_token token) -> std::coroutine_handle<>
274273
{
275274
token_ = std::move(token);
276-
ios_.get().write_some(h, ex, buffers_, token_, &ec_, &bytes_transferred_);
277-
return std::noop_coroutine();
275+
return ios_.get().write_some(h, ex, buffers_, token_, &ec_, &bytes_transferred_);
278276
}
279277
};
280278

@@ -288,7 +286,7 @@ class BOOST_COROSIO_DECL io_stream : public io_object
288286
struct io_stream_impl : io_object_impl
289287
{
290288
/// Initiate platform read operation.
291-
virtual void read_some(
289+
virtual std::coroutine_handle<> read_some(
292290
std::coroutine_handle<>,
293291
capy::executor_ref,
294292
io_buffer_param,
@@ -297,7 +295,7 @@ class BOOST_COROSIO_DECL io_stream : public io_object
297295
std::size_t*) = 0;
298296

299297
/// Initiate platform write operation.
300-
virtual void write_some(
298+
virtual std::coroutine_handle<> write_some(
301299
std::coroutine_handle<>,
302300
capy::executor_ref,
303301
io_buffer_param,

‎src/corosio/src/detail/epoll/sockets.cpp‎

Lines changed: 13 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -209,7 +209,7 @@ connect(
209209
svc_.post(&op);
210210
}
211211

212-
void
212+
std::coroutine_handle<>
213213
epoll_socket_impl::
214214
read_some(
215215
std::coroutine_handle<> h,
@@ -237,7 +237,7 @@ read_some(
237237
op.complete(0, 0);
238238
op.impl_ptr = shared_from_this();
239239
svc_.post(&op);
240-
return;
240+
return std::noop_coroutine();
241241
}
242242

243243
for (int i = 0; i < op.iovec_count; ++i)
@@ -254,7 +254,7 @@ read_some(
254254
op.complete(0, static_cast<std::size_t>(n));
255255
op.impl_ptr = shared_from_this();
256256
svc_.post(&op);
257-
return;
257+
return std::noop_coroutine();
258258
}
259259

260260
if (n == 0)
@@ -263,7 +263,7 @@ read_some(
263263
op.complete(0, 0);
264264
op.impl_ptr = shared_from_this();
265265
svc_.post(&op);
266-
return;
266+
return std::noop_coroutine();
267267
}
268268

269269
if (errno == EAGAIN || errno == EWOULDBLOCK)
@@ -289,7 +289,7 @@ read_some(
289289
svc_.post(claimed);
290290
svc_.work_finished();
291291
}
292-
return;
292+
return std::noop_coroutine();
293293
}
294294
}
295295

@@ -302,15 +302,16 @@ read_some(
302302
svc_.work_finished();
303303
}
304304
}
305-
return;
305+
return std::noop_coroutine();
306306
}
307307

308308
op.complete(errno, 0);
309309
op.impl_ptr = shared_from_this();
310310
svc_.post(&op);
311+
return std::noop_coroutine();
311312
}
312313

313-
void
314+
std::coroutine_handle<>
314315
epoll_socket_impl::
315316
write_some(
316317
std::coroutine_handle<> h,
@@ -337,7 +338,7 @@ write_some(
337338
op.complete(0, 0);
338339
op.impl_ptr = shared_from_this();
339340
svc_.post(&op);
340-
return;
341+
return std::noop_coroutine();
341342
}
342343

343344
for (int i = 0; i < op.iovec_count; ++i)
@@ -358,7 +359,7 @@ write_some(
358359
op.complete(0, static_cast<std::size_t>(n));
359360
op.impl_ptr = shared_from_this();
360361
svc_.post(&op);
361-
return;
362+
return std::noop_coroutine();
362363
}
363364

364365
if (errno == EAGAIN || errno == EWOULDBLOCK)
@@ -384,7 +385,7 @@ write_some(
384385
svc_.post(claimed);
385386
svc_.work_finished();
386387
}
387-
return;
388+
return std::noop_coroutine();
388389
}
389390
}
390391

@@ -397,12 +398,13 @@ write_some(
397398
svc_.work_finished();
398399
}
399400
}
400-
return;
401+
return std::noop_coroutine();
401402
}
402403

403404
op.complete(errno ? errno : EIO, 0);
404405
op.impl_ptr = shared_from_this();
405406
svc_.post(&op);
407+
return std::noop_coroutine();
406408
}
407409

408410
std::error_code

‎src/corosio/src/detail/epoll/sockets.hpp‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -102,15 +102,15 @@ class epoll_socket_impl
102102
std::stop_token,
103103
std::error_code*) override;
104104

105-
void read_some(
105+
std::coroutine_handle<> read_some(
106106
std::coroutine_handle<>,
107107
capy::executor_ref,
108108
io_buffer_param,
109109
std::stop_token,
110110
std::error_code*,
111111
std::size_t*) override;
112112

113-
void write_some(
113+
std::coroutine_handle<> write_some(
114114
std::coroutine_handle<>,
115115
capy::executor_ref,
116116
io_buffer_param,

0 commit comments

Comments
 (0)