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