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_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_op.hpp>
26 : #include <boost/corosio/native/detail/reactor/reactor_op_complete.hpp>
27 : #include <boost/capy/buffers.hpp>
28 :
29 : #include <coroutine>
30 : #include <memory>
31 : #include <mutex>
32 : #include <utility>
33 :
34 : #include <errno.h>
35 : #include <sys/uio.h>
36 : #include <unistd.h>
37 :
38 : /* Reactor-backed implementation of posix_descriptor.
39 :
40 : Deliberately does not derive from reactor_basic_socket: that base
41 : is rooted in native_socket_base, which overrides local_endpoint(),
42 : set_option() and get_option() on its ImplBase. posix_descriptor has
43 : no socket verbs, so the shared logic (init_and_register, register_op,
44 : the cancel/close/release op sweeps) is carried here against five op
45 : slots instead of eight.
46 :
47 : The one behavior that is genuinely new: O_NONBLOCK is armed lazily,
48 : on the first read_some/write_some and never from assign() or wait().
49 : The flag lives on the shared open file description, so arming it is
50 : visible to every other holder of that description -- which is why a
51 : wait()-only user must never trigger it.
52 :
53 : PARALLEL COPY: init_and_register, register_op, cancel_single_op and
54 : the cancel/close/release op sweeps here mirror the socket versions in
55 : reactor_basic_socket.hpp (init_and_register, register_op,
56 : cancel_single_op, do_cancel, do_close_socket, do_release_socket).
57 : The two are separate code because that base carries native_socket_base
58 : and its socket verbs; they are not separate protocols. A fix to the
59 : cancel/park protocol -- slot claiming under desc_state_.mutex, the
60 : cached-edge replay in register_op, the impl_ref_ pinning during
61 : teardown -- belongs in both files.
62 : */
63 :
64 : namespace boost::corosio::detail {
65 :
66 : // ============================================================
67 : // Op types
68 : // ============================================================
69 :
70 : /* Descriptor op family.
71 :
72 : Mirrors reactor_stream_ops.hpp. The Acceptor parameter is a
73 : placeholder: reactor_op is parameterized on both a socket and an
74 : acceptor impl type, and the shared completion helpers name
75 : acceptor_impl_ in a branch that a descriptor op never takes but the
76 : compiler still instantiates. Passing the backend's acceptor type (as
77 : reactor_dgram_socket_impl already does) keeps those helpers shared
78 : rather than duplicated here.
79 :
80 : @tparam Traits Backend traits (epoll_traits, kqueue_traits, ...).
81 : @tparam Descriptor The concrete descriptor type (forward-declared).
82 : @tparam Acceptor Placeholder acceptor type for the op base.
83 : */
84 :
85 : template<class Traits, class Descriptor, class Acceptor>
86 : struct reactor_descriptor_base_op : reactor_op<Descriptor, Acceptor>
87 : {
88 : void operator()() override;
89 : void cancel() noexcept override;
90 : };
91 :
92 : template<class Traits, class Descriptor, class Acceptor>
93 : struct reactor_descriptor_read_op final
94 : : reactor_read_op<reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>
95 : {};
96 :
97 : template<class Traits, class Descriptor, class Acceptor>
98 : struct reactor_descriptor_write_op final
99 : : reactor_write_op<
100 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>,
101 : typename Traits::descriptor_write_policy>
102 : {};
103 :
104 : template<class Traits, class Descriptor, class Acceptor>
105 : struct reactor_descriptor_wait_op final
106 : : reactor_wait_op<reactor_descriptor_base_op<Traits, Descriptor, Acceptor>>
107 : {
108 : void operator()() override;
109 : };
110 :
111 : // --- Deferred implementations (instantiated when Descriptor is complete) ---
112 :
113 : template<class Traits, class Descriptor, class Acceptor>
114 : void
115 HIT 25 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>::operator()()
116 : {
117 25 : complete_io_op(*this);
118 25 : }
119 :
120 : template<class Traits, class Descriptor, class Acceptor>
121 : void
122 2 : reactor_descriptor_base_op<Traits, Descriptor, Acceptor>::cancel() noexcept
123 : {
124 : // A descriptor op is only ever started against a descriptor impl, so
125 : // the acceptor arm of the stream op's cancel() has no counterpart.
126 2 : if (this->socket_impl_)
127 2 : this->socket_impl_->cancel_single_op(*this);
128 : else
129 MIS 0 : this->request_cancel();
130 HIT 2 : }
131 :
132 : template<class Traits, class Descriptor, class Acceptor>
133 : void
134 16 : reactor_descriptor_wait_op<Traits, Descriptor, Acceptor>::operator()()
135 : {
136 16 : complete_wait_op(*this);
137 16 : }
138 :
139 : // ============================================================
140 : // Descriptor implementation
141 : // ============================================================
142 :
143 : /** CRTP base for reactor-backed posix_descriptor implementations.
144 :
145 : Holds the adopted descriptor, its reactor registration state, and
146 : the five op slots (read, write, and one wait per direction).
147 :
148 : @tparam Derived The named final class (CRTP self).
149 : @tparam Traits Backend traits (epoll_traits, kqueue_traits, ...).
150 : @tparam Service The backend's descriptor service type.
151 : @tparam Acceptor Placeholder acceptor type for the op base.
152 : */
153 : template<class Derived, class Traits, class Service, class Acceptor>
154 : class reactor_descriptor
155 : : public posix_descriptor::implementation
156 : , public std::enable_shared_from_this<Derived>
157 : , public intrusive_list<Derived>::node
158 : {
159 : using base_op = reactor_descriptor_base_op<Traits, Derived, Acceptor>;
160 : using read_op = reactor_descriptor_read_op<Traits, Derived, Acceptor>;
161 : using write_op = reactor_descriptor_write_op<Traits, Derived, Acceptor>;
162 : using wait_op = reactor_descriptor_wait_op<Traits, Derived, Acceptor>;
163 :
164 : protected:
165 : // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
166 59 : explicit reactor_descriptor(Service& svc) noexcept : svc_(svc) {}
167 :
168 : public:
169 59 : ~reactor_descriptor() override = default;
170 :
171 : /// Per-descriptor state for persistent reactor registration.
172 : typename Traits::desc_state_type desc_state_;
173 :
174 : // --- Virtual method overrides ---
175 :
176 23 : std::coroutine_handle<> read_some(
177 : std::coroutine_handle<> h,
178 : capy::executor_ref ex,
179 : buffer_param param,
180 : std::stop_token token,
181 : std::error_code* ec,
182 : std::size_t* bytes_out) override
183 : {
184 23 : return do_read_some(h, ex, param, token, ec, bytes_out);
185 : }
186 :
187 6 : std::coroutine_handle<> write_some(
188 : std::coroutine_handle<> h,
189 : capy::executor_ref ex,
190 : buffer_param param,
191 : std::stop_token token,
192 : std::error_code* ec,
193 : std::size_t* bytes_out) override
194 : {
195 6 : return do_write_some(h, ex, param, token, ec, bytes_out);
196 : }
197 :
198 16 : std::coroutine_handle<> wait(
199 : std::coroutine_handle<> h,
200 : capy::executor_ref ex,
201 : wait_type w,
202 : std::stop_token token,
203 : std::error_code* ec) override
204 : {
205 16 : return do_wait(h, ex, w, token, ec);
206 : }
207 :
208 139 : native_handle_type native_handle() const noexcept override
209 : {
210 139 : return fd_;
211 : }
212 :
213 : native_handle_type release_descriptor() noexcept override;
214 :
215 3 : void cancel() noexcept override
216 : {
217 3 : do_cancel();
218 3 : }
219 :
220 : // --- Service-facing (non-virtual) ---
221 :
222 : /** Adopt the fd, initialize descriptor state, and register it.
223 :
224 : @param fd The descriptor to adopt.
225 :
226 : @return The error if the reactor rejects the descriptor, in
227 : which case the implementation is left closed and the caller
228 : retains ownership of @a fd; otherwise a default constructed
229 : error code.
230 : */
231 : std::error_code init_and_register(int fd) noexcept;
232 :
233 : /// Close the descriptor and cancel pending operations.
234 : void close_descriptor() noexcept;
235 :
236 : /// Cancel a single pending operation, claiming it from its slot.
237 : template<class Op>
238 : void cancel_single_op(Op& op) noexcept;
239 :
240 : private:
241 : /** Arm O_NONBLOCK, once, before the first speculative syscall.
242 :
243 : Reports an errno rather than an error_code because that is what
244 : the op result model records; the round trip is lossless here
245 : because fcntl only fails with codes make_err passes through.
246 : */
247 29 : int arm_nonblocking() noexcept
248 : {
249 29 : if (nonblocking_)
250 MIS 0 : return 0;
251 HIT 29 : if (auto ec = ensure_nonblocking(fd_))
252 MIS 0 : return ec.value();
253 HIT 29 : nonblocking_ = true;
254 29 : return 0;
255 : }
256 :
257 : std::coroutine_handle<> do_read_some(
258 : std::coroutine_handle<>,
259 : capy::executor_ref,
260 : buffer_param,
261 : std::stop_token const&,
262 : std::error_code*,
263 : std::size_t*);
264 :
265 : std::coroutine_handle<> do_write_some(
266 : std::coroutine_handle<>,
267 : capy::executor_ref,
268 : buffer_param,
269 : std::stop_token const&,
270 : std::error_code*,
271 : std::size_t*);
272 :
273 : std::coroutine_handle<> do_wait(
274 : std::coroutine_handle<>,
275 : capy::executor_ref,
276 : wait_type,
277 : std::stop_token const&,
278 : std::error_code*);
279 :
280 : void do_cancel() noexcept;
281 :
282 : /// Register an op with the reactor, handling cached edge events.
283 : template<class Op>
284 : void register_op(
285 : Op& op,
286 : reactor_op_base*& desc_slot,
287 : bool& ready_flag,
288 : bool is_write_direction = false) noexcept;
289 :
290 : /// Apply @a fn to each of the five op slots.
291 : template<class Fn>
292 213 : void for_each_op(Fn fn) noexcept
293 : {
294 213 : fn(rd_);
295 213 : fn(wr_);
296 213 : fn(wait_rd_);
297 213 : fn(wait_wr_);
298 213 : fn(wait_er_);
299 213 : }
300 :
301 : /** Claim every parked op out of its descriptor_state slot.
302 :
303 : @param claimed Receives the claimed ops; must hold five.
304 : @param teardown Also clear the cached edge flags and, if the
305 : state is queued in the scheduler, pin the impl alive.
306 : @param self Keepalive used by @a teardown.
307 : @return The number of ops claimed.
308 : */
309 : int claim_parked_ops(
310 : reactor_op_base** claimed,
311 : bool teardown,
312 : std::shared_ptr<Derived> const& self) noexcept;
313 :
314 : /// Post claimed ops to the scheduler, keeping the impl alive.
315 : void post_claimed_ops(
316 : reactor_op_base** claimed,
317 : int count,
318 : std::shared_ptr<Derived> const& self) noexcept;
319 :
320 : /// Sweep every op slot, then drop the reactor registration.
321 : void quiesce() noexcept;
322 :
323 : reactor_op_base** op_to_desc_slot(base_op& op) noexcept;
324 :
325 : Service& svc_;
326 : int fd_ = -1;
327 : bool nonblocking_ = false;
328 :
329 : read_op rd_;
330 : write_op wr_;
331 : wait_op wait_rd_;
332 : wait_op wait_wr_;
333 : wait_op wait_er_;
334 : };
335 :
336 : // ============================================================
337 : // Registration and teardown
338 : // ============================================================
339 :
340 : template<class Derived, class Traits, class Service, class Acceptor>
341 : std::error_code
342 47 : reactor_descriptor<Derived, Traits, Service, Acceptor>::init_and_register(
343 : int fd) noexcept
344 : {
345 47 : fd_ = fd;
346 47 : desc_state_.fd = fd;
347 : {
348 : // Every slot this type owns; connect_op is deliberately absent,
349 : // a descriptor has no connect operation to park there.
350 47 : std::lock_guard lock(desc_state_.mutex);
351 47 : desc_state_.read_op = nullptr;
352 47 : desc_state_.write_op = nullptr;
353 47 : desc_state_.wait_read_op = nullptr;
354 47 : desc_state_.wait_write_op = nullptr;
355 47 : desc_state_.wait_error_op = nullptr;
356 47 : }
357 47 : if (auto ec = svc_.scheduler().register_descriptor(fd, &desc_state_))
358 : {
359 : // Undo the partial state so a failed adopt is
360 : // indistinguishable from a closed implementation.
361 MIS 0 : fd_ = -1;
362 0 : desc_state_.fd = -1;
363 0 : desc_state_.registered_events = 0;
364 0 : return ec;
365 : }
366 HIT 47 : return {};
367 : }
368 :
369 : template<class Derived, class Traits, class Service, class Acceptor>
370 : int
371 213 : reactor_descriptor<Derived, Traits, Service, Acceptor>::claim_parked_ops(
372 : reactor_op_base** claimed,
373 : bool teardown,
374 : std::shared_ptr<Derived> const& self) noexcept
375 : {
376 213 : int count = 0;
377 213 : std::lock_guard lock(desc_state_.mutex);
378 1278 : for (auto** slot :
379 213 : {&desc_state_.read_op, &desc_state_.write_op,
380 213 : &desc_state_.wait_read_op, &desc_state_.wait_write_op,
381 213 : &desc_state_.wait_error_op})
382 : {
383 1065 : if (auto* c = std::exchange(*slot, nullptr))
384 5 : claimed[count++] = c;
385 : }
386 213 : if (teardown)
387 : {
388 210 : desc_state_.read_ready = false;
389 210 : desc_state_.write_ready = false;
390 :
391 : // Must be set under the same lock that invoke_deferred_io clears
392 : // is_enqueued_ under, or the impl could be destroyed while the
393 : // scheduler still holds the queued descriptor_state.
394 210 : if (desc_state_.is_enqueued_.load(std::memory_order_acquire))
395 33 : desc_state_.impl_ref_ = self;
396 : }
397 426 : return count;
398 213 : }
399 :
400 : template<class Derived, class Traits, class Service, class Acceptor>
401 : void
402 213 : reactor_descriptor<Derived, Traits, Service, Acceptor>::post_claimed_ops(
403 : reactor_op_base** claimed,
404 : int count,
405 : std::shared_ptr<Derived> const& self) noexcept
406 : {
407 218 : for (int i = 0; i < count; ++i)
408 : {
409 5 : claimed[i]->impl_ptr = self;
410 5 : svc_.post(claimed[i]);
411 5 : svc_.work_finished();
412 : }
413 213 : }
414 :
415 : template<class Derived, class Traits, class Service, class Acceptor>
416 : void
417 3 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_cancel() noexcept
418 : {
419 3 : auto self = this->weak_from_this().lock();
420 3 : if (!self)
421 MIS 0 : return;
422 :
423 HIT 18 : for_each_op([](auto& op) { op.request_cancel(); });
424 :
425 : reactor_op_base* claimed[5];
426 3 : int const count = claim_parked_ops(claimed, /*teardown=*/false, self);
427 3 : post_claimed_ops(claimed, count, self);
428 3 : }
429 :
430 : template<class Derived, class Traits, class Service, class Acceptor>
431 : void
432 210 : reactor_descriptor<Derived, Traits, Service, Acceptor>::quiesce() noexcept
433 : {
434 210 : auto self = this->weak_from_this().lock();
435 210 : if (self)
436 : {
437 1260 : for_each_op([](auto& op) { op.request_cancel(); });
438 :
439 : reactor_op_base* claimed[5];
440 210 : int const count = claim_parked_ops(claimed, /*teardown=*/true, self);
441 210 : post_claimed_ops(claimed, count, self);
442 : }
443 :
444 210 : if (fd_ >= 0 && desc_state_.registered_events != 0)
445 47 : svc_.scheduler().deregister_descriptor(fd_);
446 :
447 210 : desc_state_.registered_events = 0;
448 : // The next adopted fd starts from an unknown flag state.
449 210 : nonblocking_ = false;
450 210 : }
451 :
452 : template<class Derived, class Traits, class Service, class Acceptor>
453 : void
454 208 : reactor_descriptor<Derived, Traits, Service, Acceptor>::
455 : close_descriptor() noexcept
456 : {
457 208 : quiesce();
458 :
459 208 : if (fd_ >= 0)
460 : {
461 45 : ::close(fd_);
462 45 : fd_ = -1;
463 : }
464 208 : desc_state_.fd = -1;
465 208 : }
466 :
467 : template<class Derived, class Traits, class Service, class Acceptor>
468 : native_handle_type
469 2 : reactor_descriptor<Derived, Traits, Service, Acceptor>::
470 : release_descriptor() noexcept
471 : {
472 2 : quiesce();
473 :
474 : // Do NOT close -- the caller takes ownership.
475 2 : native_handle_type released = fd_;
476 2 : fd_ = -1;
477 2 : desc_state_.fd = -1;
478 2 : return released;
479 : }
480 :
481 : // ============================================================
482 : // Op registration and per-op cancellation
483 : // ============================================================
484 :
485 : template<class Derived, class Traits, class Service, class Acceptor>
486 : template<class Op>
487 : void
488 16 : reactor_descriptor<Derived, Traits, Service, Acceptor>::register_op(
489 : Op& op,
490 : reactor_op_base*& desc_slot,
491 : bool& ready_flag,
492 : bool is_write_direction) noexcept
493 : {
494 16 : svc_.work_started();
495 :
496 16 : std::lock_guard lock(desc_state_.mutex);
497 16 : bool io_done = false;
498 16 : if (ready_flag)
499 : {
500 8 : ready_flag = false;
501 8 : op.perform_io();
502 8 : io_done = (op.errn != EAGAIN && op.errn != EWOULDBLOCK);
503 8 : if (!io_done)
504 8 : op.errn = 0;
505 : }
506 :
507 16 : if (io_done || op.cancelled.load(std::memory_order_acquire))
508 : {
509 MIS 0 : svc_.post(&op);
510 0 : svc_.work_finished();
511 : }
512 : else
513 : {
514 HIT 16 : desc_slot = &op;
515 :
516 : // Select must rebuild its fd_sets when a write-direction op
517 : // is parked, so select() watches for writability. Compiled
518 : // away to nothing for epoll and kqueue.
519 : if constexpr (Service::needs_write_notification)
520 : {
521 8 : if (is_write_direction)
522 2 : svc_.scheduler().notify_reactor();
523 : }
524 : }
525 16 : }
526 :
527 : template<class Derived, class Traits, class Service, class Acceptor>
528 : reactor_op_base**
529 2 : reactor_descriptor<Derived, Traits, Service, Acceptor>::op_to_desc_slot(
530 : base_op& op) noexcept
531 : {
532 2 : if (&op == static_cast<void*>(&rd_))
533 2 : return &desc_state_.read_op;
534 MIS 0 : if (&op == static_cast<void*>(&wr_))
535 0 : return &desc_state_.write_op;
536 0 : if (&op == static_cast<void*>(&wait_rd_))
537 0 : return &desc_state_.wait_read_op;
538 0 : if (&op == static_cast<void*>(&wait_wr_))
539 0 : return &desc_state_.wait_write_op;
540 0 : if (&op == static_cast<void*>(&wait_er_))
541 0 : return &desc_state_.wait_error_op;
542 0 : return nullptr;
543 : }
544 :
545 : template<class Derived, class Traits, class Service, class Acceptor>
546 : template<class Op>
547 : void
548 HIT 2 : reactor_descriptor<Derived, Traits, Service, Acceptor>::cancel_single_op(
549 : Op& op) noexcept
550 : {
551 2 : auto self = this->weak_from_this().lock();
552 2 : if (!self)
553 MIS 0 : return;
554 :
555 HIT 2 : op.request_cancel();
556 :
557 2 : reactor_op_base** desc_op_ptr = op_to_desc_slot(op);
558 2 : if (!desc_op_ptr)
559 MIS 0 : return;
560 :
561 HIT 2 : reactor_op_base* claimed = nullptr;
562 : {
563 2 : std::lock_guard lock(desc_state_.mutex);
564 2 : if (*desc_op_ptr == &op)
565 MIS 0 : claimed = std::exchange(*desc_op_ptr, nullptr);
566 : // Not in the slot: request_cancel() above already set
567 : // op.cancelled, which register_op consults before parking
568 : // and the completion decode consults on delivery. Latching
569 : // a descriptor flag here instead would outlive this op and
570 : // cancel the next wait in the same direction.
571 HIT 2 : }
572 2 : if (claimed)
573 : {
574 MIS 0 : op.impl_ptr = self;
575 0 : svc_.post(&op);
576 0 : svc_.work_finished();
577 : }
578 HIT 2 : }
579 :
580 : // ============================================================
581 : // I/O dispatch
582 : // ============================================================
583 :
584 : template<class Derived, class Traits, class Service, class Acceptor>
585 : std::coroutine_handle<>
586 23 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_read_some(
587 : std::coroutine_handle<> h,
588 : capy::executor_ref ex,
589 : buffer_param param,
590 : std::stop_token const& token,
591 : std::error_code* ec,
592 : std::size_t* bytes_out)
593 : {
594 23 : auto& op = rd_;
595 23 : op.reset();
596 23 : op.h = h;
597 23 : op.ex = ex;
598 23 : op.ec_out = ec;
599 23 : op.bytes_out = bytes_out;
600 :
601 : // Closed-object contract: complete with bad_file_descriptor without
602 : // touching the kernel or the unregistered descriptor state.
603 23 : if (fd_ < 0)
604 : {
605 MIS 0 : op.start(token, static_cast<Derived*>(this));
606 0 : op.impl_ptr = this->shared_from_this();
607 0 : op.complete(EBADF, 0);
608 0 : svc_.post(&op);
609 0 : return std::noop_coroutine();
610 : }
611 :
612 HIT 23 : capy::mutable_buffer bufs[read_op::max_buffers];
613 23 : op.iovec_count =
614 23 : static_cast<int>(param.copy_to(bufs, read_op::max_buffers));
615 :
616 23 : if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
617 : {
618 MIS 0 : op.empty_buffer_read = true;
619 0 : op.start(token, static_cast<Derived*>(this));
620 0 : op.impl_ptr = this->shared_from_this();
621 0 : op.complete(0, 0);
622 0 : svc_.post(&op);
623 0 : return std::noop_coroutine();
624 : }
625 :
626 : // The first transferring operation is what arms O_NONBLOCK; assign()
627 : // and wait() never do.
628 HIT 23 : if (int const nerr = arm_nonblocking())
629 : {
630 MIS 0 : op.start(token, static_cast<Derived*>(this));
631 0 : op.impl_ptr = this->shared_from_this();
632 0 : op.complete(nerr, 0);
633 0 : svc_.post(&op);
634 0 : return std::noop_coroutine();
635 : }
636 :
637 HIT 48 : for (int i = 0; i < op.iovec_count; ++i)
638 : {
639 25 : op.iovecs[i].iov_base = bufs[i].data();
640 25 : op.iovecs[i].iov_len = bufs[i].size();
641 : }
642 :
643 : // Speculative read; the single-buffer case uses read() so the kernel
644 : // skips the readv iov_iter setup.
645 : ssize_t n;
646 23 : if (op.iovec_count == 1)
647 : {
648 : do
649 : {
650 21 : n = ::read(fd_, bufs[0].data(), bufs[0].size());
651 : }
652 21 : while (n < 0 && errno == EINTR);
653 : }
654 : else
655 : {
656 : do
657 : {
658 2 : n = ::readv(fd_, op.iovecs, op.iovec_count);
659 : }
660 2 : while (n < 0 && errno == EINTR);
661 : }
662 :
663 23 : if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
664 : {
665 17 : int err = (n < 0) ? errno : 0;
666 17 : auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
667 :
668 17 : if (svc_.scheduler().try_consume_inline_budget())
669 : {
670 4 : if (err)
671 MIS 0 : *ec = make_err(err);
672 HIT 4 : else if (n == 0)
673 MIS 0 : *ec = capy::error::eof;
674 : else
675 HIT 4 : *ec = {};
676 4 : *bytes_out = bytes;
677 4 : op.cont.h = h;
678 4 : return dispatch_coro(ex, op.cont);
679 : }
680 13 : op.start(token, static_cast<Derived*>(this));
681 13 : op.impl_ptr = this->shared_from_this();
682 13 : op.complete(err, bytes);
683 13 : svc_.post(&op);
684 13 : return std::noop_coroutine();
685 : }
686 :
687 : // EAGAIN — register with reactor
688 6 : op.fd = fd_;
689 6 : op.start(token, static_cast<Derived*>(this));
690 6 : op.impl_ptr = this->shared_from_this();
691 :
692 6 : register_op(op, desc_state_.read_op, desc_state_.read_ready);
693 6 : return std::noop_coroutine();
694 : }
695 :
696 : template<class Derived, class Traits, class Service, class Acceptor>
697 : std::coroutine_handle<>
698 6 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_write_some(
699 : std::coroutine_handle<> h,
700 : capy::executor_ref ex,
701 : buffer_param param,
702 : std::stop_token const& token,
703 : std::error_code* ec,
704 : std::size_t* bytes_out)
705 : {
706 6 : auto& op = wr_;
707 6 : op.reset();
708 6 : op.h = h;
709 6 : op.ex = ex;
710 6 : op.ec_out = ec;
711 6 : op.bytes_out = bytes_out;
712 :
713 6 : if (fd_ < 0)
714 : {
715 MIS 0 : op.start(token, static_cast<Derived*>(this));
716 0 : op.impl_ptr = this->shared_from_this();
717 0 : op.complete(EBADF, 0);
718 0 : svc_.post(&op);
719 0 : return std::noop_coroutine();
720 : }
721 :
722 HIT 6 : capy::mutable_buffer bufs[write_op::max_buffers];
723 6 : op.iovec_count =
724 6 : static_cast<int>(param.copy_to(bufs, write_op::max_buffers));
725 :
726 6 : if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
727 : {
728 MIS 0 : op.start(token, static_cast<Derived*>(this));
729 0 : op.impl_ptr = this->shared_from_this();
730 0 : op.complete(0, 0);
731 0 : svc_.post(&op);
732 0 : return std::noop_coroutine();
733 : }
734 :
735 HIT 6 : if (int const nerr = arm_nonblocking())
736 : {
737 MIS 0 : op.start(token, static_cast<Derived*>(this));
738 0 : op.impl_ptr = this->shared_from_this();
739 0 : op.complete(nerr, 0);
740 0 : svc_.post(&op);
741 0 : return std::noop_coroutine();
742 : }
743 :
744 HIT 14 : for (int i = 0; i < op.iovec_count; ++i)
745 : {
746 8 : op.iovecs[i].iov_base = bufs[i].data();
747 8 : op.iovecs[i].iov_len = bufs[i].size();
748 : }
749 :
750 : // Speculative write; the single-buffer case skips the iov_iter setup.
751 : ssize_t n;
752 6 : if (op.iovec_count == 1)
753 : {
754 8 : n = write_op::write_policy::write_one(
755 4 : fd_, bufs[0].data(), bufs[0].size());
756 : }
757 : else
758 : {
759 2 : n = write_op::write_policy::write(fd_, op.iovecs, op.iovec_count);
760 : }
761 :
762 6 : if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
763 : {
764 4 : int err = (n < 0) ? errno : 0;
765 4 : auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
766 :
767 4 : if (svc_.scheduler().try_consume_inline_budget())
768 : {
769 MIS 0 : *ec = err ? make_err(err) : std::error_code{};
770 0 : *bytes_out = bytes;
771 0 : op.cont.h = h;
772 0 : return dispatch_coro(ex, op.cont);
773 : }
774 HIT 4 : op.start(token, static_cast<Derived*>(this));
775 4 : op.impl_ptr = this->shared_from_this();
776 4 : op.complete(err, bytes);
777 4 : svc_.post(&op);
778 4 : return std::noop_coroutine();
779 : }
780 :
781 : // EAGAIN — register with reactor
782 2 : op.fd = fd_;
783 2 : op.start(token, static_cast<Derived*>(this));
784 2 : op.impl_ptr = this->shared_from_this();
785 :
786 2 : register_op(op, desc_state_.write_op, desc_state_.write_ready, true);
787 2 : return std::noop_coroutine();
788 : }
789 :
790 : template<class Derived, class Traits, class Service, class Acceptor>
791 : std::coroutine_handle<>
792 16 : reactor_descriptor<Derived, Traits, Service, Acceptor>::do_wait(
793 : std::coroutine_handle<> h,
794 : capy::executor_ref ex,
795 : wait_type w,
796 : std::stop_token const& token,
797 : std::error_code* ec)
798 : {
799 : // Pick refs up-front to avoid duplicating the register_op call.
800 : wait_op* op_ptr;
801 : reactor_op_base** desc_slot_ptr;
802 : std::uint32_t event;
803 :
804 16 : if (w == wait_type::read)
805 : {
806 12 : op_ptr = &wait_rd_;
807 12 : desc_slot_ptr = &desc_state_.wait_read_op;
808 12 : event = reactor_event_read;
809 : }
810 4 : else if (w == wait_type::write)
811 : {
812 2 : op_ptr = &wait_wr_;
813 2 : desc_slot_ptr = &desc_state_.wait_write_op;
814 2 : event = reactor_event_write;
815 : }
816 : else // wait_type::error
817 : {
818 2 : op_ptr = &wait_er_;
819 2 : desc_slot_ptr = &desc_state_.wait_error_op;
820 2 : event = reactor_event_error;
821 : }
822 :
823 16 : auto& op = *op_ptr;
824 :
825 : // Speculative probe: an edge-triggered reactor cannot report a
826 : // condition that already holds, so a wait initiated on an already
827 : // ready descriptor would otherwise park forever. No syscall here
828 : // modifies the descriptor -- in particular O_NONBLOCK is untouched.
829 16 : int perr = 0;
830 16 : if (wait_op::probe(fd_, event, perr))
831 : {
832 8 : if (svc_.scheduler().try_consume_inline_budget())
833 : {
834 MIS 0 : *ec = perr ? make_err(perr) : std::error_code{};
835 0 : op.cont.h = h;
836 0 : return dispatch_coro(ex, op.cont);
837 : }
838 HIT 8 : op.reset();
839 8 : op.wait_event = event;
840 8 : op.h = h;
841 8 : op.ex = ex;
842 8 : op.ec_out = ec;
843 8 : op.fd = fd_;
844 8 : op.start(token, static_cast<Derived*>(this));
845 8 : op.impl_ptr = this->shared_from_this();
846 8 : op.complete(perr, 0);
847 8 : svc_.post(&op);
848 8 : return std::noop_coroutine();
849 : }
850 :
851 8 : op.reset();
852 8 : op.wait_event = event;
853 8 : op.h = h;
854 8 : op.ex = ex;
855 8 : op.ec_out = ec;
856 8 : op.fd = fd_;
857 8 : op.start(token, static_cast<Derived*>(this));
858 8 : op.impl_ptr = this->shared_from_this();
859 :
860 : // Force register_op's ready path so the wait op re-probes under the
861 : // descriptor mutex before parking. A stale write_ready latched at
862 : // registration would otherwise report a full pipe as writable.
863 8 : bool force_probe = true;
864 8 : register_op(op, *desc_slot_ptr, force_probe, event == reactor_event_write);
865 8 : return std::noop_coroutine();
866 : }
867 :
868 : } // namespace boost::corosio::detail
869 :
870 : #endif // BOOST_COROSIO_POSIX
871 :
872 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_DESCRIPTOR_HPP
|