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_POSIX_POSIX_RESOLVER_SERVICE_HPP
12 : #define BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RESOLVER_SERVICE_HPP
13 :
14 : #include <boost/corosio/detail/platform.hpp>
15 :
16 : #if BOOST_COROSIO_POSIX
17 :
18 : #include <boost/corosio/native/detail/posix/posix_resolver.hpp>
19 : #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp>
20 : #include <boost/corosio/detail/thread_pool.hpp>
21 :
22 : #include <unordered_map>
23 :
24 : namespace boost::corosio::detail {
25 :
26 : /** Resolver service for POSIX backends.
27 :
28 : Owns all posix_resolver instances. Thread lifecycle is managed
29 : by the thread_pool service.
30 : */
31 : class BOOST_COROSIO_DECL posix_resolver_service final
32 : : public capy::execution_context::service
33 : , public io_object::io_service
34 : {
35 : public:
36 : using key_type = posix_resolver_service;
37 :
38 HIT 65 : explicit posix_resolver_service(capy::execution_context& ctx)
39 195 : : sched_(&get_scheduler(ctx))
40 65 : , pool_(ctx)
41 : {
42 65 : }
43 :
44 130 : ~posix_resolver_service() override = default;
45 :
46 : posix_resolver_service(posix_resolver_service const&) = delete;
47 : posix_resolver_service& operator=(posix_resolver_service const&) = delete;
48 :
49 : io_object::implementation* construct() override;
50 :
51 65 : void destroy(io_object::implementation* p) override
52 : {
53 65 : auto& impl = static_cast<posix_resolver&>(*p);
54 65 : impl.cancel();
55 65 : destroy_impl(impl);
56 65 : }
57 :
58 : void shutdown() override;
59 : void destroy_impl(posix_resolver& impl);
60 :
61 : void post(scheduler_op* op);
62 :
63 : /** Return the resolver thread pool.
64 :
65 : The pool's service is created on first use, so this can fail
66 : where a plain accessor could not. Its workers start later, on
67 : the first post, and a thread the system refuses there is
68 : reported by that post rather than thrown here.
69 :
70 : @throws std::bad_alloc If the service cannot be allocated.
71 :
72 : @return The context's shared blocking-I/O pool.
73 :
74 : @see thread_pool_ref::get
75 : */
76 52 : thread_pool& pool()
77 : {
78 52 : return pool_.get();
79 : }
80 :
81 : /// True when the resolver thread pool is unavailable: the `unsafe` tier,
82 : /// whose lockless scheduler cannot accept the pool's cross-thread
83 : /// completions.
84 54 : bool resolver_unavailable() const noexcept
85 : {
86 54 : return sched_->scheduler_locking_disabled();
87 : }
88 :
89 : private:
90 : scheduler* sched_;
91 : thread_pool_ref pool_;
92 : std::mutex mutex_;
93 : intrusive_list<posix_resolver> resolver_list_;
94 : std::unordered_map<posix_resolver*, std::shared_ptr<posix_resolver>>
95 : resolver_ptrs_;
96 : };
97 :
98 : // ---------------------------------------------------------------------------
99 : // Inline implementation
100 : // ---------------------------------------------------------------------------
101 :
102 : // posix_resolver_detail helpers
103 :
104 : inline int
105 33 : posix_resolver_detail::flags_to_hints(resolve_flags flags)
106 : {
107 33 : int hints = 0;
108 :
109 33 : if ((flags & resolve_flags::passive) != resolve_flags::none)
110 1 : hints |= AI_PASSIVE;
111 33 : if ((flags & resolve_flags::numeric_host) != resolve_flags::none)
112 18 : hints |= AI_NUMERICHOST;
113 33 : if ((flags & resolve_flags::numeric_service) != resolve_flags::none)
114 12 : hints |= AI_NUMERICSERV;
115 33 : if ((flags & resolve_flags::address_configured) != resolve_flags::none)
116 1 : hints |= AI_ADDRCONFIG;
117 33 : if ((flags & resolve_flags::v4_mapped) != resolve_flags::none)
118 1 : hints |= AI_V4MAPPED;
119 33 : if ((flags & resolve_flags::all_matching) != resolve_flags::none)
120 1 : hints |= AI_ALL;
121 :
122 33 : return hints;
123 : }
124 :
125 : inline int
126 17 : posix_resolver_detail::flags_to_ni_flags(reverse_flags flags)
127 : {
128 17 : int ni_flags = 0;
129 :
130 17 : if ((flags & reverse_flags::numeric_host) != reverse_flags::none)
131 7 : ni_flags |= NI_NUMERICHOST;
132 17 : if ((flags & reverse_flags::numeric_service) != reverse_flags::none)
133 7 : ni_flags |= NI_NUMERICSERV;
134 17 : if ((flags & reverse_flags::name_required) != reverse_flags::none)
135 1 : ni_flags |= NI_NAMEREQD;
136 17 : if ((flags & reverse_flags::datagram_service) != reverse_flags::none)
137 1 : ni_flags |= NI_DGRAM;
138 :
139 17 : return ni_flags;
140 : }
141 :
142 : inline std::vector<endpoint>
143 21 : posix_resolver_detail::convert_results(struct addrinfo* ai)
144 : {
145 21 : std::vector<endpoint> endpoints;
146 21 : endpoints.reserve(4); // Most lookups return 1-4 addresses
147 :
148 42 : for (auto* p = ai; p != nullptr; p = p->ai_next)
149 : {
150 21 : if (p->ai_family == AF_INET)
151 : {
152 18 : auto* addr = reinterpret_cast<sockaddr_in*>(p->ai_addr);
153 18 : endpoints.push_back(from_sockaddr_in(*addr));
154 : }
155 3 : else if (p->ai_family == AF_INET6)
156 : {
157 3 : auto* addr = reinterpret_cast<sockaddr_in6*>(p->ai_addr);
158 3 : endpoints.push_back(from_sockaddr_in6(*addr));
159 : }
160 : }
161 :
162 21 : return endpoints;
163 MIS 0 : }
164 :
165 : inline std::error_code
166 HIT 26 : posix_resolver_detail::make_gai_error(int gai_err)
167 : {
168 : // Map GAI errors to appropriate generic error codes
169 26 : switch (gai_err)
170 : {
171 1 : case EAI_AGAIN:
172 : // Temporary failure - try again later
173 1 : return std::error_code(
174 : static_cast<int>(std::errc::resource_unavailable_try_again),
175 1 : std::generic_category());
176 :
177 1 : case EAI_BADFLAGS:
178 : // Invalid flags
179 1 : return std::error_code(
180 : static_cast<int>(std::errc::invalid_argument),
181 1 : std::generic_category());
182 :
183 11 : case EAI_FAIL:
184 : // Non-recoverable failure
185 11 : return std::error_code(
186 11 : static_cast<int>(std::errc::io_error), std::generic_category());
187 :
188 1 : case EAI_FAMILY:
189 : // Address family not supported
190 1 : return std::error_code(
191 : static_cast<int>(std::errc::address_family_not_supported),
192 1 : std::generic_category());
193 :
194 1 : case EAI_MEMORY:
195 : // Memory allocation failure
196 1 : return std::error_code(
197 : static_cast<int>(std::errc::not_enough_memory),
198 1 : std::generic_category());
199 :
200 7 : case EAI_NONAME:
201 : // Host or service not found
202 7 : return std::error_code(
203 : static_cast<int>(std::errc::no_such_device_or_address),
204 7 : std::generic_category());
205 :
206 1 : case EAI_SERVICE:
207 : // Service not supported for socket type
208 1 : return std::error_code(
209 : static_cast<int>(std::errc::invalid_argument),
210 1 : std::generic_category());
211 :
212 1 : case EAI_SOCKTYPE:
213 : // Socket type not supported
214 1 : return std::error_code(
215 : static_cast<int>(std::errc::not_supported),
216 1 : std::generic_category());
217 :
218 1 : case EAI_SYSTEM:
219 : // System error - use errno
220 1 : return std::error_code(errno, std::generic_category());
221 :
222 1 : default:
223 : // Unknown error
224 1 : return std::error_code(
225 1 : static_cast<int>(std::errc::io_error), std::generic_category());
226 : }
227 : }
228 :
229 : // posix_resolver
230 :
231 66 : inline posix_resolver::posix_resolver(posix_resolver_service& svc) noexcept
232 66 : : svc_(svc)
233 : {
234 66 : }
235 :
236 : // posix_resolver::resolve_op implementation
237 :
238 : inline void
239 34 : posix_resolver::resolve_op::reset() noexcept
240 : {
241 34 : host.clear();
242 34 : service.clear();
243 34 : flags = resolve_flags::none;
244 34 : stored_results = std::vector<endpoint>{};
245 34 : gai_error = 0;
246 34 : cancelled.store(false, std::memory_order_relaxed);
247 34 : stop_cb.reset();
248 34 : ec_out = nullptr;
249 34 : out = nullptr;
250 34 : }
251 :
252 : inline void
253 32 : posix_resolver::resolve_op::operator()()
254 : {
255 32 : stop_cb.reset(); // Disconnect stop callback
256 :
257 32 : bool const was_cancelled = cancelled.load(std::memory_order_acquire);
258 :
259 32 : if (ec_out)
260 : {
261 32 : if (was_cancelled)
262 MIS 0 : *ec_out = capy::error::canceled;
263 HIT 32 : else if (gai_error != 0)
264 11 : *ec_out = posix_resolver_detail::make_gai_error(gai_error);
265 : else
266 21 : *ec_out = {}; // Clear on success
267 : }
268 :
269 32 : if (out && !was_cancelled && gai_error == 0)
270 21 : *out = std::move(stored_results);
271 :
272 : // Hold the keepalive across the dispatch: it may be the last
273 : // reference to the implementation this op is embedded in.
274 32 : auto prevent_destroy = std::move(impl_ptr);
275 32 : ex.on_work_finished();
276 32 : cont.h = h;
277 32 : dispatch_coro(ex, cont).resume();
278 32 : }
279 :
280 : inline void
281 1 : posix_resolver::resolve_op::destroy()
282 : {
283 1 : stop_cb.reset();
284 1 : auto local_ex = ex;
285 : // May destroy the implementation, and with it this op.
286 1 : impl_ptr.reset();
287 1 : local_ex.on_work_finished();
288 1 : }
289 :
290 : // posix_resolver::reverse_resolve_op implementation
291 :
292 : inline void
293 18 : posix_resolver::reverse_resolve_op::reset() noexcept
294 : {
295 18 : ep = endpoint{};
296 18 : flags = reverse_flags::none;
297 18 : stored_host.clear();
298 18 : stored_service.clear();
299 18 : gai_error = 0;
300 18 : cancelled.store(false, std::memory_order_relaxed);
301 18 : stop_cb.reset();
302 18 : ec_out = nullptr;
303 18 : result_out = nullptr;
304 18 : }
305 :
306 : inline void
307 16 : posix_resolver::reverse_resolve_op::operator()()
308 : {
309 16 : stop_cb.reset(); // Disconnect stop callback
310 :
311 16 : bool const was_cancelled = cancelled.load(std::memory_order_acquire);
312 :
313 16 : if (ec_out)
314 : {
315 16 : if (was_cancelled)
316 MIS 0 : *ec_out = capy::error::canceled;
317 HIT 16 : else if (gai_error != 0)
318 6 : *ec_out = posix_resolver_detail::make_gai_error(gai_error);
319 : else
320 10 : *ec_out = {}; // Clear on success
321 : }
322 :
323 16 : if (result_out && !was_cancelled && gai_error == 0)
324 : {
325 10 : *result_out =
326 10 : endpoint_name{std::move(stored_host), std::move(stored_service)};
327 : }
328 :
329 : // Hold the keepalive across the dispatch: it may be the last
330 : // reference to the implementation this op is embedded in.
331 16 : auto prevent_destroy = std::move(impl_ptr);
332 16 : ex.on_work_finished();
333 16 : cont.h = h;
334 16 : dispatch_coro(ex, cont).resume();
335 16 : }
336 :
337 : inline void
338 1 : posix_resolver::reverse_resolve_op::destroy()
339 : {
340 1 : stop_cb.reset();
341 1 : auto local_ex = ex;
342 : // May destroy the implementation, and with it this op.
343 1 : impl_ptr.reset();
344 1 : local_ex.on_work_finished();
345 1 : }
346 :
347 : // posix_resolver implementation
348 :
349 : inline std::coroutine_handle<>
350 35 : posix_resolver::resolve(
351 : std::coroutine_handle<> h,
352 : capy::executor_ref ex,
353 : std::string_view host,
354 : std::string_view service,
355 : resolve_flags flags,
356 : std::stop_token token,
357 : std::error_code* ec,
358 : std::vector<endpoint>* out)
359 : {
360 35 : if (svc_.resolver_unavailable())
361 : {
362 1 : *ec = std::make_error_code(std::errc::operation_not_supported);
363 1 : op_.cont.h = h;
364 1 : return dispatch_coro(ex, op_.cont);
365 : }
366 :
367 34 : auto& op = op_;
368 34 : op.reset();
369 34 : op.h = h;
370 34 : op.ex = ex;
371 34 : op.ec_out = ec;
372 34 : op.out = out;
373 34 : op.host = host;
374 34 : op.service = service;
375 34 : op.flags = flags;
376 34 : op.start(token);
377 :
378 : // Keep io_context alive while resolution is pending
379 34 : op.ex.on_work_started();
380 :
381 : // Prevent impl destruction while work is in flight
382 34 : resolve_pool_op_.resolver_ = this;
383 34 : resolve_pool_op_.ref_ = this->shared_from_this();
384 34 : resolve_pool_op_.func_ = &posix_resolver::do_resolve_work;
385 34 : if (auto pec = svc_.pool().post(&resolve_pool_op_))
386 : {
387 : // The pool is shutting down, or the system refused it a thread.
388 : // Nothing of this resolve went cross-thread, so it answers here
389 : // like the no-resolver exit above rather than through a
390 : // completion the scheduler has to carry back.
391 1 : resolve_pool_op_.ref_.reset();
392 1 : op.stop_cb.reset();
393 1 : op.ex.on_work_finished();
394 1 : *ec = pec;
395 1 : op.cont.h = h;
396 1 : return dispatch_coro(ex, op.cont);
397 : }
398 33 : return std::noop_coroutine();
399 : }
400 :
401 : inline std::coroutine_handle<>
402 19 : posix_resolver::reverse_resolve(
403 : std::coroutine_handle<> h,
404 : capy::executor_ref ex,
405 : endpoint const& ep,
406 : reverse_flags flags,
407 : std::stop_token token,
408 : std::error_code* ec,
409 : endpoint_name* result_out)
410 : {
411 19 : if (svc_.resolver_unavailable())
412 : {
413 1 : *ec = std::make_error_code(std::errc::operation_not_supported);
414 1 : reverse_op_.cont.h = h;
415 1 : return dispatch_coro(ex, reverse_op_.cont);
416 : }
417 :
418 18 : auto& op = reverse_op_;
419 18 : op.reset();
420 18 : op.h = h;
421 18 : op.ex = ex;
422 18 : op.ec_out = ec;
423 18 : op.result_out = result_out;
424 18 : op.ep = ep;
425 18 : op.flags = flags;
426 18 : op.start(token);
427 :
428 : // Keep io_context alive while resolution is pending
429 18 : op.ex.on_work_started();
430 :
431 : // Prevent impl destruction while work is in flight
432 18 : reverse_pool_op_.resolver_ = this;
433 18 : reverse_pool_op_.ref_ = this->shared_from_this();
434 18 : reverse_pool_op_.func_ = &posix_resolver::do_reverse_resolve_work;
435 18 : if (auto pec = svc_.pool().post(&reverse_pool_op_))
436 : {
437 : // The pool is shutting down, or the system refused it a thread.
438 : // Nothing of this resolve went cross-thread, so it answers here
439 : // like the no-resolver exit above rather than through a
440 : // completion the scheduler has to carry back.
441 1 : reverse_pool_op_.ref_.reset();
442 1 : op.stop_cb.reset();
443 1 : op.ex.on_work_finished();
444 1 : *ec = pec;
445 1 : op.cont.h = h;
446 1 : return dispatch_coro(ex, op.cont);
447 : }
448 17 : return std::noop_coroutine();
449 : }
450 :
451 : inline void
452 73 : posix_resolver::cancel() noexcept
453 : {
454 73 : op_.request_cancel();
455 73 : reverse_op_.request_cancel();
456 73 : }
457 :
458 : inline void
459 33 : posix_resolver::do_resolve_work(pool_work_item* w) noexcept
460 : {
461 33 : auto* pw = static_cast<pool_op*>(w);
462 33 : auto* self = pw->resolver_;
463 :
464 33 : struct addrinfo hints{};
465 33 : hints.ai_family = AF_UNSPEC;
466 33 : hints.ai_socktype = SOCK_STREAM;
467 33 : hints.ai_flags = posix_resolver_detail::flags_to_hints(self->op_.flags);
468 :
469 33 : struct addrinfo* ai = nullptr;
470 99 : int result = ::getaddrinfo(
471 66 : self->op_.host.empty() ? nullptr : self->op_.host.c_str(),
472 61 : self->op_.service.empty() ? nullptr : self->op_.service.c_str(), &hints,
473 : &ai);
474 :
475 33 : if (!self->op_.cancelled.load(std::memory_order_acquire))
476 : {
477 32 : if (result == 0 && ai)
478 : {
479 : self->op_.stored_results =
480 21 : posix_resolver_detail::convert_results(ai);
481 21 : self->op_.gai_error = 0;
482 : }
483 : else
484 : {
485 11 : self->op_.gai_error = result;
486 : }
487 : }
488 :
489 33 : if (ai)
490 22 : ::freeaddrinfo(ai);
491 :
492 : // Hand the keepalive to the op: the completion waits in the
493 : // scheduler's queue, and the implementation embedding it must
494 : // outlive that wait. Nothing may touch *self after the post.
495 33 : self->op_.impl_ptr = std::move(pw->ref_);
496 33 : self->svc_.post(&self->op_);
497 33 : }
498 :
499 : inline void
500 17 : posix_resolver::do_reverse_resolve_work(pool_work_item* w) noexcept
501 : {
502 17 : auto* pw = static_cast<pool_op*>(w);
503 17 : auto* self = pw->resolver_;
504 :
505 17 : sockaddr_storage ss{};
506 : socklen_t ss_len;
507 :
508 17 : if (self->reverse_op_.ep.is_v4())
509 : {
510 15 : auto sa = to_sockaddr_in(self->reverse_op_.ep);
511 15 : std::memcpy(&ss, &sa, sizeof(sa));
512 15 : ss_len = sizeof(sockaddr_in);
513 : }
514 : else
515 : {
516 2 : auto sa = to_sockaddr_in6(self->reverse_op_.ep);
517 2 : std::memcpy(&ss, &sa, sizeof(sa));
518 2 : ss_len = sizeof(sockaddr_in6);
519 : }
520 :
521 : char host[NI_MAXHOST];
522 : char service[NI_MAXSERV];
523 :
524 17 : int result = ::getnameinfo(
525 : reinterpret_cast<sockaddr*>(&ss), ss_len, host, sizeof(host), service,
526 : sizeof(service),
527 : posix_resolver_detail::flags_to_ni_flags(self->reverse_op_.flags));
528 :
529 17 : if (!self->reverse_op_.cancelled.load(std::memory_order_acquire))
530 : {
531 16 : if (result == 0)
532 : {
533 10 : self->reverse_op_.stored_host = host;
534 10 : self->reverse_op_.stored_service = service;
535 10 : self->reverse_op_.gai_error = 0;
536 : }
537 : else
538 : {
539 6 : self->reverse_op_.gai_error = result;
540 : }
541 : }
542 :
543 : // Hand the keepalive to the op: the completion waits in the
544 : // scheduler's queue, and the implementation embedding it must
545 : // outlive that wait. Nothing may touch *self after the post.
546 17 : self->reverse_op_.impl_ptr = std::move(pw->ref_);
547 17 : self->svc_.post(&self->reverse_op_);
548 17 : }
549 :
550 : // posix_resolver_service implementation
551 :
552 : inline void
553 65 : posix_resolver_service::shutdown()
554 : {
555 65 : std::lock_guard<std::mutex> lock(mutex_);
556 :
557 : // Cancel all resolvers (sets cancelled flag checked by pool threads)
558 66 : for (auto* impl = resolver_list_.pop_front(); impl != nullptr;
559 1 : impl = resolver_list_.pop_front())
560 : {
561 1 : impl->cancel();
562 : }
563 :
564 : // Clear the map which releases shared_ptrs.
565 : // The thread pool service shuts down separately via
566 : // execution_context service ordering.
567 65 : resolver_ptrs_.clear();
568 65 : }
569 :
570 : inline io_object::implementation*
571 66 : posix_resolver_service::construct()
572 : {
573 66 : auto ptr = std::make_shared<posix_resolver>(*this);
574 66 : auto* impl = ptr.get();
575 :
576 : {
577 66 : std::lock_guard<std::mutex> lock(mutex_);
578 66 : resolver_list_.push_back(impl);
579 66 : resolver_ptrs_[impl] = std::move(ptr);
580 66 : }
581 :
582 66 : return impl;
583 66 : }
584 :
585 : inline void
586 65 : posix_resolver_service::destroy_impl(posix_resolver& impl)
587 : {
588 65 : std::lock_guard<std::mutex> lock(mutex_);
589 65 : resolver_list_.remove(&impl);
590 65 : resolver_ptrs_.erase(&impl);
591 65 : }
592 :
593 : inline void
594 50 : posix_resolver_service::post(scheduler_op* op)
595 : {
596 50 : sched_->post(op);
597 50 : }
598 :
599 : } // namespace boost::corosio::detail
600 :
601 : #endif // BOOST_COROSIO_POSIX
602 :
603 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RESOLVER_SERVICE_HPP
|