TLA Line data Source code
1 : //
2 : // Copyright (c) 2026 Michael Vandeberg
3 : //
4 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
5 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
6 : //
7 : // Official repository: https://github.com/cppalliance/corosio
8 : //
9 :
10 : #ifndef BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
11 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
12 :
13 : #include <boost/corosio/detail/platform.hpp>
14 :
15 : #if BOOST_COROSIO_POSIX
16 :
17 : #include <boost/corosio/posix_stream_descriptor.hpp>
18 : #include <boost/corosio/wait_type.hpp>
19 : #include <boost/corosio/detail/dispatch_coro.hpp>
20 : #include <boost/corosio/detail/intrusive.hpp>
21 : #include <boost/corosio/detail/native_handle.hpp>
22 : #include <boost/corosio/native/detail/make_err.hpp>
23 : #include <boost/corosio/native/detail/validate_fd.hpp>
24 : #include <boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp>
25 : #include <boost/corosio/native/detail/reactor/reactor_io_core.hpp>
26 : #include <boost/corosio/native/detail/reactor/reactor_op.hpp>
27 : #include <boost/corosio/native/detail/reactor/reactor_op_complete.hpp>
28 : #include <boost/capy/buffers.hpp>
29 :
30 : #include <atomic>
31 : #include <coroutine>
32 : #include <memory>
33 : #include <mutex>
34 : #include <utility>
35 :
36 : #include <errno.h>
37 : #include <sys/stat.h>
38 : #include <sys/uio.h>
39 : #include <unistd.h>
40 :
41 : /* Reactor-backed implementation of posix_stream_descriptor.
42 :
43 : The register/park/cancel/teardown protocol lives in reactor_io_core.
44 :
45 : The one behavior that is genuinely new: O_NONBLOCK is armed lazily,
46 : on the first read_some/write_some and never from assign() or wait().
47 : The flag lives on the shared open file description, so arming it is
48 : visible to every other holder of that description -- which is why a
49 : wait()-only user must never trigger it.
50 : */
51 :
52 : namespace boost::corosio::detail {
53 :
54 : // ============================================================
55 : // Op types
56 : // ============================================================
57 :
58 : /* Descriptor op family.
59 :
60 : Mirrors reactor_stream_ops.hpp. The Acceptor parameter is a
61 : placeholder: reactor_op is parameterized on both a socket and an
62 : acceptor impl type, and the shared completion helpers name
63 : acceptor_impl_ in a branch that a descriptor op never takes but the
64 : compiler still instantiates. Passing the backend's acceptor type (as
65 : reactor_dgram_socket_impl already does) keeps those helpers shared
66 : rather than duplicated here.
67 :
68 : @tparam Traits Backend traits (epoll_traits, kqueue_traits, ...).
69 : @tparam Descriptor The concrete descriptor type (forward-declared).
70 : @tparam Acceptor Placeholder acceptor type for the op base.
71 : */
72 :
73 : template<class Traits, class Descriptor, class Acceptor>
74 : struct reactor_descriptor_base_op : reactor_op<Descriptor, Acceptor>
75 : {
76 : void operator()() override;
77 : void cancel() noexcept override;
78 : };
79 :
80 : template<class Traits, class Descriptor, class Acceptor>
81 : struct reactor_descriptor_read_op final
82 : : reactor_read_op<reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>
83 : {};
84 :
85 : template<class Traits, class Descriptor, class Acceptor>
86 : struct reactor_descriptor_write_op final
87 : : reactor_write_op<
88 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>,
89 : typename Traits::descriptor_write_policy>
90 : {};
91 :
92 : template<class Traits, class Descriptor, class Acceptor>
93 : struct reactor_descriptor_wait_op final
94 : : reactor_wait_op<reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>
95 : {
96 : using base_type = reactor_wait_op<
97 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>;
98 :
99 : /** Probe like reactor_wait_op::probe, but name a code for an error
100 : wait that sees POLLERR or POLLHUP. A non-socket has no SO_ERROR
101 : to consult, and success would claim no error condition exists.
102 : */
103 HIT 72 : static bool probe(int fd, std::uint32_t event, int& err) noexcept
104 : {
105 72 : if (event != reactor_event_error || fd < 0)
106 48 : return base_type::probe(fd, event, err);
107 :
108 24 : pollfd pfd{};
109 24 : pfd.fd = fd;
110 24 : pfd.events = POLLPRI;
111 : int r;
112 : do
113 : {
114 24 : r = ::poll(&pfd, 1, 0);
115 : }
116 24 : while (r < 0 && errno == EINTR);
117 24 : if (r < 0)
118 : {
119 MIS 0 : err = (errno == EAGAIN || errno == EWOULDBLOCK) ? ENOMEM : errno;
120 0 : return true;
121 : }
122 HIT 24 : if (r == 0)
123 22 : return false;
124 2 : if (pfd.revents & POLLNVAL)
125 MIS 0 : err = EBADF;
126 HIT 2 : else if (pfd.revents & (POLLERR | POLLHUP))
127 2 : err = EIO;
128 MIS 0 : else if (is_fifo(fd))
129 : // Darwin reports POLLPRI on a pipe that merely holds data;
130 : // a FIFO has no exceptional condition to signal.
131 0 : return false;
132 HIT 2 : return true;
133 : }
134 :
135 MIS 0 : static bool is_fifo(int fd) noexcept
136 : {
137 : struct stat st;
138 0 : return ::fstat(fd, &st) == 0 && S_ISFIFO(st.st_mode);
139 : }
140 :
141 HIT 35 : void perform_io() noexcept override
142 : {
143 35 : int err = 0;
144 35 : if (probe(this->fd, this->wait_event, err))
145 10 : this->complete(err, 0);
146 : else
147 25 : this->complete(EAGAIN, 0);
148 35 : }
149 :
150 : void operator()() override;
151 : };
152 :
153 : // --- Deferred implementations (instantiated when Descriptor is complete) ---
154 :
155 : template<class Traits, class Descriptor, class Acceptor>
156 : void
157 55 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>::operator()()
158 : {
159 55 : complete_io_op(*this);
160 55 : }
161 :
162 : template<class Traits, class Descriptor, class Acceptor>
163 : void
164 6 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>::cancel() noexcept
165 : {
166 : // A descriptor op is only ever started against a descriptor impl, so
167 : // the acceptor arm of the stream op's cancel() has no counterpart.
168 6 : if (this->socket_impl_)
169 6 : this->socket_impl_->cancel_single_op(*this);
170 : else
171 MIS 0 : this->request_cancel();
172 HIT 6 : }
173 :
174 : template<class Traits, class Descriptor, class Acceptor>
175 : void
176 27 : reactor_descriptor_wait_op<Traits, Descriptor, Acceptor>::operator()()
177 : {
178 27 : complete_wait_op(*this);
179 27 : }
180 :
181 : // ============================================================
182 : // Descriptor implementation
183 : // ============================================================
184 :
185 : /** CRTP base for reactor-backed posix_stream_descriptor implementations.
186 :
187 : Holds the adopted descriptor, its reactor registration state, and
188 : the five op slots (read, write, and one wait per direction).
189 :
190 : @tparam Derived The named final class (CRTP self).
191 : @tparam Traits Backend traits (epoll_traits, kqueue_traits, ...).
192 : @tparam Service The backend's descriptor service type.
193 : @tparam Acceptor Placeholder acceptor type for the op base.
194 : */
195 : template<class Derived, class Traits, class Service, class Acceptor>
196 : class reactor_descriptor
197 : : public posix_stream_descriptor::implementation
198 : , public std::enable_shared_from_this<Derived>
199 : , public reactor_io_core<Derived, Service, typename Traits::desc_state_type>
200 : , public intrusive_list<Derived>::node
201 : {
202 : using base_op = reactor_descriptor_base_op<Traits, Derived, Acceptor>;
203 : using read_op = reactor_descriptor_read_op<Traits, Derived, Acceptor>;
204 : using write_op = reactor_descriptor_write_op<Traits, Derived, Acceptor>;
205 : using wait_op = reactor_descriptor_wait_op<Traits, Derived, Acceptor>;
206 :
207 : using core_type =
208 : reactor_io_core<Derived, Service, typename Traits::desc_state_type>;
209 : friend core_type;
210 :
211 : protected:
212 : // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
213 121 : explicit reactor_descriptor(Service& svc) noexcept : core_type(svc) {}
214 :
215 : using core_type::svc_;
216 :
217 : public:
218 121 : ~reactor_descriptor() override = default;
219 :
220 : using core_type::desc_state_;
221 :
222 : // --- Virtual method overrides ---
223 :
224 49 : std::coroutine_handle<> read_some(
225 : std::coroutine_handle<> h,
226 : capy::executor_ref ex,
227 : buffer_param param,
228 : std::stop_token token,
229 : std::error_code* ec,
230 : std::size_t* bytes_out) override
231 : {
232 49 : return do_read_some(h, ex, param, token, ec, bytes_out);
233 : }
234 :
235 20 : std::coroutine_handle<> write_some(
236 : std::coroutine_handle<> h,
237 : capy::executor_ref ex,
238 : buffer_param param,
239 : std::stop_token token,
240 : std::error_code* ec,
241 : std::size_t* bytes_out) override
242 : {
243 20 : return do_write_some(h, ex, param, token, ec, bytes_out);
244 : }
245 :
246 37 : std::coroutine_handle<> wait(
247 : std::coroutine_handle<> h,
248 : capy::executor_ref ex,
249 : wait_type w,
250 : std::stop_token token,
251 : std::error_code* ec) override
252 : {
253 37 : return do_wait(h, ex, w, token, ec);
254 : }
255 :
256 262 : native_handle_type native_handle() const noexcept override
257 : {
258 262 : return fd_;
259 : }
260 :
261 : native_handle_type release_descriptor() noexcept override;
262 :
263 7 : void cancel() noexcept override
264 : {
265 7 : this->cancel_all();
266 7 : }
267 :
268 : // --- Service-facing (non-virtual) ---
269 :
270 : /** Adopt the fd, initialize descriptor state, and register it.
271 :
272 : @param fd The descriptor to adopt.
273 :
274 : @return The error if the reactor rejects the descriptor, in
275 : which case the implementation is left closed and the caller
276 : retains ownership of @a fd; otherwise a default constructed
277 : error code.
278 : */
279 : std::error_code init_and_register(int fd) noexcept;
280 :
281 : /// Close the descriptor and cancel pending operations.
282 : void close_descriptor() noexcept;
283 :
284 : private:
285 : /** Arm O_NONBLOCK, once, before the first speculative syscall.
286 :
287 : Reports an errno rather than an error_code because that is what
288 : the op result model records; the round trip is lossless here
289 : because fcntl only fails with codes make_err passes through.
290 : */
291 67 : int arm_nonblocking() noexcept
292 : {
293 : // Relaxed: ensure_nonblocking is idempotent, so a duplicate
294 : // fcntl from a racing first read and write is harmless.
295 67 : if (nonblocking_.load(std::memory_order_relaxed))
296 6 : return 0;
297 61 : if (auto ec = ensure_nonblocking(fd_))
298 MIS 0 : return ec.value();
299 HIT 61 : nonblocking_.store(true, std::memory_order_relaxed);
300 61 : return 0;
301 : }
302 :
303 : std::coroutine_handle<> do_read_some(
304 : std::coroutine_handle<>,
305 : capy::executor_ref,
306 : buffer_param,
307 : std::stop_token const&,
308 : std::error_code*,
309 : std::size_t*);
310 :
311 : std::coroutine_handle<> do_write_some(
312 : std::coroutine_handle<>,
313 : capy::executor_ref,
314 : buffer_param,
315 : std::stop_token const&,
316 : std::error_code*,
317 : std::size_t*);
318 :
319 : std::coroutine_handle<> do_wait(
320 : std::coroutine_handle<>,
321 : capy::executor_ref,
322 : wait_type,
323 : std::stop_token const&,
324 : std::error_code*);
325 :
326 : /// Apply @a fn to each of the five op slots.
327 : template<class Fn>
328 327 : void for_each_op(Fn fn) noexcept
329 : {
330 327 : fn(rd_);
331 327 : fn(wr_);
332 327 : fn(wait_rd_);
333 327 : fn(wait_wr_);
334 327 : fn(wait_er_);
335 327 : }
336 :
337 : /// Apply @a fn to each op and the descriptor_state slot it parks in.
338 : template<class Fn>
339 433 : void for_each_desc_entry(Fn fn) noexcept
340 : {
341 433 : fn(rd_, desc_state_.read_op);
342 433 : fn(wr_, desc_state_.write_op);
343 433 : fn(wait_rd_, desc_state_.wait_read_op);
344 433 : fn(wait_wr_, desc_state_.wait_write_op);
345 433 : fn(wait_er_, desc_state_.wait_error_op);
346 433 : }
347 :
348 : /// Sweep every op slot, then drop the reactor registration.
349 : void quiesce() noexcept;
350 :
351 : reactor_op_base** op_to_desc_slot(base_op& op) noexcept;
352 :
353 : int fd_ = -1;
354 : std::atomic<bool> nonblocking_{false};
355 :
356 : read_op rd_;
357 : write_op wr_;
358 : wait_op wait_rd_;
359 : wait_op wait_wr_;
360 : wait_op wait_er_;
361 : };
362 :
363 : // ============================================================
364 : // Registration and teardown
365 : // ============================================================
366 :
367 : template<class Derived, class Traits, class Service, class Acceptor>
368 : std::error_code
369 106 : reactor_descriptor<Derived, Traits, Service, Acceptor>::init_and_register(
370 : int fd) noexcept
371 : {
372 106 : fd_ = fd;
373 106 : if (auto ec = this->register_fd(fd))
374 : {
375 MIS 0 : fd_ = -1;
376 0 : return ec;
377 : }
378 HIT 106 : return {};
379 : }
380 :
381 : template<class Derived, class Traits, class Service, class Acceptor>
382 : void
383 320 : reactor_descriptor<Derived, Traits, Service, Acceptor>::quiesce() noexcept
384 : {
385 320 : this->abandon_all();
386 320 : this->unregister_fd(fd_);
387 : // The next adopted fd starts from an unknown flag state.
388 320 : nonblocking_.store(false, std::memory_order_relaxed);
389 320 : }
390 :
391 : template<class Derived, class Traits, class Service, class Acceptor>
392 : void
393 316 : reactor_descriptor<Derived, Traits, Service, Acceptor>::
394 : close_descriptor() noexcept
395 : {
396 316 : quiesce();
397 :
398 316 : if (fd_ >= 0)
399 : {
400 102 : ::close(fd_);
401 102 : fd_ = -1;
402 : }
403 316 : }
404 :
405 : template<class Derived, class Traits, class Service, class Acceptor>
406 : native_handle_type
407 4 : reactor_descriptor<Derived, Traits, Service, Acceptor>::
408 : release_descriptor() noexcept
409 : {
410 4 : quiesce();
411 :
412 : // Do NOT close -- the caller takes ownership.
413 4 : native_handle_type released = fd_;
414 4 : fd_ = -1;
415 4 : return released;
416 : }
417 :
418 : // ============================================================
419 : // Op slot lookup
420 : // ============================================================
421 :
422 : template<class Derived, class Traits, class Service, class Acceptor>
423 : reactor_op_base**
424 6 : reactor_descriptor<Derived, Traits, Service, Acceptor>::op_to_desc_slot(
425 : base_op& op) noexcept
426 : {
427 6 : if (&op == static_cast<void*>(&rd_))
428 6 : return &desc_state_.read_op;
429 MIS 0 : if (&op == static_cast<void*>(&wr_))
430 0 : return &desc_state_.write_op;
431 0 : if (&op == static_cast<void*>(&wait_rd_))
432 0 : return &desc_state_.wait_read_op;
433 0 : if (&op == static_cast<void*>(&wait_wr_))
434 0 : return &desc_state_.wait_write_op;
435 0 : if (&op == static_cast<void*>(&wait_er_))
436 0 : return &desc_state_.wait_error_op;
437 0 : return nullptr;
438 : }
439 :
440 : // ============================================================
441 : // I/O dispatch
442 : // ============================================================
443 :
444 : template<class Derived, class Traits, class Service, class Acceptor>
445 : std::coroutine_handle<>
446 HIT 49 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_read_some(
447 : std::coroutine_handle<> h,
448 : capy::executor_ref ex,
449 : buffer_param param,
450 : std::stop_token const& token,
451 : std::error_code* ec,
452 : std::size_t* bytes_out)
453 : {
454 49 : auto& op = rd_;
455 49 : op.reset();
456 49 : op.h = h;
457 49 : op.ex = ex;
458 49 : op.ec_out = ec;
459 49 : op.bytes_out = bytes_out;
460 :
461 : // Closed-object contract: complete with bad_file_descriptor without
462 : // touching the kernel or the unregistered descriptor state.
463 49 : if (fd_ < 0)
464 : {
465 MIS 0 : op.start(token, static_cast<Derived*>(this));
466 0 : op.impl_ptr = this->shared_from_this();
467 0 : op.complete(EBADF, 0);
468 0 : svc_.post(&op);
469 0 : return std::noop_coroutine();
470 : }
471 :
472 HIT 49 : capy::mutable_buffer bufs[read_op::max_buffers];
473 49 : op.iovec_count =
474 49 : static_cast<int>(param.copy_to(bufs, read_op::max_buffers));
475 :
476 49 : if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
477 : {
478 2 : op.empty_buffer_read = true;
479 2 : op.start(token, static_cast<Derived*>(this));
480 2 : op.impl_ptr = this->shared_from_this();
481 2 : op.complete(0, 0);
482 2 : svc_.post(&op);
483 2 : return std::noop_coroutine();
484 : }
485 :
486 : // The first transferring operation is what arms O_NONBLOCK; assign()
487 : // and wait() never do.
488 47 : if (int const nerr = arm_nonblocking())
489 : {
490 MIS 0 : op.start(token, static_cast<Derived*>(this));
491 0 : op.impl_ptr = this->shared_from_this();
492 0 : op.complete(nerr, 0);
493 0 : svc_.post(&op);
494 0 : return std::noop_coroutine();
495 : }
496 :
497 HIT 96 : for (int i = 0; i < op.iovec_count; ++i)
498 : {
499 49 : op.iovecs[i].iov_base = bufs[i].data();
500 49 : op.iovecs[i].iov_len = bufs[i].size();
501 : }
502 :
503 : // Speculative read; the single-buffer case uses read() so the kernel
504 : // skips the readv iov_iter setup.
505 : ssize_t n;
506 47 : if (op.iovec_count == 1)
507 : {
508 : do
509 : {
510 45 : n = ::read(fd_, bufs[0].data(), bufs[0].size());
511 : }
512 45 : while (n < 0 && errno == EINTR);
513 : }
514 : else
515 : {
516 : do
517 : {
518 2 : n = ::readv(fd_, op.iovecs, op.iovec_count);
519 : }
520 2 : while (n < 0 && errno == EINTR);
521 : }
522 :
523 47 : if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
524 : {
525 23 : int err = (n < 0) ? errno : 0;
526 23 : auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
527 :
528 23 : if (svc_.scheduler().try_consume_inline_budget())
529 : {
530 6 : if (err)
531 MIS 0 : *ec = make_err(err);
532 HIT 6 : else if (n == 0)
533 2 : *ec = capy::error::eof;
534 : else
535 4 : *ec = {};
536 6 : *bytes_out = bytes;
537 6 : op.cont.h = h;
538 6 : return dispatch_coro(ex, op.cont);
539 : }
540 17 : op.start(token, static_cast<Derived*>(this));
541 17 : op.impl_ptr = this->shared_from_this();
542 17 : op.complete(err, bytes);
543 17 : svc_.post(&op);
544 17 : return std::noop_coroutine();
545 : }
546 :
547 : // EAGAIN — register with reactor
548 24 : op.fd = fd_;
549 24 : op.start(token, static_cast<Derived*>(this));
550 24 : op.impl_ptr = this->shared_from_this();
551 :
552 24 : this->register_op(op, desc_state_.read_op, desc_state_.read_ready);
553 24 : return std::noop_coroutine();
554 : }
555 :
556 : template<class Derived, class Traits, class Service, class Acceptor>
557 : std::coroutine_handle<>
558 20 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_write_some(
559 : std::coroutine_handle<> h,
560 : capy::executor_ref ex,
561 : buffer_param param,
562 : std::stop_token const& token,
563 : std::error_code* ec,
564 : std::size_t* bytes_out)
565 : {
566 20 : auto& op = wr_;
567 20 : op.reset();
568 20 : op.h = h;
569 20 : op.ex = ex;
570 20 : op.ec_out = ec;
571 20 : op.bytes_out = bytes_out;
572 :
573 20 : if (fd_ < 0)
574 : {
575 MIS 0 : op.start(token, static_cast<Derived*>(this));
576 0 : op.impl_ptr = this->shared_from_this();
577 0 : op.complete(EBADF, 0);
578 0 : svc_.post(&op);
579 0 : return std::noop_coroutine();
580 : }
581 :
582 HIT 20 : capy::mutable_buffer bufs[write_op::max_buffers];
583 20 : op.iovec_count =
584 20 : static_cast<int>(param.copy_to(bufs, write_op::max_buffers));
585 :
586 20 : if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
587 : {
588 MIS 0 : op.start(token, static_cast<Derived*>(this));
589 0 : op.impl_ptr = this->shared_from_this();
590 0 : op.complete(0, 0);
591 0 : svc_.post(&op);
592 0 : return std::noop_coroutine();
593 : }
594 :
595 HIT 20 : if (int const nerr = arm_nonblocking())
596 : {
597 MIS 0 : op.start(token, static_cast<Derived*>(this));
598 0 : op.impl_ptr = this->shared_from_this();
599 0 : op.complete(nerr, 0);
600 0 : svc_.post(&op);
601 0 : return std::noop_coroutine();
602 : }
603 :
604 HIT 42 : for (int i = 0; i < op.iovec_count; ++i)
605 : {
606 22 : op.iovecs[i].iov_base = bufs[i].data();
607 22 : op.iovecs[i].iov_len = bufs[i].size();
608 : }
609 :
610 : // Speculative write; the single-buffer case skips the iov_iter setup.
611 : ssize_t n;
612 20 : if (op.iovec_count == 1)
613 : {
614 36 : n = write_op::write_policy::write_one(
615 18 : fd_, bufs[0].data(), bufs[0].size());
616 : }
617 : else
618 : {
619 2 : n = write_op::write_policy::write(fd_, op.iovecs, op.iovec_count);
620 : }
621 :
622 20 : if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
623 : {
624 14 : int err = (n < 0) ? errno : 0;
625 14 : auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
626 :
627 14 : if (svc_.scheduler().try_consume_inline_budget())
628 : {
629 2 : *ec = err ? make_err(err) : std::error_code{};
630 2 : *bytes_out = bytes;
631 2 : op.cont.h = h;
632 2 : return dispatch_coro(ex, op.cont);
633 : }
634 12 : op.start(token, static_cast<Derived*>(this));
635 12 : op.impl_ptr = this->shared_from_this();
636 12 : op.complete(err, bytes);
637 12 : svc_.post(&op);
638 12 : return std::noop_coroutine();
639 : }
640 :
641 : // EAGAIN — register with reactor
642 6 : op.fd = fd_;
643 6 : op.start(token, static_cast<Derived*>(this));
644 6 : op.impl_ptr = this->shared_from_this();
645 :
646 6 : this->register_op(op, desc_state_.write_op, desc_state_.write_ready, true);
647 6 : return std::noop_coroutine();
648 : }
649 :
650 : template<class Derived, class Traits, class Service, class Acceptor>
651 : std::coroutine_handle<>
652 37 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_wait(
653 : std::coroutine_handle<> h,
654 : capy::executor_ref ex,
655 : wait_type w,
656 : std::stop_token const& token,
657 : std::error_code* ec)
658 : {
659 : // Pick refs up-front to avoid duplicating the register_op call.
660 : wait_op* op_ptr;
661 : reactor_op_base** desc_slot_ptr;
662 : std::uint32_t event;
663 :
664 37 : if (w == wait_type::read)
665 : {
666 16 : op_ptr = &wait_rd_;
667 16 : desc_slot_ptr = &desc_state_.wait_read_op;
668 16 : event = reactor_event_read;
669 : }
670 21 : else if (w == wait_type::write)
671 : {
672 8 : op_ptr = &wait_wr_;
673 8 : desc_slot_ptr = &desc_state_.wait_write_op;
674 8 : event = reactor_event_write;
675 : }
676 : else // wait_type::error
677 : {
678 13 : op_ptr = &wait_er_;
679 13 : desc_slot_ptr = &desc_state_.wait_error_op;
680 13 : event = reactor_event_error;
681 : }
682 :
683 37 : auto& op = *op_ptr;
684 :
685 : // Speculative probe: an edge-triggered reactor cannot report a
686 : // condition that already holds, so a wait initiated on an already
687 : // ready descriptor would otherwise park forever. No syscall here
688 : // modifies the descriptor -- in particular O_NONBLOCK is untouched.
689 37 : int perr = 0;
690 37 : if (wait_op::probe(fd_, event, perr))
691 : {
692 12 : if (svc_.scheduler().try_consume_inline_budget())
693 : {
694 2 : *ec = perr ? make_err(perr) : std::error_code{};
695 2 : op.cont.h = h;
696 2 : return dispatch_coro(ex, op.cont);
697 : }
698 10 : op.reset();
699 10 : op.wait_event = event;
700 10 : op.h = h;
701 10 : op.ex = ex;
702 10 : op.ec_out = ec;
703 10 : op.fd = fd_;
704 10 : op.start(token, static_cast<Derived*>(this));
705 10 : op.impl_ptr = this->shared_from_this();
706 10 : op.complete(perr, 0);
707 10 : svc_.post(&op);
708 10 : return std::noop_coroutine();
709 : }
710 :
711 25 : op.reset();
712 25 : op.wait_event = event;
713 25 : op.h = h;
714 25 : op.ex = ex;
715 25 : op.ec_out = ec;
716 25 : op.fd = fd_;
717 25 : op.start(token, static_cast<Derived*>(this));
718 25 : op.impl_ptr = this->shared_from_this();
719 :
720 : // Force register_op's ready path so the wait op re-probes under the
721 : // descriptor mutex before parking. A stale write_ready latched at
722 : // registration would otherwise report a full pipe as writable.
723 25 : bool force_probe = true;
724 25 : this->register_op(
725 : op, *desc_slot_ptr, force_probe, event == reactor_event_write);
726 25 : return std::noop_coroutine();
727 : }
728 :
729 : } // namespace boost::corosio::detail
730 :
731 : #endif // BOOST_COROSIO_POSIX
732 :
733 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
|