TLA Line data Source code
1 : //
2 : // Copyright (c) 2026 Steve Gerbino
3 : // Copyright (c) 2026 Michael Vandeberg
4 : //
5 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
6 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7 : //
8 : // Official repository: https://github.com/cppalliance/corosio
9 : //
10 :
11 : #ifndef BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_STREAM_SOCKET_HPP
12 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_STREAM_SOCKET_HPP
13 :
14 : #include <boost/corosio/tcp_socket.hpp>
15 : #include <boost/corosio/shutdown_type.hpp>
16 : #include <boost/corosio/wait_type.hpp>
17 : #include <boost/corosio/native/detail/reactor/reactor_basic_socket.hpp>
18 : #include <boost/corosio/native/detail/reactor/reactor_descriptor_state.hpp>
19 : #include <boost/corosio/detail/dispatch_coro.hpp>
20 : #include <boost/capy/buffers.hpp>
21 :
22 : #include <coroutine>
23 :
24 : #include <errno.h>
25 : #include <sys/socket.h>
26 : #include <sys/uio.h>
27 :
28 : namespace boost::corosio::detail {
29 :
30 : /** CRTP base for reactor-backed stream socket implementations.
31 :
32 : Inherits shared data members and cancel/close/register logic
33 : from reactor_basic_socket. Adds the stream-specific remote
34 : endpoint, shutdown, and I/O dispatch (connect, read, write, wait).
35 :
36 : @tparam Derived The concrete socket type (CRTP).
37 : @tparam Service The backend's socket service type.
38 : @tparam ConnOp The backend's connect op type.
39 : @tparam ReadOp The backend's read op type.
40 : @tparam WriteOp The backend's write op type.
41 : @tparam WaitOp The backend's wait op type.
42 : @tparam DescState The backend's descriptor_state type.
43 : @tparam ImplBase The public vtable base
44 : (tcp_socket::implementation or
45 : local_stream_socket::implementation).
46 : @tparam Endpoint The endpoint type (endpoint or local_endpoint).
47 : */
48 : template<
49 : class Derived,
50 : class Service,
51 : class ConnOp,
52 : class ReadOp,
53 : class WriteOp,
54 : class WaitOp,
55 : class DescState,
56 : class ImplBase = tcp_socket::implementation,
57 : class Endpoint = endpoint>
58 : class reactor_stream_socket
59 : : public reactor_basic_socket<
60 : Derived,
61 : ImplBase,
62 : Service,
63 : DescState,
64 : Endpoint>
65 : {
66 : using base_type =
67 : reactor_basic_socket<Derived, ImplBase, Service, DescState, Endpoint>;
68 : using self_type = reactor_stream_socket<
69 : Derived,
70 : Service,
71 : ConnOp,
72 : ReadOp,
73 : WriteOp,
74 : WaitOp,
75 : DescState,
76 : ImplBase,
77 : Endpoint>;
78 : friend base_type;
79 : friend reactor_io_core<Derived, Service, DescState>;
80 : friend Derived;
81 :
82 : protected:
83 : // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
84 HIT 14466 : explicit reactor_stream_socket(Service& svc) noexcept : base_type(svc) {}
85 :
86 : protected:
87 : Endpoint remote_endpoint_;
88 :
89 : public:
90 : /// Pending connect operation slot.
91 : ConnOp conn_;
92 :
93 : /// Pending read operation slot.
94 : ReadOp rd_;
95 :
96 : /// Pending write operation slot.
97 : WriteOp wr_;
98 :
99 : /// Pending wait-for-read operation slot.
100 : WaitOp wait_rd_;
101 :
102 : /// Pending wait-for-write operation slot.
103 : WaitOp wait_wr_;
104 :
105 : /// Pending wait-for-error operation slot.
106 : WaitOp wait_er_;
107 :
108 14466 : ~reactor_stream_socket() override = default;
109 :
110 : /// Return the cached remote endpoint.
111 62 : Endpoint remote_endpoint() const noexcept override
112 : {
113 62 : return remote_endpoint_;
114 : }
115 :
116 : // --- Virtual method overrides (satisfy ImplBase pure virtuals) ---
117 :
118 4667 : std::coroutine_handle<> connect(
119 : std::coroutine_handle<> h,
120 : capy::executor_ref ex,
121 : Endpoint ep,
122 : std::stop_token token,
123 : std::error_code* ec) override
124 : {
125 4667 : return do_connect(h, ex, ep, token, ec);
126 : }
127 :
128 210443 : std::coroutine_handle<> read_some(
129 : std::coroutine_handle<> h,
130 : capy::executor_ref ex,
131 : buffer_param param,
132 : std::stop_token token,
133 : std::error_code* ec,
134 : std::size_t* bytes_out) override
135 : {
136 210443 : return do_read_some(h, ex, param, token, ec, bytes_out);
137 : }
138 :
139 209742 : std::coroutine_handle<> write_some(
140 : std::coroutine_handle<> h,
141 : capy::executor_ref ex,
142 : buffer_param param,
143 : std::stop_token token,
144 : std::error_code* ec,
145 : std::size_t* bytes_out) override
146 : {
147 209742 : return do_write_some(h, ex, param, token, ec, bytes_out);
148 : }
149 :
150 96 : std::coroutine_handle<> wait(
151 : std::coroutine_handle<> h,
152 : capy::executor_ref ex,
153 : wait_type w,
154 : std::stop_token token,
155 : std::error_code* ec) override
156 : {
157 96 : return do_wait(h, ex, w, token, ec);
158 : }
159 :
160 25 : std::error_code shutdown(corosio::shutdown_type what) noexcept override
161 : {
162 25 : return do_shutdown(static_cast<int>(what));
163 : }
164 :
165 226 : void cancel() noexcept override
166 : {
167 226 : this->do_cancel();
168 226 : }
169 :
170 : // --- End virtual overrides ---
171 :
172 : /// Close the socket (non-virtual, called by the service).
173 : void close_socket() noexcept
174 : {
175 : this->do_close_socket();
176 : }
177 :
178 : /** Shut down part or all of the full-duplex connection.
179 :
180 : @param what 0 = receive, 1 = send, 2 = both.
181 : */
182 25 : std::error_code do_shutdown(int what) noexcept
183 : {
184 : int how;
185 25 : switch (what)
186 : {
187 4 : case 0: // shutdown_receive
188 4 : how = SHUT_RD;
189 4 : break;
190 17 : case 1: // shutdown_send
191 17 : how = SHUT_WR;
192 17 : break;
193 4 : case 2: // shutdown_both
194 4 : how = SHUT_RDWR;
195 4 : break;
196 MIS 0 : default:
197 0 : return make_err(EINVAL);
198 : }
199 HIT 25 : if (::shutdown(this->fd_, how) != 0)
200 2 : return make_err(errno);
201 23 : return {};
202 : }
203 :
204 : /// Cache local and remote endpoints.
205 9392 : void set_endpoints(Endpoint local, Endpoint remote) noexcept
206 : {
207 9392 : this->local_endpoint_ = std::move(local);
208 9392 : remote_endpoint_ = std::move(remote);
209 9392 : }
210 :
211 : /** Shared connect dispatch.
212 :
213 : Tries the connect syscall speculatively. On synchronous
214 : completion, returns via inline budget or posts through queue.
215 : On EINPROGRESS, registers with the reactor.
216 : */
217 : std::coroutine_handle<> do_connect(
218 : std::coroutine_handle<>,
219 : capy::executor_ref,
220 : Endpoint const&,
221 : std::stop_token const&,
222 : std::error_code*);
223 :
224 : /** Shared scatter-read dispatch.
225 :
226 : Tries readv() speculatively. On success or hard error,
227 : returns via inline budget or posts through queue.
228 : On EAGAIN, registers with the reactor.
229 : */
230 : std::coroutine_handle<> do_read_some(
231 : std::coroutine_handle<>,
232 : capy::executor_ref,
233 : buffer_param,
234 : std::stop_token const&,
235 : std::error_code*,
236 : std::size_t*);
237 :
238 : /** Shared gather-write dispatch.
239 :
240 : Tries the write via WriteOp::write_policy speculatively.
241 : On success or hard error, returns via inline budget or
242 : posts through queue. On EAGAIN, registers with the reactor.
243 : */
244 : std::coroutine_handle<> do_write_some(
245 : std::coroutine_handle<>,
246 : capy::executor_ref,
247 : buffer_param,
248 : std::stop_token const&,
249 : std::error_code*,
250 : std::size_t*);
251 :
252 : /** Shared readiness-wait dispatch.
253 :
254 : Every wait type probes the descriptor with a zero-timeout
255 : `poll()` and completes at once if the condition already
256 : holds; otherwise the op re-probes under the descriptor mutex
257 : and parks, completing when a reactor event arrives and a
258 : fresh probe confirms the condition. A write wait therefore
259 : completes only while a non-blocking write can make progress.
260 : */
261 : std::coroutine_handle<> do_wait(
262 : std::coroutine_handle<>,
263 : capy::executor_ref,
264 : wait_type,
265 : std::stop_token const&,
266 : std::error_code*);
267 :
268 : /** Close the socket and cancel pending operations.
269 :
270 : Extends the base do_close_socket() to also reset
271 : the remote endpoint.
272 : */
273 43218 : void do_close_socket() noexcept
274 : {
275 43218 : base_type::do_close_socket();
276 43218 : remote_endpoint_ = Endpoint{};
277 43218 : }
278 :
279 : /// Release ownership of the descriptor and drop the cached peer.
280 8 : native_handle_type do_release_socket() noexcept
281 : {
282 8 : auto fd = base_type::do_release_socket();
283 8 : remote_endpoint_ = Endpoint{};
284 8 : return fd;
285 : }
286 :
287 : private:
288 : // CRTP callbacks for reactor_io_core cancel/close
289 :
290 : template<class Op>
291 219 : reactor_op_base** op_to_desc_slot(Op& op) noexcept
292 : {
293 219 : if (&op == static_cast<void*>(&conn_))
294 MIS 0 : return &this->desc_state_.connect_op;
295 HIT 219 : if (&op == static_cast<void*>(&rd_))
296 208 : return &this->desc_state_.read_op;
297 11 : if (&op == static_cast<void*>(&wr_))
298 6 : return &this->desc_state_.write_op;
299 5 : if (&op == static_cast<void*>(&wait_rd_))
300 3 : return &this->desc_state_.wait_read_op;
301 2 : if (&op == static_cast<void*>(&wait_wr_))
302 MIS 0 : return &this->desc_state_.wait_write_op;
303 HIT 2 : if (&op == static_cast<void*>(&wait_er_))
304 2 : return &this->desc_state_.wait_error_op;
305 MIS 0 : return nullptr;
306 : }
307 :
308 : template<class Fn>
309 HIT 43452 : void for_each_op(Fn fn) noexcept
310 : {
311 43452 : fn(conn_);
312 43452 : fn(rd_);
313 43452 : fn(wr_);
314 43452 : fn(wait_rd_);
315 43452 : fn(wait_wr_);
316 43452 : fn(wait_er_);
317 43452 : }
318 :
319 : template<class Fn>
320 48374 : void for_each_desc_entry(Fn fn) noexcept
321 : {
322 48374 : fn(conn_, this->desc_state_.connect_op);
323 48374 : fn(rd_, this->desc_state_.read_op);
324 48374 : fn(wr_, this->desc_state_.write_op);
325 48374 : fn(wait_rd_, this->desc_state_.wait_read_op);
326 48374 : fn(wait_wr_, this->desc_state_.wait_write_op);
327 48374 : fn(wait_er_, this->desc_state_.wait_error_op);
328 48374 : }
329 : };
330 :
331 : template<
332 : class Derived,
333 : class Service,
334 : class ConnOp,
335 : class ReadOp,
336 : class WriteOp,
337 : class WaitOp,
338 : class DescState,
339 : class ImplBase,
340 : class Endpoint>
341 : std::coroutine_handle<>
342 4667 : reactor_stream_socket<
343 : Derived,
344 : Service,
345 : ConnOp,
346 : ReadOp,
347 : WriteOp,
348 : WaitOp,
349 : DescState,
350 : ImplBase,
351 : Endpoint>::
352 : do_connect(
353 : std::coroutine_handle<> h,
354 : capy::executor_ref ex,
355 : Endpoint const& ep,
356 : std::stop_token const& token,
357 : std::error_code* ec)
358 : {
359 4667 : auto& op = conn_;
360 :
361 4667 : sockaddr_storage storage{};
362 4667 : socklen_t addrlen = to_sockaddr(ep, socket_family(this->fd_), storage);
363 : int result =
364 4667 : ::connect(this->fd_, reinterpret_cast<sockaddr*>(&storage), addrlen);
365 :
366 4667 : if (result == 0)
367 : {
368 29 : sockaddr_storage local_storage{};
369 29 : socklen_t local_len = sizeof(local_storage);
370 29 : if (::getsockname(
371 : this->fd_, reinterpret_cast<sockaddr*>(&local_storage),
372 29 : &local_len) == 0)
373 MIS 0 : this->local_endpoint_ =
374 HIT 29 : from_sockaddr_as(local_storage, local_len, Endpoint{});
375 29 : remote_endpoint_ = ep;
376 : }
377 :
378 4667 : if (result == 0 || errno != EINPROGRESS)
379 : {
380 37 : int err = (result < 0) ? errno : 0;
381 37 : if (this->svc_.scheduler().try_consume_inline_budget())
382 : {
383 MIS 0 : *ec = err ? make_err(err) : std::error_code{};
384 0 : op.cont.h = h;
385 0 : return dispatch_coro(ex, op.cont);
386 : }
387 HIT 37 : op.reset();
388 37 : op.h = h;
389 37 : op.ex = ex;
390 37 : op.ec_out = ec;
391 37 : op.fd = this->fd_;
392 37 : op.target_endpoint = ep;
393 37 : op.start(token, static_cast<Derived*>(this));
394 37 : op.impl_ptr = this->shared_from_this();
395 37 : op.complete(err, 0);
396 37 : this->svc_.post(&op);
397 37 : return std::noop_coroutine();
398 : }
399 :
400 : // EINPROGRESS — register with reactor
401 4630 : op.reset();
402 4630 : op.h = h;
403 4630 : op.ex = ex;
404 4630 : op.ec_out = ec;
405 4630 : op.fd = this->fd_;
406 4630 : op.target_endpoint = ep;
407 4630 : op.start(token, static_cast<Derived*>(this));
408 4630 : op.impl_ptr = this->shared_from_this();
409 :
410 4630 : this->register_op(
411 4630 : op, this->desc_state_.connect_op, this->desc_state_.write_ready, true);
412 4630 : return std::noop_coroutine();
413 : }
414 :
415 : template<
416 : class Derived,
417 : class Service,
418 : class ConnOp,
419 : class ReadOp,
420 : class WriteOp,
421 : class WaitOp,
422 : class DescState,
423 : class ImplBase,
424 : class Endpoint>
425 : std::coroutine_handle<>
426 210443 : reactor_stream_socket<
427 : Derived,
428 : Service,
429 : ConnOp,
430 : ReadOp,
431 : WriteOp,
432 : WaitOp,
433 : DescState,
434 : ImplBase,
435 : Endpoint>::
436 : do_read_some(
437 : std::coroutine_handle<> h,
438 : capy::executor_ref ex,
439 : buffer_param param,
440 : std::stop_token const& token,
441 : std::error_code* ec,
442 : std::size_t* bytes_out)
443 : {
444 210443 : auto& op = rd_;
445 210443 : op.reset();
446 :
447 : // Closed-object contract: complete with bad_file_descriptor without
448 : // touching the kernel or the unregistered descriptor state.
449 210443 : if (this->fd_ < 0)
450 : {
451 8 : op.h = h;
452 8 : op.ex = ex;
453 8 : op.ec_out = ec;
454 8 : op.bytes_out = bytes_out;
455 8 : op.start(token, static_cast<Derived*>(this));
456 8 : op.impl_ptr = this->shared_from_this();
457 8 : op.complete(EBADF, 0);
458 8 : this->svc_.post(&op);
459 8 : return std::noop_coroutine();
460 : }
461 :
462 210435 : capy::mutable_buffer bufs[ReadOp::max_buffers];
463 210435 : op.iovec_count = static_cast<int>(param.copy_to(bufs, ReadOp::max_buffers));
464 :
465 210435 : if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
466 : {
467 4 : op.empty_buffer_read = true;
468 4 : op.h = h;
469 4 : op.ex = ex;
470 4 : op.ec_out = ec;
471 4 : op.bytes_out = bytes_out;
472 4 : op.start(token, static_cast<Derived*>(this));
473 4 : op.impl_ptr = this->shared_from_this();
474 4 : op.complete(0, 0);
475 4 : this->svc_.post(&op);
476 4 : return std::noop_coroutine();
477 : }
478 :
479 420880 : for (int i = 0; i < op.iovec_count; ++i)
480 : {
481 210449 : op.iovecs[i].iov_base = bufs[i].data();
482 210449 : op.iovecs[i].iov_len = bufs[i].size();
483 : }
484 :
485 : // Speculative read; for the single-buffer case use recv() so the
486 : // kernel skips the readv iov_iter setup.
487 : ssize_t n;
488 210431 : if (op.iovec_count == 1)
489 : {
490 : do
491 : {
492 210419 : n = ::recv(this->fd_, bufs[0].data(), bufs[0].size(), 0);
493 : }
494 210419 : while (n < 0 && errno == EINTR);
495 : }
496 : else
497 : {
498 : do
499 : {
500 16 : n = ::readv(this->fd_, op.iovecs, op.iovec_count);
501 : }
502 16 : while (n < 0 && errno == EINTR);
503 : }
504 :
505 210431 : if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
506 : {
507 209650 : int err = (n < 0) ? errno : 0;
508 209650 : auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
509 :
510 209650 : if (this->svc_.scheduler().try_consume_inline_budget())
511 : {
512 167752 : if (err)
513 8 : *ec = make_err(err);
514 167744 : else if (n == 0)
515 17 : *ec = capy::error::eof;
516 : else
517 167727 : *ec = {};
518 167752 : *bytes_out = bytes;
519 167752 : op.cont.h = h;
520 167752 : return dispatch_coro(ex, op.cont);
521 : }
522 41898 : op.h = h;
523 41898 : op.ex = ex;
524 41898 : op.ec_out = ec;
525 41898 : op.bytes_out = bytes_out;
526 41898 : op.start(token, static_cast<Derived*>(this));
527 41898 : op.impl_ptr = this->shared_from_this();
528 41898 : op.complete(err, bytes);
529 41898 : this->svc_.post(&op);
530 41898 : return std::noop_coroutine();
531 : }
532 :
533 : // EAGAIN — register with reactor
534 781 : op.h = h;
535 781 : op.ex = ex;
536 781 : op.ec_out = ec;
537 781 : op.bytes_out = bytes_out;
538 781 : op.fd = this->fd_;
539 781 : op.start(token, static_cast<Derived*>(this));
540 781 : op.impl_ptr = this->shared_from_this();
541 :
542 781 : this->register_op(
543 781 : op, this->desc_state_.read_op, this->desc_state_.read_ready);
544 781 : return std::noop_coroutine();
545 : }
546 :
547 : template<
548 : class Derived,
549 : class Service,
550 : class ConnOp,
551 : class ReadOp,
552 : class WriteOp,
553 : class WaitOp,
554 : class DescState,
555 : class ImplBase,
556 : class Endpoint>
557 : std::coroutine_handle<>
558 209742 : reactor_stream_socket<
559 : Derived,
560 : Service,
561 : ConnOp,
562 : ReadOp,
563 : WriteOp,
564 : WaitOp,
565 : DescState,
566 : ImplBase,
567 : Endpoint>::
568 : do_write_some(
569 : std::coroutine_handle<> h,
570 : capy::executor_ref ex,
571 : buffer_param param,
572 : std::stop_token const& token,
573 : std::error_code* ec,
574 : std::size_t* bytes_out)
575 : {
576 209742 : auto& op = wr_;
577 209742 : op.reset();
578 :
579 : // Closed-object contract: complete with bad_file_descriptor without
580 : // touching the kernel or the unregistered descriptor state.
581 209742 : if (this->fd_ < 0)
582 : {
583 10 : op.h = h;
584 10 : op.ex = ex;
585 10 : op.ec_out = ec;
586 10 : op.bytes_out = bytes_out;
587 10 : op.start(token, static_cast<Derived*>(this));
588 10 : op.impl_ptr = this->shared_from_this();
589 10 : op.complete(EBADF, 0);
590 10 : this->svc_.post(&op);
591 10 : return std::noop_coroutine();
592 : }
593 :
594 209732 : capy::mutable_buffer bufs[WriteOp::max_buffers];
595 209732 : op.iovec_count =
596 209732 : static_cast<int>(param.copy_to(bufs, WriteOp::max_buffers));
597 :
598 209732 : if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
599 : {
600 4 : op.h = h;
601 4 : op.ex = ex;
602 4 : op.ec_out = ec;
603 4 : op.bytes_out = bytes_out;
604 4 : op.start(token, static_cast<Derived*>(this));
605 4 : op.impl_ptr = this->shared_from_this();
606 4 : op.complete(0, 0);
607 4 : this->svc_.post(&op);
608 4 : return std::noop_coroutine();
609 : }
610 :
611 419472 : for (int i = 0; i < op.iovec_count; ++i)
612 : {
613 209744 : op.iovecs[i].iov_base = bufs[i].data();
614 209744 : op.iovecs[i].iov_len = bufs[i].size();
615 : }
616 :
617 : // Speculative write; the single-buffer case dispatches to a
618 : // backend-specific fast path so the kernel skips msghdr/iov_iter
619 : // setup (and so each backend can pick the right SIGPIPE strategy).
620 : ssize_t n;
621 209728 : if (op.iovec_count == 1)
622 : {
623 419432 : n = WriteOp::write_policy::write_one(
624 209716 : this->fd_, bufs[0].data(), bufs[0].size());
625 : }
626 : else
627 : {
628 12 : n = WriteOp::write_policy::write(this->fd_, op.iovecs, op.iovec_count);
629 : }
630 :
631 209728 : if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
632 : {
633 209584 : int err = (n < 0) ? errno : 0;
634 209584 : auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
635 :
636 209584 : if (this->svc_.scheduler().try_consume_inline_budget())
637 : {
638 167611 : *ec = err ? make_err(err) : std::error_code{};
639 167611 : *bytes_out = bytes;
640 167611 : op.cont.h = h;
641 167611 : return dispatch_coro(ex, op.cont);
642 : }
643 41973 : op.h = h;
644 41973 : op.ex = ex;
645 41973 : op.ec_out = ec;
646 41973 : op.bytes_out = bytes_out;
647 41973 : op.start(token, static_cast<Derived*>(this));
648 41973 : op.impl_ptr = this->shared_from_this();
649 41973 : op.complete(err, bytes);
650 41973 : this->svc_.post(&op);
651 41973 : return std::noop_coroutine();
652 : }
653 :
654 : // EAGAIN — register with reactor
655 144 : op.h = h;
656 144 : op.ex = ex;
657 144 : op.ec_out = ec;
658 144 : op.bytes_out = bytes_out;
659 144 : op.fd = this->fd_;
660 144 : op.start(token, static_cast<Derived*>(this));
661 144 : op.impl_ptr = this->shared_from_this();
662 :
663 144 : this->register_op(
664 144 : op, this->desc_state_.write_op, this->desc_state_.write_ready, true);
665 144 : return std::noop_coroutine();
666 : }
667 :
668 : template<
669 : class Derived,
670 : class Service,
671 : class ConnOp,
672 : class ReadOp,
673 : class WriteOp,
674 : class WaitOp,
675 : class DescState,
676 : class ImplBase,
677 : class Endpoint>
678 : std::coroutine_handle<>
679 96 : reactor_stream_socket<
680 : Derived,
681 : Service,
682 : ConnOp,
683 : ReadOp,
684 : WriteOp,
685 : WaitOp,
686 : DescState,
687 : ImplBase,
688 : Endpoint>::
689 : do_wait(
690 : std::coroutine_handle<> h,
691 : capy::executor_ref ex,
692 : wait_type w,
693 : std::stop_token const& token,
694 : std::error_code* ec)
695 : {
696 : // Pick refs up-front to avoid duplicating the register_op call.
697 : WaitOp* op_ptr;
698 : reactor_op_base** desc_slot_ptr;
699 : std::uint32_t event;
700 :
701 96 : if (w == wait_type::read)
702 : {
703 55 : op_ptr = &wait_rd_;
704 55 : desc_slot_ptr = &this->desc_state_.wait_read_op;
705 55 : event = reactor_event_read;
706 : }
707 41 : else if (w == wait_type::write)
708 : {
709 23 : op_ptr = &wait_wr_;
710 23 : desc_slot_ptr = &this->desc_state_.wait_write_op;
711 23 : event = reactor_event_write;
712 : }
713 : else // wait_type::error
714 : {
715 18 : op_ptr = &wait_er_;
716 18 : desc_slot_ptr = &this->desc_state_.wait_error_op;
717 18 : event = reactor_event_error;
718 : }
719 :
720 96 : auto& op = *op_ptr;
721 :
722 : // Speculative probe, mirroring the speculative read: an
723 : // edge-triggered reactor cannot report a condition that already
724 : // holds, so a wait initiated on an already-ready socket would
725 : // otherwise park forever.
726 96 : int perr = 0;
727 96 : if (WaitOp::probe(this->fd_, event, perr))
728 : {
729 40 : if (this->svc_.scheduler().try_consume_inline_budget())
730 : {
731 8 : *ec = perr ? make_err(perr) : std::error_code{};
732 8 : op.cont.h = h;
733 8 : return dispatch_coro(ex, op.cont);
734 : }
735 32 : op.reset();
736 32 : op.wait_event = event;
737 32 : op.h = h;
738 32 : op.ex = ex;
739 32 : op.ec_out = ec;
740 32 : op.fd = this->fd_;
741 32 : op.start(token, static_cast<Derived*>(this));
742 32 : op.impl_ptr = this->shared_from_this();
743 32 : op.complete(perr, 0);
744 32 : this->svc_.post(&op);
745 32 : return std::noop_coroutine();
746 : }
747 :
748 56 : op.reset();
749 56 : op.wait_event = event;
750 56 : op.h = h;
751 56 : op.ex = ex;
752 56 : op.ec_out = ec;
753 56 : op.fd = this->fd_;
754 56 : op.start(token, static_cast<Derived*>(this));
755 56 : op.impl_ptr = this->shared_from_this();
756 :
757 : // Force register_op's ready path so the wait op re-probes under
758 : // the descriptor mutex before parking. An edge consumed between
759 : // the speculative probe above and the park (a concurrent short
760 : // read, or an error event dispatched to an empty slot) would
761 : // otherwise leave the wait parked on a ready socket.
762 56 : bool force_probe = true;
763 56 : this->register_op(
764 : op, *desc_slot_ptr, force_probe, event == reactor_event_write);
765 56 : return std::noop_coroutine();
766 : }
767 :
768 : } // namespace boost::corosio::detail
769 :
770 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_STREAM_SOCKET_HPP
|