1#include <xrpl/basics/ResolverAsio.h>
3#include <xrpl/basics/Log.h>
4#include <xrpl/basics/Resolver.h>
5#include <xrpl/beast/net/IPAddressConversion.h>
6#include <xrpl/beast/net/IPEndpoint.h>
7#include <xrpl/beast/utility/Journal.h>
8#include <xrpl/beast/utility/instrumentation.h>
10#include <boost/asio/bind_executor.hpp>
11#include <boost/asio/dispatch.hpp>
12#include <boost/asio/error.hpp>
13#include <boost/asio/io_context.hpp>
14#include <boost/asio/ip/tcp.hpp>
15#include <boost/asio/post.hpp>
16#include <boost/asio/strand.hpp>
17#include <boost/system/detail/error_code.hpp>
41template <
class Derived>
52 XRPL_ASSERT(
pending_.load() == 0,
"xrpl::AsyncObject::~AsyncObject : nothing pending");
75 if (--
owner_->pending_ == 0)
76 owner_->asyncHandlersComplete();
96 (
static_cast<Derived*
>(
this))->asyncHandlersComplete();
114 boost::asio::strand<boost::asio::io_context::executor_type>
strand;
130 template <
class StringSequence>
133 names.reserve(inNames.size());
153 XRPL_ASSERT(
work.empty(),
"xrpl::ResolverAsioImpl::~ResolverAsioImpl : no pending work");
154 XRPL_ASSERT(
stopped,
"xrpl::ResolverAsioImpl::~ResolverAsioImpl : stopped");
176 XRPL_ASSERT(
stopped ==
true,
"xrpl::ResolverAsioImpl::start : stopped");
177 XRPL_ASSERT(
stopCalled ==
false,
"xrpl::ResolverAsioImpl::start : not stopping");
194 boost::asio::dispatch(
196 boost::asio::bind_executor(
197 strand, [
this, counter = CompletionCounter(
this)] {
doStop(counter); }));
199 JLOG(
journal.debug()) <<
"Queued a stop request";
208 JLOG(
journal.debug()) <<
"Waiting to stop";
212 JLOG(
journal.debug()) <<
"Stopped";
218 XRPL_ASSERT(
stopCalled ==
false,
"xrpl::ResolverAsioImpl::resolve : not stopping");
219 XRPL_ASSERT(!names.
empty(),
"xrpl::ResolverAsioImpl::resolve : names non-empty");
223 boost::asio::dispatch(
225 boost::asio::bind_executor(
226 strand, [
this, names, handler, counter = CompletionCounter(
this)] {
236 XRPL_ASSERT(
stopCalled ==
true,
"xrpl::ResolverAsioImpl::doStop : stopping");
250 boost::system::error_code
const& ec,
252 boost::asio::ip::tcp::resolver::results_type results,
255 if (ec == boost::asio::error::operation_aborted)
259 auto iter = results.
begin();
265 while (iter != results.end())
272 handler(name, addresses);
276 boost::asio::bind_executor(
277 strand, [
this, counter = CompletionCounter(
this)] {
doWork(counter); }));
288 return make_pair(result->address().to_string(),
std::to_string(result->port()));
297 auto const findWhitespace = [&loc](std::string::value_type c) {
298 return std::isspace<std::string::value_type>(c, loc);
307 if (hostFirst >= portLast)
311 auto const findPortSeparator = [](
char const c) ->
bool {
312 if (std::isspace(
static_cast<unsigned char>(c)))
321 auto hostLast =
std::find_if(hostFirst, portLast, findPortSeparator);
341 work.front().names.pop_back();
343 if (
work.front().names.empty())
346 auto const [host, port] =
parseName(name);
350 JLOG(
journal.error()) <<
"Unable to parse '" << name <<
"'";
354 boost::asio::bind_executor(
355 strand, [
this, counter = CompletionCounter(
this)] {
doWork(counter); }));
363 [
this, name, handler, counter = CompletionCounter(
this)](
364 boost::system::error_code
const& ec,
365 boost::asio::ip::tcp::resolver::results_type results) {
366 doFinish(name, ec, handler, results, counter);
373 XRPL_ASSERT(!names.
empty(),
"xrpl::ResolverAsioImpl::doResolve : names non-empty");
377 work.emplace_back(names, handler);
379 JLOG(
journal.debug()) <<
"Queued new job with " << names.
size() <<
" tasks. "
380 <<
work.size() <<
" jobs outstanding.";
386 boost::asio::bind_executor(
387 strand, [
this, counter = CompletionCounter(
this)] {
doWork(counter); }));
T back_inserter(T... args)
A generic endpoint for log messages.
static std::optional< Endpoint > fromStringChecked(std::string const &s)
Create an Endpoint from a string.
RAII container that maintains the count of pending I/O.
CompletionCounter & operator=(CompletionCounter const &)=delete
CompletionCounter(Derived *owner)
CompletionCounter(CompletionCounter const &other)
std::atomic< int > pending_
void resolve(std::vector< std::string > const &names, HandlerType const &handler) override
void stopAsync() override
Issue an asynchronous stop request.
void start() override
Issue a synchronous start request.
ResolverAsioImpl(boost::asio::io_context &ioContext, beast::Journal journal)
std::condition_variable cv
std::atomic< bool > stopCalled
void doWork(CompletionCounter)
void stop() override
Issue a synchronous stop request.
~ResolverAsioImpl() override
bool asyncHandlersCompleted
static HostAndPort parseName(std::string const &str)
void doFinish(std::string name, boost::system::error_code const &ec, HandlerType handler, boost::asio::ip::tcp::resolver::results_type results, CompletionCounter)
boost::asio::ip::tcp::resolver resolver
boost::asio::strand< boost::asio::io_context::executor_type > strand
std::pair< std::string, std::string > HostAndPort
boost::asio::io_context & ioContext
std::atomic< bool > stopped
void doStop(CompletionCounter)
void asyncHandlersComplete()
void doResolve(std::vector< std::string > const &names, HandlerType const &handler, CompletionCounter)
static std::unique_ptr< ResolverAsio > make(boost::asio::io_context &, beast::Journal)
std::function< void(std::string, std::vector< beast::ip::Endpoint >)> HandlerType
Use hash_* containers for keys that do not need a cryptographically secure hashing algorithm.
T reverse_copy(T... args)
static ip::Endpoint fromAsio(boost::asio::ip::address const &address)
Work(StringSequence const &inNames, HandlerType handler)
std::vector< std::string > names