diff --git a/docker/telemetry/xrpld-telemetry.cfg b/docker/telemetry/xrpld-telemetry.cfg index f03af1cbde..12de7576c3 100644 --- a/docker/telemetry/xrpld-telemetry.cfg +++ b/docker/telemetry/xrpld-telemetry.cfg @@ -65,9 +65,8 @@ trace_ledger=1 # endpoint and prefix are informational only. OTelCollector records on the # global MeterProvider that [telemetry] configures, and formatName() does not # apply the prefix, so metric names are bare and lowercase. -# Known limitation: the collector is built before the MeterProvider is -# registered, so beast::insight instruments bind to a no-op meter and are not -# exported yet. Tracing and the [telemetry] metrics pipeline are unaffected. +# Requires [telemetry] enabled=1: the MeterProvider these instruments record on +# is owned by the telemetry module, and without it they are discarded. [insight] server=otel endpoint=http://localhost:4318/v1/metrics diff --git a/include/xrpl/beast/insight/Collector.h b/include/xrpl/beast/insight/Collector.h index 9da2a8bb74..d3e7f4d5e5 100644 --- a/include/xrpl/beast/insight/Collector.h +++ b/include/xrpl/beast/insight/Collector.h @@ -30,6 +30,29 @@ public: virtual ~Collector() = 0; + /** + * Called once the services that hook handlers read are constructed. + * + * Implementations that poll their producers must not do so before this: + * hook handlers read live application state. Default is a no-op, for + * collectors that only push. + */ + virtual void + onCollectionReady() + { + } + + /** + * Called before those services are shut down. + * + * Polling must have stopped by the time this returns. Paired with + * onCollectionReady(). + */ + virtual void + onCollectionStopping() + { + } + /** * Create a hook. * diff --git a/include/xrpl/server/Wallet.h b/include/xrpl/server/Wallet.h index 95486cc468..af6c92b83d 100644 --- a/include/xrpl/server/Wallet.h +++ b/include/xrpl/server/Wallet.h @@ -16,6 +16,7 @@ #include #include +#include #include #include #include @@ -88,6 +89,19 @@ addValidatorManifest(soci::session& session, std::string const& serialized); void clearNodeIdentity(soci::session& session); +/** + * Returns this node's stored keypair, if the database holds a valid one. + * + * Read-only: unlike getNodeIdentity(), never generates or persists a key. A row + * counts only when its public and secret keys are a pair. + * + * @param session Session with the database. + * + * @return The stored keypair, or std::nullopt. + */ +std::optional> +readNodeIdentity(soci::session& session); + /** * Returns a stable public and private key for this node. * diff --git a/src/libxrpl/beast/insight/OTelCollector.cpp b/src/libxrpl/beast/insight/OTelCollector.cpp index 2251c64504..c74b1738a3 100644 --- a/src/libxrpl/beast/insight/OTelCollector.cpp +++ b/src/libxrpl/beast/insight/OTelCollector.cpp @@ -230,13 +230,9 @@ public: * @param name Export-ready metric name, already run through * formatName() by the collector: prefix prepended * and dots replaced with underscores. - * @param meter OTel Meter used to create the observable gauge. * @param collector Owning collector, used to invoke hooks before reads. */ - OTelGaugeImpl( - std::string const& name, - opentelemetry::nostd::shared_ptr const& meter, - std::shared_ptr const& collector); + OTelGaugeImpl(std::string name, std::shared_ptr const& collector); ~OTelGaugeImpl() override; @@ -274,6 +270,25 @@ public: static void gaugeCallback(opentelemetry::metrics::ObserverResult result, void* state); + /** + * Create the observable instrument and register the callback, once. + * + * Called when the collector is told collection is ready, because the + * callback reads live application state. + */ + void + arm(); + + /** + * Remove the callback, so the reader thread stops observing this gauge. + * + * RemoveCallback is synchronous: the SDK guards its callback list and the + * observe pass with the same mutex, so no callback is running once this + * returns. Idempotent. + */ + void + disarm(); + private: /** * Current gauge value, updated atomically by set()/increment(). @@ -281,10 +296,20 @@ private: std::atomic value_{0}; /** - * OTel observable gauge handle (prevents deregistration). + * Export-ready metric name, held until arm() creates the instrument. + */ + std::string const name_; + + /** + * OTel observable gauge handle, null until arm() runs. */ opentelemetry::nostd::shared_ptr gauge_; + /** + * Guards gauge_ against concurrent arm()/disarm(). + */ + std::mutex armMutex_; + /** * Owning collector, used to invoke hooks before reading gauge values. */ @@ -445,6 +470,12 @@ public: Gauge makeGauge(std::string const& name) override; + void + onCollectionReady() override; + + void + onCollectionStopping() override; + Meter makeMeter(std::string const& name) override; /** @} */ @@ -626,16 +657,37 @@ OTelEventImpl::notify(value_type const& value) // OTelGaugeImpl //------------------------------------------------------------------------------ -OTelGaugeImpl::OTelGaugeImpl( - std::string const& name, - opentelemetry::nostd::shared_ptr const& meter, - std::shared_ptr const& collector) - : gauge_(meter->CreateInt64ObservableGauge(name)), collector_(collector) +OTelGaugeImpl::OTelGaugeImpl(std::string name, std::shared_ptr const& collector) + : name_(std::move(name)), collector_(collector) { collector_->addGauge(this); +} + +void +OTelGaugeImpl::arm() +{ + // AddCallback arms the SDK reader thread against this gauge, and the + // callback runs hook handlers that read application services. The registry + // does not de-duplicate callbacks, so arm at most once. + std::scoped_lock const lock(armMutex_); + if (gauge_) + return; + + gauge_ = collector_->otelMeter()->CreateInt64ObservableGauge(name_); gauge_->AddCallback(gaugeCallback, this); } +void +OTelGaugeImpl::disarm() +{ + std::scoped_lock const lock(armMutex_); + if (!gauge_) + return; + + gauge_->RemoveCallback(gaugeCallback, this); + gauge_ = nullptr; +} + void OTelGaugeImpl::gaugeCallback(opentelemetry::metrics::ObserverResult result, void* state) { @@ -656,7 +708,8 @@ OTelGaugeImpl::~OTelGaugeImpl() // The SDK's ObservableRegistry guards its callback list and the Observe() // pass with the same mutex, so RemoveCallback cannot return while a // callback for this instrument is in flight — removal is synchronous. - gauge_->RemoveCallback(gaugeCallback, this); + // A no-op when never armed, or already disarmed at shutdown. + disarm(); collector_->removeGauge(this); } @@ -785,7 +838,7 @@ OTelCollectorImp::makeEvent(std::string const& name) Gauge OTelCollectorImp::makeGauge(std::string const& name) { - return Gauge(std::make_shared(formatName(name), otelMeter_, shared_from_this())); + return Gauge(std::make_shared(formatName(name), shared_from_this())); } Meter @@ -851,6 +904,58 @@ OTelCollectorImp::removeGauge(OTelGaugeImpl* gauge) std::erase(gauges_, gauge); } +void +OTelCollectorImp::onCollectionReady() +{ + // Snapshot under the lock, arm outside it. arm() enters the SDK's + // observable registry lock, and the reader thread takes that lock before + // calling callHooks(), which wants mutex_. callHooks() copies its hook list + // for the same reason. + std::vector gauges; + { + std::scoped_lock const lock(mutex_); + gauges = gauges_; + } + + std::size_t armed = 0; + for (auto* gauge : gauges) + { + // Telemetry must never stop the node, so one bad instrument costs only + // its own metric. + try + { + gauge->arm(); + ++armed; + } + catch (std::exception const& e) + { + JLOG(journal_.error()) << "OTelCollector: could not register an observable gauge, " + "so that metric will not be exported: " + << e.what(); + } + } + + JLOG(journal_.info()) << "OTelCollector: registered " << armed << " of " << gauges.size() + << " observable gauges"; +} + +void +OTelCollectorImp::onCollectionStopping() +{ + // Same lock discipline as onCollectionReady(): snapshot, then act outside + // the lock, because disarm() enters the SDK's observable registry lock. + std::vector gauges; + { + std::scoped_lock const lock(mutex_); + gauges = gauges_; + } + + for (auto* gauge : gauges) + gauge->disarm(); + + JLOG(journal_.info()) << "OTelCollector: stopped observing " << gauges.size() << " gauges"; +} + opentelemetry::nostd::shared_ptr const& OTelCollectorImp::otelMeter() const { diff --git a/src/libxrpl/beast/insight/StatsDCollector.cpp b/src/libxrpl/beast/insight/StatsDCollector.cpp index dc19aa8953..bc2640ca77 100644 --- a/src/libxrpl/beast/insight/StatsDCollector.cpp +++ b/src/libxrpl/beast/insight/StatsDCollector.cpp @@ -23,6 +23,7 @@ #include #include +#include #include #include #include @@ -218,6 +219,13 @@ private: std::recursive_mutex metricsLock_; List metrics_; + /** + * Whether hook handlers may be called. False until onCollectionReady(), + * because the handlers read application services that are still being + * constructed while this collector exists. + */ + std::atomic polling_{false}; + // Must come last for order of init std::thread thread_; @@ -255,6 +263,22 @@ public: thread_.join(); } + void + onCollectionReady() override + { + polling_.store(true, std::memory_order_release); + } + + void + onCollectionStopping() override + { + polling_.store(false, std::memory_order_release); + + // onTimer holds metricsLock_ across the handler loop, so acquiring it + // here waits for a handler that is already running. + std::scoped_lock const _(metricsLock_); + } + Hook makeHook(HookImpl::HandlerType const& handler) override { @@ -437,12 +461,15 @@ public: return; } - std::scoped_lock const _(metricsLock_); + if (polling_.load(std::memory_order_acquire)) + { + std::scoped_lock const _(metricsLock_); - for (auto& m : metrics_) - m.doProcess(); + for (auto& m : metrics_) + m.doProcess(); - sendBuffers(); + sendBuffers(); + } setTimer(); } diff --git a/src/libxrpl/server/Wallet.cpp b/src/libxrpl/server/Wallet.cpp index 56d0db67d4..92317d40f6 100644 --- a/src/libxrpl/server/Wallet.cpp +++ b/src/libxrpl/server/Wallet.cpp @@ -32,6 +32,7 @@ #include #include #include +#include #include #include #include @@ -147,27 +148,34 @@ clearNodeIdentity(soci::session& session) session << "DELETE FROM NodeIdentity;"; } +std::optional> +readNodeIdentity(soci::session& session) +{ + // SOCI requires boost::optional (not std::optional) as the parameter. + boost::optional pubKO, priKO; + soci::statement st = + (session.prepare << "SELECT PublicKey, PrivateKey FROM NodeIdentity;", + soci::into(pubKO), + soci::into(priKO)); + st.execute(); + while (st.fetch()) + { + auto const sk = parseBase58(TokenType::NodePrivate, priKO.value_or("")); + auto const pk = parseBase58(TokenType::NodePublic, pubKO.value_or("")); + + // Only use if the public and secret keys are a pair + if (sk && pk && (*pk == derivePublicKey(KeyType::Secp256k1, *sk))) + return std::pair{*pk, *sk}; + } + + return std::nullopt; +} + std::pair getNodeIdentity(soci::session& session) { - { - // SOCI requires boost::optional (not std::optional) as the parameter. - boost::optional pubKO, priKO; - soci::statement st = - (session.prepare << "SELECT PublicKey, PrivateKey FROM NodeIdentity;", - soci::into(pubKO), - soci::into(priKO)); - st.execute(); - while (st.fetch()) - { - auto const sk = parseBase58(TokenType::NodePrivate, priKO.value_or("")); - auto const pk = parseBase58(TokenType::NodePublic, pubKO.value_or("")); - - // Only use if the public and secret keys are a pair - if (sk && pk && (*pk == derivePublicKey(KeyType::Secp256k1, *sk))) - return {*pk, *sk}; - } - } + if (auto const stored = readNodeIdentity(session)) + return *stored; // If a valid identity wasn't found, we randomly generate a new one: auto [newpublicKey, newsecretKey] = randomKeyPair(KeyType::Secp256k1); diff --git a/src/libxrpl/telemetry/Telemetry.cpp b/src/libxrpl/telemetry/Telemetry.cpp index d285c090a7..bc0135124e 100644 --- a/src/libxrpl/telemetry/Telemetry.cpp +++ b/src/libxrpl/telemetry/Telemetry.cpp @@ -66,6 +66,7 @@ #include #include +#include #include #include #include @@ -317,11 +318,154 @@ class TelemetryImpl : public Telemetry */ opentelemetry::nostd::shared_ptr contextStorage_; + /** + * Set by stop(), so a second call does nothing. + */ + bool stopped_{false}; + + /** + * Build the OTel resource shared by the tracer and meter providers. + * + * Both pipelines must report the same resource identity, so this is the + * single place the attributes are named. Called twice because the two + * providers are now built at different times: metrics in the constructor, + * traces in start(). + * + * @return The resource carrying service and network identity. + */ + [[nodiscard]] resource::Resource + makeResource() const + { + return resource::Resource::Create({ + {opentelemetry::semconv::service::kServiceName, setup_.serviceName}, + {opentelemetry::semconv::service::kServiceVersion, setup_.serviceVersion}, + {opentelemetry::semconv::service::kServiceInstanceId, setup_.serviceInstanceId}, + {std::string(attr::networkId), + static_cast(setup_.networkId)}, // LCOV_EXCL_LINE + {std::string(attr::networkType), setup_.networkType}, // LCOV_EXCL_LINE + }); + } + + /** + * Build and publish the metrics pipeline (MeterProvider + periodic reader + * + OTLP exporter + histogram view). + * + * Called from the constructor so the provider is published before any + * subsystem creates an instrument. opentelemetry-cpp 1.28.0 has no proxy + * MeterProvider: a meter is a point-in-time copy and is never rebound, so + * an instrument created before this would hold a noop meter for the process + * lifetime. + * + * The reader is attached here too. The SDK notes a reader added later "may + * not receive any in-flight meter data". + * + * Observable instruments are registered later, once the services their + * callbacks read exist. See Collector::onCollectionReady(). + * + * @note Throws whatever the SDK factories throw; the constructor catches. + */ + void + initMetrics() + { + // Derive the metrics endpoint from the trace endpoint by swapping + // the trailing "/v1/traces" path for "/v1/metrics". Any other URL + // shape is used as-is. + std::string metricsEndpoint = setup_.exporterEndpoint; + constexpr std::string_view tracesPath{"/v1/traces"}; + if (metricsEndpoint.ends_with(tracesPath)) + { + metricsEndpoint.replace( + metricsEndpoint.size() - tracesPath.size(), tracesPath.size(), "/v1/metrics"); + } + + // Configure OTLP HTTP metric exporter, honoring the same TLS + // options as the trace exporter. + otlp_http::OtlpHttpMetricExporterOptions metricExporterOpts; + metricExporterOpts.url = metricsEndpoint; + if (setup_.useTls) + { + metricExporterOpts.ssl_ca_cert_path = setup_.tlsCertPath; + metricExporterOpts.ssl_client_cert_path = setup_.tlsClientCertPath; + metricExporterOpts.ssl_client_key_path = setup_.tlsClientKeyPath; + } + + auto metricExporter = otlp_http::OtlpHttpMetricExporterFactory::Create(metricExporterOpts); + + // Configure periodic metric reader (1-second export interval, + // matching the beast OTelCollector path). + metrics_sdk::PeriodicExportingMetricReaderOptions readerOpts; + readerOpts.export_interval_millis = std::chrono::milliseconds(1000); + readerOpts.export_timeout_millis = std::chrono::milliseconds(500); + + auto reader = metrics_sdk::PeriodicExportingMetricReaderFactory::Create( + std::move(metricExporter), readerOpts); + + // Create MeterProvider with the shared resource, then attach reader. + meterProvider_ = metrics_sdk::MeterProviderFactory::Create( + std::make_unique(), makeResource()); + meterProvider_->AddMetricReader(std::move(reader)); + + // Histogram view: SpanMetrics-compatible bucket boundaries (ms) so + // histogram instruments align with the collector's SpanMetrics. + auto histogramSelector = metrics_sdk::InstrumentSelectorFactory::Create( + metrics_sdk::InstrumentType::kHistogram, "*", "ms"); + + // Must match the meter name used by getMeter() and the beast + // OTelCollector, or the view never applies. + auto meterSelector = + metrics_sdk::MeterSelectorFactory::Create(std::string(kMeterName), "", ""); + + auto histogramConfig = std::make_shared(); + histogramConfig->boundaries_ = + std::vector{1.0, 5.0, 10.0, 25.0, 50.0, 100.0, 250.0, 500.0, 1000.0, 5000.0}; + + // An empty view name applies the buckets without renaming. A name here + // would collapse every matching histogram into one series. + auto histogramView = metrics_sdk::ViewFactory::Create( + "", + "SpanMetrics-compatible histogram buckets", + metrics_sdk::AggregationType::kHistogram, + std::move(histogramConfig)); + + meterProvider_->AddView( + std::move(histogramSelector), std::move(meterSelector), std::move(histogramView)); + + // Publish as the global meter provider so developers (and the beast + // OTelCollector shim) reach the same pipeline. + metrics_api::Provider::SetMeterProvider( + opentelemetry::nostd::shared_ptr(meterProvider_)); + } + public: TelemetryImpl(Setup setup, beast::Journal journal) : setup_(std::move(setup)), journal_(journal) { + // Publish the MeterProvider before any subsystem is constructed; see + // initMetrics(). setup_.serviceInstanceId is already resolved by the + // caller, so the resource is complete. + // + // A failure must never stop the node starting: the global provider + // stays noop and every instrument call remains valid. + try + { + initMetrics(); + } + catch (std::exception const& e) + { + JLOG(journal_.error()) << "Telemetry metrics pipeline failed to initialise, " + "continuing without metrics: " + << e.what(); + } } + /** + * Override the service instance id, for callers that learn it late. + * + * Affects only the tracer resource, which start() builds. The metrics + * resource is built by the constructor and is immutable, so supply the id + * through Setup to have it on both. + * + * @param id The instance id to report on spans. + */ void setServiceInstanceId(std::string const& id) override { @@ -361,15 +505,9 @@ public: // marked with kDiscardedAttr (via SpanGuard::discard()). auto processor = std::make_unique(std::move(batchProcessor)); - // Configure resource attributes - auto resourceAttrs = resource::Resource::Create({ - {opentelemetry::semconv::service::kServiceName, setup_.serviceName}, - {opentelemetry::semconv::service::kServiceVersion, setup_.serviceVersion}, - {opentelemetry::semconv::service::kServiceInstanceId, setup_.serviceInstanceId}, - {std::string(attr::networkId), - static_cast(setup_.networkId)}, // LCOV_EXCL_LINE - {std::string(attr::networkType), setup_.networkType}, // LCOV_EXCL_LINE - }); + // Configure resource attributes. Shared with the metrics pipeline the + // constructor already published, so both report one identity. + auto resourceAttrs = makeResource(); // Configure sampler. Head sampling is fixed at 1.0 (sample everything); // setup_.samplingRatio is not config-driven. Wrap the ratio sampler in a @@ -407,69 +545,7 @@ public: trace_api::Provider::SetTracerProvider( opentelemetry::nostd::shared_ptr(sdkProvider_)); - // Build the metrics pipeline, parallel to the tracer above and - // reusing the same resourceAttrs so metrics and traces share one - // resource identity. - - // Derive the metrics endpoint from the trace endpoint by swapping - // the trailing "/v1/traces" path for "/v1/metrics". Any other URL - // shape is used as-is. - std::string metricsEndpoint = setup_.exporterEndpoint; - constexpr std::string_view tracesPath{"/v1/traces"}; - if (metricsEndpoint.ends_with(tracesPath)) - { - metricsEndpoint.replace( - metricsEndpoint.size() - tracesPath.size(), tracesPath.size(), "/v1/metrics"); - } - - // Configure OTLP HTTP metric exporter, honoring the same TLS - // options as the trace exporter. - otlp_http::OtlpHttpMetricExporterOptions metricExporterOpts; - metricExporterOpts.url = metricsEndpoint; - if (setup_.useTls) - { - metricExporterOpts.ssl_ca_cert_path = setup_.tlsCertPath; - metricExporterOpts.ssl_client_cert_path = setup_.tlsClientCertPath; - metricExporterOpts.ssl_client_key_path = setup_.tlsClientKeyPath; - } - - auto metricExporter = otlp_http::OtlpHttpMetricExporterFactory::Create(metricExporterOpts); - - // Configure periodic metric reader (1-second export interval, - // matching the beast OTelCollector path). - metrics_sdk::PeriodicExportingMetricReaderOptions readerOpts; - readerOpts.export_interval_millis = std::chrono::milliseconds(1000); - readerOpts.export_timeout_millis = std::chrono::milliseconds(500); - - auto reader = metrics_sdk::PeriodicExportingMetricReaderFactory::Create( - std::move(metricExporter), readerOpts); - - // Create MeterProvider with the shared resource, then attach reader. - meterProvider_ = metrics_sdk::MeterProviderFactory::Create( - std::make_unique(), resourceAttrs); - meterProvider_->AddMetricReader(std::move(reader)); - - // Histogram view: SpanMetrics-compatible bucket boundaries (ms) so - // histogram instruments align with the collector's SpanMetrics. - auto histogramSelector = metrics_sdk::InstrumentSelectorFactory::Create( - metrics_sdk::InstrumentType::kHistogram, "*", "ms"); - auto meterSelector = metrics_sdk::MeterSelectorFactory::Create("xrpld_metrics", "", ""); - auto histogramConfig = std::make_shared(); - histogramConfig->boundaries_ = - std::vector{1.0, 5.0, 10.0, 25.0, 50.0, 100.0, 250.0, 500.0, 1000.0, 5000.0}; - auto histogramView = metrics_sdk::ViewFactory::Create( - "default_histogram", - "Default histogram view with SpanMetrics-compatible buckets", - metrics_sdk::AggregationType::kHistogram, - std::move(histogramConfig)); - - meterProvider_->AddView( - std::move(histogramSelector), std::move(meterSelector), std::move(histogramView)); - - // Publish as the global meter provider so developers (and the beast - // OTelCollector shim) reach the same pipeline. - metrics_api::Provider::SetMeterProvider( - opentelemetry::nostd::shared_ptr(meterProvider_)); + // The metrics pipeline was built by initMetrics() in the constructor. // Register as the global Telemetry instance so SpanGuard factory // methods can access it without callers passing a reference. @@ -481,10 +557,16 @@ public: void stop() override { + if (stopped_) + return; + stopped_ = true; + JLOG(journal_.info()) << "Telemetry stopping"; - // Unregister global instance before tearing down the pipeline. - Telemetry::setInstance(nullptr); + // Unregister global instance before tearing down the pipeline, but only + // if this object is the one that published it. + if (Telemetry::getInstance() == this) + Telemetry::setInstance(nullptr); if (sdkProvider_) { diff --git a/src/xrpld/app/main/Application.cpp b/src/xrpld/app/main/Application.cpp index 83bbaa8fba..2a184b8fc6 100644 --- a/src/xrpld/app/main/Application.cpp +++ b/src/xrpld/app/main/Application.cpp @@ -316,7 +316,8 @@ public: ApplicationImp( std::unique_ptr config, std::unique_ptr logs, - std::unique_ptr timeKeeper) + std::unique_ptr timeKeeper, + std::optional const& nodePublicKey) : BasicApp(numberOfThreads(*config)) , config_(std::move(config)) , logs_(std::move(logs)) @@ -330,11 +331,15 @@ public: *this, logs_->journal("PerfLog"), [this] { signalStop("PerfLog"); })) + // Telemetry publishes the MeterProvider on construction, so it must + // precede collectorManager_ below and every subsystem that creates an + // instrument. Its resource is immutable, so the instance id has to be + // supplied now; empty means this run reports none. , telemetry_( telemetry::makeTelemetry( telemetry::makeTelemetrySetup( config_->section("telemetry"), - "", // Updated later via setServiceInstanceId() + nodePublicKey.value_or(""), build_info::getVersionString(), config_->networkId), logs_->journal("Telemetry"))) @@ -515,6 +520,33 @@ public: add(ledgerCleaner_.get()); } + /** + * Stop observing and stop telemetry before the members are destroyed. + * + * The metrics reader thread runs callbacks that read the services member + * destruction is about to tear down. telemetry_ is declared early because + * the collector needs its MeterProvider, so reverse-order member destruction + * would take it down last. + * + * run() does both on the normal path; this covers the paths that never + * reach it -- every `return false` in setup(), and the unit tests. Both + * calls are idempotent. + */ + ~ApplicationImp() override + { + // A shutdown diagnostic must never terminate the process, and a + // destructor is implicitly noexcept. + try + { + collectorManager_->collector()->onCollectionStopping(); + telemetry_->stop(); + } + catch (std::exception const& e) + { + JLOG(journal_.error()) << "Error stopping telemetry: " << e.what(); + } + } + //-------------------------------------------------------------------------- bool @@ -1167,9 +1199,8 @@ private: * global Telemetry instance is not yet live, and the first consensus * round runs inside setup(). * - * @pre nodeIdentity_ is populated, so setServiceInstanceId() has - * already supplied the service.instance.id resource attribute - * (the Telemetry resource is fixed once start() builds it). + * The resource attributes, including service.instance.id, were supplied at + * construction. */ void startTelemetry() const; @@ -1271,11 +1302,8 @@ ApplicationImp::setup(boost::program_options::variables_map const& cmdline) nodeIdentity_ = getNodeIdentity(*this, cmdline); - // Now that the node identity is known, inject it into the telemetry - // resource attributes — but only if the user didn't already set a - // custom service_instance_id in [telemetry]. The Telemetry object - // was constructed with an empty serviceInstanceId because - // nodeIdentity_ is not available in the member initializer list. + // The metrics resource was fixed at construction, but the tracer resource is + // built by start() below, so a key minted just now can still reach spans. if (!config_->section("telemetry").exists("service_instance_id")) telemetry_->setServiceInstanceId(toBase58(TokenType::NodePublic, nodeIdentity_->first)); @@ -1452,6 +1480,12 @@ ApplicationImp::setup(boost::program_options::variables_map const& cmdline) collectorManager_->collector()); add(*overlay_); // add to PropertyStream + // Register the collector's observable instruments. Their callbacks run hook + // handlers that read ledgerMaster_, networkOPs_, the peer finder, the job + // queue and overlay_ -- the last of these to be built. Still before the + // first consensus round below, so that round is covered. + collectorManager_->collector()->onCollectionReady(); + // start first consensus round if (!networkOPs_->beginConsensus(ledgerMaster_->getClosedLedger()->header().hash, {})) { @@ -1669,6 +1703,11 @@ ApplicationImp::run() return getValidators().trustedPublisher(pubKey); }); + // Stop observing before any service below is stopped: the collector's gauge + // callbacks run hook handlers that read ledgerMaster_, networkOPs_, the peer + // finder, the job queue and overlay_. Returns once no callback is running. + collectorManager_->collector()->onCollectionStopping(); + // The order of these stop calls is delicate. // Re-ordering them risks undefined behavior. loadManager_->stop(); @@ -2252,9 +2291,19 @@ makeApplication( std::unique_ptr config, std::unique_ptr logs, std::unique_ptr timeKeeper) +{ + return makeApplication(std::move(config), std::move(logs), std::move(timeKeeper), std::nullopt); +} + +std::unique_ptr +makeApplication( + std::unique_ptr config, + std::unique_ptr logs, + std::unique_ptr timeKeeper, + std::optional const& nodePublicKey) { return std::make_unique( - std::move(config), std::move(logs), std::move(timeKeeper)); + std::move(config), std::move(logs), std::move(timeKeeper), nodePublicKey); } void diff --git a/src/xrpld/app/main/Application.h b/src/xrpld/app/main/Application.h index 225275afe4..1d7125cd64 100644 --- a/src/xrpld/app/main/Application.h +++ b/src/xrpld/app/main/Application.h @@ -174,4 +174,19 @@ makeApplication( std::unique_ptr logs, std::unique_ptr timeKeeper); +/** + * Construct the application with a known node public key. + * + * Telemetry builds its resource attributes during construction and they are + * immutable, so the base58 node public key must be supplied here. Pass + * std::nullopt when it is unknown; that run reports no instance id. See + * resolveNodePublicKey(). + */ +std::unique_ptr +makeApplication( + std::unique_ptr config, + std::unique_ptr logs, + std::unique_ptr timeKeeper, + std::optional const& nodePublicKey); + } // namespace xrpl diff --git a/src/xrpld/app/main/Main.cpp b/src/xrpld/app/main/Main.cpp index ba6520db5f..82b959f7bc 100644 --- a/src/xrpld/app/main/Main.cpp +++ b/src/xrpld/app/main/Main.cpp @@ -1,4 +1,5 @@ #include +#include #include #include #include @@ -36,6 +37,7 @@ #include #include #include +#include #include #include #include @@ -804,8 +806,30 @@ run(int argc, char** argv) if (vm.contains("debug")) setDebugLogSink(logs->makeSink("Debug", beast::Severity::Trace)); - auto app = - makeApplication(std::move(config), std::move(logs), std::make_unique()); + // Telemetry needs the node public key at construction, so read it here + // where a config error can still be reported and the process can exit + // cleanly. getNodeIdentity() in setup() stays authoritative. + std::optional nodePublicKey; + try + { + nodePublicKey = resolveNodePublicKey(*config, vm, logs->journal("Application")); + } + catch (std::exception const& e) + { + std::cerr << "Unable to start " << systemName() << ": " << e.what() << std::endl; + return -1; + } + + if (!nodePublicKey) + { + JLOG(logs->journal("Application").warn()) + << "Telemetry: no node identity available yet, so this run reports an empty " + "service.instance.id. Set [telemetry] service_instance_id, or restart once " + "the node key exists."; + } + + auto app = makeApplication( + std::move(config), std::move(logs), std::make_unique(), nodePublicKey); if (!app->setup(vm)) return -1; diff --git a/src/xrpld/app/main/NodeIdentity.cpp b/src/xrpld/app/main/NodeIdentity.cpp index 8198c43af7..fcf460f185 100644 --- a/src/xrpld/app/main/NodeIdentity.cpp +++ b/src/xrpld/app/main/NodeIdentity.cpp @@ -8,13 +8,18 @@ #include #include #include +#include +#include #include #include +#include +#include #include #include #include +#include #include namespace xrpl { @@ -58,4 +63,82 @@ getNodeIdentity(Application& app, boost::program_options::variables_map const& c return getNodeIdentity(*db); } +std::optional +resolveNodePublicKey( + Config const& config, + boost::program_options::variables_map const& cmdline, + beast::Journal journal) +{ + std::optional seed; + bool seedConfigured = false; + + if (cmdline.contains("nodeid")) + { + seedConfigured = true; + seed = parseGenericSeed(cmdline["nodeid"].as(), false); + } + else if (config.exists(Sections::kNodeSeed)) + { + seedConfigured = true; + if (auto const& lines = config.section(Sections::kNodeSeed).lines(); !lines.empty()) + seed = parseBase58(lines.front()); + } + + // A configured seed decides the identity outright. A malformed or missing + // one is reported by getNodeIdentity(), which runs later. + if (seedConfigured) + { + if (!seed) + return std::nullopt; + + auto const secretKey = generateSecretKey(KeyType::Secp256k1, *seed); + return toBase58(TokenType::NodePublic, derivePublicKey(KeyType::Secp256k1, secretKey)); + } + + // --newnodeid discards whatever is stored. + if (cmdline.contains("newnodeid")) + return std::nullopt; + + try + { + auto setup = setupDatabaseCon(config, journal); + + // Standalone uses a temporary database, so nothing is persisted and this + // run will mint a fresh key. + if (setup.standAlone && setup.startUp != StartUpType::Load && + setup.startUp != StartUpType::LoadFile && setup.startUp != StartUpType::Replay) + { + return std::nullopt; + } + + // The global pragmas include journal_mode, which rewrites the database + // header. The wallet is opened without them everywhere else. + setup.useGlobalPragma = false; + + // Only read an existing file: SQLite would otherwise create one. + if (std::error_code ec; !std::filesystem::exists(setup.dataDir / kWalletDbName, ec)) + { + return std::nullopt; + } + + // Empty init SQL: open the existing schema, never create it. + DatabaseCon walletDb{ + setup, + kWalletDbName, + std::array{}, + std::array{}, + journal}; + + auto db = walletDb.checkoutDb(); + if (auto const stored = readNodeIdentity(*db)) + return toBase58(TokenType::NodePublic, stored->first); + } + catch (std::exception const& e) + { + JLOG(journal.warn()) << "Could not read the node identity: " << e.what(); + } + + return std::nullopt; +} + } // namespace xrpl diff --git a/src/xrpld/app/main/NodeIdentity.h b/src/xrpld/app/main/NodeIdentity.h index 117acffdb1..7309f6007a 100644 --- a/src/xrpld/app/main/NodeIdentity.h +++ b/src/xrpld/app/main/NodeIdentity.h @@ -1,12 +1,16 @@ #pragma once #include +#include +#include #include #include #include +#include +#include #include namespace xrpl { @@ -20,4 +24,26 @@ namespace xrpl { std::pair getNodeIdentity(Application& app, boost::program_options::variables_map const& cmdline); +/** + * This server's public key, read without creating or modifying anything. + * + * For callers that need the identity before the Application exists, such as + * telemetry building its resource attributes in the member-init list. Derives + * from a configured seed when there is one, otherwise reads the wallet database + * only if it already exists. + * + * getNodeIdentity() remains authoritative and mints a key when none exists. + * + * @param config The server configuration. + * @param cmdline The command line parameters passed into the application. + * @param journal Journal for reporting an unreadable database. + * @return The base58-encoded node public key, or std::nullopt if none can be + * read. + */ +std::optional +resolveNodePublicKey( + Config const& config, + boost::program_options::variables_map const& cmdline, + beast::Journal journal); + } // namespace xrpl