#include "etl/LoadBalancer.hpp" #include "data/BackendInterface.hpp" #include "etl/ETLState.hpp" #include "etl/InitialLoadObserverInterface.hpp" #include "etl/LoadBalancerInterface.hpp" #include "etl/NetworkValidatedLedgersInterface.hpp" #include "etl/Source.hpp" #include "feed/SubscriptionManagerInterface.hpp" #include "rpc/Errors.hpp" #include "util/Assert.hpp" #include "util/CoroutineGroup.hpp" #include "util/Profiler.hpp" #include "util/Random.hpp" #include "util/ResponseExpirationCache.hpp" #include "util/config/ArrayView.hpp" #include "util/config/ConfigDefinition.hpp" #include "util/config/ObjectView.hpp" #include "util/log/Logger.hpp" #include "util/prometheus/Label.hpp" #include "util/prometheus/Prometheus.hpp" #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include using namespace util::config; using util::prometheus::Labels; namespace etl { std::shared_ptr LoadBalancer::makeLoadBalancer( ClioConfigDefinition const& config, boost::asio::io_context& ioc, std::shared_ptr backend, std::shared_ptr subscriptions, std::unique_ptr randomGenerator, std::shared_ptr validatedLedgers, SourceFactory sourceFactory ) { return std::make_shared( config, ioc, std::move(backend), std::move(subscriptions), std::move(randomGenerator), std::move(validatedLedgers), std::move(sourceFactory) ); } LoadBalancer::LoadBalancer( ClioConfigDefinition const& config, boost::asio::io_context& ioc, std::shared_ptr backend, std::shared_ptr subscriptions, std::unique_ptr randomGenerator, std::shared_ptr validatedLedgers, SourceFactory sourceFactory ) : randomGenerator_(std::move(randomGenerator)) , forwardingCounters_{ .successDuration = PrometheusService::counterInt( "forwarding_duration_milliseconds_counter", Labels({util::prometheus::Label{"status", "success"}}), "The duration of processing successful forwarded requests" ), .failDuration = PrometheusService::counterInt( "forwarding_duration_milliseconds_counter", Labels({util::prometheus::Label{"status", "fail"}}), "The duration of processing failed forwarded requests" ), .retries = PrometheusService::counterInt( "forwarding_retries_counter", Labels(), "The number of retries before a forwarded request was successful. Initial attempt " "excluded" ), .cacheHit = PrometheusService::counterInt( "forwarding_cache_hit_counter", Labels(), "The number of requests that we served from the cache" ), .cacheMiss = PrometheusService::counterInt( "forwarding_cache_miss_counter", Labels(), "The number of requests that were not served from the cache" ) } { auto const forwardingCacheTimeout = config.get("forwarding.cache_timeout"); if (forwardingCacheTimeout > 0.f) { forwardingCache_ = util::ResponseExpirationCache{ util::config::ClioConfigDefinition::toMilliseconds(forwardingCacheTimeout), {"server_info", "server_state", "server_definitions", "fee", "ledger_closed"} }; } auto const numMarkers = config.getValueView("num_markers"); if (numMarkers.hasValue()) { auto const value = numMarkers.asIntType(); downloadRanges_ = value; } else if (backend->fetchLedgerRange()) { downloadRanges_ = 4; } auto const allowNoEtl = config.get("allow_no_etl"); auto const checkOnETLFailure = [this, allowNoEtl](std::string const& log) { LOG(log_.warn()) << log; if (!allowNoEtl) { LOG( log_.error() ) << "Set allow_no_etl as true in config to allow clio run without valid ETL sources."; throw std::logic_error("ETL configuration error."); } }; auto const forwardingTimeout = ClioConfigDefinition::toMilliseconds(config.get("forwarding.request_timeout")); auto const etlArray = config.getArray("etl_sources"); for (auto it = etlArray.begin(); it != etlArray.end(); ++it) { auto source = sourceFactory( *it, ioc, subscriptions, validatedLedgers, forwardingTimeout, [this]() { if (not hasForwardingSource_.lock().get()) chooseForwardingSource(); }, [this](bool wasForwarding) { if (wasForwarding) chooseForwardingSource(); }, [this]() { if (forwardingCache_.has_value()) forwardingCache_->invalidate(); } ); // checking etl node validity auto const stateOpt = ETLState::fetchETLStateFromSource(*source); if (!stateOpt) { LOG(log_.warn()) << "Failed to fetch ETL state from source = " << source->toString() << " Please check the configuration and network"; } else if (etlState_ && etlState_->networkID != stateOpt->networkID) { checkOnETLFailure( fmt::format( "ETL sources must be on the same network. Source network id = {} does not " "match others network id " "= {}", stateOpt->networkID, etlState_->networkID ) ); } else { etlState_ = stateOpt; } sources_.push_back(std::move(source)); LOG(log_.info()) << "Added etl source - " << sources_.back()->toString(); } if (!etlState_) { checkOnETLFailure( "Failed to fetch ETL state from any source. Please check the configuration and network" ); } if (sources_.empty()) checkOnETLFailure("No ETL sources configured. Please check the configuration"); // This is made separate from source creation to prevent UB in case one of the sources will call // chooseForwardingSource while we are still filling the sources_ vector for (auto const& source : sources_) { source->run(); } } InitialLedgerLoadResult LoadBalancer::loadInitialLedger( uint32_t sequence, InitialLoadObserverInterface& loadObserver, std::chrono::steady_clock::duration retryAfter ) { InitialLedgerLoadResult response; execute( [this, &response, &sequence, &loadObserver](auto& source) { auto res = source->loadInitialLedger(sequence, downloadRanges_, loadObserver); if (not res.has_value() and res.error() == InitialLedgerLoadError::Errored) { LOG(log_.error()) << "Failed to download initial ledger." << " Sequence = " << sequence << " source = " << source->toString(); return false; // should retry on error } response = std::move(res); // cancelled or data received return true; }, sequence, retryAfter ); return response; } LoadBalancer::OptionalGetLedgerResponseType LoadBalancer::fetchLedger( uint32_t ledgerSequence, bool getObjects, bool getObjectNeighbors, std::chrono::steady_clock::duration retryAfter ) { GetLedgerResponseType response; execute( [&response, ledgerSequence, getObjects, getObjectNeighbors, log = log_](auto& source) { auto [status, data] = source->fetchLedger(ledgerSequence, getObjects, getObjectNeighbors); response = std::move(data); if (status.ok() && response.validated()) { LOG(log.info()) << "Successfully fetched ledger = " << ledgerSequence << " from source = " << source->toString(); return true; } LOG(log.warn()) << "Could not fetch ledger " << ledgerSequence << ", Reply: " << response.DebugString() << ", error_code: " << status.error_code() << ", error_msg: " << status.error_message() << ", source = " << source->toString(); return false; }, ledgerSequence, retryAfter ); return response; } std::expected LoadBalancer::forwardToRippled( boost::json::object const& request, std::optional const& clientIp, bool isAdmin, boost::asio::yield_context yield ) { if (not request.contains("command")) return std::unexpected{rpc::ClioError::RpcCommandIsMissing}; auto const cmd = boost::json::value_to(request.at("command")); if (shouldUseCache(isAdmin)) { // NOLINTNEXTLINE(bugprone-unchecked-optional-access) if (auto cachedResponse = forwardingCache_->get(cmd, request); cachedResponse) { forwardingCounters_.cacheHit.get() += 1; return *std::move(cachedResponse); } } forwardingCounters_.cacheMiss.get() += 1; ASSERT(not sources_.empty(), "ETL sources must be configured to forward requests."); std::size_t sourceIdx = randomGenerator_->uniform(0ul, sources_.size() - 1); auto numAttempts = 0u; auto xUserValue = isAdmin ? kAdminForwardingXUserValue : kUserForwardingXUserValue; std::optional response; rpc::ClioError error = rpc::ClioError::EtlConnectionError; while (numAttempts < sources_.size()) { auto [res, duration] = util::timed([&]() { return sources_[sourceIdx]->forwardToRippled(request, clientIp, xUserValue, yield); }); if (res) { forwardingCounters_.successDuration.get() += duration; response = std::move(res).value(); break; } forwardingCounters_.failDuration.get() += duration; ++forwardingCounters_.retries.get(); error = std::max(error, res.error()); // Choose the best result between all sources sourceIdx = (sourceIdx + 1) % sources_.size(); ++numAttempts; } if (response) { if (shouldUseCache(isAdmin) and not response->contains("error")) { // NOLINTNEXTLINE(bugprone-unchecked-optional-access) forwardingCache_->put(cmd, request, *response); } return *std::move(response); } return std::unexpected{error}; } boost::json::value LoadBalancer::toJson() const { boost::json::array ret; for (auto& src : sources_) ret.push_back(src->toJson()); return ret; } template void LoadBalancer::execute( Func f, uint32_t ledgerSequence, std::chrono::steady_clock::duration retryAfter ) { ASSERT(not sources_.empty(), "ETL sources must be configured to execute functions."); size_t sourceIdx = randomGenerator_->uniform(0ul, sources_.size() - 1); size_t numAttempts = 0; while (true) { auto& source = sources_[sourceIdx]; LOG(log_.debug()) << "Attempting to execute func. ledger sequence = " << ledgerSequence << " - source = " << source->toString(); // Originally, it was (source->hasLedger(ledgerSequence) || true) /* Sometimes rippled has ledger but doesn't actually know. However, but this does NOT happen in the normal case and is safe to remove This || true is only needed when loading full history standalone */ if (source->hasLedger(ledgerSequence)) { bool const res = f(source); if (res) { LOG(log_.debug()) << "Successfully executed func at source = " << source->toString() << " - ledger sequence = " << ledgerSequence; break; } LOG(log_.warn()) << "Failed to execute func at source = " << source->toString() << " - ledger sequence = " << ledgerSequence; } else { LOG(log_.warn()) << "Ledger not present at source = " << source->toString() << " - ledger sequence = " << ledgerSequence; } sourceIdx = (sourceIdx + 1) % sources_.size(); numAttempts++; if (numAttempts % sources_.size() == 0) { LOG( log_.info() ) << "Ledger sequence " << ledgerSequence << " is not yet available from any configured sources. Sleeping and trying again"; std::this_thread::sleep_for(retryAfter); } } } std::optional LoadBalancer::getETLState() noexcept { if (!etlState_) { // retry ETLState fetch etlState_ = ETLState::fetchETLStateFromSource(*this); } return etlState_; } void LoadBalancer::stop(boost::asio::yield_context yield) { util::CoroutineGroup group{yield}; std::ranges::for_each(sources_, [&group, yield](auto& source) { group.spawn(yield, [&source](boost::asio::yield_context innerYield) { source->stop(innerYield); }); }); group.asyncWait(yield); } void LoadBalancer::chooseForwardingSource() { LOG(log_.info()) << "Choosing a new source to forward subscriptions"; auto hasForwardingSourceLock = hasForwardingSource_.lock(); hasForwardingSourceLock.get() = false; for (auto& source : sources_) { if (not hasForwardingSourceLock.get() and source->isConnected()) { source->setForwarding(true); hasForwardingSourceLock.get() = true; } else { source->setForwarding(false); } } } bool LoadBalancer::shouldUseCache(bool isAdmin) const { return forwardingCache_.has_value() and not isAdmin; } } // namespace etl