TLA Line data Source code
1 : //
2 : // Copyright (c) 2026 Steve Gerbino
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_STREAM_SOCKET_HPP
11 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_STREAM_SOCKET_HPP
12 :
13 : #include <boost/corosio/tcp_socket.hpp>
14 : #include <boost/corosio/shutdown_type.hpp>
15 : #include <boost/corosio/wait_type.hpp>
16 : #include <boost/corosio/detail/config.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 : #include <cstring>
24 :
25 : #include <errno.h>
26 : #include <sys/socket.h>
27 : #include <sys/uio.h>
28 :
29 : namespace boost::corosio::detail {
30 :
31 : /** CRTP base for reactor-backed stream socket implementations.
32 :
33 : Inherits shared data members and cancel/close/register logic
34 : from reactor_basic_socket. Adds the stream-specific remote
35 : endpoint, shutdown, and I/O dispatch (connect, read, write, wait).
36 :
37 : @tparam Derived The concrete socket type (CRTP).
38 : @tparam Service The backend's socket service type.
39 : @tparam ConnOp The backend's connect op type.
40 : @tparam ReadOp The backend's read op type.
41 : @tparam WriteOp The backend's write op type.
42 : @tparam WaitOp The backend's wait op type.
43 : @tparam DescState The backend's descriptor_state type.
44 : @tparam ImplBase The public vtable base
45 : (tcp_socket::implementation or
46 : local_stream_socket::implementation).
47 : @tparam Endpoint The endpoint type (endpoint or local_endpoint).
48 : */
49 : template<
50 : class Derived,
51 : class Service,
52 : class ConnOp,
53 : class ReadOp,
54 : class WriteOp,
55 : class WaitOp,
56 : class DescState,
57 : class ImplBase = tcp_socket::implementation,
58 : class Endpoint = endpoint>
59 : class reactor_stream_socket
60 : : public reactor_basic_socket<
61 : Derived,
62 : ImplBase,
63 : Service,
64 : DescState,
65 : Endpoint>
66 : {
67 : using base_type =
68 : reactor_basic_socket<Derived, ImplBase, Service, DescState, Endpoint>;
69 : using self_type = reactor_stream_socket<
70 : Derived,
71 : Service,
72 : ConnOp,
73 : ReadOp,
74 : WriteOp,
75 : WaitOp,
76 : DescState,
77 : ImplBase,
78 : Endpoint>;
79 : friend base_type;
80 : friend Derived;
81 :
82 : protected:
83 : // NOLINTNEXTLINE(bugprone-crtp-constructor-accessibility)
84 HIT 2133 : 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 2133 : ~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 4919 : 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 4919 : return do_connect(h, ex, ep, token, ec);
126 : }
127 :
128 229162 : 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 229162 : return do_read_some(h, ex, param, token, ec, bytes_out);
137 : }
138 :
139 228472 : 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 228472 : return do_write_some(h, ex, param, token, ec, bytes_out);
148 : }
149 :
150 90 : 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 90 : 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 220 : void cancel() noexcept override
166 : {
167 220 : this->do_cancel();
168 220 : }
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 9893 : void set_endpoints(Endpoint local, Endpoint remote) noexcept
206 : {
207 9893 : this->local_endpoint_ = std::move(local);
208 9893 : remote_endpoint_ = std::move(remote);
209 9893 : }
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 45958 : void do_close_socket() noexcept
274 : {
275 45958 : base_type::do_close_socket();
276 45958 : remote_endpoint_ = Endpoint{};
277 45958 : }
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 : /** Reset op slots and the cached peer for recycling.
288 :
289 : Each op's own `reset()` runs again at its next `do_*` call
290 : before any field is read, so only `remote_endpoint_` needs an
291 : active reset here; `stop_cb` is checked instead of reset — a
292 : still-armed callback would mean a stop_token outlived the op's
293 : completion, which `reset()` always disengages, so asserting it
294 : is what actually catches a missed completion rather than
295 : papering over it.
296 :
297 : @pre refs_ == 0, fd closed and deregistered, no op in flight.
298 : */
299 13144 : void reuse() noexcept
300 : {
301 13144 : base_type::reuse();
302 13144 : remote_endpoint_ = Endpoint{};
303 13144 : BOOST_COROSIO_ASSERT(!conn_.stop_cb);
304 13144 : BOOST_COROSIO_ASSERT(!rd_.stop_cb);
305 13144 : BOOST_COROSIO_ASSERT(!wr_.stop_cb);
306 13144 : BOOST_COROSIO_ASSERT(!wait_rd_.stop_cb);
307 13144 : BOOST_COROSIO_ASSERT(!wait_wr_.stop_cb);
308 13144 : BOOST_COROSIO_ASSERT(!wait_er_.stop_cb);
309 13144 : }
310 :
311 : #if !defined(NDEBUG)
312 : /// Extend base_type::poison() to the cached peer endpoint.
313 15272 : void poison() noexcept
314 : {
315 15272 : base_type::poison();
316 15272 : std::memset(
317 15272 : static_cast<void*>(&remote_endpoint_), 0xDB,
318 : sizeof(remote_endpoint_));
319 15272 : }
320 : #endif
321 :
322 : private:
323 : // CRTP callbacks for reactor_basic_socket cancel/close
324 :
325 : template<class Op>
326 217 : reactor_op_base** op_to_desc_slot(Op& op) noexcept
327 : {
328 217 : if (&op == static_cast<void*>(&conn_))
329 MIS 0 : return &this->desc_state_.connect_op;
330 HIT 217 : if (&op == static_cast<void*>(&rd_))
331 206 : return &this->desc_state_.read_op;
332 11 : if (&op == static_cast<void*>(&wr_))
333 6 : return &this->desc_state_.write_op;
334 5 : if (&op == static_cast<void*>(&wait_rd_))
335 3 : return &this->desc_state_.wait_read_op;
336 2 : if (&op == static_cast<void*>(&wait_wr_))
337 MIS 0 : return &this->desc_state_.wait_write_op;
338 HIT 2 : if (&op == static_cast<void*>(&wait_er_))
339 2 : return &this->desc_state_.wait_error_op;
340 MIS 0 : return nullptr;
341 : }
342 :
343 : template<class Fn>
344 HIT 46186 : void for_each_op(Fn fn) noexcept
345 : {
346 46186 : fn(conn_);
347 46186 : fn(rd_);
348 46186 : fn(wr_);
349 46186 : fn(wait_rd_);
350 46186 : fn(wait_wr_);
351 46186 : fn(wait_er_);
352 46186 : }
353 :
354 : template<class Fn>
355 46186 : void for_each_desc_entry(Fn fn) noexcept
356 : {
357 46186 : fn(conn_, this->desc_state_.connect_op);
358 46186 : fn(rd_, this->desc_state_.read_op);
359 46186 : fn(wr_, this->desc_state_.write_op);
360 46186 : fn(wait_rd_, this->desc_state_.wait_read_op);
361 46186 : fn(wait_wr_, this->desc_state_.wait_write_op);
362 46186 : fn(wait_er_, this->desc_state_.wait_error_op);
363 46186 : }
364 : };
365 :
366 : template<
367 : class Derived,
368 : class Service,
369 : class ConnOp,
370 : class ReadOp,
371 : class WriteOp,
372 : class WaitOp,
373 : class DescState,
374 : class ImplBase,
375 : class Endpoint>
376 : std::coroutine_handle<>
377 4919 : reactor_stream_socket<
378 : Derived,
379 : Service,
380 : ConnOp,
381 : ReadOp,
382 : WriteOp,
383 : WaitOp,
384 : DescState,
385 : ImplBase,
386 : Endpoint>::
387 : do_connect(
388 : std::coroutine_handle<> h,
389 : capy::executor_ref ex,
390 : Endpoint const& ep,
391 : std::stop_token const& token,
392 : std::error_code* ec)
393 : {
394 4919 : auto& op = conn_;
395 :
396 4919 : sockaddr_storage storage{};
397 4919 : socklen_t addrlen = to_sockaddr(ep, socket_family(this->fd_), storage);
398 : int result =
399 4919 : ::connect(this->fd_, reinterpret_cast<sockaddr*>(&storage), addrlen);
400 :
401 4919 : if (result == 0)
402 : {
403 31 : sockaddr_storage local_storage{};
404 31 : socklen_t local_len = sizeof(local_storage);
405 31 : if (::getsockname(
406 : this->fd_, reinterpret_cast<sockaddr*>(&local_storage),
407 31 : &local_len) == 0)
408 MIS 0 : this->local_endpoint_ =
409 HIT 31 : from_sockaddr_as(local_storage, local_len, Endpoint{});
410 31 : remote_endpoint_ = ep;
411 : }
412 :
413 4919 : if (result == 0 || errno != EINPROGRESS)
414 : {
415 39 : int err = (result < 0) ? errno : 0;
416 39 : if (this->svc_.scheduler().try_consume_inline_budget())
417 : {
418 MIS 0 : *ec = err ? make_err(err) : std::error_code{};
419 0 : op.cont.h = h;
420 0 : return dispatch_coro(ex, op.cont);
421 : }
422 HIT 39 : op.reset();
423 39 : op.h = h;
424 39 : op.ex = ex;
425 39 : op.ec_out = ec;
426 39 : op.fd = this->fd_;
427 39 : op.target_endpoint = ep;
428 39 : op.start(token, static_cast<Derived*>(this));
429 39 : op.object_ref_ = detail::object_ref(this);
430 39 : op.complete(err, 0);
431 39 : this->svc_.post(&op);
432 39 : return std::noop_coroutine();
433 : }
434 :
435 : // EINPROGRESS — register with reactor
436 4880 : op.reset();
437 4880 : op.h = h;
438 4880 : op.ex = ex;
439 4880 : op.ec_out = ec;
440 4880 : op.fd = this->fd_;
441 4880 : op.target_endpoint = ep;
442 4880 : op.start(token, static_cast<Derived*>(this));
443 4880 : op.object_ref_ = detail::object_ref(this);
444 :
445 4880 : this->register_op(
446 4880 : op, this->desc_state_.connect_op, this->desc_state_.write_ready, true);
447 4880 : return std::noop_coroutine();
448 : }
449 :
450 : template<
451 : class Derived,
452 : class Service,
453 : class ConnOp,
454 : class ReadOp,
455 : class WriteOp,
456 : class WaitOp,
457 : class DescState,
458 : class ImplBase,
459 : class Endpoint>
460 : std::coroutine_handle<>
461 229162 : reactor_stream_socket<
462 : Derived,
463 : Service,
464 : ConnOp,
465 : ReadOp,
466 : WriteOp,
467 : WaitOp,
468 : DescState,
469 : ImplBase,
470 : Endpoint>::
471 : do_read_some(
472 : std::coroutine_handle<> h,
473 : capy::executor_ref ex,
474 : buffer_param param,
475 : std::stop_token const& token,
476 : std::error_code* ec,
477 : std::size_t* bytes_out)
478 : {
479 229162 : auto& op = rd_;
480 229162 : op.reset();
481 :
482 : // Closed-object contract: complete with bad_file_descriptor without
483 : // touching the kernel or the unregistered descriptor state.
484 229162 : if (this->fd_ < 0)
485 : {
486 8 : op.h = h;
487 8 : op.ex = ex;
488 8 : op.ec_out = ec;
489 8 : op.bytes_out = bytes_out;
490 8 : op.start(token, static_cast<Derived*>(this));
491 8 : op.object_ref_ = detail::object_ref(this);
492 8 : op.complete(EBADF, 0);
493 8 : this->svc_.post(&op);
494 8 : return std::noop_coroutine();
495 : }
496 :
497 229154 : capy::mutable_buffer bufs[ReadOp::max_buffers];
498 229154 : op.iovec_count = static_cast<int>(param.copy_to(bufs, ReadOp::max_buffers));
499 :
500 229154 : if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
501 : {
502 4 : op.empty_buffer_read = true;
503 4 : op.h = h;
504 4 : op.ex = ex;
505 4 : op.ec_out = ec;
506 4 : op.bytes_out = bytes_out;
507 4 : op.start(token, static_cast<Derived*>(this));
508 4 : op.object_ref_ = detail::object_ref(this);
509 4 : op.complete(0, 0);
510 4 : this->svc_.post(&op);
511 4 : return std::noop_coroutine();
512 : }
513 :
514 458318 : for (int i = 0; i < op.iovec_count; ++i)
515 : {
516 229168 : op.iovecs[i].iov_base = bufs[i].data();
517 229168 : op.iovecs[i].iov_len = bufs[i].size();
518 : }
519 :
520 : // Speculative read; for the single-buffer case use recv() so the
521 : // kernel skips the readv iov_iter setup.
522 : ssize_t n;
523 229150 : if (op.iovec_count == 1)
524 : {
525 : do
526 : {
527 229138 : n = ::recv(this->fd_, bufs[0].data(), bufs[0].size(), 0);
528 : }
529 229138 : while (n < 0 && errno == EINTR);
530 : }
531 : else
532 : {
533 : do
534 : {
535 16 : n = ::readv(this->fd_, op.iovecs, op.iovec_count);
536 : }
537 16 : while (n < 0 && errno == EINTR);
538 : }
539 :
540 229150 : if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
541 : {
542 228347 : int err = (n < 0) ? errno : 0;
543 228347 : auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
544 :
545 228347 : if (this->svc_.scheduler().try_consume_inline_budget())
546 : {
547 182684 : if (err)
548 4 : *ec = make_err(err);
549 182680 : else if (n == 0)
550 15 : *ec = capy::error::eof;
551 : else
552 182665 : *ec = {};
553 182684 : *bytes_out = bytes;
554 182684 : op.cont.h = h;
555 182684 : return dispatch_coro(ex, op.cont);
556 : }
557 45663 : op.h = h;
558 45663 : op.ex = ex;
559 45663 : op.ec_out = ec;
560 45663 : op.bytes_out = bytes_out;
561 45663 : op.start(token, static_cast<Derived*>(this));
562 45663 : op.object_ref_ = detail::object_ref(this);
563 45663 : op.complete(err, bytes);
564 45663 : this->svc_.post(&op);
565 45663 : return std::noop_coroutine();
566 : }
567 :
568 : // EAGAIN — register with reactor
569 803 : op.h = h;
570 803 : op.ex = ex;
571 803 : op.ec_out = ec;
572 803 : op.bytes_out = bytes_out;
573 803 : op.fd = this->fd_;
574 803 : op.start(token, static_cast<Derived*>(this));
575 803 : op.object_ref_ = detail::object_ref(this);
576 :
577 803 : this->register_op(
578 803 : op, this->desc_state_.read_op, this->desc_state_.read_ready);
579 803 : return std::noop_coroutine();
580 : }
581 :
582 : template<
583 : class Derived,
584 : class Service,
585 : class ConnOp,
586 : class ReadOp,
587 : class WriteOp,
588 : class WaitOp,
589 : class DescState,
590 : class ImplBase,
591 : class Endpoint>
592 : std::coroutine_handle<>
593 228472 : reactor_stream_socket<
594 : Derived,
595 : Service,
596 : ConnOp,
597 : ReadOp,
598 : WriteOp,
599 : WaitOp,
600 : DescState,
601 : ImplBase,
602 : Endpoint>::
603 : do_write_some(
604 : std::coroutine_handle<> h,
605 : capy::executor_ref ex,
606 : buffer_param param,
607 : std::stop_token const& token,
608 : std::error_code* ec,
609 : std::size_t* bytes_out)
610 : {
611 228472 : auto& op = wr_;
612 228472 : op.reset();
613 :
614 : // Closed-object contract: complete with bad_file_descriptor without
615 : // touching the kernel or the unregistered descriptor state.
616 228472 : if (this->fd_ < 0)
617 : {
618 10 : op.h = h;
619 10 : op.ex = ex;
620 10 : op.ec_out = ec;
621 10 : op.bytes_out = bytes_out;
622 10 : op.start(token, static_cast<Derived*>(this));
623 10 : op.object_ref_ = detail::object_ref(this);
624 10 : op.complete(EBADF, 0);
625 10 : this->svc_.post(&op);
626 10 : return std::noop_coroutine();
627 : }
628 :
629 228462 : capy::mutable_buffer bufs[WriteOp::max_buffers];
630 228462 : op.iovec_count =
631 228462 : static_cast<int>(param.copy_to(bufs, WriteOp::max_buffers));
632 :
633 228462 : if (op.iovec_count == 0 || (op.iovec_count == 1 && bufs[0].size() == 0))
634 : {
635 4 : op.h = h;
636 4 : op.ex = ex;
637 4 : op.ec_out = ec;
638 4 : op.bytes_out = bytes_out;
639 4 : op.start(token, static_cast<Derived*>(this));
640 4 : op.object_ref_ = detail::object_ref(this);
641 4 : op.complete(0, 0);
642 4 : this->svc_.post(&op);
643 4 : return std::noop_coroutine();
644 : }
645 :
646 456932 : for (int i = 0; i < op.iovec_count; ++i)
647 : {
648 228474 : op.iovecs[i].iov_base = bufs[i].data();
649 228474 : op.iovecs[i].iov_len = bufs[i].size();
650 : }
651 :
652 : // Speculative write; the single-buffer case dispatches to a
653 : // backend-specific fast path so the kernel skips msghdr/iov_iter
654 : // setup (and so each backend can pick the right SIGPIPE strategy).
655 : ssize_t n;
656 228458 : if (op.iovec_count == 1)
657 : {
658 456892 : n = WriteOp::write_policy::write_one(
659 228446 : this->fd_, bufs[0].data(), bufs[0].size());
660 : }
661 : else
662 : {
663 12 : n = WriteOp::write_policy::write(this->fd_, op.iovecs, op.iovec_count);
664 : }
665 :
666 228458 : if (n >= 0 || (errno != EAGAIN && errno != EWOULDBLOCK))
667 : {
668 228318 : int err = (n < 0) ? errno : 0;
669 228318 : auto bytes = (n > 0) ? static_cast<std::size_t>(n) : std::size_t(0);
670 :
671 228318 : if (this->svc_.scheduler().try_consume_inline_budget())
672 : {
673 182614 : *ec = err ? make_err(err) : std::error_code{};
674 182614 : *bytes_out = bytes;
675 182614 : op.cont.h = h;
676 182614 : return dispatch_coro(ex, op.cont);
677 : }
678 45704 : op.h = h;
679 45704 : op.ex = ex;
680 45704 : op.ec_out = ec;
681 45704 : op.bytes_out = bytes_out;
682 45704 : op.start(token, static_cast<Derived*>(this));
683 45704 : op.object_ref_ = detail::object_ref(this);
684 45704 : op.complete(err, bytes);
685 45704 : this->svc_.post(&op);
686 45704 : return std::noop_coroutine();
687 : }
688 :
689 : // EAGAIN — register with reactor
690 140 : op.h = h;
691 140 : op.ex = ex;
692 140 : op.ec_out = ec;
693 140 : op.bytes_out = bytes_out;
694 140 : op.fd = this->fd_;
695 140 : op.start(token, static_cast<Derived*>(this));
696 140 : op.object_ref_ = detail::object_ref(this);
697 :
698 140 : this->register_op(
699 140 : op, this->desc_state_.write_op, this->desc_state_.write_ready, true);
700 140 : return std::noop_coroutine();
701 : }
702 :
703 : template<
704 : class Derived,
705 : class Service,
706 : class ConnOp,
707 : class ReadOp,
708 : class WriteOp,
709 : class WaitOp,
710 : class DescState,
711 : class ImplBase,
712 : class Endpoint>
713 : std::coroutine_handle<>
714 90 : reactor_stream_socket<
715 : Derived,
716 : Service,
717 : ConnOp,
718 : ReadOp,
719 : WriteOp,
720 : WaitOp,
721 : DescState,
722 : ImplBase,
723 : Endpoint>::
724 : do_wait(
725 : std::coroutine_handle<> h,
726 : capy::executor_ref ex,
727 : wait_type w,
728 : std::stop_token const& token,
729 : std::error_code* ec)
730 : {
731 : // Pick refs up-front to avoid duplicating the register_op call.
732 : WaitOp* op_ptr;
733 : reactor_op_base** desc_slot_ptr;
734 : std::uint32_t event;
735 :
736 90 : if (w == wait_type::read)
737 : {
738 51 : op_ptr = &wait_rd_;
739 51 : desc_slot_ptr = &this->desc_state_.wait_read_op;
740 51 : event = reactor_event_read;
741 : }
742 39 : else if (w == wait_type::write)
743 : {
744 21 : op_ptr = &wait_wr_;
745 21 : desc_slot_ptr = &this->desc_state_.wait_write_op;
746 21 : event = reactor_event_write;
747 : }
748 : else // wait_type::error
749 : {
750 18 : op_ptr = &wait_er_;
751 18 : desc_slot_ptr = &this->desc_state_.wait_error_op;
752 18 : event = reactor_event_error;
753 : }
754 :
755 90 : auto& op = *op_ptr;
756 :
757 : // Speculative probe, mirroring the speculative read: an
758 : // edge-triggered reactor cannot report a condition that already
759 : // holds, so a wait initiated on an already-ready socket would
760 : // otherwise park forever.
761 90 : int perr = 0;
762 90 : if (WaitOp::probe(this->fd_, event, perr))
763 : {
764 36 : if (this->svc_.scheduler().try_consume_inline_budget())
765 : {
766 8 : *ec = perr ? make_err(perr) : std::error_code{};
767 8 : op.cont.h = h;
768 8 : return dispatch_coro(ex, op.cont);
769 : }
770 28 : op.reset();
771 28 : op.wait_event = event;
772 28 : op.h = h;
773 28 : op.ex = ex;
774 28 : op.ec_out = ec;
775 28 : op.fd = this->fd_;
776 28 : op.start(token, static_cast<Derived*>(this));
777 28 : op.object_ref_ = detail::object_ref(this);
778 28 : op.complete(perr, 0);
779 28 : this->svc_.post(&op);
780 28 : return std::noop_coroutine();
781 : }
782 :
783 54 : op.reset();
784 54 : op.wait_event = event;
785 54 : op.h = h;
786 54 : op.ex = ex;
787 54 : op.ec_out = ec;
788 54 : op.fd = this->fd_;
789 54 : op.start(token, static_cast<Derived*>(this));
790 54 : op.object_ref_ = detail::object_ref(this);
791 :
792 : // Force register_op's ready path so the wait op re-probes under
793 : // the descriptor mutex before parking. An edge consumed between
794 : // the speculative probe above and the park (a concurrent short
795 : // read, or an error event dispatched to an empty slot) would
796 : // otherwise leave the wait parked on a ready socket.
797 54 : bool force_probe = true;
798 54 : this->register_op(
799 : op, *desc_slot_ptr, force_probe, event == reactor_event_write);
800 54 : return std::noop_coroutine();
801 : }
802 :
803 : } // namespace boost::corosio::detail
804 :
805 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_STREAM_SOCKET_HPP
|