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 1225 : posix_resolver_service(capy::execution_context& ctx, scheduler& sched)
39 2450 : : sched_(&sched)
40 1225 : , pool_(ctx.use_service<thread_pool>())
41 : {
42 1225 : }
43 :
44 2450 : ~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 43 : void destroy(io_object::implementation* p) override
52 : {
53 43 : auto& impl = static_cast<posix_resolver&>(*p);
54 43 : impl.cancel();
55 43 : destroy_impl(impl);
56 43 : }
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 34 : thread_pool& pool() noexcept
67 : {
68 34 : 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 36 : bool resolver_unavailable() const noexcept
75 : {
76 36 : 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 22 : posix_resolver_detail::flags_to_hints(resolve_flags flags)
108 : {
109 22 : int hints = 0;
110 :
111 22 : if ((flags & resolve_flags::passive) != resolve_flags::none)
112 1 : hints |= AI_PASSIVE;
113 22 : if ((flags & resolve_flags::numeric_host) != resolve_flags::none)
114 12 : hints |= AI_NUMERICHOST;
115 22 : if ((flags & resolve_flags::numeric_service) != resolve_flags::none)
116 9 : hints |= AI_NUMERICSERV;
117 22 : if ((flags & resolve_flags::address_configured) != resolve_flags::none)
118 1 : hints |= AI_ADDRCONFIG;
119 22 : if ((flags & resolve_flags::v4_mapped) != resolve_flags::none)
120 1 : hints |= AI_V4MAPPED;
121 22 : if ((flags & resolve_flags::all_matching) != resolve_flags::none)
122 1 : hints |= AI_ALL;
123 :
124 22 : 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 17 : posix_resolver_detail::convert_results(
146 : struct addrinfo* ai, std::string_view host, std::string_view service)
147 : {
148 17 : std::vector<resolver_entry> entries;
149 17 : entries.reserve(4); // Most lookups return 1-4 addresses
150 :
151 34 : for (auto* p = ai; p != nullptr; p = p->ai_next)
152 : {
153 17 : if (p->ai_family == AF_INET)
154 : {
155 15 : auto* addr = reinterpret_cast<sockaddr_in*>(p->ai_addr);
156 15 : auto ep = from_sockaddr_in(*addr);
157 15 : 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 17 : 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 43 : inline posix_resolver::posix_resolver(posix_resolver_service& svc) noexcept
237 43 : : svc_(svc)
238 : {
239 43 : }
240 :
241 : // posix_resolver::resolve_op implementation
242 :
243 : inline void
244 22 : posix_resolver::resolve_op::reset() noexcept
245 : {
246 22 : host.clear();
247 22 : service.clear();
248 22 : flags = resolve_flags::none;
249 22 : stored_results = resolver_results{};
250 22 : gai_error = 0;
251 22 : cancelled.store(false, std::memory_order_relaxed);
252 22 : stop_cb.reset();
253 22 : ec_out = nullptr;
254 22 : out = nullptr;
255 22 : }
256 :
257 : inline void
258 22 : posix_resolver::resolve_op::operator()()
259 : {
260 22 : stop_cb.reset(); // Disconnect stop callback
261 :
262 22 : bool const was_cancelled = cancelled.load(std::memory_order_acquire);
263 :
264 22 : if (ec_out)
265 : {
266 22 : if (was_cancelled)
267 1 : *ec_out = capy::error::canceled;
268 21 : else if (gai_error != 0)
269 4 : *ec_out = posix_resolver_detail::make_gai_error(gai_error);
270 : else
271 17 : *ec_out = {}; // Clear on success
272 : }
273 :
274 22 : if (out && !was_cancelled && gai_error == 0)
275 17 : *out = std::move(stored_results);
276 :
277 22 : impl->svc_.work_finished();
278 22 : cont.h = h;
279 22 : dispatch_coro(ex, cont).resume();
280 22 : }
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 48 : posix_resolver::resolve_op::request_cancel() noexcept
290 : {
291 48 : cancelled.store(true, std::memory_order_release);
292 48 : }
293 :
294 : inline void
295 22 : posix_resolver::resolve_op::start(std::stop_token const& token)
296 : {
297 22 : cancelled.store(false, std::memory_order_release);
298 22 : stop_cb.reset();
299 :
300 22 : if (token.stop_possible())
301 1 : stop_cb.emplace(token, canceller{this});
302 22 : }
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 48 : posix_resolver::reverse_resolve_op::request_cancel() noexcept
356 : {
357 48 : cancelled.store(true, std::memory_order_release);
358 48 : }
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 23 : 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 23 : 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 22 : auto& op = op_;
391 22 : op.reset();
392 22 : op.h = h;
393 22 : op.ex = ex;
394 22 : op.impl = this;
395 22 : op.ec_out = ec;
396 22 : op.out = out;
397 22 : op.host = host;
398 22 : op.service = service;
399 22 : op.flags = flags;
400 22 : op.start(token);
401 :
402 : // Keep io_context alive while resolution is pending
403 22 : op.ex.on_work_started();
404 :
405 : // Prevent impl destruction while work is in flight
406 22 : resolve_pool_op_.resolver_ = this;
407 22 : resolve_pool_op_.ref_ = this->shared_from_this();
408 22 : resolve_pool_op_.func_ = &posix_resolver::do_resolve_work;
409 22 : 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 22 : 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 47 : posix_resolver::cancel() noexcept
466 : {
467 47 : op_.request_cancel();
468 47 : reverse_op_.request_cancel();
469 47 : }
470 :
471 : inline void
472 22 : posix_resolver::do_resolve_work(pool_work_item* w) noexcept
473 : {
474 22 : auto* pw = static_cast<pool_op*>(w);
475 22 : auto* self = pw->resolver_;
476 :
477 22 : struct addrinfo hints{};
478 22 : hints.ai_family = AF_UNSPEC;
479 22 : hints.ai_socktype = SOCK_STREAM;
480 22 : hints.ai_flags = posix_resolver_detail::flags_to_hints(self->op_.flags);
481 :
482 22 : struct addrinfo* ai = nullptr;
483 66 : int result = ::getaddrinfo(
484 44 : self->op_.host.empty() ? nullptr : self->op_.host.c_str(),
485 44 : self->op_.service.empty() ? nullptr : self->op_.service.c_str(), &hints,
486 : &ai);
487 :
488 22 : if (!self->op_.cancelled.load(std::memory_order_acquire))
489 : {
490 21 : if (result == 0 && ai)
491 : {
492 34 : self->op_.stored_results = posix_resolver_detail::convert_results(
493 17 : ai, self->op_.host, self->op_.service);
494 17 : self->op_.gai_error = 0;
495 : }
496 : else
497 : {
498 4 : self->op_.gai_error = result;
499 : }
500 : }
501 :
502 22 : if (ai)
503 18 : ::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 22 : auto ref = std::move(pw->ref_);
508 22 : self->svc_.post(&self->op_);
509 22 : }
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 1225 : posix_resolver_service::shutdown()
565 : {
566 1225 : std::lock_guard<std::mutex> lock(mutex_);
567 :
568 : // Cancel all resolvers (sets cancelled flag checked by pool threads)
569 1225 : 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 1225 : resolver_ptrs_.clear();
579 1225 : }
580 :
581 : inline io_object::implementation*
582 43 : posix_resolver_service::construct()
583 : {
584 43 : auto ptr = std::make_shared<posix_resolver>(*this);
585 43 : auto* impl = ptr.get();
586 :
587 : {
588 43 : std::lock_guard<std::mutex> lock(mutex_);
589 43 : resolver_list_.push_back(impl);
590 43 : resolver_ptrs_[impl] = std::move(ptr);
591 43 : }
592 :
593 43 : return impl;
594 43 : }
595 :
596 : inline void
597 43 : posix_resolver_service::destroy_impl(posix_resolver& impl)
598 : {
599 43 : std::lock_guard<std::mutex> lock(mutex_);
600 43 : resolver_list_.remove(&impl);
601 43 : resolver_ptrs_.erase(&impl);
602 43 : }
603 :
604 : inline void
605 34 : posix_resolver_service::post(scheduler_op* op)
606 : {
607 34 : sched_->post(op);
608 34 : }
609 :
610 : inline void
611 : posix_resolver_service::work_started() noexcept
612 : {
613 : sched_->work_started();
614 : }
615 :
616 : inline void
617 34 : posix_resolver_service::work_finished() noexcept
618 : {
619 34 : sched_->work_finished();
620 34 : }
621 :
622 : // Free function to get/create the resolver service
623 :
624 : inline posix_resolver_service&
625 1225 : get_resolver_service(capy::execution_context& ctx, scheduler& sched)
626 : {
627 1225 : 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
|