Skip to content

Commit 4bd1735

Browse files
committed
refactor: share the native op envelope, decode, socket accessors, and io_uring service plumbing
Deduplicate the backend-independent parts of the native backends without sharing the scheduler loop: - coro_op: the shared non-template op envelope (reactor_op_base, io_uring_op, overlapped_op derive from it) + coro_op_complete (drain/resume tail and decode_io_result); each backend keeps its own native-error -> error_code conversion. - native_socket_base: socket accessors shared by reactor and io_uring sockets. - io_uring_socket_service_base / io_uring_file_service_base: shared io_uring service construct/destroy/shutdown/close plumbing. The scheduler loop is deliberately not shared — io_uring keeps its leader/follower loop, the reactors share reactor_scheduler, IOCP stays standalone.
1 parent 7c636ac commit 4bd1735

22 files changed

Lines changed: 1081 additions & 1019 deletions
Lines changed: 134 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,134 @@
1+
//
2+
// Copyright (c) 2026 Michael Vandeberg
3+
//
4+
// Distributed under the Boost Software License, Version 1.0. (See accompanying
5+
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
6+
//
7+
// Official repository: https://github.com/cppalliance/corosio
8+
//
9+
10+
#ifndef BOOST_COROSIO_NATIVE_DETAIL_CORO_OP_HPP
11+
#define BOOST_COROSIO_NATIVE_DETAIL_CORO_OP_HPP
12+
13+
#include <boost/corosio/detail/config.hpp>
14+
#include <boost/corosio/detail/continuation_op.hpp>
15+
#include <boost/corosio/detail/scheduler_op.hpp>
16+
#include <boost/capy/ex/executor_ref.hpp>
17+
18+
#include <atomic>
19+
#include <coroutine>
20+
#include <cstddef>
21+
#include <memory>
22+
#include <optional>
23+
#include <stop_token>
24+
#include <system_error>
25+
26+
/*
27+
Shared, non-template op envelope for every native backend — the readiness
28+
reactors (epoll/kqueue/select), io_uring, and IOCP. It captures the part of
29+
an async operation that is identical regardless of how completion is
30+
reported: the coroutine to resume, the executor it dispatches on, the
31+
output pointers, the stop_token wiring, the cancelled flag, and the
32+
keepalive that holds the owning impl alive while the op is in flight.
33+
34+
What is deliberately NOT here (it differs by backend and stays in the
35+
derived op layer):
36+
- the result model: the reactors re-run the syscall and record
37+
`errn`/`bytes_transferred` (reactor_op_base); io_uring stores the raw
38+
`res`/`cqe_flags`; IOCP stores `dwError`/`bytes_transferred`. Each
39+
decodes its own result.
40+
- the submission + the kernel cancel action. Cancellation is unified only
41+
at the call site via the virtual `on_cancel()` hook: the stop_callback
42+
always targets `coro_op`, and each backend overrides `on_cancel()` —
43+
the reactors route to the owning impl's cancel(), io_uring submits an
44+
ASYNC_CANCEL SQE, IOCP calls the stored cancel_func_/CancelIoEx.
45+
46+
See tasks/proactor-dedup-decisions.md and coro-op-unification-scope.md.
47+
*/
48+
49+
namespace boost::corosio::detail {
50+
51+
/** Non-template op envelope shared by every native backend's operations.
52+
53+
`reactor_op_base`, `io_uring_op`, and `overlapped_op` all derive from this.
54+
Derives from scheduler_op so ops queue intrusively and dispatch through the
55+
function-pointer (io_uring/IOCP) or virtual (reactors) completion path —
56+
hence both a default and a func_type constructor.
57+
58+
@note For IOCP, the concrete op multiply-inherits `OVERLAPPED` as its
59+
first base (so `static_cast<OVERLAPPED*>` round-trips); `coro_op`
60+
follows it.
61+
*/
62+
struct coro_op : scheduler_op
63+
{
64+
/** Stop-callback handler: routes a stop_token firing to `on_cancel()`.
65+
66+
A single canceller type for both backends keeps `stop_cb` (and thus
67+
`start()`) in this shared base; the backend-specific action lives
68+
behind the `on_cancel()` virtual.
69+
*/
70+
struct canceller
71+
{
72+
coro_op* op;
73+
void operator()() const noexcept { op->on_cancel(); }
74+
};
75+
76+
std::coroutine_handle<> h;
77+
detail::continuation_op cont_op;
78+
capy::executor_ref ex;
79+
std::error_code* ec_out = nullptr;
80+
std::size_t* bytes_out = nullptr;
81+
82+
/// True for receive/read ops (drives the zero-byte == EOF decision).
83+
bool is_read = false;
84+
/// True when the submitted buffer was zero-length (suppresses EOF).
85+
bool empty_buffer = false;
86+
87+
std::atomic<bool> cancelled{false};
88+
std::optional<std::stop_callback<canceller>> stop_cb;
89+
90+
/// Keeps the owning impl alive while the op is in flight (the kernel
91+
/// owns user buffers until completion). Dropped in the handler's resume
92+
/// tail (see coro_op_complete.hpp).
93+
std::shared_ptr<void> impl_ptr;
94+
95+
/// Default-construct for virtual-dispatch backends (the reactors, which
96+
/// override operator()/destroy() and leave func_ null).
97+
coro_op() noexcept = default;
98+
99+
/// Construct with the completion function for func-pointer dispatch
100+
/// (io_uring / IOCP completion handlers).
101+
explicit coro_op(func_type func) noexcept : scheduler_op(func) {}
102+
103+
/** Arm the stop-token callback. Call before the op is submitted.
104+
105+
Resets the cancellation flag and (re)arms `stop_cb` against @a token.
106+
Derived ops that carry extra pre-submit state (e.g. io_uring's
107+
`sqe_set`) extend this.
108+
*/
109+
void start(std::stop_token const& token)
110+
{
111+
cancelled.store(false, std::memory_order_relaxed);
112+
stop_cb.reset();
113+
if (token.stop_possible())
114+
stop_cb.emplace(token, canceller{this});
115+
}
116+
117+
/// Mark this op cancellation-requested. Shared by every backend.
118+
void request_cancel() noexcept
119+
{
120+
cancelled.store(true, std::memory_order_release);
121+
}
122+
123+
/** Backend cancellation hook, invoked when the stop_token fires.
124+
125+
The default just records the request. Backends override to also
126+
drive the kernel: io_uring submits an ASYNC_CANCEL SQE; IOCP calls
127+
its stored cancel_func_ (CancelIoEx / wait-reactor deregister).
128+
*/
129+
virtual void on_cancel() noexcept { request_cancel(); }
130+
};
131+
132+
} // namespace boost::corosio::detail
133+
134+
#endif
Lines changed: 136 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,136 @@
1+
//
2+
// Copyright (c) 2026 Michael Vandeberg
3+
//
4+
// Distributed under the Boost Software License, Version 1.0. (See accompanying
5+
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
6+
//
7+
// Official repository: https://github.com/cppalliance/corosio
8+
//
9+
10+
#ifndef BOOST_COROSIO_NATIVE_DETAIL_CORO_OP_COMPLETE_HPP
11+
#define BOOST_COROSIO_NATIVE_DETAIL_CORO_OP_COMPLETE_HPP
12+
13+
#include <boost/corosio/detail/dispatch_coro.hpp>
14+
#include <boost/corosio/native/detail/coro_op.hpp>
15+
#include <boost/capy/error.hpp>
16+
17+
#include <cstddef>
18+
#include <memory>
19+
#include <system_error>
20+
21+
/*
22+
Shared completion-tail helpers for proactor ops. Every IOCP and io_uring
23+
I/O handler ends the same way once its backend-specific result has been
24+
decoded into ec_out/bytes_out:
25+
26+
1. disarm the stop_callback,
27+
2. on the shutdown-drain path (owner == nullptr) just break the
28+
impl_ptr keepalive cycle and return without resuming,
29+
3. otherwise resume the coroutine on its executor, dropping the
30+
keepalive only after the continuation has been handed off.
31+
32+
The *decode* step (raw DWORD/res -> {ec, bytes, eof, canceled}) stays
33+
backend-specific because the raw encodings differ; in Phase 3 it is
34+
formalized as `Traits::decode_result`. These two helpers capture the
35+
backend-agnostic prologue and resume tail so the per-op handlers shrink to
36+
"drain-or-decode, then resume".
37+
*/
38+
39+
namespace boost::corosio::detail {
40+
41+
/** Translate a decoded I/O result into `*ec_out` using the cancelled /
42+
error / EOF / success priority shared by every native backend.
43+
44+
The raw error encodings differ per backend (reactor positive `errno`,
45+
io_uring negative `res`, IOCP `DWORD`), so the native-error -> error_code
46+
step stays backend-local: the caller passes @a err already converted
47+
(an empty error_code means "no error"). This helper owns only the
48+
priority logic, which is byte-for-byte identical everywhere:
49+
50+
cancelled -> operation_canceled
51+
err set -> err
52+
is_read && bytes == 0 && !empty -> end_of_file
53+
otherwise -> success
54+
55+
Writes nothing when @a ec_out is null. Does not touch bytes_out — callers
56+
that report a byte count write it separately (connect/wait carry none).
57+
58+
@param ec_out Destination (may be null).
59+
@param cancelled The op's cancellation flag.
60+
@param err Backend error already converted to error_code, or a
61+
default-constructed error_code on success.
62+
@param is_read True only for reads that should map a 0-byte
63+
completion to EOF — false for writes, connect, wait,
64+
and datagrams (a 0-byte datagram is success, not EOF).
65+
@param bytes Bytes transferred (consulted only for the EOF test).
66+
@param empty_buffer True when the submitted buffer was zero-length,
67+
which suppresses the otherwise-spurious EOF.
68+
*/
69+
inline void
70+
decode_io_result(
71+
std::error_code* ec_out,
72+
bool cancelled,
73+
std::error_code err,
74+
bool is_read,
75+
std::size_t bytes,
76+
bool empty_buffer) noexcept
77+
{
78+
if (!ec_out)
79+
return;
80+
if (cancelled)
81+
*ec_out = capy::error::canceled;
82+
else if (err)
83+
*ec_out = err;
84+
else if (is_read && bytes == 0 && !empty_buffer)
85+
*ec_out = capy::error::eof;
86+
else
87+
*ec_out = {};
88+
}
89+
90+
/** Completion prologue shared by every proactor handler.
91+
92+
Disarms the stop_callback, then detects the shutdown-drain path.
93+
94+
@param owner The scheduler pointer (nullptr during shutdown drain).
95+
@param self The completing op.
96+
@return True if this was a shutdown drain — the caller must `return`
97+
immediately without decoding or resuming. On that path the
98+
impl_ptr keepalive is dropped here (which may destroy the impl,
99+
and with it the op storage).
100+
*/
101+
inline bool
102+
coro_drain_if_shutdown(void* owner, coro_op* self) noexcept
103+
{
104+
self->stop_cb.reset();
105+
if (owner == nullptr)
106+
{
107+
auto suicide = std::move(self->impl_ptr);
108+
return true;
109+
}
110+
return false;
111+
}
112+
113+
/** Resume tail shared by every proactor handler.
114+
115+
Resumes the op's coroutine on its executor and then drops the impl_ptr
116+
keepalive. The keepalive is moved into a local that is released *after*
117+
`resume()` returns, matching the existing io_uring ordering: the impl (and
118+
therefore this op's storage) may be destroyed as the local goes out of
119+
scope, so nothing may touch `*self` after the resume.
120+
121+
@pre `self->ec_out`/`bytes_out` have already been written by the
122+
backend's decode step.
123+
*/
124+
inline void
125+
coro_resume(coro_op* self) noexcept
126+
{
127+
self->cont_op.cont.h = self->h;
128+
auto next = dispatch_coro(self->ex, self->cont_op.cont);
129+
auto suicide = std::move(self->impl_ptr);
130+
next.resume();
131+
// suicide drops here; may destroy impl + self.
132+
}
133+
134+
} // namespace boost::corosio::detail
135+
136+
#endif

‎include/boost/corosio/native/detail/io_uring/io_uring_dgram_ops.hpp‎

Lines changed: 24 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818

1919
#include <boost/corosio/detail/dispatch_coro.hpp>
2020
#include <boost/corosio/native/detail/io_uring/io_uring_op.hpp>
21+
#include <boost/corosio/native/detail/coro_op_complete.hpp>
2122
#include <boost/corosio/native/detail/speculative_state.hpp>
2223
#include <boost/corosio/native/detail/io_uring/io_uring_socket_ops.hpp>
2324
#include <boost/corosio/native/detail/make_err.hpp>
@@ -132,23 +133,18 @@ struct uring_dgram_send_op : io_uring_op
132133
std::uint32_t /*bytes*/, std::uint32_t /*error*/) noexcept
133134
{
134135
auto* self = static_cast<uring_dgram_send_op*>(base);
135-
self->stop_cb.reset();
136-
137-
if (owner == nullptr)
138-
{
139-
auto suicide = std::move(self->impl_ptr);
136+
if (coro_drain_if_shutdown(owner, self))
140137
return;
141-
}
142138

143-
if (self->ec_out)
144-
{
145-
if (self->cancelled.load(std::memory_order_acquire))
146-
*self->ec_out = capy::error::canceled;
147-
else if (self->res < 0)
148-
*self->ec_out = make_err(-self->res);
149-
else
150-
*self->ec_out = {};
151-
}
139+
if (self->sched_)
140+
self->sched_->reset_inline_budget();
141+
142+
// Datagram send: no EOF (a 0-byte send is success).
143+
decode_io_result(
144+
self->ec_out,
145+
self->cancelled.load(std::memory_order_acquire),
146+
self->res < 0 ? make_err(-self->res) : std::error_code{},
147+
/*is_read=*/false, /*bytes=*/0, /*empty_buffer=*/false);
152148
if (self->bytes_out)
153149
*self->bytes_out = (self->res >= 0)
154150
? static_cast<std::size_t>(self->res) : 0;
@@ -159,10 +155,7 @@ struct uring_dgram_send_op : io_uring_op
159155
self->spec_state->on_async_write_ready();
160156
}
161157

162-
self->cont_op.cont.h = self->h;
163-
auto next = dispatch_coro(self->ex, self->cont_op.cont);
164-
auto suicide = std::move(self->impl_ptr);
165-
next.resume();
158+
coro_resume(self);
166159
}
167160
};
168161

@@ -299,23 +292,19 @@ struct uring_dgram_recv_op : io_uring_op
299292
std::uint32_t /*bytes*/, std::uint32_t /*error*/) noexcept
300293
{
301294
auto* self = static_cast<uring_dgram_recv_op*>(base);
302-
self->stop_cb.reset();
303-
304-
if (owner == nullptr)
305-
{
306-
auto suicide = std::move(self->impl_ptr);
295+
if (coro_drain_if_shutdown(owner, self))
307296
return;
308-
}
309297

310-
if (self->ec_out)
311-
{
312-
if (self->cancelled.load(std::memory_order_acquire))
313-
*self->ec_out = capy::error::canceled;
314-
else if (self->res < 0)
315-
*self->ec_out = make_err(-self->res);
316-
else
317-
*self->ec_out = {}; // zero-byte datagram is success, not EOF
318-
}
298+
if (self->sched_)
299+
self->sched_->reset_inline_budget();
300+
301+
// Datagram recv: a 0-byte datagram is success, not EOF — is_read
302+
// stays false so the shared decode never maps it to end_of_file.
303+
decode_io_result(
304+
self->ec_out,
305+
self->cancelled.load(std::memory_order_acquire),
306+
self->res < 0 ? make_err(-self->res) : std::error_code{},
307+
/*is_read=*/false, /*bytes=*/0, /*empty_buffer=*/false);
319308
if (self->bytes_out)
320309
*self->bytes_out = (self->res >= 0)
321310
? static_cast<std::size_t>(self->res) : 0;
@@ -332,10 +321,7 @@ struct uring_dgram_recv_op : io_uring_op
332321
self->source_writer(self->source_writer_ctx,
333322
self->source_storage, self->source_len);
334323

335-
self->cont_op.cont.h = self->h;
336-
auto next = dispatch_coro(self->ex, self->cont_op.cont);
337-
auto suicide = std::move(self->impl_ptr);
338-
next.resume();
324+
coro_resume(self);
339325
}
340326
};
341327

0 commit comments

Comments
 (0)