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_BASIC_SOCKET_HPP
11 : #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_BASIC_SOCKET_HPP
12 :
13 : #include <boost/corosio/detail/config.hpp>
14 : #include <boost/corosio/detail/intrusive.hpp>
15 : #include <boost/corosio/detail/native_handle.hpp>
16 : #include <boost/corosio/endpoint.hpp>
17 : #include <boost/corosio/native/detail/native_socket_base.hpp>
18 : #include <boost/corosio/native/detail/reactor/reactor_op_base.hpp>
19 : #include <boost/corosio/native/detail/make_err.hpp>
20 : #include <boost/corosio/native/detail/endpoint_convert.hpp>
21 :
22 : #include <cstring>
23 : #include <memory>
24 : #include <mutex>
25 : #include <utility>
26 :
27 : #include <errno.h>
28 : #include <netinet/in.h>
29 : #include <sys/socket.h>
30 : #include <unistd.h>
31 :
32 : namespace boost::corosio::detail {
33 :
34 : /** CRTP base for reactor-backed socket implementations.
35 :
36 : Extracts the shared data members, virtual overrides, and
37 : cancel/close/register logic that is identical across TCP
38 : (reactor_stream_socket) and UDP (reactor_datagram_socket).
39 :
40 : Derived classes provide CRTP callbacks that enumerate their
41 : specific op slots so cancel/close can iterate them generically.
42 :
43 : @tparam Derived The concrete socket type (CRTP).
44 : @tparam ImplBase The public vtable base (tcp_socket::implementation
45 : or udp_socket::implementation).
46 : @tparam Service The backend's service type.
47 : @tparam DescState The backend's descriptor_state type.
48 : @tparam Endpoint The endpoint type (endpoint or local_endpoint).
49 : */
50 : template<
51 : class Derived,
52 : class ImplBase,
53 : class Service,
54 : class DescState,
55 : class Endpoint = endpoint>
56 : class reactor_basic_socket
57 : : public native_socket_base<Derived, ImplBase, Endpoint>
58 : , public intrusive_list<Derived>::node
59 : {
60 : friend Derived;
61 :
62 : template<class, class, class, class, class, class, class, class, class>
63 : friend class reactor_stream_socket;
64 :
65 : template<
66 : class,
67 : class,
68 : class,
69 : class,
70 : class,
71 : class,
72 : class,
73 : class,
74 : class,
75 : class,
76 : class>
77 : friend class reactor_datagram_socket;
78 :
79 HIT 2765 : explicit reactor_basic_socket(Service& svc) noexcept : svc_(svc) {}
80 :
81 : protected:
82 : // fd_ / local_endpoint_ and the synchronous accessors (native_handle,
83 : // is_open, set_option/get_option, set_socket/set_local_endpoint, do_bind)
84 : // live in native_socket_base — the readiness/completion-agnostic base
85 : // shared with io_uring's sockets. The using-declarations make the
86 : // inherited members visible to this template's own unqualified
87 : // references below (two-phase lookup).
88 : using native_socket_base<Derived, ImplBase, Endpoint>::fd_;
89 : using native_socket_base<Derived, ImplBase, Endpoint>::local_endpoint_;
90 :
91 : Service& svc_;
92 :
93 : public:
94 : /// Per-descriptor state for persistent reactor registration.
95 : DescState desc_state_;
96 :
97 2765 : ~reactor_basic_socket() override = default;
98 :
99 : /** Assign the fd, initialize descriptor state, and register with
100 : the reactor.
101 :
102 : @param fd The descriptor to adopt.
103 :
104 : @return The error if the reactor rejects the descriptor, in
105 : which case the implementation is left closed and the caller
106 : retains ownership of @a fd; otherwise a default constructed
107 : error code.
108 : */
109 5753 : std::error_code init_and_register(int fd) noexcept
110 : {
111 5753 : fd_ = fd;
112 5753 : desc_state_.fd = fd;
113 : {
114 5753 : std::lock_guard lock(desc_state_.mutex);
115 5753 : desc_state_.read_op = nullptr;
116 5753 : desc_state_.write_op = nullptr;
117 5753 : desc_state_.connect_op = nullptr;
118 5753 : }
119 5753 : if (auto ec = svc_.scheduler().register_descriptor(fd, &desc_state_))
120 : {
121 : // Undo the partial state so a failed adopt is
122 : // indistinguishable from a closed implementation.
123 3 : fd_ = -1;
124 3 : desc_state_.fd = -1;
125 3 : desc_state_.registered_events = 0;
126 3 : return ec;
127 : }
128 5750 : return {};
129 : }
130 :
131 : /** Register an op with the reactor.
132 :
133 : Handles cached edge events. Called on the EAGAIN/EINPROGRESS
134 : path when speculative I/O failed.
135 : */
136 : template<class Op>
137 : void register_op(
138 : Op& op,
139 : reactor_op_base*& desc_slot,
140 : bool& ready_flag,
141 : bool is_write_direction = false) noexcept;
142 :
143 : /** Cancel a single pending operation.
144 :
145 : Claims the operation from its descriptor_state slot under
146 : the mutex and posts it to the scheduler as cancelled.
147 : Derived must implement:
148 : op_to_desc_slot(Op&) -> reactor_op_base**
149 : */
150 : template<class Op>
151 : void cancel_single_op(Op& op) noexcept;
152 :
153 : /** Cancel all pending operations.
154 :
155 : Invoked by the derived class's cancel() override.
156 : Derived must implement:
157 : for_each_op(auto fn)
158 : for_each_desc_entry(auto fn)
159 : */
160 : void do_cancel() noexcept;
161 :
162 : /** Close the socket and cancel pending operations.
163 :
164 : Invoked by the derived class's close_socket(). The
165 : derived class may add backend-specific cleanup after
166 : calling this method.
167 : Derived must implement:
168 : for_each_op(auto fn)
169 : for_each_desc_entry(auto fn)
170 : */
171 : void do_close_socket() noexcept;
172 :
173 : /** Release the socket without closing the fd.
174 :
175 : Like do_close_socket() but does not call ::close().
176 : Returns the fd so the caller can take ownership.
177 : */
178 : native_handle_type do_release_socket() noexcept;
179 :
180 : /** Reset descriptor state for recycling.
181 :
182 : Called by the owning service's `construct()` when popping this
183 : impl back off the free list. Actively re-initializes every
184 : field it owns rather than trusting `close_socket()`'s prior
185 : writes to still be there — `poison()` below depends on that,
186 : and a future field added to either side without the other
187 : would otherwise go unnoticed until this one re-sets it.
188 : `object_ref_` and `is_enqueued_` are asserted instead: `poison()`
189 : never touches them, so they must already hold their
190 : zero-action-complete state.
191 :
192 : @pre refs_ == 0, fd closed and deregistered, no op in flight.
193 : */
194 13287 : void reuse() noexcept
195 : {
196 13287 : fd_ = -1;
197 13287 : local_endpoint_ = Endpoint{};
198 13287 : desc_state_.fd = -1;
199 13287 : desc_state_.registered_events = 0;
200 13287 : desc_state_.read_op = nullptr;
201 13287 : desc_state_.write_op = nullptr;
202 13287 : desc_state_.connect_op = nullptr;
203 13287 : desc_state_.wait_read_op = nullptr;
204 13287 : desc_state_.wait_write_op = nullptr;
205 13287 : desc_state_.wait_error_op = nullptr;
206 13287 : desc_state_.read_ready = false;
207 13287 : desc_state_.write_ready = false;
208 13287 : BOOST_COROSIO_ASSERT(!desc_state_.object_ref_);
209 13287 : BOOST_COROSIO_ASSERT(!desc_state_.is_enqueued_.load(
210 : std::memory_order_relaxed));
211 13287 : }
212 :
213 : #if !defined(NDEBUG)
214 : /** Poison the fields `reuse()` re-initializes, before the pool
215 : parks this impl on the free list.
216 :
217 : A field `reuse()` forgets to reset then reads back as this
218 : unmistakable byte pattern instead of silently matching
219 : whatever `close_socket()` happened to leave, turning a missed
220 : reset into an immediate, loud failure rather than a latent
221 : stale-state bug. Does not touch the mutex, `ready_events_` /
222 : `is_enqueued_` (atomics), `scheduler_`, `object_ref_`, the
223 : `intrusive_list` hook, or `refs_`: `recycle()` still reads
224 : some of those after this call returns, and smashing the rest
225 : would be undefined rather than diagnostic.
226 :
227 : @pre refs_ == 0, fd closed and deregistered, no op in flight.
228 : */
229 16047 : void poison() noexcept
230 : {
231 192564 : auto smash = [](auto& field) {
232 192564 : std::memset(static_cast<void*>(&field), 0xDB, sizeof(field));
233 : };
234 16047 : smash(fd_);
235 16047 : smash(local_endpoint_);
236 16047 : smash(desc_state_.fd);
237 16047 : smash(desc_state_.registered_events);
238 16047 : smash(desc_state_.read_op);
239 16047 : smash(desc_state_.write_op);
240 16047 : smash(desc_state_.connect_op);
241 16047 : smash(desc_state_.wait_read_op);
242 16047 : smash(desc_state_.wait_write_op);
243 16047 : smash(desc_state_.wait_error_op);
244 16047 : smash(desc_state_.read_ready);
245 16047 : smash(desc_state_.write_ready);
246 16047 : }
247 : #endif
248 : };
249 :
250 : template<
251 : class Derived,
252 : class ImplBase,
253 : class Service,
254 : class DescState,
255 : class Endpoint>
256 : template<class Op>
257 : void
258 6140 : reactor_basic_socket<Derived, ImplBase, Service, DescState, Endpoint>::
259 : register_op(
260 : Op& op,
261 : reactor_op_base*& desc_slot,
262 : bool& ready_flag,
263 : bool is_write_direction) noexcept
264 : {
265 6140 : svc_.work_started();
266 :
267 6140 : std::lock_guard lock(desc_state_.mutex);
268 6140 : bool io_done = false;
269 6140 : if (ready_flag)
270 : {
271 285 : ready_flag = false;
272 285 : op.perform_io();
273 285 : io_done = (op.errn != EAGAIN && op.errn != EWOULDBLOCK);
274 285 : if (!io_done)
275 279 : op.errn = 0;
276 : }
277 :
278 6140 : if (io_done || op.cancelled.load(std::memory_order_acquire))
279 : {
280 6 : svc_.post(&op);
281 6 : svc_.work_finished();
282 : }
283 : else
284 : {
285 6134 : desc_slot = &op;
286 :
287 : // Select must rebuild its fd_sets when a write-direction op
288 : // is parked, so select() watches for writability. Compiled
289 : // away to nothing for epoll and kqueue.
290 : if constexpr (requires { Service::needs_write_notification; })
291 : {
292 : if constexpr (Service::needs_write_notification)
293 : {
294 2817 : if (is_write_direction)
295 2282 : svc_.scheduler().notify_reactor();
296 : }
297 : }
298 : }
299 6140 : }
300 :
301 : template<
302 : class Derived,
303 : class ImplBase,
304 : class Service,
305 : class DescState,
306 : class Endpoint>
307 : template<class Op>
308 : void
309 239 : reactor_basic_socket<Derived, ImplBase, Service, DescState, Endpoint>::
310 : cancel_single_op(Op& op) noexcept
311 : {
312 239 : op.request_cancel();
313 :
314 239 : auto* d = static_cast<Derived*>(this);
315 239 : reactor_op_base** desc_op_ptr = d->op_to_desc_slot(op);
316 :
317 239 : if (desc_op_ptr)
318 : {
319 239 : reactor_op_base* claimed = nullptr;
320 : {
321 239 : std::lock_guard lock(desc_state_.mutex);
322 239 : if (*desc_op_ptr == &op)
323 229 : claimed = std::exchange(*desc_op_ptr, nullptr);
324 : // Not in the slot: request_cancel() above already set
325 : // op.cancelled, which register_op consults before parking
326 : // and the completion decode consults on delivery. Latching
327 : // a descriptor flag here instead would outlive this op and
328 : // cancel the next wait in the same direction.
329 239 : }
330 239 : if (claimed)
331 : {
332 229 : op.object_ref_ = detail::object_ref(this);
333 229 : svc_.post(&op);
334 229 : svc_.work_finished();
335 : }
336 : }
337 239 : }
338 :
339 : template<
340 : class Derived,
341 : class ImplBase,
342 : class Service,
343 : class DescState,
344 : class Endpoint>
345 : void
346 283 : reactor_basic_socket<Derived, ImplBase, Service, DescState, Endpoint>::
347 : do_cancel() noexcept
348 : {
349 283 : auto* d = static_cast<Derived*>(this);
350 :
351 2107 : d->for_each_op([](auto& op) { op.request_cancel(); });
352 :
353 : // Claim ops under a single lock acquisition
354 : struct claimed_entry
355 : {
356 : reactor_op_base* op = nullptr;
357 : reactor_op_base* base = nullptr;
358 : };
359 : // Max 8 ops: conn, rd, wr, wait_rd, wait_wr, wait_er, recv_rd, send_wr
360 283 : claimed_entry claimed[8];
361 283 : int count = 0;
362 :
363 : {
364 283 : std::lock_guard lock(desc_state_.mutex);
365 3931 : d->for_each_desc_entry([&](auto& op, reactor_op_base*& desc_slot) {
366 1824 : if (desc_slot == &op)
367 : {
368 196 : claimed[count].op = std::exchange(desc_slot, nullptr);
369 196 : claimed[count].base = &op;
370 196 : ++count;
371 : }
372 : });
373 283 : }
374 :
375 479 : for (int i = 0; i < count; ++i)
376 : {
377 196 : claimed[i].base->object_ref_ = detail::object_ref(this);
378 196 : svc_.post(claimed[i].base);
379 196 : svc_.work_finished();
380 : }
381 283 : }
382 :
383 : template<
384 : class Derived,
385 : class ImplBase,
386 : class Service,
387 : class DescState,
388 : class Endpoint>
389 : void
390 48764 : reactor_basic_socket<Derived, ImplBase, Service, DescState, Endpoint>::
391 : do_close_socket() noexcept
392 : {
393 48764 : auto* d = static_cast<Derived*>(this);
394 :
395 346960 : d->for_each_op([](auto& op) { op.request_cancel(); });
396 :
397 : struct claimed_entry
398 : {
399 : reactor_op_base* base = nullptr;
400 : };
401 48764 : claimed_entry claimed[8];
402 48764 : int count = 0;
403 :
404 : {
405 48764 : std::lock_guard lock(desc_state_.mutex);
406 48764 : d->for_each_desc_entry(
407 596392 : [&](auto& /*op*/, reactor_op_base*& desc_slot) {
408 298196 : auto* c = std::exchange(desc_slot, nullptr);
409 298196 : if (c)
410 : {
411 58 : claimed[count].base = c;
412 58 : ++count;
413 : }
414 : });
415 48764 : desc_state_.read_ready = false;
416 48764 : desc_state_.write_ready = false;
417 :
418 48764 : if (desc_state_.is_enqueued_.load(std::memory_order_acquire))
419 1333 : desc_state_.object_ref_ = detail::object_ref(this);
420 48764 : }
421 :
422 48822 : for (int i = 0; i < count; ++i)
423 : {
424 58 : claimed[i].base->object_ref_ = detail::object_ref(this);
425 58 : svc_.post(claimed[i].base);
426 58 : svc_.work_finished();
427 : }
428 :
429 48764 : if (fd_ >= 0)
430 : {
431 10602 : if (desc_state_.registered_events != 0)
432 10600 : svc_.scheduler().deregister_descriptor(fd_);
433 10602 : ::close(fd_);
434 10602 : fd_ = -1;
435 : }
436 :
437 48764 : desc_state_.fd = -1;
438 48764 : desc_state_.registered_events = 0;
439 :
440 48764 : local_endpoint_ = Endpoint{};
441 48764 : }
442 :
443 : template<
444 : class Derived,
445 : class ImplBase,
446 : class Service,
447 : class DescState,
448 : class Endpoint>
449 : native_handle_type
450 16 : reactor_basic_socket<Derived, ImplBase, Service, DescState, Endpoint>::
451 : do_release_socket() noexcept
452 : {
453 : // Cancel pending ops (same as do_close_socket)
454 16 : auto* d = static_cast<Derived*>(this);
455 :
456 128 : d->for_each_op([](auto& op) { op.request_cancel(); });
457 :
458 : struct claimed_entry
459 : {
460 : reactor_op_base* base = nullptr;
461 : };
462 16 : claimed_entry claimed[8];
463 16 : int count = 0;
464 :
465 : {
466 16 : std::lock_guard lock(desc_state_.mutex);
467 16 : d->for_each_desc_entry(
468 224 : [&](auto& /*op*/, reactor_op_base*& desc_slot) {
469 112 : auto* c = std::exchange(desc_slot, nullptr);
470 112 : if (c)
471 : {
472 12 : claimed[count].base = c;
473 12 : ++count;
474 : }
475 : });
476 16 : desc_state_.read_ready = false;
477 16 : desc_state_.write_ready = false;
478 :
479 16 : if (desc_state_.is_enqueued_.load(std::memory_order_acquire))
480 3 : desc_state_.object_ref_ = detail::object_ref(this);
481 16 : }
482 :
483 28 : for (int i = 0; i < count; ++i)
484 : {
485 12 : claimed[i].base->object_ref_ = detail::object_ref(this);
486 12 : svc_.post(claimed[i].base);
487 12 : svc_.work_finished();
488 : }
489 :
490 16 : native_handle_type released = fd_;
491 :
492 16 : if (fd_ >= 0)
493 : {
494 16 : if (desc_state_.registered_events != 0)
495 16 : svc_.scheduler().deregister_descriptor(fd_);
496 : // Do NOT close -- caller takes ownership
497 16 : fd_ = -1;
498 : }
499 :
500 16 : desc_state_.fd = -1;
501 16 : desc_state_.registered_events = 0;
502 :
503 16 : local_endpoint_ = Endpoint{};
504 :
505 16 : return released;
506 : }
507 :
508 : } // namespace boost::corosio::detail
509 :
510 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_BASIC_SOCKET_HPP
|