diff --git a/doc/.vale/styles/config/vocabularies/Corosio/accept.txt b/doc/.vale/styles/config/vocabularies/Corosio/accept.txt
index 3275f5036..631a2772a 100644
--- a/doc/.vale/styles/config/vocabularies/Corosio/accept.txt
+++ b/doc/.vale/styles/config/vocabularies/Corosio/accept.txt
@@ -89,6 +89,8 @@
(?i)kqueue
(?i)io_uring
(?i)iovec
+(?i)ttys?
+(?i)inodes?
(?i)syscalls?
(?i)datagrams?
(?i)wakeups?
@@ -133,6 +135,7 @@
# --- Coined adjectives and nouns the guide uses ---------------------------------
(?i)joinable
+(?i)pollable
(?i)schedulable
(?i)launchable
(?i)buildable
@@ -278,6 +281,7 @@
(?i)wolfssl
(?i)backend[’']s
(?i)scheduler[’']s
+(?i)select[’']s
(?i)IP[’']s
(?i)UDP[’']s
(?i)URL[’']s
diff --git a/doc/modules/ROOT/nav.adoc b/doc/modules/ROOT/nav.adoc
index 2005b0d12..b52330919 100644
--- a/doc/modules/ROOT/nav.adoc
+++ b/doc/modules/ROOT/nav.adoc
@@ -40,6 +40,7 @@
** xref:4.guide/4p.unix-sockets.adoc[Unix Domain Sockets]
** xref:4.guide/4q.udp.adoc[UDP Sockets]
** xref:4.guide/4r.wait.adoc[Readiness Wait]
+** xref:4.guide/4s.native-descriptors.adoc[Native Descriptors]
* xref:5.testing/5.intro.adoc[Testing]
** xref:5.testing/5a.mocket.adoc[Mock Sockets]
** xref:5.testing/5b.socket-pair.adoc[Socket Pairs]
diff --git a/doc/modules/ROOT/pages/4.guide/4o.file-io.adoc b/doc/modules/ROOT/pages/4.guide/4o.file-io.adoc
index dbbd52107..93b8d0247 100644
--- a/doc/modules/ROOT/pages/4.guide/4o.file-io.adoc
+++ b/doc/modules/ROOT/pages/4.guide/4o.file-io.adoc
@@ -88,7 +88,7 @@ Both file types accept a bitmask of `file_base::flags` when opening:
| `create` | Create the file if it does not exist
| `exclusive` | Fail if the file already exists (requires `create`)
| `truncate` | Truncate the file to zero length on open
-| `append` | Seek to end on open (stream_file only)
+| `append` | Seek to end on open (cpp:stream_file[] only)
| `sync_all_on_write` | Synchronize data to disk on each write
|===
@@ -137,6 +137,20 @@ file object is already associated and cannot be re-adopted there. On
POSIX platforms no such restriction exists.
====
+On POSIX, `assign()` accepts only what a file object can position:
+regular files, block devices, and character devices. A pipe or socket is
+rejected with `errc::operation_not_supported` — adopt those into a
+xref:4.guide/4s.native-descriptors.adoc[`posix_descriptor`] instead; a
+directory is not adoptable by either type.
+
+A character device that cannot seek, such as a tty, passes adoption on
+every backend. What happens next is backend-specific: the POSIX
+backends issue `preadv`/`pwritev` and fail at the first read or write
+with `ESPIPE`. io_uring submits `READV`/`WRITEV` at offset `-1` for a
+`stream_file` and reads the tty successfully. Adopt one into a
+xref:4.guide/4s.native-descriptors.adoc[`posix_descriptor`] if you want
+the same behavior everywhere.
+
== Error Handling
File operations follow the same error model as sockets. Reads past
diff --git a/doc/modules/ROOT/pages/4.guide/4r.wait.adoc b/doc/modules/ROOT/pages/4.guide/4r.wait.adoc
index f1bb0c514..8b240a82b 100644
--- a/doc/modules/ROOT/pages/4.guide/4r.wait.adoc
+++ b/doc/modules/ROOT/pages/4.guide/4r.wait.adoc
@@ -70,6 +70,12 @@ library's socket directly and `release()` it before the library needs
exclusive ownership again. Or create a true duplicate with
`WSADuplicateSocketW` and adopt that.
+This applies to sockets. A library that owns a non-socket descriptor —
+a pipe, a tty, `inotify`, `eventfd` — needs the same readiness
+notification.
+xref:4.guide/4s.native-descriptors.adoc[`posix_descriptor`] provides it
+with the same `wait()` on the same terms.
+
== Acceptors
cpp:tcp_acceptor[] and cpp:local_stream_acceptor[] expose the same `wait()`.
@@ -150,3 +156,4 @@ uniform across platforms.
* xref:4.guide/4d.sockets.adoc[Sockets]
* xref:4.guide/4e.tcp-acceptor.adoc[Acceptors]
* xref:4.guide/4q.udp.adoc[UDP Sockets]
+* xref:4.guide/4s.native-descriptors.adoc[Native Descriptors]
diff --git a/doc/modules/ROOT/pages/4.guide/4s.native-descriptors.adoc b/doc/modules/ROOT/pages/4.guide/4s.native-descriptors.adoc
new file mode 100644
index 000000000..6c30e41da
--- /dev/null
+++ b/doc/modules/ROOT/pages/4.guide/4s.native-descriptors.adoc
@@ -0,0 +1,224 @@
+//
+// Copyright (c) 2026 Michael Vandeberg
+//
+// Distributed under the Boost Software License, Version 1.0. (See accompanying
+// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
+//
+// Official repository: https://github.com/cppalliance/corosio
+//
+
+= Native Descriptors
+:page-mode: how-to
+
+cpp:posix_descriptor[] adopts a file descriptor you already have and
+drives it from an `io_context`, giving it `read_some()`, `write_some()`
+and `wait()`. It exists on POSIX platforms only.
+
+[NOTE]
+====
+Code snippets assume:
+[source,cpp]
+----
+include::example$snippets/4s_native_descriptors.cpp[tag=assume]
+----
+====
+
+== Overview
+
+Corosio calls a descriptor pollable when a reactor can wait on it for
+readiness. Anything pollable that corosio does not already wrap is in
+scope. That includes character devices, `inotify`, `eventfd`,
+`timerfd`, `pidfd`, pipes, ttys, and socket kinds that have no
+dedicated corosio type.
+
+The type never creates a descriptor. Open it with whichever platform
+call suits it — `eventfd()`, `inotify_init1()`, `open()` on a device
+node — and hand the result to `assign()`. corosio supplies the event
+loop, not the constructor.
+
+== Why the Name Says POSIX
+
+A single portable `native_descriptor` spanning POSIX and Windows was
+considered and rejected. The platform gaps here are not edge cases
+around a shared core; they are the type's semantics. `O_NONBLOCK` on a
+shared open file description, `dup()`, and a file-type reject list
+expressed in `st_mode` bits are the entire contract below. An open
+file description is the kernel-side object a descriptor refers to.
+Every `dup()` of a descriptor shares the same one. None of them has a
+Windows counterpart. A type that named both would have to either
+document each rule twice or say nothing precise about either.
+
+Portability lives one layer up instead. A cpp:posix_descriptor[] is an
+cpp:io_object[], cpp:io_read_stream[], cpp:io_write_stream[] and
+cpp:io_stream[] — the same bases `tcp_socket` has — and, like every
+corosio stream, it satisfies `capy::Stream`. That concept is what the
+generic algorithms are written against, so they run on a descriptor, a
+socket and a `tls_stream` alike:
+
+[source,cpp]
+----
+include::example$snippets/4s_native_descriptors.cpp[tag=layering,indent=0]
+----
+
+Only the handful of lines that produce the descriptor are
+platform-specific. xref:4.guide/4l.tls.adoc[TLS] layers over it on the
+same terms.
+
+== Adopting a Descriptor
+
+`assign()` takes ownership: `close()` and the destructor close the
+descriptor. It returns a `std::error_code` and is pass:[[[nodiscard]]]; a
+failed `assign()` leaves the descriptor with you, so close it yourself.
+
+An `eventfd` as a cross-thread wakeup:
+
+[source,cpp]
+----
+include::example$snippets/4s_native_descriptors.cpp[tag=adopt_eventfd,indent=0]
+----
+
+[NOTE]
+====
+Distinct cpp:posix_descriptor[] objects are safe to use from different
+threads. A shared object must not run two operations of the same kind
+at once. One read and one write may overlap.
+====
+
+An `inotify` watch. The descriptor is a stream of variable-length
+records, so `read_some()` is the whole interface you need:
+
+[source,cpp]
+----
+include::example$snippets/4s_native_descriptors.cpp[tag=adopt_inotify,indent=0]
+----
+
+`release()` hands the descriptor back, cancelling pending operations
+and leaving the object not-open.
+
+== Ownership and the `dup()` Rule
+
+Another party sometimes owns the descriptor — a C library that does
+its own I/O on it, or a process-wide descriptor such as
+`STDIN_FILENO`. When that happens, adopt a `dup()` of it rather than
+the descriptor itself:
+
+[source,cpp]
+----
+include::example$snippets/4s_native_descriptors.cpp[tag=dup_for_foreign_fd,indent=0]
+----
+
+Both descriptors refer to one open file description, so the duplicate
+reports exactly the original's readiness. Corosio closing the
+duplicate can never close the original. This is the same rule
+xref:4.guide/4r.wait.adoc[Readiness Wait] states for adopted sockets.
+
+== Descriptor Flags
+
+[WARNING]
+====
+`O_NONBLOCK` is set on the first `read_some()` or `write_some()`, never
+by `assign()`, and it is never restored.
+
+The flag lives on the shared open file description, not on the
+descriptor, so every other holder of that description sees it.
+Restoring it on close would race whoever else is holding it. Permanent
+is the only safe choice.
+
+A `dup()` does not shield the other holder from this: the duplicate
+shares the same description, so the flag change reaches them anyway. It
+separates the lifetimes, nothing more. When another party owns the
+descriptor and cannot tolerate `O_NONBLOCK`, the way out is `wait()`
+and doing the I/O yourself. `wait()` never modifies the descriptor at
+all, flags included.
+====
+
+That is what makes standard input safe to adopt for readiness alone.
+Flipping `O_NONBLOCK` on it would change the terminal the parent shell
+is still using.
+
+[source,cpp]
+----
+include::example$snippets/4s_native_descriptors.cpp[tag=wait_only,indent=0]
+----
+
+== What Is Rejected
+
+`assign()` rejects regular files, block devices, and directories with a
+code comparing equal to `errc::operation_not_supported`. A reactor
+cannot report readiness for them, and they already have a home:
+xref:4.guide/4o.file-io.adoc[`stream_file` and `random_access_file`]
+adopt exactly those kinds.
+
+The test is a reject list, not an accept list. The reason: the
+flagship descriptor kinds — `eventfd`, `timerfd`, `inotify`, `pidfd` —
+are anonymous inodes whose `st_mode` type bits are all zero. An accept list
+would reject the descriptors this type exists to carry.
+
+A negative or closed descriptor fails with `errc::bad_file_descriptor`,
+and re-assigning the descriptor the object already holds fails with
+`errc::invalid_argument`.
+
+== Where Errors Surface
+
+Validation runs before anything is mutated. A rejected descriptor
+leaves the object holding whatever it held before — pending operations
+included — and leaves you owning the descriptor.
+
+A refusal from the *kernel* is the one exception, and where it appears
+depends on the backend:
+
+epoll, kqueue, select:: These register the descriptor with the reactor
+during `assign()`, so a refusal fails `assign()`. The old descriptor
+has already been closed by then, so the object is left closed. On
+select this is reachable in normal use: `select()` cannot monitor a
+descriptor at or above `FD_SETSIZE`, and such a descriptor is rejected
+with `EMFILE`.
+
+io_uring:: There is no adopt-time registration syscall, so `assign()`
+succeeds and a refusal appears at the first operation instead.
+
+Whichever way `assign()` fails, the descriptor you passed is still
+yours to close.
+
+Not every character device can be adopted. `/dev/null`, `/dev/zero` and
+`/dev/urandom` are not pollable: epoll refuses them with `EPERM` and
+kqueue with `EINVAL`, so `assign()` fails on those backends. select and
+io_uring have nothing to refuse them with, so `assign()` succeeds and the
+descriptor works — a `read_some()` on `/dev/zero` returns zeros. Character
+devices backed by a real driver, a tty among them, are pollable and
+adopt everywhere.
+
+`wait(wait_type::error)` is the one verb that is not uniform. epoll and
+io_uring report a pipe or FIFO hangup as an error condition and name a
+code. kqueue and select do not. kqueue raises an error event only for
+`EV_ERROR` or for `EV_EOF` with `fflags != 0`, and a hangup sets
+neither. select's exceptional set does not cover it. On those two backends
+the wait never completes; end it with `cancel()` or a stop token. Prefer
+`wait(wait_type::read)`, which is uniform — the hangup surfaces there as
+readiness, and the read that follows names the real failure.
+
+== `SIGPIPE`
+
+[WARNING]
+====
+Writing to a descriptor whose peer has closed raises `SIGPIPE` in the
+default disposition, which terminates the process. The socket types
+suppress this; cpp:posix_descriptor[] cannot.
+
+The suppression sockets get has no general form. `MSG_NOSIGNAL` is a
+`send()` flag and there is no `writev()` equivalent. `SO_NOSIGPIPE`
+is a socket option. Neither applies to an arbitrary descriptor.
+
+Install `SIG_IGN` for `SIGPIPE` — or handle it through a
+xref:4.guide/4i.signals.adoc[`signal_set`] — before writing to an
+adopted descriptor. The write then fails with `EPIPE` instead.
+====
+
+Asio's `posix::stream_descriptor` behaves the same way, for the same
+reason. Code ported from it needs no change here.
+
+== See Also
+
+* xref:4.guide/4r.wait.adoc[Readiness Wait]
+* xref:4.guide/4o.file-io.adoc[File I/O]
+* xref:reference:boost/corosio/posix_descriptor.adoc[`posix_descriptor` reference]
diff --git a/include/boost/corosio.hpp b/include/boost/corosio.hpp
index a6bcb41aa..4337cc84e 100644
--- a/include/boost/corosio.hpp
+++ b/include/boost/corosio.hpp
@@ -1,5 +1,6 @@
//
// Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com)
+// Copyright (c) 2026 Michael Vandeberg
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
@@ -22,6 +23,7 @@
#include
#include
#include
+#include // POSIX-only; self-guarded
#include
#include
#include
diff --git a/include/boost/corosio/backend.hpp b/include/boost/corosio/backend.hpp
index 1bfdbac39..1cb1caf07 100644
--- a/include/boost/corosio/backend.hpp
+++ b/include/boost/corosio/backend.hpp
@@ -1,5 +1,6 @@
//
// Copyright (c) 2026 Steve Gerbino
+// Copyright (c) 2026 Michael Vandeberg
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
@@ -39,6 +40,8 @@ class epoll_local_stream_acceptor;
class epoll_local_stream_acceptor_service;
class epoll_local_datagram_socket;
class epoll_local_datagram_service;
+class epoll_descriptor;
+class epoll_descriptor_service;
class epoll_scheduler;
class posix_signal;
@@ -84,6 +87,11 @@ struct epoll_t
/// The service that owns the Unix domain datagram implementations.
using local_datagram_service_type = detail::epoll_local_datagram_service;
+ /// The concrete adopted-descriptor type.
+ using descriptor_type = detail::epoll_descriptor;
+ /// The service that owns the descriptor implementations.
+ using descriptor_service_type = detail::epoll_descriptor_service;
+
/// The concrete signal set type.
using signal_type = detail::posix_signal;
/// The service that owns the signal set implementations.
@@ -141,6 +149,8 @@ class select_local_stream_acceptor;
class select_local_stream_acceptor_service;
class select_local_datagram_socket;
class select_local_datagram_service;
+class select_descriptor;
+class select_descriptor_service;
class select_scheduler;
class posix_signal;
@@ -186,6 +196,11 @@ struct select_t
/// The service that owns the Unix domain datagram implementations.
using local_datagram_service_type = detail::select_local_datagram_service;
+ /// The concrete adopted-descriptor type.
+ using descriptor_type = detail::select_descriptor;
+ /// The service that owns the descriptor implementations.
+ using descriptor_service_type = detail::select_descriptor_service;
+
/// The concrete signal set type.
using signal_type = detail::posix_signal;
/// The service that owns the signal set implementations.
@@ -243,6 +258,8 @@ class kqueue_local_stream_acceptor;
class kqueue_local_stream_acceptor_service;
class kqueue_local_datagram_socket;
class kqueue_local_datagram_service;
+class kqueue_descriptor;
+class kqueue_descriptor_service;
class kqueue_scheduler;
class posix_signal;
@@ -288,6 +305,11 @@ struct kqueue_t
/// The service that owns the Unix domain datagram implementations.
using local_datagram_service_type = detail::kqueue_local_datagram_service;
+ /// The concrete adopted-descriptor type.
+ using descriptor_type = detail::kqueue_descriptor;
+ /// The service that owns the descriptor implementations.
+ using descriptor_service_type = detail::kqueue_descriptor_service;
+
/// The concrete signal set type.
using signal_type = detail::posix_signal;
/// The service that owns the signal set implementations.
@@ -345,6 +367,8 @@ class uring_local_stream_acceptor;
class uring_local_stream_acceptor_service;
class uring_local_datagram_socket;
class uring_local_datagram_service;
+class uring_descriptor;
+class uring_descriptor_service;
class uring_stream_file;
class uring_stream_file_service;
class uring_random_access_file;
@@ -390,6 +414,11 @@ struct uring_t
/// The service that owns the Unix domain datagram implementations.
using local_datagram_service_type = detail::uring_local_datagram_service;
+ /// The concrete adopted-descriptor type.
+ using descriptor_type = detail::uring_descriptor;
+ /// The service that owns the descriptor implementations.
+ using descriptor_service_type = detail::uring_descriptor_service;
+
/// The concrete signal set type.
using signal_type = detail::posix_signal;
/// The service that owns the signal set implementations.
diff --git a/include/boost/corosio/detail/descriptor_service.hpp b/include/boost/corosio/detail/descriptor_service.hpp
new file mode 100644
index 000000000..b3aa4da28
--- /dev/null
+++ b/include/boost/corosio/detail/descriptor_service.hpp
@@ -0,0 +1,66 @@
+//
+// Copyright (c) 2026 Michael Vandeberg
+//
+// Distributed under the Boost Software License, Version 1.0. (See accompanying
+// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
+//
+// Official repository: https://github.com/cppalliance/corosio
+//
+
+#ifndef BOOST_COROSIO_DETAIL_DESCRIPTOR_SERVICE_HPP
+#define BOOST_COROSIO_DETAIL_DESCRIPTOR_SERVICE_HPP
+
+#include
+#include
+
+#if BOOST_COROSIO_POSIX
+
+#include
+#include
+
+#include
+
+namespace boost::corosio::detail {
+
+/* Abstract descriptor service base class.
+
+ Concrete implementations (epoll, select, kqueue, uring) inherit
+ from this class. The three reactor backends register the adopted
+ fd with their reactor; uring has no adopt-time registration and
+ only takes ownership of it.
+ The context constructor installs whichever backend via
+ make_service, and posix_descriptor.cpp retrieves it via
+ create_handle().
+*/
+class BOOST_COROSIO_DECL descriptor_service
+ : public capy::execution_context::service
+ , public io_object::io_service
+{
+public:
+ /// Identifies this service for execution_context lookup.
+ using key_type = descriptor_service;
+
+ /** Adopt an existing native descriptor.
+
+ Validates before mutating: on failure the implementation
+ keeps its previous descriptor and pending operations, and
+ the caller retains ownership of @a fd. On success the
+ implementation takes ownership and will close it.
+
+ @param impl The descriptor implementation to assign to.
+ @param fd The native descriptor to adopt.
+ @return Error code on failure, empty on success.
+ */
+ virtual std::error_code assign_descriptor(
+ posix_descriptor::implementation& impl, native_handle_type fd) = 0;
+
+protected:
+ descriptor_service() = default;
+ ~descriptor_service() override = default;
+};
+
+} // namespace boost::corosio::detail
+
+#endif // BOOST_COROSIO_POSIX
+
+#endif
diff --git a/include/boost/corosio/native/detail/epoll/epoll_scheduler.hpp b/include/boost/corosio/native/detail/epoll/epoll_scheduler.hpp
index 574c25665..b8d05a934 100644
--- a/include/boost/corosio/native/detail/epoll/epoll_scheduler.hpp
+++ b/include/boost/corosio/native/detail/epoll/epoll_scheduler.hpp
@@ -387,7 +387,29 @@ epoll_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
auto* desc =
static_cast(event_buffer_[i].data.ptr);
- desc->add_ready_events(event_buffer_[i].events);
+
+ // A pipe or tty whose peer closed reports EPOLLHUP on its own --
+ // no EPOLLIN, no EPOLLERR -- and EPOLLHUP maps to no
+ // reactor_event_* bit, so invoke_deferred_io() would take no
+ // branch and, the registration being edge-triggered, never get
+ // another chance.
+ //
+ // Sockets are unaffected because they never report EPOLLHUP
+ // alone. tcp_poll() and unix_poll() raise it only once
+ // sk_shutdown is SHUTDOWN_MASK (or the state is TCP_CLOSE), and
+ // both also report the socket readable and writable there --
+ // tcp_poll() takes an explicit `else mask |= EPOLLOUT` branch
+ // once SEND_SHUTDOWN is set, because a send on a shut-down
+ // socket fails fast rather than blocking. Measured on every
+ // state that produces EPOLLHUP -- peer close plus local
+ // SHUT_WR/SHUT_RDWR, RST, RST with the send buffer full, and
+ // the AF_UNIX equivalents -- the mask is always IN|OUT|HUP
+ // (0x15), or IN|OUT|ERR|HUP (0x1d) for a reset. The forced bits
+ // are therefore already set on every socket path.
+ std::uint32_t ev = event_buffer_[i].events;
+ if (ev & EPOLLHUP)
+ ev |= EPOLLIN | EPOLLOUT;
+ desc->add_ready_events(ev);
bool expected = false;
if (desc->is_enqueued_.compare_exchange_strong(
diff --git a/include/boost/corosio/native/detail/epoll/epoll_traits.hpp b/include/boost/corosio/native/detail/epoll/epoll_traits.hpp
index 8b57e9b0d..22d61cc1b 100644
--- a/include/boost/corosio/native/detail/epoll/epoll_traits.hpp
+++ b/include/boost/corosio/native/detail/epoll/epoll_traits.hpp
@@ -23,6 +23,8 @@
#include
#include
#include
+#include
+#include
/* epoll backend traits.
@@ -92,6 +94,39 @@ struct epoll_traits
}
};
+ // Descriptors are not sockets: sendmsg() fails with ENOTSOCK on a
+ // pipe or character device, so the write path is writev()/write()
+ // and SIGPIPE suppression is structurally unavailable -- MSG_NOSIGNAL
+ // is a send() flag and SO_NOSIGPIPE a socket option. A write to a
+ // pipe whose read end has closed raises SIGPIPE, exactly as a plain
+ // write(2) would; callers install SIG_IGN.
+ struct descriptor_write_policy
+ {
+ static ssize_t write(int fd, iovec* iovecs, int count) noexcept
+ {
+ ssize_t n;
+ do
+ {
+ n = ::writev(fd, iovecs, count);
+ }
+ while (n < 0 && errno == EINTR);
+ return n;
+ }
+
+ // Single-buffer fast path: skips the kernel's iov_iter setup.
+ static ssize_t
+ write_one(int fd, void const* data, std::size_t size) noexcept
+ {
+ ssize_t n;
+ do
+ {
+ n = ::write(fd, data, size);
+ }
+ while (n < 0 && errno == EINTR);
+ return n;
+ }
+ };
+
struct accept_policy
{
static int
diff --git a/include/boost/corosio/native/detail/epoll/epoll_types.hpp b/include/boost/corosio/native/detail/epoll/epoll_types.hpp
index 5ea24afba..ad1623172 100644
--- a/include/boost/corosio/native/detail/epoll/epoll_types.hpp
+++ b/include/boost/corosio/native/detail/epoll/epoll_types.hpp
@@ -25,6 +25,7 @@
#include
#include
#include
+#include
namespace boost::corosio::detail {
@@ -41,6 +42,8 @@ class epoll_local_stream_acceptor;
class epoll_local_stream_acceptor_service;
class epoll_local_datagram_socket;
class epoll_local_datagram_service;
+class epoll_descriptor;
+class epoll_descriptor_service;
// --- Stream sockets ---
@@ -235,6 +238,29 @@ class epoll_local_stream_acceptor final
}
};
+// --- Descriptors ---
+
+class epoll_descriptor final
+ : public reactor_descriptor<
+ epoll_descriptor,
+ epoll_traits,
+ epoll_descriptor_service,
+ epoll_tcp_acceptor>
+{
+ using base_type = reactor_descriptor<
+ epoll_descriptor,
+ epoll_traits,
+ epoll_descriptor_service,
+ epoll_tcp_acceptor>;
+ friend epoll_descriptor_service;
+
+public:
+ explicit epoll_descriptor(epoll_descriptor_service& svc) noexcept
+ : base_type(svc)
+ {
+ }
+};
+
// --- Services ---
class BOOST_COROSIO_DECL epoll_tcp_service final
@@ -351,6 +377,24 @@ class BOOST_COROSIO_DECL epoll_local_stream_acceptor_service final
}
};
+class BOOST_COROSIO_DECL epoll_descriptor_service final
+ : public reactor_descriptor_service<
+ epoll_descriptor_service,
+ epoll_traits,
+ epoll_descriptor>
+{
+ using base_type = reactor_descriptor_service<
+ epoll_descriptor_service,
+ epoll_traits,
+ epoll_descriptor>;
+
+public:
+ explicit epoll_descriptor_service(capy::execution_context& ctx)
+ : base_type(ctx)
+ {
+ }
+};
+
} // namespace boost::corosio::detail
#endif // BOOST_COROSIO_HAS_EPOLL
diff --git a/include/boost/corosio/native/detail/kqueue/kqueue_traits.hpp b/include/boost/corosio/native/detail/kqueue/kqueue_traits.hpp
index cfbb163a4..62114097f 100644
--- a/include/boost/corosio/native/detail/kqueue/kqueue_traits.hpp
+++ b/include/boost/corosio/native/detail/kqueue/kqueue_traits.hpp
@@ -103,6 +103,39 @@ struct kqueue_traits
}
};
+ // Descriptors are not sockets: sendmsg() fails with ENOTSOCK on a
+ // pipe or character device, so the write path is writev()/write()
+ // and SIGPIPE suppression is structurally unavailable -- MSG_NOSIGNAL
+ // is a send() flag and SO_NOSIGPIPE a socket option. A write to a
+ // pipe whose read end has closed raises SIGPIPE, exactly as a plain
+ // write(2) would; callers install SIG_IGN.
+ struct descriptor_write_policy
+ {
+ static ssize_t write(int fd, iovec* iovecs, int count) noexcept
+ {
+ ssize_t n;
+ do
+ {
+ n = ::writev(fd, iovecs, count);
+ }
+ while (n < 0 && errno == EINTR);
+ return n;
+ }
+
+ // Single-buffer fast path: skips the kernel's iov_iter setup.
+ static ssize_t
+ write_one(int fd, void const* data, std::size_t size) noexcept
+ {
+ ssize_t n;
+ do
+ {
+ n = ::write(fd, data, size);
+ }
+ while (n < 0 && errno == EINTR);
+ return n;
+ }
+ };
+
struct accept_policy
{
static int
diff --git a/include/boost/corosio/native/detail/kqueue/kqueue_types.hpp b/include/boost/corosio/native/detail/kqueue/kqueue_types.hpp
index 071a2f128..3e18accc2 100644
--- a/include/boost/corosio/native/detail/kqueue/kqueue_types.hpp
+++ b/include/boost/corosio/native/detail/kqueue/kqueue_types.hpp
@@ -25,6 +25,7 @@
#include
#include
#include
+#include
namespace boost::corosio::detail {
@@ -41,6 +42,8 @@ class kqueue_local_stream_acceptor;
class kqueue_local_stream_acceptor_service;
class kqueue_local_datagram_socket;
class kqueue_local_datagram_service;
+class kqueue_descriptor;
+class kqueue_descriptor_service;
// --- Stream sockets ---
@@ -238,6 +241,29 @@ class kqueue_local_stream_acceptor final
}
};
+// --- Descriptors ---
+
+class kqueue_descriptor final
+ : public reactor_descriptor<
+ kqueue_descriptor,
+ kqueue_traits,
+ kqueue_descriptor_service,
+ kqueue_tcp_acceptor>
+{
+ using base_type = reactor_descriptor<
+ kqueue_descriptor,
+ kqueue_traits,
+ kqueue_descriptor_service,
+ kqueue_tcp_acceptor>;
+ friend kqueue_descriptor_service;
+
+public:
+ explicit kqueue_descriptor(kqueue_descriptor_service& svc) noexcept
+ : base_type(svc)
+ {
+ }
+};
+
// --- Services ---
class BOOST_COROSIO_DECL kqueue_tcp_service final
@@ -358,6 +384,24 @@ class BOOST_COROSIO_DECL kqueue_local_stream_acceptor_service final
}
};
+class BOOST_COROSIO_DECL kqueue_descriptor_service final
+ : public reactor_descriptor_service<
+ kqueue_descriptor_service,
+ kqueue_traits,
+ kqueue_descriptor>
+{
+ using base_type = reactor_descriptor_service<
+ kqueue_descriptor_service,
+ kqueue_traits,
+ kqueue_descriptor>;
+
+public:
+ explicit kqueue_descriptor_service(capy::execution_context& ctx)
+ : base_type(ctx)
+ {
+ }
+};
+
} // namespace boost::corosio::detail
#endif // BOOST_COROSIO_HAS_KQUEUE
diff --git a/include/boost/corosio/native/detail/posix/posix_random_access_file.hpp b/include/boost/corosio/native/detail/posix/posix_random_access_file.hpp
index c43431f47..169c2ec35 100644
--- a/include/boost/corosio/native/detail/posix/posix_random_access_file.hpp
+++ b/include/boost/corosio/native/detail/posix/posix_random_access_file.hpp
@@ -25,6 +25,7 @@
#include
#include
#include
+#include
#include
#include
#include
@@ -269,6 +270,23 @@ posix_random_access_file::release()
inline std::error_code
posix_random_access_file::assign(native_handle_type handle) noexcept
{
+ // handle >= 0 guard: an unset impl reports native_handle() == -1, and a
+ // caller-supplied -1 must fail as a bad fd, not a self-assign.
+ if (handle >= 0 && handle == fd_)
+ return std::make_error_code(std::errc::invalid_argument);
+
+ // Validate before touching the held fd: a failed assign must leave
+ // this object unchanged and the caller still owning handle.
+ if (auto ec = validate_file_fd(handle))
+ return ec;
+
+ // cancel() first: an in-flight read_at/write_at's pool-thread
+ // completion reads fd_ at execution time, not at post time, so
+ // without this a pending op silently completes against the newly
+ // adopted file instead of being cancelled. The service's
+ // close(handle) / destroy() normally pair cancel()+close_file();
+ // assign() bypasses that path and must do the same pairing itself.
+ cancel();
close_file();
fd_ = handle;
return {};
diff --git a/include/boost/corosio/native/detail/posix/posix_stream_file.hpp b/include/boost/corosio/native/detail/posix/posix_stream_file.hpp
index 0f0d6c092..9c7d4981e 100644
--- a/include/boost/corosio/native/detail/posix/posix_stream_file.hpp
+++ b/include/boost/corosio/native/detail/posix/posix_stream_file.hpp
@@ -26,6 +26,7 @@
#include
#include
#include
+#include
#include
#include
#include
@@ -328,6 +329,23 @@ posix_stream_file::release()
inline std::error_code
posix_stream_file::assign(native_handle_type handle) noexcept
{
+ // handle >= 0 guard: an unset impl reports native_handle() == -1, and a
+ // caller-supplied -1 must fail as a bad fd, not a self-assign.
+ if (handle >= 0 && handle == fd_)
+ return std::make_error_code(std::errc::invalid_argument);
+
+ // Validate before touching the held fd: a failed assign must leave
+ // this object unchanged and the caller still owning handle.
+ if (auto ec = validate_file_fd(handle))
+ return ec;
+
+ // cancel() first: an in-flight read/write's pool-thread completion
+ // reads fd_/offset_ at execution time, not at post time, so without
+ // this a pending op silently completes against the newly adopted
+ // file instead of being cancelled. The service's close(handle) /
+ // destroy() normally pair cancel()+close_file(); assign() bypasses
+ // that path and must do the same pairing itself.
+ cancel();
close_file();
fd_ = handle;
offset_ = 0;
diff --git a/include/boost/corosio/native/detail/reactor/reactor_descriptor.hpp b/include/boost/corosio/native/detail/reactor/reactor_descriptor.hpp
new file mode 100644
index 000000000..76168ee66
--- /dev/null
+++ b/include/boost/corosio/native/detail/reactor/reactor_descriptor.hpp
@@ -0,0 +1,872 @@
+//
+// Copyright (c) 2026 Michael Vandeberg
+//
+// Distributed under the Boost Software License, Version 1.0. (See accompanying
+// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
+//
+// Official repository: https://github.com/cppalliance/corosio
+//
+
+#ifndef BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
+#define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
+
+#include
+
+#if BOOST_COROSIO_POSIX
+
+#include
+#include
+#include
+#include
+#include
+#include
+#include
+#include
+#include
+#include
+#include
+
+#include
+#include
+#include
+#include
+
+#include
+#include
+#include
+
+/* Reactor-backed implementation of posix_descriptor.
+
+ Deliberately does not derive from reactor_basic_socket: that base
+ is rooted in native_socket_base, which overrides local_endpoint(),
+ set_option() and get_option() on its ImplBase. posix_descriptor has
+ no socket verbs, so the shared logic (init_and_register, register_op,
+ the cancel/close/release op sweeps) is carried here against five op
+ slots instead of eight.
+
+ The one behavior that is genuinely new: O_NONBLOCK is armed lazily,
+ on the first read_some/write_some and never from assign() or wait().
+ The flag lives on the shared open file description, so arming it is
+ visible to every other holder of that description -- which is why a
+ wait()-only user must never trigger it.
+
+ PARALLEL COPY: init_and_register, register_op, cancel_single_op and
+ the cancel/close/release op sweeps here mirror the socket versions in
+ reactor_basic_socket.hpp (init_and_register, register_op,
+ cancel_single_op, do_cancel, do_close_socket, do_release_socket).
+ The two are separate code because that base carries native_socket_base
+ and its socket verbs; they are not separate protocols. A fix to the
+ cancel/park protocol -- slot claiming under desc_state_.mutex, the
+ cached-edge replay in register_op, the impl_ref_ pinning during
+ teardown -- belongs in both files.
+*/
+
+namespace boost::corosio::detail {
+
+// ============================================================
+// Op types
+// ============================================================
+
+/* Descriptor op family.
+
+ Mirrors reactor_stream_ops.hpp. The Acceptor parameter is a
+ placeholder: reactor_op is parameterized on both a socket and an
+ acceptor impl type, and the shared completion helpers name
+ acceptor_impl_ in a branch that a descriptor op never takes but the
+ compiler still instantiates. Passing the backend's acceptor type (as
+ reactor_dgram_socket_impl already does) keeps those helpers shared
+ rather than duplicated here.
+
+ @tparam Traits Backend traits (epoll_traits, kqueue_traits, ...).
+ @tparam Descriptor The concrete descriptor type (forward-declared).
+ @tparam Acceptor Placeholder acceptor type for the op base.
+*/
+
+template
+struct reactor_descriptor_base_op : reactor_op
+{
+ void operator()() override;
+ void cancel() noexcept override;
+};
+
+template
+struct reactor_descriptor_read_op final
+ : reactor_read_op>
+{};
+
+template
+struct reactor_descriptor_write_op final
+ : reactor_write_op<
+ reactor_descriptor_base_op,
+ typename Traits::descriptor_write_policy>
+{};
+
+template
+struct reactor_descriptor_wait_op final
+ : reactor_wait_op>
+{
+ void operator()() override;
+};
+
+// --- Deferred implementations (instantiated when Descriptor is complete) ---
+
+template
+void
+reactor_descriptor_base_op::operator()()
+{
+ complete_io_op(*this);
+}
+
+template
+void
+reactor_descriptor_base_op::cancel() noexcept
+{
+ // A descriptor op is only ever started against a descriptor impl, so
+ // the acceptor arm of the stream op's cancel() has no counterpart.
+ if (this->socket_impl_)
+ this->socket_impl_->cancel_single_op(*this);
+ else
+ this->request_cancel();
+}
+
+template
+void
+reactor_descriptor_wait_op::operator()()
+{
+ complete_wait_op(*this);
+}
+
+// ============================================================
+// Descriptor implementation
+// ============================================================
+
+/** CRTP base for reactor-backed posix_descriptor implementations.
+
+ Holds the adopted descriptor, its reactor registration state, and
+ the five op slots (read, write, and one wait per direction).
+
+ @tparam Derived The named final class (CRTP self).
+ @tparam Traits Backend traits (epoll_traits, kqueue_traits, ...).
+ @tparam Service The backend's descriptor service type.
+ @tparam Acceptor Placeholder acceptor type for the op base.
+*/
+template
+class reactor_descriptor
+ : public posix_descriptor::implementation
+ , public std::enable_shared_from_this
+ , public intrusive_list::node
+{
+ using base_op = reactor_descriptor_base_op;
+ using read_op = reactor_descriptor_read_op;
+ using write_op = reactor_descriptor_write_op;
+ using wait_op = reactor_descriptor_wait_op;
+
+protected:
+ // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
+ explicit reactor_descriptor(Service& svc) noexcept : svc_(svc) {}
+
+public:
+ ~reactor_descriptor() override = default;
+
+ /// Per-descriptor state for persistent reactor registration.
+ typename Traits::desc_state_type desc_state_;
+
+ // --- Virtual method overrides ---
+
+ std::coroutine_handle<> read_some(
+ std::coroutine_handle<> h,
+ capy::executor_ref ex,
+ buffer_param param,
+ std::stop_token token,
+ std::error_code* ec,
+ std::size_t* bytes_out) override
+ {
+ return do_read_some(h, ex, param, token, ec, bytes_out);
+ }
+
+ std::coroutine_handle<> write_some(
+ std::coroutine_handle<> h,
+ capy::executor_ref ex,
+ buffer_param param,
+ std::stop_token token,
+ std::error_code* ec,
+ std::size_t* bytes_out) override
+ {
+ return do_write_some(h, ex, param, token, ec, bytes_out);
+ }
+
+ std::coroutine_handle<> wait(
+ std::coroutine_handle<> h,
+ capy::executor_ref ex,
+ wait_type w,
+ std::stop_token token,
+ std::error_code* ec) override
+ {
+ return do_wait(h, ex, w, token, ec);
+ }
+
+ native_handle_type native_handle() const noexcept override
+ {
+ return fd_;
+ }
+
+ native_handle_type release_descriptor() noexcept override;
+
+ void cancel() noexcept override
+ {
+ do_cancel();
+ }
+
+ // --- Service-facing (non-virtual) ---
+
+ /** Adopt the fd, initialize descriptor state, and register it.
+
+ @param fd The descriptor to adopt.
+
+ @return The error if the reactor rejects the descriptor, in
+ which case the implementation is left closed and the caller
+ retains ownership of @a fd; otherwise a default constructed
+ error code.
+ */
+ std::error_code init_and_register(int fd) noexcept;
+
+ /// Close the descriptor and cancel pending operations.
+ void close_descriptor() noexcept;
+
+ /// Cancel a single pending operation, claiming it from its slot.
+ template
+ void cancel_single_op(Op& op) noexcept;
+
+private:
+ /** Arm O_NONBLOCK, once, before the first speculative syscall.
+
+ Reports an errno rather than an error_code because that is what
+ the op result model records; the round trip is lossless here
+ because fcntl only fails with codes make_err passes through.
+ */
+ int arm_nonblocking() noexcept
+ {
+ if (nonblocking_)
+ return 0;
+ if (auto ec = ensure_nonblocking(fd_))
+ return ec.value();
+ nonblocking_ = true;
+ return 0;
+ }
+
+ std::coroutine_handle<> do_read_some(
+ std::coroutine_handle<>,
+ capy::executor_ref,
+ buffer_param,
+ std::stop_token const&,
+ std::error_code*,
+ std::size_t*);
+
+ std::coroutine_handle<> do_write_some(
+ std::coroutine_handle<>,
+ capy::executor_ref,
+ buffer_param,
+ std::stop_token const&,
+ std::error_code*,
+ std::size_t*);
+
+ std::coroutine_handle<> do_wait(
+ std::coroutine_handle<>,
+ capy::executor_ref,
+ wait_type,
+ std::stop_token const&,
+ std::error_code*);
+
+ void do_cancel() noexcept;
+
+ /// Register an op with the reactor, handling cached edge events.
+ template
+ void register_op(
+ Op& op,
+ reactor_op_base*& desc_slot,
+ bool& ready_flag,
+ bool is_write_direction = false) noexcept;
+
+ /// Apply @a fn to each of the five op slots.
+ template
+ void for_each_op(Fn fn) noexcept
+ {
+ fn(rd_);
+ fn(wr_);
+ fn(wait_rd_);
+ fn(wait_wr_);
+ fn(wait_er_);
+ }
+
+ /** Claim every parked op out of its descriptor_state slot.
+
+ @param claimed Receives the claimed ops; must hold five.
+ @param teardown Also clear the cached edge flags and, if the
+ state is queued in the scheduler, pin the impl alive.
+ @param self Keepalive used by @a teardown.
+ @return The number of ops claimed.
+ */
+ int claim_parked_ops(
+ reactor_op_base** claimed,
+ bool teardown,
+ std::shared_ptr const& self) noexcept;
+
+ /// Post claimed ops to the scheduler, keeping the impl alive.
+ void post_claimed_ops(
+ reactor_op_base** claimed,
+ int count,
+ std::shared_ptr const& self) noexcept;
+
+ /// Sweep every op slot, then drop the reactor registration.
+ void quiesce() noexcept;
+
+ reactor_op_base** op_to_desc_slot(base_op& op) noexcept;
+
+ Service& svc_;
+ int fd_ = -1;
+ bool nonblocking_ = false;
+
+ read_op rd_;
+ write_op wr_;
+ wait_op wait_rd_;
+ wait_op wait_wr_;
+ wait_op wait_er_;
+};
+
+// ============================================================
+// Registration and teardown
+// ============================================================
+
+template
+std::error_code
+reactor_descriptor::init_and_register(
+ int fd) noexcept
+{
+ fd_ = fd;
+ desc_state_.fd = fd;
+ {
+ // Every slot this type owns; connect_op is deliberately absent,
+ // a descriptor has no connect operation to park there.
+ std::lock_guard lock(desc_state_.mutex);
+ desc_state_.read_op = nullptr;
+ desc_state_.write_op = nullptr;
+ desc_state_.wait_read_op = nullptr;
+ desc_state_.wait_write_op = nullptr;
+ desc_state_.wait_error_op = nullptr;
+ }
+ if (auto ec = svc_.scheduler().register_descriptor(fd, &desc_state_))
+ {
+ // Undo the partial state so a failed adopt is
+ // indistinguishable from a closed implementation.
+ fd_ = -1;
+ desc_state_.fd = -1;
+ desc_state_.registered_events = 0;
+ return ec;
+ }
+ return {};
+}
+
+template
+int
+reactor_descriptor::claim_parked_ops(
+ reactor_op_base** claimed,
+ bool teardown,
+ std::shared_ptr const& self) noexcept
+{
+ int count = 0;
+ std::lock_guard lock(desc_state_.mutex);
+ for (auto** slot :
+ {&desc_state_.read_op, &desc_state_.write_op,
+ &desc_state_.wait_read_op, &desc_state_.wait_write_op,
+ &desc_state_.wait_error_op})
+ {
+ if (auto* c = std::exchange(*slot, nullptr))
+ claimed[count++] = c;
+ }
+ if (teardown)
+ {
+ desc_state_.read_ready = false;
+ desc_state_.write_ready = false;
+
+ // Must be set under the same lock that invoke_deferred_io clears
+ // is_enqueued_ under, or the impl could be destroyed while the
+ // scheduler still holds the queued descriptor_state.
+ if (desc_state_.is_enqueued_.load(std::memory_order_acquire))
+ desc_state_.impl_ref_ = self;
+ }
+ return count;
+}
+
+template
+void
+reactor_descriptor::post_claimed_ops(
+ reactor_op_base** claimed,
+ int count,
+ std::shared_ptr const& self) noexcept
+{
+ for (int i = 0; i < count; ++i)
+ {
+ claimed[i]->impl_ptr = self;
+ svc_.post(claimed[i]);
+ svc_.work_finished();
+ }
+}
+
+template
+void
+reactor_descriptor::do_cancel() noexcept
+{
+ auto self = this->weak_from_this().lock();
+ if (!self)
+ return;
+
+ for_each_op([](auto& op) { op.request_cancel(); });
+
+ reactor_op_base* claimed[5];
+ int const count = claim_parked_ops(claimed, /*teardown=*/false, self);
+ post_claimed_ops(claimed, count, self);
+}
+
+template
+void
+reactor_descriptor::quiesce() noexcept
+{
+ auto self = this->weak_from_this().lock();
+ if (self)
+ {
+ for_each_op([](auto& op) { op.request_cancel(); });
+
+ reactor_op_base* claimed[5];
+ int const count = claim_parked_ops(claimed, /*teardown=*/true, self);
+ post_claimed_ops(claimed, count, self);
+ }
+
+ if (fd_ >= 0 && desc_state_.registered_events != 0)
+ svc_.scheduler().deregister_descriptor(fd_);
+
+ desc_state_.registered_events = 0;
+ // The next adopted fd starts from an unknown flag state.
+ nonblocking_ = false;
+}
+
+template
+void
+reactor_descriptor::
+ close_descriptor() noexcept
+{
+ quiesce();
+
+ if (fd_ >= 0)
+ {
+ ::close(fd_);
+ fd_ = -1;
+ }
+ desc_state_.fd = -1;
+}
+
+template
+native_handle_type
+reactor_descriptor::
+ release_descriptor() noexcept
+{
+ quiesce();
+
+ // Do NOT close -- the caller takes ownership.
+ native_handle_type released = fd_;
+ fd_ = -1;
+ desc_state_.fd = -1;
+ return released;
+}
+
+// ============================================================
+// Op registration and per-op cancellation
+// ============================================================
+
+template
+template
+void
+reactor_descriptor::register_op(
+ Op& op,
+ reactor_op_base*& desc_slot,
+ bool& ready_flag,
+ bool is_write_direction) noexcept
+{
+ svc_.work_started();
+
+ std::lock_guard lock(desc_state_.mutex);
+ bool io_done = false;
+ if (ready_flag)
+ {
+ ready_flag = false;
+ op.perform_io();
+ io_done = (op.errn != EAGAIN && op.errn != EWOULDBLOCK);
+ if (!io_done)
+ op.errn = 0;
+ }
+
+ if (io_done || op.cancelled.load(std::memory_order_acquire))
+ {
+ svc_.post(&op);
+ svc_.work_finished();
+ }
+ else
+ {
+ desc_slot = &op;
+
+ // Select must rebuild its fd_sets when a write-direction op
+ // is parked, so select() watches for writability. Compiled
+ // away to nothing for epoll and kqueue.
+ if constexpr (Service::needs_write_notification)
+ {
+ if (is_write_direction)
+ svc_.scheduler().notify_reactor();
+ }
+ }
+}
+
+template
+reactor_op_base**
+reactor_descriptor::op_to_desc_slot(
+ base_op& op) noexcept
+{
+ if (&op == static_cast(&rd_))
+ return &desc_state_.read_op;
+ if (&op == static_cast(&wr_))
+ return &desc_state_.write_op;
+ if (&op == static_cast(&wait_rd_))
+ return &desc_state_.wait_read_op;
+ if (&op == static_cast(&wait_wr_))
+ return &desc_state_.wait_write_op;
+ if (&op == static_cast(&wait_er_))
+ return &desc_state_.wait_error_op;
+ return nullptr;
+}
+
+template
+template
+void
+reactor_descriptor::cancel_single_op(
+ Op& op) noexcept
+{
+ auto self = this->weak_from_this().lock();
+ if (!self)
+ return;
+
+ op.request_cancel();
+
+ reactor_op_base** desc_op_ptr = op_to_desc_slot(op);
+ if (!desc_op_ptr)
+ return;
+
+ reactor_op_base* claimed = nullptr;
+ {
+ std::lock_guard lock(desc_state_.mutex);
+ if (*desc_op_ptr == &op)
+ claimed = std::exchange(*desc_op_ptr, nullptr);
+ // Not in the slot: request_cancel() above already set
+ // op.cancelled, which register_op consults before parking
+ // and the completion decode consults on delivery. Latching
+ // a descriptor flag here instead would outlive this op and
+ // cancel the next wait in the same direction.
+ }
+ if (claimed)
+ {
+ op.impl_ptr = self;
+ svc_.post(&op);
+ svc_.work_finished();
+ }
+}
+
+// ============================================================
+// I/O dispatch
+// ============================================================
+
+template
+std::coroutine_handle<>
+reactor_descriptor::do_read_some(
+ std::coroutine_handle<> h,
+ capy::executor_ref ex,
+ buffer_param param,
+ std::stop_token const& token,
+ std::error_code* ec,
+ std::size_t* bytes_out)
+{
+ auto& op = rd_;
+ op.reset();
+ op.h = h;
+ op.ex = ex;
+ op.ec_out = ec;
+ op.bytes_out = bytes_out;
+
+ // Closed-object contract: complete with bad_file_descriptor without
+ // touching the kernel or the unregistered descriptor state.
+ if (fd_ < 0)
+ {
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+ op.complete(EBADF, 0);
+ svc_.post(&op);
+ return std::noop_coroutine();
+ }
+
+ capy::mutable_buffer bufs[read_op::max_buffers];
+ op.iovec_count =
+ static_cast(param.copy_to(bufs, read_op::max_buffers));
+
+ if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
+ {
+ op.empty_buffer_read = true;
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+ op.complete(0, 0);
+ svc_.post(&op);
+ return std::noop_coroutine();
+ }
+
+ // The first transferring operation is what arms O_NONBLOCK; assign()
+ // and wait() never do.
+ if (int const nerr = arm_nonblocking())
+ {
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+ op.complete(nerr, 0);
+ svc_.post(&op);
+ return std::noop_coroutine();
+ }
+
+ for (int i = 0; i < op.iovec_count; ++i)
+ {
+ op.iovecs[i].iov_base = bufs[i].data();
+ op.iovecs[i].iov_len = bufs[i].size();
+ }
+
+ // Speculative read; the single-buffer case uses read() so the kernel
+ // skips the readv iov_iter setup.
+ ssize_t n;
+ if (op.iovec_count == 1)
+ {
+ do
+ {
+ n = ::read(fd_, bufs[0].data(), bufs[0].size());
+ }
+ while (n < 0 && errno == EINTR);
+ }
+ else
+ {
+ do
+ {
+ n = ::readv(fd_, op.iovecs, op.iovec_count);
+ }
+ while (n < 0 && errno == EINTR);
+ }
+
+ if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
+ {
+ int err = (n < 0) ? errno : 0;
+ auto bytes = (n > 0) ? static_cast(n) : std::size_t(0);
+
+ if (svc_.scheduler().try_consume_inline_budget())
+ {
+ if (err)
+ *ec = make_err(err);
+ else if (n == 0)
+ *ec = capy::error::eof;
+ else
+ *ec = {};
+ *bytes_out = bytes;
+ op.cont.h = h;
+ return dispatch_coro(ex, op.cont);
+ }
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+ op.complete(err, bytes);
+ svc_.post(&op);
+ return std::noop_coroutine();
+ }
+
+ // EAGAIN — register with reactor
+ op.fd = fd_;
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+
+ register_op(op, desc_state_.read_op, desc_state_.read_ready);
+ return std::noop_coroutine();
+}
+
+template
+std::coroutine_handle<>
+reactor_descriptor::do_write_some(
+ std::coroutine_handle<> h,
+ capy::executor_ref ex,
+ buffer_param param,
+ std::stop_token const& token,
+ std::error_code* ec,
+ std::size_t* bytes_out)
+{
+ auto& op = wr_;
+ op.reset();
+ op.h = h;
+ op.ex = ex;
+ op.ec_out = ec;
+ op.bytes_out = bytes_out;
+
+ if (fd_ < 0)
+ {
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+ op.complete(EBADF, 0);
+ svc_.post(&op);
+ return std::noop_coroutine();
+ }
+
+ capy::mutable_buffer bufs[write_op::max_buffers];
+ op.iovec_count =
+ static_cast(param.copy_to(bufs, write_op::max_buffers));
+
+ if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
+ {
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+ op.complete(0, 0);
+ svc_.post(&op);
+ return std::noop_coroutine();
+ }
+
+ if (int const nerr = arm_nonblocking())
+ {
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+ op.complete(nerr, 0);
+ svc_.post(&op);
+ return std::noop_coroutine();
+ }
+
+ for (int i = 0; i < op.iovec_count; ++i)
+ {
+ op.iovecs[i].iov_base = bufs[i].data();
+ op.iovecs[i].iov_len = bufs[i].size();
+ }
+
+ // Speculative write; the single-buffer case skips the iov_iter setup.
+ ssize_t n;
+ if (op.iovec_count == 1)
+ {
+ n = write_op::write_policy::write_one(
+ fd_, bufs[0].data(), bufs[0].size());
+ }
+ else
+ {
+ n = write_op::write_policy::write(fd_, op.iovecs, op.iovec_count);
+ }
+
+ if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
+ {
+ int err = (n < 0) ? errno : 0;
+ auto bytes = (n > 0) ? static_cast(n) : std::size_t(0);
+
+ if (svc_.scheduler().try_consume_inline_budget())
+ {
+ *ec = err ? make_err(err) : std::error_code{};
+ *bytes_out = bytes;
+ op.cont.h = h;
+ return dispatch_coro(ex, op.cont);
+ }
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+ op.complete(err, bytes);
+ svc_.post(&op);
+ return std::noop_coroutine();
+ }
+
+ // EAGAIN — register with reactor
+ op.fd = fd_;
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+
+ register_op(op, desc_state_.write_op, desc_state_.write_ready, true);
+ return std::noop_coroutine();
+}
+
+template
+std::coroutine_handle<>
+reactor_descriptor::do_wait(
+ std::coroutine_handle<> h,
+ capy::executor_ref ex,
+ wait_type w,
+ std::stop_token const& token,
+ std::error_code* ec)
+{
+ // Pick refs up-front to avoid duplicating the register_op call.
+ wait_op* op_ptr;
+ reactor_op_base** desc_slot_ptr;
+ std::uint32_t event;
+
+ if (w == wait_type::read)
+ {
+ op_ptr = &wait_rd_;
+ desc_slot_ptr = &desc_state_.wait_read_op;
+ event = reactor_event_read;
+ }
+ else if (w == wait_type::write)
+ {
+ op_ptr = &wait_wr_;
+ desc_slot_ptr = &desc_state_.wait_write_op;
+ event = reactor_event_write;
+ }
+ else // wait_type::error
+ {
+ op_ptr = &wait_er_;
+ desc_slot_ptr = &desc_state_.wait_error_op;
+ event = reactor_event_error;
+ }
+
+ auto& op = *op_ptr;
+
+ // Speculative probe: an edge-triggered reactor cannot report a
+ // condition that already holds, so a wait initiated on an already
+ // ready descriptor would otherwise park forever. No syscall here
+ // modifies the descriptor -- in particular O_NONBLOCK is untouched.
+ int perr = 0;
+ if (wait_op::probe(fd_, event, perr))
+ {
+ if (svc_.scheduler().try_consume_inline_budget())
+ {
+ *ec = perr ? make_err(perr) : std::error_code{};
+ op.cont.h = h;
+ return dispatch_coro(ex, op.cont);
+ }
+ op.reset();
+ op.wait_event = event;
+ op.h = h;
+ op.ex = ex;
+ op.ec_out = ec;
+ op.fd = fd_;
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+ op.complete(perr, 0);
+ svc_.post(&op);
+ return std::noop_coroutine();
+ }
+
+ op.reset();
+ op.wait_event = event;
+ op.h = h;
+ op.ex = ex;
+ op.ec_out = ec;
+ op.fd = fd_;
+ op.start(token, static_cast(this));
+ op.impl_ptr = this->shared_from_this();
+
+ // Force register_op's ready path so the wait op re-probes under the
+ // descriptor mutex before parking. A stale write_ready latched at
+ // registration would otherwise report a full pipe as writable.
+ bool force_probe = true;
+ register_op(op, *desc_slot_ptr, force_probe, event == reactor_event_write);
+ return std::noop_coroutine();
+}
+
+} // namespace boost::corosio::detail
+
+#endif // BOOST_COROSIO_POSIX
+
+#endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
diff --git a/include/boost/corosio/native/detail/reactor/reactor_descriptor_service.hpp b/include/boost/corosio/native/detail/reactor/reactor_descriptor_service.hpp
new file mode 100644
index 000000000..d1ae1a617
--- /dev/null
+++ b/include/boost/corosio/native/detail/reactor/reactor_descriptor_service.hpp
@@ -0,0 +1,171 @@
+//
+// Copyright (c) 2026 Michael Vandeberg
+//
+// Distributed under the Boost Software License, Version 1.0. (See accompanying
+// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
+//
+// Official repository: https://github.com/cppalliance/corosio
+//
+
+#ifndef BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_SERVICE_HPP
+#define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_SERVICE_HPP
+
+#include
+
+#if BOOST_COROSIO_POSIX
+
+#include
+#include
+#include
+#include
+#include
+#include
+
+#include
+#include
+#include
+
+/* Reactor-backed descriptor_service.
+
+ assign_descriptor is the validate-before-mutate core the public
+ assign() contract rests on, modelled on do_assign_fd in
+ reactor_service_finals.hpp.
+
+ PARALLEL COPY: construct, destroy, close and shutdown here mirror the
+ same four members of reactor_socket_service.hpp, which this cannot
+ reuse because it calls close_socket() by name. A fix to the service
+ lifecycle -- the construct/destroy bookkeeping under state_->mutex_,
+ or shutdown's deliberate retention of impl_ptrs_ so impls outlive the
+ scheduler's drain -- belongs in both files.
+*/
+
+namespace boost::corosio::detail {
+
+/** CRTP base for reactor-backed descriptor services.
+
+ @tparam Derived The named final service type (CRTP self).
+ @tparam Traits Backend traits (epoll_traits, kqueue_traits, ...).
+ @tparam DescFinal The named final descriptor impl type.
+*/
+template
+class reactor_descriptor_service : public descriptor_service
+{
+ using scheduler_type = typename Traits::scheduler_type;
+ using state_type = reactor_service_state;
+
+ friend Derived;
+
+protected:
+ // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
+ explicit reactor_descriptor_service(capy::execution_context& ctx)
+ : state_(
+ std::make_unique(
+ ctx.template use_service()))
+ {
+ }
+
+public:
+ /// True when a parked write-direction op must wake the reactor.
+ static constexpr bool needs_write_notification =
+ Traits::needs_write_notification;
+
+ ~reactor_descriptor_service() override = default;
+
+ std::error_code assign_descriptor(
+ posix_descriptor::implementation& impl, native_handle_type fd) override;
+
+ void shutdown() override
+ {
+ std::lock_guard lock(state_->mutex_);
+
+ while (auto* impl = state_->impl_list_.pop_front())
+ impl->close_descriptor();
+
+ // Don't clear impl_ptrs_ here: the scheduler shuts down after us
+ // and drains completed_ops_, so every impl must outlive that.
+ }
+
+ io_object::implementation* construct() override
+ {
+ auto impl = std::make_shared(static_cast(*this));
+ auto* raw = impl.get();
+
+ {
+ std::lock_guard lock(state_->mutex_);
+ state_->impl_ptrs_.emplace(raw, std::move(impl));
+ state_->impl_list_.push_back(raw);
+ }
+
+ return raw;
+ }
+
+ void destroy(io_object::implementation* impl) override
+ {
+ auto* typed = static_cast(impl);
+ typed->close_descriptor();
+ std::lock_guard lock(state_->mutex_);
+ state_->impl_list_.remove(typed);
+ state_->impl_ptrs_.erase(typed);
+ }
+
+ void close(io_object::handle& h) override
+ {
+ static_cast(h.get())->close_descriptor();
+ }
+
+ scheduler_type& scheduler() const noexcept
+ {
+ return state_->sched_;
+ }
+
+ void post(scheduler_op* op)
+ {
+ state_->sched_.post(op);
+ }
+
+ void work_started() noexcept
+ {
+ state_->sched_.work_started();
+ }
+
+ void work_finished() noexcept
+ {
+ state_->sched_.work_finished();
+ }
+
+protected:
+ std::unique_ptr state_;
+
+private:
+ reactor_descriptor_service(reactor_descriptor_service const&) = delete;
+ reactor_descriptor_service&
+ operator=(reactor_descriptor_service const&) = delete;
+};
+
+template
+std::error_code
+reactor_descriptor_service::assign_descriptor(
+ posix_descriptor::implementation& impl_base, native_handle_type fd)
+{
+ auto* impl = static_cast(&impl_base);
+
+ // fd >= 0 guard: an unset impl reports native_handle() == -1, and a
+ // caller-supplied -1 must fail as a bad fd, not a self-assign.
+ if (fd >= 0 && fd == impl->native_handle())
+ return std::make_error_code(std::errc::invalid_argument);
+
+ // Validate before touching the held descriptor: a failed assign
+ // must leave the object unchanged and the caller owning the fd.
+ if (auto ec = validate_descriptor_fd(fd))
+ return ec;
+
+ impl->close_descriptor();
+
+ return impl->init_and_register(fd);
+}
+
+} // namespace boost::corosio::detail
+
+#endif // BOOST_COROSIO_POSIX
+
+#endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_SERVICE_HPP
diff --git a/include/boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp b/include/boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp
index 4dda0eb76..c8a8fddd4 100644
--- a/include/boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp
+++ b/include/boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp
@@ -1,5 +1,6 @@
//
// Copyright (c) 2026 Steve Gerbino
+// Copyright (c) 2026 Michael Vandeberg
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
@@ -152,7 +153,31 @@ reactor_descriptor_state::invoke_deferred_io()
{
socklen_t len = sizeof(err);
if (::getsockopt(fd, SOL_SOCKET, SO_ERROR, &err, &len) < 0)
- err = errno;
+ {
+ if (errno == ENOTSOCK)
+ {
+ // Non-socket fd (pipe, chardev, ...): no SO_ERROR, so
+ // let the op's own syscall name the real failure.
+ // Also force the read/write dispatch below to run: an
+ // edge-triggered EPOLLERR can arrive alone, and
+ // without this a parked op never calls perform_io()
+ // and, the edge being one-shot, never gets another
+ // chance -- a permanent hang.
+ //
+ // Assumes at least one parked op's own syscall makes
+ // non-EAGAIN progress; if every op re-parks with
+ // EAGAIN this sticky error is never redelivered and
+ // they hang. No such case is known for pipes -- a
+ // future non-socket type that hits one should be
+ // handled here.
+ err = 0;
+ ev |= reactor_event_read | reactor_event_write;
+ }
+ else
+ {
+ err = errno;
+ }
+ }
// select raises its exceptional set for out-of-band/urgent
// data as well as for genuine faults; on a healthy socket the
// probe then reads SO_ERROR == 0. Faulting a pending read or
diff --git a/include/boost/corosio/native/detail/select/select_traits.hpp b/include/boost/corosio/native/detail/select/select_traits.hpp
index 3b08f6faa..08d5fbd01 100644
--- a/include/boost/corosio/native/detail/select/select_traits.hpp
+++ b/include/boost/corosio/native/detail/select/select_traits.hpp
@@ -25,6 +25,7 @@
#include
#include
#include
+#include
#include
/* select backend traits.
@@ -111,6 +112,39 @@ struct select_traits
}
};
+ // Descriptors are not sockets: sendmsg() fails with ENOTSOCK on a
+ // pipe or character device, so the write path is writev()/write()
+ // and SIGPIPE suppression is structurally unavailable -- MSG_NOSIGNAL
+ // is a send() flag and SO_NOSIGPIPE a socket option. A write to a
+ // pipe whose read end has closed raises SIGPIPE, exactly as a plain
+ // write(2) would; callers install SIG_IGN.
+ struct descriptor_write_policy
+ {
+ static ssize_t write(int fd, iovec* iovecs, int count) noexcept
+ {
+ ssize_t n;
+ do
+ {
+ n = ::writev(fd, iovecs, count);
+ }
+ while (n < 0 && errno == EINTR);
+ return n;
+ }
+
+ // Single-buffer fast path: skips the kernel's iov_iter setup.
+ static ssize_t
+ write_one(int fd, void const* data, std::size_t size) noexcept
+ {
+ ssize_t n;
+ do
+ {
+ n = ::write(fd, data, size);
+ }
+ while (n < 0 && errno == EINTR);
+ return n;
+ }
+ };
+
struct accept_policy
{
static int
diff --git a/include/boost/corosio/native/detail/select/select_types.hpp b/include/boost/corosio/native/detail/select/select_types.hpp
index cb919a7a4..df21ee80b 100644
--- a/include/boost/corosio/native/detail/select/select_types.hpp
+++ b/include/boost/corosio/native/detail/select/select_types.hpp
@@ -25,6 +25,7 @@
#include
#include
#include
+#include
namespace boost::corosio::detail {
@@ -41,6 +42,8 @@ class select_local_stream_acceptor;
class select_local_stream_acceptor_service;
class select_local_datagram_socket;
class select_local_datagram_service;
+class select_descriptor;
+class select_descriptor_service;
// --- Stream sockets ---
@@ -238,6 +241,29 @@ class select_local_stream_acceptor final
}
};
+// --- Descriptors ---
+
+class select_descriptor final
+ : public reactor_descriptor<
+ select_descriptor,
+ select_traits,
+ select_descriptor_service,
+ select_tcp_acceptor>
+{
+ using base_type = reactor_descriptor<
+ select_descriptor,
+ select_traits,
+ select_descriptor_service,
+ select_tcp_acceptor>;
+ friend select_descriptor_service;
+
+public:
+ explicit select_descriptor(select_descriptor_service& svc) noexcept
+ : base_type(svc)
+ {
+ }
+};
+
// --- Services ---
class BOOST_COROSIO_DECL select_tcp_service final
@@ -358,6 +384,24 @@ class BOOST_COROSIO_DECL select_local_stream_acceptor_service final
}
};
+class BOOST_COROSIO_DECL select_descriptor_service final
+ : public reactor_descriptor_service<
+ select_descriptor_service,
+ select_traits,
+ select_descriptor>
+{
+ using base_type = reactor_descriptor_service<
+ select_descriptor_service,
+ select_traits,
+ select_descriptor>;
+
+public:
+ explicit select_descriptor_service(capy::execution_context& ctx)
+ : base_type(ctx)
+ {
+ }
+};
+
} // namespace boost::corosio::detail
#endif // BOOST_COROSIO_HAS_SELECT
diff --git a/include/boost/corosio/native/detail/uring/uring_descriptor.hpp b/include/boost/corosio/native/detail/uring/uring_descriptor.hpp
new file mode 100644
index 000000000..c7f9167ea
--- /dev/null
+++ b/include/boost/corosio/native/detail/uring/uring_descriptor.hpp
@@ -0,0 +1,569 @@
+//
+// Copyright (c) 2026 Michael Vandeberg
+//
+// Distributed under the Boost Software License, Version 1.0. (See accompanying
+// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
+//
+// Official repository: https://github.com/cppalliance/corosio
+//
+
+#ifndef BOOST_COROSIO_NATIVE_DETAIL_URING_URING_DESCRIPTOR_HPP
+#define BOOST_COROSIO_NATIVE_DETAIL_URING_URING_DESCRIPTOR_HPP
+
+#include
+
+#if BOOST_COROSIO_HAS_URING
+
+#include
+#include
+#include
+#include
+#include
+#include
+
+#include
+#include
+#include
+#include
+#include
+#include
+
+#include
+#include
+#include
+
+/* io_uring-backed implementation of posix_descriptor.
+
+ Three things differ from the reactor backends and from the other
+ io_uring services:
+
+ Transfers submit READV/WRITEV at offset -1 so the kernel uses (and
+ advances) the descriptor's own file position. The file services
+ pass a real offset; a pipe, tty or character device has none.
+
+ O_NONBLOCK is armed lazily, on the first read_some/write_some and
+ never from assign() or wait(). Here that is a cancellability
+ requirement, not only the public contract: a transfer the kernel
+ would have to block on is punted to an io-wq worker, and a request
+ running on a worker cannot be cancelled. With the flag set the
+ kernel reports readiness instead, so every operation stays
+ cancellable on every kernel version.
+
+ That reporting is what the transfer ops' two-phase shape handles.
+ An O_NONBLOCK descriptor the kernel cannot retry internally
+ completes with -EAGAIN; the op then re-arms itself as a poll_add on
+ the same descriptor and re-submits the transfer when the poll says
+ ready. The handler makes that decision *before* coro_drain_if_shutdown,
+ which disarms stop_cb: an op going round again keeps its
+ cancellation wiring, so a stop_token firing between phases still
+ reaches the kernel.
+
+ The gap between a CQE and its dispatch is the whole difficulty of
+ that shape, and two epoch counters close it. Nothing of a transfer
+ waiting in that gap is in the ring, so neither cancel-by-fd nor a
+ one-shot stop_callback can reach it; cancel() bumps cancel_epoch_
+ and every descriptor change bumps desc_epoch_, and an op whose
+ snapshot no longer matches completes instead of re-arming. The
+ descriptor epoch is also what makes the staleness check exact
+ rather than heuristic: a closed fd number the next assign() gets
+ back would satisfy a bare fd comparison.
+
+ There is no adopt-time registration. assign() validates and takes
+ the descriptor; a kernel that refuses it says so at the first
+ operation, as the public docstring promises.
+*/
+
+namespace boost::corosio::detail {
+
+class uring_descriptor;
+
+/** Advance a two-phase transfer op, or report that it is finished.
+
+ A kernel `EAGAIN` becomes a `poll_add` on the same descriptor, and
+ the poll's completion re-submits the transfer. Every path that
+ stops instead leaves @a op carrying a result the completion decode
+ can read as terminal.
+
+ @param op The op whose CQE just arrived.
+ @return True when a fresh SQE was submitted, in which case the
+ caller must neither complete nor disarm the op.
+*/
+template
+bool uring_descriptor_continue(Op& op) noexcept;
+
+/** Scatter read via `IORING_OP_READV` at the descriptor's own offset.
+
+ @see uring_descriptor_continue for the `polling` phase.
+*/
+struct uring_descriptor_read_op final : uring_file_read_op_base
+{
+ uring_descriptor* desc = nullptr;
+ /// True while the submitted SQE is the readiness poll, not the read.
+ bool polling = false;
+ /// Owner epochs snapshotted at submission; see uring_descriptor_continue.
+ std::uint32_t cancel_epoch = 0;
+ std::uint32_t desc_epoch = 0;
+
+ uring_descriptor_read_op() noexcept : uring_file_read_op_base(&do_handler)
+ {
+ prep_func = &do_prep;
+ }
+
+ static void do_prep(uring_op* base, ::io_uring_sqe* sqe) noexcept
+ {
+ auto* self = static_cast(base);
+ if (self->polling)
+ ::io_uring_prep_poll_add(sqe, self->fd, POLLIN);
+ else
+ uring_file_read_op_base::do_prep(base, sqe);
+ }
+
+ static void do_handler(
+ void* owner,
+ scheduler_op* base,
+ std::uint32_t bytes,
+ std::uint32_t error) noexcept;
+};
+
+/// Gather write via `IORING_OP_WRITEV` at the descriptor's own offset.
+struct uring_descriptor_write_op final : uring_file_write_op_base
+{
+ uring_descriptor* desc = nullptr;
+ /// True while the submitted SQE is the readiness poll, not the write.
+ bool polling = false;
+ /// Owner epochs snapshotted at submission; see uring_descriptor_continue.
+ std::uint32_t cancel_epoch = 0;
+ std::uint32_t desc_epoch = 0;
+
+ uring_descriptor_write_op() noexcept : uring_file_write_op_base(&do_handler)
+ {
+ prep_func = &do_prep;
+ }
+
+ static void do_prep(uring_op* base, ::io_uring_sqe* sqe) noexcept
+ {
+ auto* self = static_cast(base);
+ if (self->polling)
+ ::io_uring_prep_poll_add(sqe, self->fd, POLLOUT);
+ else
+ uring_file_write_op_base::do_prep(base, sqe);
+ }
+
+ static void do_handler(
+ void* owner,
+ scheduler_op* base,
+ std::uint32_t bytes,
+ std::uint32_t error) noexcept;
+};
+
+/** Native io_uring implementation of @ref posix_descriptor.
+
+ Holds the adopted descriptor and the five embedded op slots: one
+ transfer per direction, and one wait per direction.
+
+ @par Thread Safety
+ Distinct objects: Safe.@n
+ Shared objects: Unsafe. Each slot carries a single pending
+ operation, so a descriptor must not have two operations of the
+ same kind in flight.
+*/
+class BOOST_COROSIO_DECL uring_descriptor final
+ : public posix_descriptor::implementation
+ , public std::enable_shared_from_this
+{
+ uring_scheduler* sched_ = nullptr;
+ int fd_ = -1;
+ bool nonblocking_ = false;
+
+ // Bumped by cancel() and by every descriptor change respectively.
+ // A transfer op parked between its EAGAIN CQE and its dispatch is
+ // invisible to the ring, so these are the only thing that can stop
+ // it re-arming against an intent it no longer belongs to.
+ std::atomic cancel_epoch_{0};
+ std::atomic desc_epoch_{0};
+
+ uring_descriptor_read_op rd_;
+ uring_descriptor_write_op wr_;
+ uring_wait_op wait_rd_;
+ uring_wait_op wait_wr_;
+ uring_wait_op wait_er_;
+
+public:
+ explicit uring_descriptor(uring_scheduler& sched) noexcept : sched_(&sched)
+ {
+ }
+
+ ~uring_descriptor() override
+ {
+ close_descriptor();
+ }
+
+ // -- io_stream::implementation --
+
+ std::coroutine_handle<> read_some(
+ std::coroutine_handle<> h,
+ capy::executor_ref ex,
+ buffer_param buffers,
+ std::stop_token token,
+ std::error_code* ec,
+ std::size_t* bytes) override
+ {
+ rd_.prepare(
+ h, ex, ec, bytes, fd_, /*file_offset=*/-1, sched_,
+ shared_from_this(), buffers, token);
+ arm_slot(rd_);
+ sched_->work_started();
+
+ // Closed-object contract outranks the zero-length no-op.
+ if (fd_ < 0)
+ {
+ rd_.empty_buffer = false;
+ rd_.res = -EBADF;
+ push_completed(&rd_);
+ return std::noop_coroutine();
+ }
+
+ if (rd_.empty_buffer || rd_.cancelled.load(std::memory_order_acquire))
+ {
+ push_completed(&rd_);
+ return std::noop_coroutine();
+ }
+
+ // The first transferring operation is what arms O_NONBLOCK;
+ // assign() and wait() never do.
+ if (int const nerr = arm_nonblocking())
+ {
+ rd_.res = -nerr;
+ push_completed(&rd_);
+ return std::noop_coroutine();
+ }
+
+ uring_submit_op(*sched_, &rd_);
+ return std::noop_coroutine();
+ }
+
+ std::coroutine_handle<> write_some(
+ std::coroutine_handle<> h,
+ capy::executor_ref ex,
+ buffer_param buffers,
+ std::stop_token token,
+ std::error_code* ec,
+ std::size_t* bytes) override
+ {
+ wr_.prepare(
+ h, ex, ec, bytes, fd_, /*file_offset=*/-1, sched_,
+ shared_from_this(), buffers, token);
+ arm_slot(wr_);
+ sched_->work_started();
+
+ if (fd_ < 0)
+ {
+ wr_.empty_buffer = false;
+ wr_.res = -EBADF;
+ push_completed(&wr_);
+ return std::noop_coroutine();
+ }
+
+ if (wr_.empty_buffer || wr_.cancelled.load(std::memory_order_acquire))
+ {
+ push_completed(&wr_);
+ return std::noop_coroutine();
+ }
+
+ if (int const nerr = arm_nonblocking())
+ {
+ wr_.res = -nerr;
+ push_completed(&wr_);
+ return std::noop_coroutine();
+ }
+
+ uring_submit_op(*sched_, &wr_);
+ return std::noop_coroutine();
+ }
+
+ // -- posix_descriptor::implementation --
+
+ std::coroutine_handle<> wait(
+ std::coroutine_handle<> h,
+ capy::executor_ref ex,
+ wait_type w,
+ std::stop_token token,
+ std::error_code* ec) override
+ {
+ uring_wait_op* op = nullptr;
+ int poll_flags = 0;
+ switch (w)
+ {
+ case wait_type::read:
+ op = &wait_rd_;
+ poll_flags = POLLIN;
+ break;
+ case wait_type::write:
+ op = &wait_wr_;
+ poll_flags = POLLOUT;
+ break;
+ case wait_type::error:
+ op = &wait_er_;
+ // POLLERR, POLLHUP and POLLNVAL are reported whether or not
+ // they are asked for, so the error wait names only POLLPRI.
+ poll_flags = POLLPRI;
+ break;
+ }
+
+ op->prepare(
+ h, ex, ec, fd_, sched_, shared_from_this(), poll_flags, token);
+ sched_->work_started();
+
+ // No arm_nonblocking() on any branch here: a wait must leave a
+ // descriptor someone else owns exactly as it found it.
+ if (fd_ < 0)
+ {
+ op->res = -EBADF;
+ push_completed(op);
+ return std::noop_coroutine();
+ }
+
+ if (op->cancelled.load(std::memory_order_acquire))
+ {
+ push_completed(op);
+ return std::noop_coroutine();
+ }
+
+ uring_submit_op(*sched_, op);
+ return std::noop_coroutine();
+ }
+
+ native_handle_type native_handle() const noexcept override
+ {
+ return fd_;
+ }
+
+ native_handle_type release_descriptor() noexcept override
+ {
+ // Flush the cancel while the fd is still open so the kernel
+ // resolves it before the caller can close and recycle the
+ // number. Do NOT close -- the caller takes ownership.
+ if (fd_ >= 0)
+ sched_->cancel_and_flush(fd_);
+ native_handle_type released = fd_;
+ fd_ = -1;
+ nonblocking_ = false;
+ desc_epoch_.fetch_add(1, std::memory_order_release);
+ return released;
+ }
+
+ void cancel() noexcept override
+ {
+ // Bump before the SQE: cancel-by-fd reaches only what the ring
+ // currently holds, and an op waiting for its handler to run
+ // holds nothing there. The epoch is what that op consults.
+ cancel_epoch_.fetch_add(1, std::memory_order_release);
+ if (fd_ >= 0)
+ sched_->submit_cancel_by_fd(fd_);
+ }
+
+ /// Epoch bumped by every @ref cancel.
+ std::uint32_t cancel_epoch() const noexcept
+ {
+ return cancel_epoch_.load(std::memory_order_acquire);
+ }
+
+ /// Epoch bumped by every change of the held descriptor.
+ std::uint32_t desc_epoch() const noexcept
+ {
+ return desc_epoch_.load(std::memory_order_acquire);
+ }
+
+ // -- Service-facing (non-virtual) --
+
+ /** Adopt an already-validated descriptor.
+
+ Resets the lazy-nonblocking latch so a freshly adopted fd is
+ not assumed to carry the flag from whatever this object held
+ before.
+
+ @param fd The descriptor to adopt.
+ */
+ void set_descriptor(int fd) noexcept
+ {
+ fd_ = fd;
+ nonblocking_ = false;
+ desc_epoch_.fetch_add(1, std::memory_order_release);
+ }
+
+ /// Cancel pending operations and close the descriptor. No-op when
+ /// already closed.
+ void close_descriptor() noexcept
+ {
+ if (fd_ < 0)
+ return;
+ // Both kernel entries below can run a queued pipe write as task
+ // work; with the reader already gone that raises SIGPIPE.
+ scoped_sigpipe_block no_sigpipe;
+ sched_->cancel_and_flush(fd_);
+ ::close(fd_);
+ fd_ = -1;
+ nonblocking_ = false;
+ desc_epoch_.fetch_add(1, std::memory_order_release);
+ }
+
+private:
+ /** Arm O_NONBLOCK, once, before the first transfer.
+
+ Reports an errno rather than an error_code because the op
+ result model records a negated errno in `res`; the round trip
+ is lossless because fcntl only fails with codes make_err
+ passes through.
+ */
+ int arm_nonblocking() noexcept
+ {
+ if (nonblocking_)
+ return 0;
+ if (auto ec = ensure_nonblocking(fd_))
+ return ec.value();
+ nonblocking_ = true;
+ return 0;
+ }
+
+ /** Bind a transfer slot to this descriptor for a fresh submission.
+
+ The epoch snapshot taken here is what a later re-arm compares
+ against; see uring_descriptor_continue.
+ */
+ template
+ void arm_slot(Op& op) noexcept
+ {
+ op.desc = this;
+ op.polling = false;
+ op.cancel_epoch = cancel_epoch_.load(std::memory_order_acquire);
+ op.desc_epoch = desc_epoch_.load(std::memory_order_acquire);
+ }
+
+ /// Queue an already-counted op for the next dispatch cycle.
+ void push_completed(scheduler_op* op) noexcept
+ {
+ uring_scheduler::lock_type lock(sched_->dispatch_mutex());
+ sched_->push_completed_locked(op);
+ }
+};
+
+// --- Deferred implementations (need uring_descriptor complete) ---
+
+template
+bool
+uring_descriptor_continue(Op& op) noexcept
+{
+ // A poll CQE carries its revents in `res` -- a small positive
+ // integer the completion decode would otherwise read as a byte
+ // count and as success. Every path that abandons the op while that
+ // value is sitting there has to overwrite it with a terminal one.
+ bool const mid_poll = op.polling && op.res >= 0;
+
+ if (!op.desc)
+ {
+ if (mid_poll)
+ op.res = -EBADF;
+ return false;
+ }
+
+ // Two epochs rather than one, because the two reasons to abandon an
+ // op name different codes to the caller.
+ if (op.cancelled.load(std::memory_order_acquire) ||
+ op.desc->cancel_epoch() != op.cancel_epoch)
+ {
+ if (mid_poll)
+ op.res = -ECANCELED;
+ return false;
+ }
+
+ if (op.desc->desc_epoch() != op.desc_epoch)
+ {
+ if (mid_poll)
+ op.res = -EBADF;
+ return false;
+ }
+
+ if (op.polling)
+ {
+ // A poll that failed or was cancelled is the operation's answer.
+ if (op.res < 0)
+ return false;
+ op.polling = false;
+ }
+ else if (op.res == -EAGAIN || op.res == -EWOULDBLOCK)
+ {
+ op.polling = true;
+ }
+ else
+ {
+ return false;
+ }
+
+ // do_one spends a work_finished() on every op it dispatches, so an
+ // op going round again has to be counted again.
+ op.sched_->work_started();
+ uring_submit_op(*op.sched_, &op);
+
+ // stop_cb is one-shot and has already fired for the SQE that just
+ // completed, and cancel-by-fd found nothing while this op was out
+ // of the ring: a cancel racing the submission above would reach no
+ // kernel request at all, so re-check and drive it here.
+ if (op.cancelled.load(std::memory_order_acquire) ||
+ op.desc->cancel_epoch() != op.cancel_epoch)
+ op.sched_->submit_cancel_by_user_data(&op);
+ return true;
+}
+
+inline void
+uring_descriptor_read_op::do_handler(
+ void* owner,
+ scheduler_op* base,
+ std::uint32_t /*bytes*/,
+ std::uint32_t /*error*/) noexcept
+{
+ auto* self = static_cast(base);
+ if (owner != nullptr && uring_descriptor_continue(*self))
+ return;
+
+ if (coro_drain_if_shutdown(owner, self))
+ return;
+
+ if (self->sched_)
+ self->sched_->reset_inline_budget();
+
+ uring_set_result(self, /*is_read=*/true, self->empty_buffer);
+ if (self->bytes_out)
+ *self->bytes_out =
+ self->res >= 0 ? static_cast(self->res) : 0u;
+ coro_resume(self);
+}
+
+inline void
+uring_descriptor_write_op::do_handler(
+ void* owner,
+ scheduler_op* base,
+ std::uint32_t /*bytes*/,
+ std::uint32_t /*error*/) noexcept
+{
+ auto* self = static_cast(base);
+ if (owner != nullptr && uring_descriptor_continue(*self))
+ return;
+
+ if (coro_drain_if_shutdown(owner, self))
+ return;
+
+ if (self->sched_)
+ self->sched_->reset_inline_budget();
+
+ uring_set_result(self, /*is_read=*/false, self->empty_buffer);
+ if (self->bytes_out)
+ *self->bytes_out =
+ self->res >= 0 ? static_cast(self->res) : 0u;
+ coro_resume(self);
+}
+
+} // namespace boost::corosio::detail
+
+#endif // BOOST_COROSIO_HAS_URING
+
+#endif // BOOST_COROSIO_NATIVE_DETAIL_URING_URING_DESCRIPTOR_HPP
diff --git a/include/boost/corosio/native/detail/uring/uring_descriptor_service.hpp b/include/boost/corosio/native/detail/uring/uring_descriptor_service.hpp
new file mode 100644
index 000000000..712d948ac
--- /dev/null
+++ b/include/boost/corosio/native/detail/uring/uring_descriptor_service.hpp
@@ -0,0 +1,143 @@
+//
+// Copyright (c) 2026 Michael Vandeberg
+//
+// Distributed under the Boost Software License, Version 1.0. (See accompanying
+// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
+//
+// Official repository: https://github.com/cppalliance/corosio
+//
+
+#ifndef BOOST_COROSIO_NATIVE_DETAIL_URING_URING_DESCRIPTOR_SERVICE_HPP
+#define BOOST_COROSIO_NATIVE_DETAIL_URING_URING_DESCRIPTOR_SERVICE_HPP
+
+#include
+
+#if BOOST_COROSIO_HAS_URING
+
+#include
+#include
+#include
+#include
+#include
+#include
+
+#include
+#include
+#include
+#include
+#include
+
+/* io_uring-backed descriptor_service.
+
+ assign_descriptor is the validate-before-mutate core the public
+ assign() contract rests on. It is the reactor version minus the
+ registration step: io_uring has no adopt-time registration syscall,
+ so adoption cannot fail once validation has passed, and a kernel
+ refusal surfaces at the first operation instead.
+
+ Neither uring_socket_service_base nor uring_file_service_base fits:
+ the first constructs impls with a (service&, scheduler&) ctor and
+ only cancels on shutdown, the second names its teardown hook
+ close_file(). The lifecycle here is small enough to carry directly.
+
+ PARALLEL COPY: construct, destroy, close and shutdown here mirror the
+ same four members of reactor_descriptor_service.hpp, which this
+ cannot reuse because that template is rooted in the reactor's
+ scheduler and service state. A fix to the descriptor service
+ lifecycle -- the construct/destroy bookkeeping, or shutdown's
+ deliberate retention of the impl map so impls outlive the
+ scheduler's drain -- belongs in both files.
+*/
+
+namespace boost::corosio::detail {
+
+/// Native io_uring descriptor service. Owns every @ref uring_descriptor
+/// the context creates.
+class BOOST_COROSIO_DECL uring_descriptor_service final
+ : public descriptor_service
+{
+ uring_scheduler* sched_ = nullptr;
+ std::mutex mutex_;
+ std::unordered_map>
+ impls_;
+
+public:
+ explicit uring_descriptor_service(capy::execution_context& ctx)
+ : sched_(&ctx.use_service())
+ {
+ }
+
+ std::error_code assign_descriptor(
+ posix_descriptor::implementation& impl_base,
+ native_handle_type fd) override
+ {
+ auto* impl = static_cast(&impl_base);
+
+ // fd >= 0 guard: an unset impl reports native_handle() == -1, and
+ // a caller-supplied -1 must fail as a bad fd, not a self-assign.
+ if (fd >= 0 && fd == impl->native_handle())
+ return std::make_error_code(std::errc::invalid_argument);
+
+ // Validate before touching the held descriptor: a failed assign
+ // must leave the object unchanged and the caller owning the fd.
+ if (auto ec = validate_descriptor_fd(fd))
+ return ec;
+
+ impl->close_descriptor();
+ impl->set_descriptor(fd);
+ return {};
+ }
+
+ io_object::implementation* construct() override
+ {
+ auto p = std::make_shared(*sched_);
+ auto* raw = p.get();
+ std::lock_guard lock(mutex_);
+ impls_.emplace(raw, std::move(p));
+ return raw;
+ }
+
+ void destroy(io_object::implementation* p) override
+ {
+ if (!p)
+ return;
+ auto* impl = static_cast(p);
+ impl->close_descriptor();
+ std::lock_guard lock(mutex_);
+ impls_.erase(impl);
+ }
+
+ void close(io_object::handle& h) override
+ {
+ if (auto* impl = static_cast(h.get()))
+ impl->close_descriptor();
+ }
+
+ void shutdown() override
+ {
+ // Snapshot, then close without the lock held. impls_ is
+ // deliberately not cleared: the scheduler shuts down after this
+ // service and drains its completed ops, so every impl must
+ // outlive that drain.
+ std::vector> live;
+ {
+ std::lock_guard lock(mutex_);
+ live.reserve(impls_.size());
+ for (auto& [raw, p] : impls_)
+ live.push_back(p);
+ }
+ for (auto& p : live)
+ p->close_descriptor();
+ }
+
+private:
+ uring_descriptor_service(uring_descriptor_service const&) = delete;
+ uring_descriptor_service&
+ operator=(uring_descriptor_service const&) = delete;
+};
+
+} // namespace boost::corosio::detail
+
+#endif // BOOST_COROSIO_HAS_URING
+
+#endif // BOOST_COROSIO_NATIVE_DETAIL_URING_URING_DESCRIPTOR_SERVICE_HPP
diff --git a/include/boost/corosio/native/detail/uring/uring_random_access_file.hpp b/include/boost/corosio/native/detail/uring/uring_random_access_file.hpp
index 773b049bf..361bc123c 100644
--- a/include/boost/corosio/native/detail/uring/uring_random_access_file.hpp
+++ b/include/boost/corosio/native/detail/uring/uring_random_access_file.hpp
@@ -1,5 +1,6 @@
//
// Copyright (c) 2026 Steve Gerbino
+// Copyright (c) 2026 Michael Vandeberg
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
@@ -20,6 +21,7 @@
#include
#include
#include
+#include
#include
#include
@@ -154,6 +156,20 @@ class BOOST_COROSIO_DECL uring_random_access_file final
std::error_code assign(native_handle_type handle) noexcept override
{
+ // handle >= 0 guard: an unset impl reports native_handle() == -1,
+ // and a caller-supplied -1 must fail as a bad fd, not a
+ // self-assign.
+ if (handle >= 0 && handle == fd_)
+ return std::make_error_code(std::errc::invalid_argument);
+
+ // Validate before touching the held fd: a failed assign must
+ // leave this object unchanged and the caller still owning handle.
+ if (auto ec = validate_file_fd(handle))
+ return ec;
+
+ // No cancel() before this, unlike the POSIX twin: close_file()
+ // calls cancel_and_flush(fd_) itself, so the mandatory
+ // cancel-before-close pairing is already there.
close_file();
fd_ = handle;
return {};
diff --git a/include/boost/corosio/native/detail/uring/uring_socket_ops.hpp b/include/boost/corosio/native/detail/uring/uring_socket_ops.hpp
index 0b564aa7d..bca726b68 100644
--- a/include/boost/corosio/native/detail/uring/uring_socket_ops.hpp
+++ b/include/boost/corosio/native/detail/uring/uring_socket_ops.hpp
@@ -636,17 +636,39 @@ struct uring_wait_op : uring_op
// way — SO_ERROR, or EIO when the kernel has none — instead of
// completing wait(error) with an empty, benign-looking code.
// OOB (POLLPRI) is a readiness signal, not an error.
+ //
+ // POLLHUP is a fault only for wait(error). A readiness wait --
+ // the POLLIN/POLLOUT flag sets -- treats it as ready: a pipe or
+ // fully shut-down socket whose peer hung up is readable-at-EOF,
+ // and the read or write that follows names the condition. The
+ // reactors' poll() probe already reports it that way, so
+ // without this io_uring alone answers the wait-then-read idiom
+ // with a spurious EIO.
+ bool const readiness_wait =
+ (self->poll_flags & (POLLIN | POLLOUT)) != 0;
+ int const fault_bits = readiness_wait
+ ? (POLLERR | POLLNVAL)
+ : (POLLERR | POLLHUP | POLLNVAL);
+
std::error_code ec{};
if (self->res < 0)
{
ec = make_err(-self->res);
}
- else if (self->res & (POLLERR | POLLHUP | POLLNVAL))
+ else if (self->res & fault_bits)
{
int so_err = 0;
socklen_t len = sizeof(so_err);
if (::getsockopt(self->fd, SOL_SOCKET, SO_ERROR, &so_err, &len) < 0)
- so_err = errno;
+ {
+ // A non-socket (pipe, chardev, ...) has no SO_ERROR and
+ // fails the probe with ENOTSOCK; reporting that would
+ // name the probe rather than the fault, so fall through
+ // to the EIO substitution below. Every other failure
+ // (EBADF from a concurrent close, say) still reports
+ // itself, unchanged on the socket hot path.
+ so_err = (errno == ENOTSOCK) ? 0 : errno;
+ }
if (so_err == 0)
so_err = EIO;
ec = make_err(so_err);
diff --git a/include/boost/corosio/native/detail/uring/uring_stream_file.hpp b/include/boost/corosio/native/detail/uring/uring_stream_file.hpp
index 05bb2ccf6..99f8509c1 100644
--- a/include/boost/corosio/native/detail/uring/uring_stream_file.hpp
+++ b/include/boost/corosio/native/detail/uring/uring_stream_file.hpp
@@ -1,5 +1,6 @@
//
// Copyright (c) 2026 Steve Gerbino
+// Copyright (c) 2026 Michael Vandeberg
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
@@ -20,6 +21,7 @@
#include
#include
#include
+#include
#include
#include
@@ -161,6 +163,20 @@ class BOOST_COROSIO_DECL uring_stream_file final
std::error_code assign(native_handle_type handle) noexcept override
{
+ // handle >= 0 guard: an unset impl reports native_handle() == -1,
+ // and a caller-supplied -1 must fail as a bad fd, not a
+ // self-assign.
+ if (handle >= 0 && handle == fd_)
+ return std::make_error_code(std::errc::invalid_argument);
+
+ // Validate before touching the held fd: a failed assign must
+ // leave this object unchanged and the caller still owning handle.
+ if (auto ec = validate_file_fd(handle))
+ return ec;
+
+ // No cancel() before this, unlike the POSIX twin: close_file()
+ // calls cancel_and_flush(fd_) itself, so the mandatory
+ // cancel-before-close pairing is already there.
close_file();
fd_ = handle;
return {};
diff --git a/include/boost/corosio/native/detail/validate_fd.hpp b/include/boost/corosio/native/detail/validate_fd.hpp
index 39c4e0c61..fac25449b 100644
--- a/include/boost/corosio/native/detail/validate_fd.hpp
+++ b/include/boost/corosio/native/detail/validate_fd.hpp
@@ -1,5 +1,6 @@
//
// Copyright (c) 2026 Steve Gerbino
+// Copyright (c) 2026 Michael Vandeberg
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
@@ -19,7 +20,9 @@
#include
#include
+#include
#include
+#include
namespace boost::corosio::detail {
@@ -65,6 +68,114 @@ validate_socket_fd(int fd, int expected_type, bool is_ip) noexcept
return {};
}
+/** Validate a caller-supplied fd for adoption by @ref posix_descriptor.
+
+ Non-mutating: interrogates the fd without changing any of its
+ flags, so a rejected fd goes back to the caller untouched. In
+ particular `O_NONBLOCK` is not applied here — see
+ @ref ensure_nonblocking.
+
+ The file-type test is a reject-list, not an accept-list. The
+ flagship descriptor kinds -- eventfd, timerfd, inotify, pidfd --
+ are anonymous inodes whose `st_mode` type bits are all zero, so
+ an accept-list would silently reject exactly the fds this type
+ exists to carry.
+
+ @param fd The descriptor to validate.
+ @return Empty on success; `EBADF` for a closed or negative fd,
+ `operation_not_supported` for a regular file, directory or
+ block device, or the `errno` reported by `fstat`.
+*/
+inline std::error_code
+validate_descriptor_fd(int fd) noexcept
+{
+ if (fd < 0)
+ return make_err(EBADF);
+
+ struct stat st{};
+ if (::fstat(fd, &st) != 0)
+ return make_err(errno);
+
+ // Regular files, block devices and directories are the province of
+ // stream_file / random_access_file, whose assign() already adopts
+ // them; a reactor cannot report readiness for them anyway.
+ switch (st.st_mode & S_IFMT)
+ {
+ case S_IFREG:
+ case S_IFBLK:
+ case S_IFDIR:
+ return std::make_error_code(std::errc::operation_not_supported);
+ default:
+ return {};
+ }
+}
+
+/** Validate a caller-supplied fd for adoption by a file object.
+
+ Non-mutating. Accepts the kinds a file object can position and
+ read: regular files, block devices, and character devices such
+ as /dev/null and /dev/zero.
+
+ This is an accept-list, the inverse of @ref validate_descriptor_fd's
+ reject-list: a file object needs a positionable fd, and the
+ anonymous inodes that motivate the descriptor reject-list are
+ exactly what a file object cannot use.
+
+ `S_IFCHR` is deliberately broad: it admits `/dev/null` and
+ `/dev/zero`, but also non-seekable character devices such as a
+ tty. Those pass here and then fail loudly at first I/O on the
+ POSIX backends, where `preadv`/`pwritev` report `ESPIPE`.
+
+ @param fd The descriptor to validate.
+ @return Empty on success; `EBADF` for a closed or negative fd,
+ `operation_not_supported` for a directory or a descriptor
+ with no file position, or the `errno` from `fstat`.
+*/
+inline std::error_code
+validate_file_fd(int fd) noexcept
+{
+ if (fd < 0)
+ return make_err(EBADF);
+
+ struct stat st{};
+ if (::fstat(fd, &st) != 0)
+ return make_err(errno);
+
+ switch (st.st_mode & S_IFMT)
+ {
+ case S_IFREG:
+ case S_IFBLK:
+ case S_IFCHR:
+ return {};
+ default:
+ return std::make_error_code(std::errc::operation_not_supported);
+ }
+}
+
+/** Put a descriptor into non-blocking mode, idempotently.
+
+ Called lazily on the first `read_some` / `write_some`, never from
+ `assign()`. The change is permanent: `O_NONBLOCK` lives on the
+ shared open file description, so restoring it later would race
+ every other holder of that description. A `wait()`-only user
+ never reaches this function and their fd is never modified.
+
+ @param fd The descriptor to modify.
+ @return Empty on success, otherwise the `errno` from `fcntl`.
+*/
+inline std::error_code
+ensure_nonblocking(int fd) noexcept
+{
+ int flags = ::fcntl(fd, F_GETFL, 0);
+ if (flags < 0)
+ return make_err(errno);
+ if (flags & O_NONBLOCK)
+ return {};
+ if (::fcntl(fd, F_SETFL, flags | O_NONBLOCK) < 0)
+ return make_err(errno);
+ return {};
+}
+
} // namespace boost::corosio::detail
#endif // BOOST_COROSIO_POSIX
diff --git a/include/boost/corosio/native/native.hpp b/include/boost/corosio/native/native.hpp
index affe8fa62..ef47ee298 100644
--- a/include/boost/corosio/native/native.hpp
+++ b/include/boost/corosio/native/native.hpp
@@ -1,5 +1,6 @@
//
// Copyright (c) 2026 Steve Gerbino
+// Copyright (c) 2026 Michael Vandeberg
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
@@ -20,6 +21,7 @@
#include
#include
#include
+#include
#include
#include
#include
diff --git a/include/boost/corosio/native/native_posix_descriptor.hpp b/include/boost/corosio/native/native_posix_descriptor.hpp
new file mode 100644
index 000000000..5275c205e
--- /dev/null
+++ b/include/boost/corosio/native/native_posix_descriptor.hpp
@@ -0,0 +1,229 @@
+//
+// Copyright (c) 2026 Michael Vandeberg
+//
+// Distributed under the Boost Software License, Version 1.0. (See accompanying
+// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
+//
+// Official repository: https://github.com/cppalliance/corosio
+//
+
+#ifndef BOOST_COROSIO_NATIVE_NATIVE_POSIX_DESCRIPTOR_HPP
+#define BOOST_COROSIO_NATIVE_NATIVE_POSIX_DESCRIPTOR_HPP
+
+#include
+#include
+#include
+
+#if BOOST_COROSIO_POSIX || defined(BOOST_COROSIO_MRDOCS)
+
+#ifndef BOOST_COROSIO_MRDOCS
+#if BOOST_COROSIO_HAS_EPOLL
+#include
+#endif
+
+#if BOOST_COROSIO_HAS_SELECT
+#include
+#endif
+
+#if BOOST_COROSIO_HAS_KQUEUE
+#include
+#endif
+
+#if BOOST_COROSIO_HAS_URING
+// uring_types.hpp does not declare the descriptor: it was added
+// after that header, in its own pair of files.
+#include
+#endif
+#endif // !BOOST_COROSIO_MRDOCS
+
+namespace boost::corosio {
+
+/** Drives an already-open POSIX descriptor, calling the backend directly.
+
+ This class template inherits from @ref posix_descriptor and
+ shadows the async operations (`read_some`, `write_some`, `wait`)
+ with versions that call the backend implementation directly.
+ This lets the compiler inline through the entire call chain.
+
+ Non-async operations (`assign`, `release`, `close`, `cancel`)
+ remain unchanged and dispatch through the compiled library.
+
+ A `native_posix_descriptor` IS-A `posix_descriptor` and can be
+ passed to any function expecting `posix_descriptor&` or
+ `io_stream&`, in which case virtual dispatch is used
+ transparently.
+
+ @tparam Backend A backend tag value (e.g., `epoll`) whose type
+ provides the concrete implementation types.
+
+ @par Thread Safety
+ Same as @ref posix_descriptor.
+
+ @par Example
+ @par !example assign_and_wait
+
+ @see posix_descriptor, epoll_t, kqueue_t
+*/
+template
+class native_posix_descriptor : public posix_descriptor
+{
+ using backend_type = decltype(Backend);
+ using impl_type = typename backend_type::descriptor_type;
+ using service_type = typename backend_type::descriptor_service_type;
+
+ impl_type& get_impl() noexcept
+ {
+ return *static_cast(h_.get());
+ }
+
+ template
+ struct native_read_awaitable
+ : detail::bytes_op_base>
+ {
+ native_posix_descriptor& self_;
+ MutableBufferSequence buffers_;
+
+ native_read_awaitable(
+ native_posix_descriptor& self,
+ MutableBufferSequence buffers) noexcept
+ : self_(self)
+ , buffers_(std::move(buffers))
+ {
+ }
+
+ std::coroutine_handle<>
+ dispatch(std::coroutine_handle<> h, capy::executor_ref ex) const
+ {
+ return self_.get_impl().read_some(
+ h, ex, buffers_, this->token_, &this->ec_, &this->bytes_);
+ }
+ };
+
+ template
+ struct native_write_awaitable
+ : detail::bytes_op_base>
+ {
+ native_posix_descriptor& self_;
+ ConstBufferSequence buffers_;
+
+ native_write_awaitable(
+ native_posix_descriptor& self, ConstBufferSequence buffers) noexcept
+ : self_(self)
+ , buffers_(std::move(buffers))
+ {
+ }
+
+ std::coroutine_handle<>
+ dispatch(std::coroutine_handle<> h, capy::executor_ref ex) const
+ {
+ return self_.get_impl().write_some(
+ h, ex, buffers_, this->token_, &this->ec_, &this->bytes_);
+ }
+ };
+
+ struct native_wait_awaitable : detail::void_op_base
+ {
+ native_posix_descriptor& self_;
+ wait_type w_;
+
+ native_wait_awaitable(
+ native_posix_descriptor& self, wait_type w) noexcept
+ : self_(self)
+ , w_(w)
+ {
+ }
+
+ std::coroutine_handle<>
+ dispatch(std::coroutine_handle<> h, capy::executor_ref ex) const
+ {
+ return self_.get_impl().wait(h, ex, w_, this->token_, &this->ec_);
+ }
+ };
+
+public:
+ /** Construct a native descriptor from an execution context.
+
+ @param ctx The execution context that owns this object.
+ */
+ explicit native_posix_descriptor(capy::execution_context& ctx)
+ : io_object(handle(ctx, ctx.use_service()))
+ {
+ }
+
+ /** Construct a native descriptor from an executor.
+
+ @param ex The executor whose context owns this object.
+ */
+ template