diff --git a/.github/scripts/telemetry/check_regression_bounds.py b/.github/scripts/telemetry/check_regression_bounds.py index 5454d5dd29..b07045ce70 100644 --- a/.github/scripts/telemetry/check_regression_bounds.py +++ b/.github/scripts/telemetry/check_regression_bounds.py @@ -43,10 +43,12 @@ baseline's own. Six rules are checked: declare, carries a reason, and has neither a threshold override nor a baseline value left behind. -Rule A subtracts ``excluded_keys`` before comparing, so a quantile removed from -the gated set does not read as a missing baseline. Rule F is what keeps that -subtraction honest: an exclusion is the one edit here that makes the gate cover -LESS, so a stale or misspelt entry must fail rather than silently widen itself. +Rule A subtracts ``excluded_keys`` from BOTH sides of its comparison, so a +quantile removed from the gated set neither reads as a missing baseline nor as +an undeclared one -- an exclusion left in the baseline is one failure, rule F's, +which names the file to edit. Rule F is what keeps that subtraction honest: an +exclusion is the one edit here that makes the gate cover LESS, so a stale or +misspelt entry must fail rather than silently widen itself. A PLACEHOLDER baseline -- ``"placeholder": true`` or an empty ``metrics`` object -- exits 0, because that is the documented bootstrap state and CI has to @@ -56,8 +58,9 @@ without having checked anything is the same green-build-that-is-not failure this script exists to prevent, so renaming or deleting one of its inputs must not silence it. -A baseline ENTRY that is not a positive finite number is rejected before any -rule runs, by ``_unusable_baseline``. Every rule does arithmetic on that value, +A baseline ENTRY that is not an object, or whose value is not a positive finite +number, is rejected before any rule runs -- see ``check_key`` and +``_unusable_baseline``. Every rule does arithmetic on that value, and a degenerate one made the script crash with a traceback (rule D divides by it) or emit advice about the wrong file (a negative value inverts rule D's comparison). Reporting malformed input is what this script is for, so it must @@ -287,7 +290,20 @@ def check_key(key, entry, thresholds, ladders): malformed input, and reporting malformed input is this script's whole job, so it is rejected up front rather than arithmetic being attempted on it -- see ``_unusable_baseline``. + + The entry's SHAPE is checked first, for the same reason. A hand edit that + writes the bare number instead of the ``{"value": .., "unit": ..}`` object + leaves no ``.get`` to call, and the script died with an AttributeError + traceback naming a line in itself rather than the key at fault. """ + if not isinstance(entry, dict): + return [ + f"{key}: baseline entry {entry!r} is not an object carrying value and " + f"unit, so no bound can be derived from it. Recapture the baseline " + f"from a CI run rather than editing it by hand -- see " + f"baselines/README.md" + ] + value, unit = entry.get("value"), entry.get("unit", "") edges = ladders.get(unit) if value is None or edges is None: @@ -381,8 +397,16 @@ def main(): failures.extend(check_exclusions(metrics_cfg, thresholds, gated, declared)) # Rule A compares against the GATED surface, so a deliberately excluded # quantile is not reported as a baseline that was never captured. - declared -= set(metrics_cfg.get("excluded_keys", {})) - for key in sorted(set(gated) - declared): + # + # The same keys come off rule A's over-coverage side too. An excluded key + # left in the baseline is rule F's finding, reported with the file to edit; + # rule A would add a second failure for the same single mistake, saying the + # key is not declared -- which is not even true, it is declared and then + # excluded. Only exclusions the surface really declares are subtracted, so a + # misspelt exclusion naming a stale baseline key still reaches rule A. + excluded_declared = set(metrics_cfg.get("excluded_keys", {})) & declared + declared -= excluded_declared + for key in sorted(set(gated) - declared - excluded_declared): failures.append( f"{key}: in the baseline but not declared by {METRICS}, so it is " f"reported every run and can never gate -- remove it (rule A)" diff --git a/.github/scripts/telemetry/test_check_regression_bounds.py b/.github/scripts/telemetry/test_check_regression_bounds.py index ba3da1848c..5e7a4b7811 100644 --- a/.github/scripts/telemetry/test_check_regression_bounds.py +++ b/.github/scripts/telemetry/test_check_regression_bounds.py @@ -46,7 +46,12 @@ class CheckerCase(unittest.TestCase): def setUp(self): self.tree = Path(tempfile.mkdtemp()) - self.addCleanup(self._cleanup) + # Bound to THIS tree, not read off self.tree when the cleanup finally + # runs. A test that calls setUp again for a fresh tree (see the + # degenerate-baseline subTests) rebinds self.tree, and a late read would + # make every registered cleanup remove the LAST tree, leaving each + # earlier one behind in /tmp. + self.addCleanup(self._cleanup, self.tree) for rel in INPUTS: dest = self.tree / rel dest.parent.mkdir(parents=True, exist_ok=True) @@ -55,11 +60,12 @@ class CheckerCase(unittest.TestCase): script.parent.mkdir(parents=True, exist_ok=True) shutil.copy(CHECKER, script) - def _cleanup(self): - for path in self.tree.rglob("*"): + def _cleanup(self, tree): + """Remove one scratch tree, restoring permissions rmtree needs first.""" + for path in tree.rglob("*"): if path.is_file(): path.chmod(stat.S_IRUSR | stat.S_IWUSR) - shutil.rmtree(self.tree, ignore_errors=True) + shutil.rmtree(tree, ignore_errors=True) def run_checker(self): """Run the checker in the scratch tree, returning (code, stdout+stderr).""" diff --git a/OpenTelemetryPlan/06-implementation-phases.md b/OpenTelemetryPlan/06-implementation-phases.md index 9be3127c15..3f606c717b 100644 --- a/OpenTelemetryPlan/06-implementation-phases.md +++ b/OpenTelemetryPlan/06-implementation-phases.md @@ -570,7 +570,7 @@ graph LR # [insight] section — new "otel" server option [insight] server=otel # NEW: uses OTel OTLP metrics exporter -prefix=xrpld # metric name prefix (preserved) +# No prefix: it applies on the StatsD path only, not this one. # Endpoint and auth inherited from [telemetry] section: [telemetry] diff --git a/docker/telemetry/workload/capture_timings.py b/docker/telemetry/workload/capture_timings.py index 3cd4bb8cd2..9c90ac5f30 100644 --- a/docker/telemetry/workload/capture_timings.py +++ b/docker/telemetry/workload/capture_timings.py @@ -122,10 +122,17 @@ def _capture_status(metrics: dict, min_ratio: float) -> dict: keys on, and it is the same predicate that decides this script's exit code — see the module docstring for why it lives in the artifact. - An empty surface is vacuously complete: there is nothing for the capture - to have fallen short of, and that is the case the exit-code check has - always passed. Defining it any other way here would make the flag and the - exit code disagree, which is the drift this block exists to remove. + An empty surface is NOT complete. ``declared == 0`` means the capture asked + Prometheus for nothing: ``build_query_plan`` returns an empty plan, without + complaining, for any config that yields no ``spans``/``rpc_methods``/ + ``job_queue`` entries -- a ``--metrics`` path pointing at the wrong file, a + truncated one, or every key excluded. Nothing about that run is evidence + the pipeline works, so treating it as vacuously complete would exit 0 and + hand the paste-me path a ``metrics: {}`` artifact to offer as baseline + material. Pasted in, it still reads as a placeholder, so the gate stays off + while the workflow reports the baseline as activated -- the silent-green + outcome the whole ``capture`` block exists to prevent. The exit code reads + this same flag, so the two still cannot disagree. """ declared = len(metrics) captured = sum(1 for entry in metrics.values() if entry["value"] is not None) @@ -133,7 +140,7 @@ def _capture_status(metrics: dict, min_ratio: float) -> dict: "declared": declared, "captured": captured, "min_ratio": min_ratio, - "complete": declared == 0 or (captured / declared) >= min_ratio, + "complete": declared > 0 and (captured / declared) >= min_ratio, } @@ -236,16 +243,28 @@ def main() -> int: logger.info("Wrote %s (%d/%d metrics captured)", args.output, captured, total) if not status["complete"]: - logger.error( - "Only %d/%d (%.0f%%) metrics captured — below the %.0f%% minimum. " - "Is Prometheus reachable at %s? The file is marked " - "capture.complete=false and must not be pasted into the baseline.", - captured, - total, - captured / total * 100, - args.min_capture_ratio * 100, - args.prometheus, - ) + if total == 0: + # No ratio to report: nothing was asked for, so the shortfall is the + # declared surface, not Prometheus. Named separately because the + # percentage below would divide by zero. + logger.error( + "No metrics were declared, so nothing was captured. Does %s " + "declare spans/rpc_methods/job_queue names, and does " + "excluded_keys leave any of them gated? The file is marked " + "capture.complete=false and must not be pasted into the baseline.", + args.metrics, + ) + else: + logger.error( + "Only %d/%d (%.0f%%) metrics captured — below the %.0f%% minimum. " + "Is Prometheus reachable at %s? The file is marked " + "capture.complete=false and must not be pasted into the baseline.", + captured, + total, + captured / total * 100, + args.min_capture_ratio * 100, + args.prometheus, + ) return 1 return 0 diff --git a/docker/telemetry/workload/test_validate_telemetry.py b/docker/telemetry/workload/test_validate_telemetry.py index 514a42b1df..84358e92a8 100644 --- a/docker/telemetry/workload/test_validate_telemetry.py +++ b/docker/telemetry/workload/test_validate_telemetry.py @@ -214,6 +214,52 @@ def test_wildcard_child_matches_any_family_member() -> None: assert report.results[0].passed, report.results[0].message +def test_wildcard_predicate_carries_no_backslash_escape() -> None: + """The TraceQL predicate must not contain a backslash escape. + + Tempo's string lexer rejects `\\.` outright -- run 33062418036 returned + HTTP 400, "invalid TraceQL query: parse error at line 1, col 68: invalid char + escape", on the predicate re.escape produced. Asserting the absence of a + backslash rather than a specific spelling keeps this test about the property + the lexer enforces instead of about one way of satisfying it. + + Note the earlier stub could not have caught this: it evaluated the pattern + with Python's re, which accepts `\\.` happily, so it modelled the regex engine + rather than the query lexer in front of it. + """ + predicate = vt._traceql_name_predicate("rpc.command.*") + assert "\\" not in predicate, f"backslash escape reaches Tempo: {predicate}" + + +def test_wildcard_predicate_matches_the_family_but_not_near_misses() -> None: + """The pattern must still mean what the glob meant. + + Dropping the escaping must not be done by making the dots match any + character: `rpc.command.*` should accept rpc.command.fee and reject a name + that differs in the separator positions, which is the looseness + _span_name_matches exists to avoid. + """ + import re as _re + + pattern = vt._traceql_name_predicate("rpc.command.*").split('"')[1] + assert _re.fullmatch(pattern, "rpc.command.fee") + assert _re.fullmatch(pattern, "rpc.command.server_info") + # Separators deliberately not dots: if the pattern left its dots bare they + # would match these too. Colons rather than a made-up letter so the spell + # checker still sees three real words. + assert not _re.fullmatch(pattern, "rpc:command:fee") + assert not _re.fullmatch(pattern, "other.command.fee") + + +def test_literal_predicate_uses_equality() -> None: + """A non-glob child must use `=`, not a regex. + + Equality is what makes a longer emitted name unable to satisfy a shorter + contract, the same guarantee _span_name_matches gives on the client side. + """ + assert vt._traceql_name_predicate("txq.accept_tx") == 'name="txq.accept_tx"' + + def main() -> int: tests = [v for k, v in sorted(globals().items()) if k.startswith("test_")] failed = 0 diff --git a/docker/telemetry/workload/validate_telemetry.py b/docker/telemetry/workload/validate_telemetry.py index 4733af8dbc..b535109897 100644 --- a/docker/telemetry/workload/validate_telemetry.py +++ b/docker/telemetry/workload/validate_telemetry.py @@ -97,6 +97,21 @@ SYNC_DIAGNOSTICS_GROUP = "sync_diagnostics" # not hammer the single-container Prometheus the harness runs. METRIC_POLL_CONCURRENCY = 8 +# Bound on ONE HTTP request to Tempo, Prometheus, Loki or Grafana. aiohttp's +# own default is total=300s, which is longer than any poll window here: a +# single wedged endpoint would blow the shared deadline above and then keep the +# run alive until the CI job's own budget killed it, losing the report and the +# artifacts with it. Failing one request fast and reporting it beats being +# killed with nothing. +# +# Derived from the poll window rather than picked: no single request may outlast +# the phase budget it sits inside, since a request that does can only ever blow +# that deadline. Connecting is held to one poll interval, because an endpoint +# that is absent or wedged should be named immediately, not waited on. +REQUEST_TIMEOUT = aiohttp.ClientTimeout( + total=METRIC_POLL_TIMEOUT_SEC, sock_connect=METRIC_POLL_INTERVAL_SEC +) + # The Prometheus exporter splits one histogram instrument into three series # names. Reverse coverage folds them back onto the base family so a contract # entry (or an accounted_patterns regex) written for the family accounts for @@ -421,11 +436,18 @@ def _traceql_name_predicate(expected_name: str) -> str: """Build the TraceQL `name` predicate that selects a contract span name. A literal contract name becomes an equality test. A glob becomes a regex - test, because TraceQL has no glob operator: `rpc.command.*` must be sent as - `name=~"rpc\\.command\\..*"`, with the dots escaped so they match literal - dots rather than any character. Sending the glob unescaped would still match - the intended spans, but would also match names differing in those positions, - which is the looseness _span_name_matches exists to avoid. + test, because TraceQL has no glob operator: `rpc.command.*` is sent as + `name=~"rpc[.]command[.].*"`. + + A literal dot is written as the character class `[.]` rather than as `\\.`, + and that is not a style choice. TraceQL's string lexer rejects a backslash + escape it does not recognise, so the re.escape spelling this replaced -- + `name=~"rpc\\.command\\..*"` -- came back as HTTP 400, "invalid TraceQL + query: parse error at line 1, col 68: invalid char escape". `[.]` carries no + backslash, so nothing reaches the lexer that it can refuse, while still + meaning a literal dot to the regex engine behind it. Leaving the dots bare + would parse but match any character in those positions, which is the + looseness _span_name_matches exists to avoid. Args: expected_name: Span name or glob from expected_spans.json. @@ -435,7 +457,18 @@ def _traceql_name_predicate(expected_name: str) -> str: """ if "*" not in expected_name: return f'name="{expected_name}"' - pattern = "".join(".*" if ch == "*" else re.escape(ch) for ch in expected_name) + # Span names are lower_snake_case segments joined by dots, so `.` and `*` are + # the only characters here that mean anything to a regex engine. Anything + # else appearing would need its own handling rather than silent passthrough. + unexpected = set(expected_name) - set("abcdefghijklmnopqrstuvwxyz0123456789_.*") + if unexpected: + raise ValueError( + f"span name {expected_name!r} contains {sorted(unexpected)}, which " + "this predicate builder does not know how to escape for TraceQL" + ) + pattern = "".join( + ".*" if ch == "*" else "[.]" if ch == "." else ch for ch in expected_name + ) return f'name=~"{pattern}"' @@ -717,6 +750,25 @@ async def _check_attributes_on_first_trace( try: trace_id = traces[0].get("traceID", "") if not trace_id: + # Recorded rather than returned on silently. Returning with no + # result would drop this span's attribute contract out of the + # report and shrink the check total, so the surface would look + # smaller with nothing saying why. Reported the same way as a + # fetched trace holding no matching span, below: being unable to + # verify is itself the finding. The caller only reaches here for a + # span that declares required_attributes, so this adds no check + # where none was expected. + report.add( + CheckResult( + name=f"span.attrs.{span_name}", + category="span", + passed=False, + message=( + f"{span_name}: newest trace carried no traceID, cannot " + "verify its attributes" + ), + ) + ) return spans = await _tempo_get_trace(session, tempo_url, trace_id) await _validate_span_attributes_otlp(spans, span_def, report) @@ -2519,7 +2571,7 @@ async def run_validation( report = ValidationReport() report.start_time = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()) - async with aiohttp.ClientSession() as session: + async with aiohttp.ClientSession(timeout=REQUEST_TIMEOUT) as session: await validate_spans(session, tempo_url, report) await validate_span_durations(session, tempo_url, report) await assert_trace_join_groups(session, tempo_url, report) diff --git a/docker/telemetry/xrpld-telemetry-mainnet.cfg b/docker/telemetry/xrpld-telemetry-mainnet.cfg index 720c9a1b6b..88a4a09fb9 100644 --- a/docker/telemetry/xrpld-telemetry-mainnet.cfg +++ b/docker/telemetry/xrpld-telemetry-mainnet.cfg @@ -135,15 +135,13 @@ data/logs/mainnet/debug.log # --- Insight (native OTel metrics via beast::insight) ----------------------- +# server is the only key that changes behaviour here. No prefix: formatName() +# ignores it, so names stay bare (jobq_job_count). service_instance_id is read +# and discarded (OTelCollector.cpp `(void)instanceId`); the label Prometheus +# shows comes from [telemetry] service_instance_id below. [insight] server=otel endpoint=http://localhost:4318/v1/metrics -prefix=xrpld -# Sets the OTel service.instance.id resource attribute, which Prometheus -# exposes as the `service_instance_id` label. Dashboards filter on it via the -# $node template variable, so without this every insight-backed panel is -# empty. Matches [telemetry] service_instance_id for a single node identity. -service_instance_id=xrpld-mainnet # --- OpenTelemetry tracing -------------------------------------------------- diff --git a/docs/telemetry-runbook.md b/docs/telemetry-runbook.md index f54664a2bf..220044a652 100644 --- a/docs/telemetry-runbook.md +++ b/docs/telemetry-runbook.md @@ -1631,12 +1631,13 @@ Add to `xrpld.cfg`: [insight] server=otel endpoint=http://localhost:4318/v1/metrics -prefix=xrpld ``` The `OTelCollector` implementation exports metrics via OTLP/HTTP to the same OTel Collector that receives traces. No separate StatsD receiver is needed. -> **Fallback**: Set `server=statsd` and `address=127.0.0.1:8125` to use the legacy StatsD UDP path. This requires re-enabling the `statsd` receiver in `otel-collector-config.yaml` and uncommenting port 8125 in `docker-compose.yml`. +Do not set `prefix` on this path. `formatName()` never applies it, so the setting is silently ignored and the exported names are bare and lowercase — `jobq_job_count`, not `xrpld_jobq_job_count`. Queries written against a prefixed name return no series. + +> **Fallback**: Set `server=statsd` and `address=127.0.0.1:8125` to use the legacy StatsD UDP path. This requires re-enabling the `statsd` receiver in `otel-collector-config.yaml` and uncommenting port 8125 in `docker-compose.yml`. On that path `prefix` **is** applied to the metric name, which is why the StatsD examples elsewhere in this document keep it. ### Metric Reference diff --git a/include/xrpl/consensus/Consensus.h b/include/xrpl/consensus/Consensus.h index f299ce2b9a..78750656a4 100644 --- a/include/xrpl/consensus/Consensus.h +++ b/include/xrpl/consensus/Consensus.h @@ -1758,16 +1758,23 @@ Consensus::updateOurPositions(std::unique_ptr const& mutableSet->erase(txId); } - auto const yaysStr = std::to_string(dispute.getYays()); - auto const naysStr = std::to_string(dispute.getNays()); - span.addEvent( - consensus::span::event::disputeResolve, - {{consensus::span::attr::txId, to_string(txId)}, - {consensus::span::attr::disputeOurVote, - dispute.getOurVote() ? std::string_view{consensus::span::val::yes} - : std::string_view{consensus::span::val::no}}, - {consensus::span::attr::disputeYays, yaysStr}, - {consensus::span::attr::disputeNays, naysStr}}); + // The event exists only for the span, so it is guarded on the + // span being active. Unguarded, every dispute that flips + // position builds a 64-character tx hash plus two number + // strings, on every establish tick. + if (span) + { + auto const yaysStr = std::to_string(dispute.getYays()); + auto const naysStr = std::to_string(dispute.getNays()); + span.addEvent( + consensus::span::event::disputeResolve, + {{consensus::span::attr::txId, to_string(txId)}, + {consensus::span::attr::disputeOurVote, + dispute.getOurVote() ? std::string_view{consensus::span::val::yes} + : std::string_view{consensus::span::val::no}}, + {consensus::span::attr::disputeYays, yaysStr}, + {consensus::span::attr::disputeNays, naysStr}}); + } } } diff --git a/include/xrpl/telemetry/Recording.h b/include/xrpl/telemetry/Recording.h new file mode 100644 index 0000000000..ecb2d9767f --- /dev/null +++ b/include/xrpl/telemetry/Recording.h @@ -0,0 +1,205 @@ +#pragma once + +/** + * Utilities for state and work that exist only to be recorded. + * + * Each type below holds real state when telemetry is compiled in and is an + * empty type with no-op methods when it is not. The member is declared in + * both configurations, so a class's member set and public API never differ + * between builds -- a difference that has previously made a test mock + * abstract. The compiled-out forms are empty types, so such a member costs a + * byte of padding rather than nothing. `[[no_unique_address]]` would remove + * even that, but MSVC ignores the standard spelling for ABI compatibility, so + * it is deliberately not used. What these types buy is work not being done, + * not a smaller struct. + * + * kEnabled ---- if constexpr ---- telemetry-only blocks + * | + * +-- Stopwatch (a clock read nobody reads when off) + * +-- Counter (an atomic nobody reads when off) + * +-- Mirror (a value kept only to be reported) + * + * @note A no-op method does NOT skip evaluation of its arguments. + * `counter.add(expensiveCount())` still calls `expensiveCount()` when + * telemetry is compiled out. Pass cheap values only; put expensive work + * inside `if constexpr (kEnabled)`. + * + * @note `if constexpr (kEnabled)` still type-checks its discarded branch in + * non-template code, so use it only where the block names no + * `opentelemetry::` type. That is the normal case, because SpanGuard exists + * to keep those types out of call sites. + * + * @note Thread safety: `Counter` is safe to update from any thread. + * `Stopwatch` and `Mirror` are not synchronized; guard them the same way + * you guard the state they sit beside. + * + * Usage: + * @code + * // Time a loop without a single #ifdef. + * telemetry::Stopwatch const timer; + * for (auto const& obj : objects) + * lookUp(obj); + * recordLookupMetrics(timer.elapsedUs()); // 0 when compiled out + * @endcode + * + * @code + * // A counter that disappears, along with its storage, when off. + * class Acquirer + * { + * telemetry::Counter<> timeouts_; + * public: + * void onTimeout() { timeouts_.add(); } + * }; + * @endcode + * + * @code + * // Edge case -- an expensive value still needs a block guard, + * // because arguments are evaluated even when the method is a no-op. + * if constexpr (telemetry::kEnabled) + * span.setAttribute( + * pathfind_span::attr::sourceAccount, redactAccount(account)); + * @endcode + */ + +#include +#include + +// Counter names std::atomic only when telemetry is compiled in, so guarding +// the include keeps misc-include-cleaner from seeing an unused one. +#ifdef XRPL_ENABLE_TELEMETRY +#include +#endif + +namespace xrpl::telemetry { + +#ifdef XRPL_ENABLE_TELEMETRY +/** + * True when telemetry code is compiled into this build. + */ +inline constexpr bool kEnabled = true; +#else +inline constexpr bool kEnabled = false; +#endif + +/** + * A monotonic elapsed-time measurement taken only for telemetry. + * + * Reads the clock on construction when telemetry is compiled in, and does + * nothing at all when it is not, so an untraced build performs no clock + * read on the measured path. + */ +class Stopwatch +{ +#ifdef XRPL_ENABLE_TELEMETRY + /** + * When the measurement started. + */ + std::chrono::steady_clock::time_point start_{std::chrono::steady_clock::now()}; +#endif + +public: + // These read start_ when telemetry is compiled in and touch no member + // when it is not, so clang-tidy asks for them to be static. Making them + // static would give the two builds different signatures. + // NOLINTBEGIN(readability-convert-member-functions-to-static) + + /** + * Begin the measurement again from now. + */ + void + restart() noexcept + { +#ifdef XRPL_ENABLE_TELEMETRY + start_ = std::chrono::steady_clock::now(); +#endif + } + + /** + * Microseconds since construction or the last restart(); 0 when off. + */ + [[nodiscard]] std::chrono::microseconds + elapsedUs() const noexcept + { +#ifdef XRPL_ENABLE_TELEMETRY + return std::chrono::duration_cast( + std::chrono::steady_clock::now() - start_); +#else + return std::chrono::microseconds{0}; +#endif + } + + // NOLINTEND(readability-convert-member-functions-to-static) +}; + +/** + * A monotonically increasing count kept only to be reported. + * + * @tparam T The counter's value type; defaults to std::uint64_t. + */ +template +class Counter +{ +#ifdef XRPL_ENABLE_TELEMETRY + /** + * The running count. Relaxed: only ever read for reporting. + */ + std::atomic value_{0}; +#endif + +public: + /** + * Default-constructed at zero. + */ + Counter() = default; + + // Copying and moving are deleted so that a class holding a Counter has the + // same copy semantics in both builds. std::atomic deletes all four + // implicitly when telemetry is compiled in; without these declarations the + // compiled-out Counter would be an empty, freely copyable type, and its + // owner would silently become copyable in that build only. + Counter(Counter const&) = delete; + Counter& + operator=(Counter const&) = delete; + Counter(Counter&&) = delete; + Counter& + operator=(Counter&&) = delete; + + // These read value_ when telemetry is compiled in and touch no member + // when it is not, so clang-tidy asks for them to be static. Making them + // static would give the two builds different signatures. + // NOLINTBEGIN(readability-convert-member-functions-to-static) + + /** + * Add to the count. A no-op, with no storage, when off. + * + * @param n How much to add; defaults to 1. + */ + void + add(T const n = 1) noexcept + { +#ifdef XRPL_ENABLE_TELEMETRY + value_.fetch_add(n, std::memory_order_relaxed); +#else + (void)n; +#endif + } + + /** + * The current count. + * + * @return The accumulated count, or T{} when telemetry is compiled out. + */ + [[nodiscard]] T + load() const noexcept + { +#ifdef XRPL_ENABLE_TELEMETRY + return value_.load(std::memory_order_relaxed); +#else + return T{}; +#endif + } + + // NOLINTEND(readability-convert-member-functions-to-static) +}; + +} // namespace xrpl::telemetry diff --git a/include/xrpl/telemetry/SpanGuard.h b/include/xrpl/telemetry/SpanGuard.h index 4758602b74..95c471f677 100644 --- a/include/xrpl/telemetry/SpanGuard.h +++ b/include/xrpl/telemetry/SpanGuard.h @@ -179,10 +179,16 @@ #include #include #include -#include #include +#include #include +#ifdef XRPL_ENABLE_TELEMETRY +// The smart-pointer members all belong to the telemetry-enabled declarations; +// the compiled-out types hold nothing, so this include is unused there. +#include +#endif + namespace protocol { class TraceContext; } // namespace protocol @@ -476,6 +482,23 @@ public: [[nodiscard]] TraceBytes getTraceBytes() const; + /** + * Report whether this thread has a live span context to propagate. + * + * Lets a caller decide whether to create an optional protobuf + * submessage before calling injectCurrentContextToProtobuf(), which + * writes nothing when no span is active. Tests exactly the two + * conditions that injector checks and allocates nothing, so it is cheap + * enough to call on every broadcast. threadLocalContext() is not + * a substitute: it builds a shared SpanContext::Impl before validity + * can be tested. + * + * @return true when a span with a valid context is active on this + * thread. + */ + [[nodiscard]] static bool + hasCurrentContext() noexcept; + /** * Inject the calling thread's currently-active OTel context into a * protobuf TraceContext message for cross-node propagation. @@ -962,7 +985,17 @@ class ScopedActivation { public: ScopedActivation() = default; - ~ScopedActivation() = default; + /** + * Written out by hand rather than defaulted, on purpose. A defaulted + * destructor on an empty class is trivial, and compilers then report every + * activation that is held only for its scope as an unused variable. Writing + * the destructor by hand matches the real ScopedActivation and keeps those + * call sites warning free. The body is empty, so it costs nothing once + * inlined. + */ + ~ScopedActivation() // NOLINT(modernize-use-equals-default) + { + } ScopedActivation(ScopedActivation&&) = delete; ScopedActivation& operator=(ScopedActivation&&) = delete; @@ -975,7 +1008,15 @@ class SpanGuard { public: SpanGuard() = default; - ~SpanGuard() = default; + /** + * Written out by hand rather than defaulted, for the same reason as + * ScopedActivation above: a trivial destructor makes a guard that is held + * only for its scope look like an unused variable. The empty body costs + * nothing once inlined. + */ + ~SpanGuard() // NOLINT(modernize-use-equals-default) + { + } SpanGuard(SpanGuard&&) noexcept = default; SpanGuard& operator=(SpanGuard&&) noexcept = default; @@ -1051,6 +1092,12 @@ public: return {}; } + [[nodiscard]] static bool + hasCurrentContext() noexcept + { + return false; + } + static void injectCurrentContextToProtobuf(protocol::TraceContext&) { @@ -1135,7 +1182,15 @@ public: ScopedSpanGuard(TraceCategory, std::string_view, std::string_view) noexcept { } - ~ScopedSpanGuard() = default; + /** + * Written out by hand rather than defaulted, for the same reason as + * ScopedActivation above: a trivial destructor makes a guard that is held + * only for its scope look like an unused variable. The empty body costs + * nothing once inlined. + */ + ~ScopedSpanGuard() // NOLINT(modernize-use-equals-default) + { + } ScopedSpanGuard(ScopedSpanGuard&&) = delete; ScopedSpanGuard& @@ -1251,4 +1306,21 @@ activateIfLive(SpanGuardHandle const& guard) #endif // XRPL_ENABLE_TELEMETRY +// These three types are held purely for their scope: callers create one and +// never read it again. A compiler only stays quiet about such a variable if +// destroying it might do something, which means the destructor must not be +// trivial. Both the real types and the compiled-out ones therefore declare a +// destructor by hand. Asserting it here fails the build immediately if one is +// ever replaced with `= default`, instead of producing an unused-variable +// error at every call site. +static_assert( + !std::is_trivially_destructible_v, + "SpanGuard must keep a hand-written destructor; see the note above"); +static_assert( + !std::is_trivially_destructible_v, + "ScopedSpanGuard must keep a hand-written destructor; see the note above"); +static_assert( + !std::is_trivially_destructible_v, + "ScopedActivation must keep a hand-written destructor; see the note above"); + } // namespace xrpl::telemetry diff --git a/include/xrpl/telemetry/Telemetry.h b/include/xrpl/telemetry/Telemetry.h index 6d13e9c8c9..bb95298c98 100644 --- a/include/xrpl/telemetry/Telemetry.h +++ b/include/xrpl/telemetry/Telemetry.h @@ -96,7 +96,6 @@ #include #include #include -#include #ifdef XRPL_ENABLE_TELEMETRY #include @@ -105,6 +104,9 @@ #include #include #include + +// std::string_view appears only in the telemetry-enabled declarations below. +#include #endif namespace xrpl::telemetry { diff --git a/src/libxrpl/basics/Log.cpp b/src/libxrpl/basics/Log.cpp index 750cdc0e8a..e49301cfdd 100644 --- a/src/libxrpl/basics/Log.cpp +++ b/src/libxrpl/basics/Log.cpp @@ -16,7 +16,6 @@ #endif // XRPL_ENABLE_TELEMETRY #include -#include #include #include #include @@ -29,6 +28,11 @@ #include #include +#ifdef XRPL_ENABLE_TELEMETRY +// std::size_t names the hex widths used when formatting a trace context. +#include +#endif // XRPL_ENABLE_TELEMETRY + namespace xrpl { Logs::Sink::Sink(std::string partition, beast::Severity thresh, Logs& logs) diff --git a/src/libxrpl/nodestore/backend/NuDBFactory.cpp b/src/libxrpl/nodestore/backend/NuDBFactory.cpp index 69173782b8..f9d56576a5 100644 --- a/src/libxrpl/nodestore/backend/NuDBFactory.cpp +++ b/src/libxrpl/nodestore/backend/NuDBFactory.cpp @@ -66,6 +66,7 @@ public: std::atomic deletePath; Scheduler& scheduler; +#ifdef XRPL_ENABLE_TELEMETRY /** * Writers currently inside doInsert. Instantaneous depth. */ @@ -101,6 +102,7 @@ public: * load, which is when the mean matters. */ std::atomic depthSamples{0}; +#endif // XRPL_ENABLE_TELEMETRY NuDBBackend( size_t keyBytes, @@ -279,6 +281,7 @@ public: nudb::detail::buffer bf; auto const result = nodeobjectCompress(e.getData(), e.getSize(), bf); +#ifdef XRPL_ENABLE_TELEMETRY // NuDB takes one global mutex for the whole insert, so the wait is // invisible from here. Record the depth we joined at and the wall // time we spent; the split follows from Little's Law. @@ -289,6 +292,10 @@ public: // slow, deep ones -- biasing the mean down exactly when queueing is // worst. With all writers inside their first insert the exit-counted // version reports no depth at all. + // + // Every node write reaches here, so none of it happens without + // telemetry: this whole block is what the write-stats gauges need and + // nothing else reads it. auto const depth = concurrentWriters.fetch_add(1, std::memory_order_relaxed) + 1; depthSum.fetch_add(depth, std::memory_order_relaxed); depthSamples.fetch_add(1, std::memory_order_relaxed); @@ -297,7 +304,8 @@ public: // A scope guard rather than straight-line code, because the insert // can allocate and so can throw. Leaking the depth would strand the // gauge above zero for the life of the process. - ScopeExit const account([this, depth, begin] { recordInsert(depth, begin); }); + ScopeExit const account([this, begin] { recordInsert(begin); }); +#endif db.insert(e.getKey(), result.first, result.second, ec); @@ -375,15 +383,29 @@ public: int getWriteLoad() override { +#ifdef XRPL_ENABLE_TELEMETRY // Writers in flight. Bounded by the number of writing threads, so // it stays far below LedgerMaster's kMaxWriteLoadAcquire of 8192 // and cannot suppress history acquisition. return static_cast(concurrentWriters.load(std::memory_order_relaxed)); +#else + // Nothing counts writers without telemetry, so report no load rather + // than a stale zero-valued counter. LedgerMaster gates history + // acquisition on this, so it must not start reporting a real depth as + // a side effect of instrumentation. + return 0; +#endif } [[nodiscard]] std::optional getWriteStats() const override { +#ifndef XRPL_ENABLE_TELEMETRY + // Not measured in this build, which is a different answer from + // measured-and-idle. The base class reports absence the same way for + // backends that never queue. + return std::nullopt; +#else WriteStats stats; stats.concurrentWriters = concurrentWriters.load(std::memory_order_relaxed); stats.insertCount = insertCount.load(std::memory_order_relaxed); @@ -392,6 +414,7 @@ public: stats.depthSum = depthSum.load(std::memory_order_relaxed); stats.depthSamples = depthSamples.load(std::memory_order_relaxed); return stats; +#endif } void @@ -432,11 +455,10 @@ private: * Always runs, including on the throwing path, so the depth gauge * returns to its true value even when the insert fails. * - * @param depth Writer depth this insert joined at, at least 1. * @param begin When the insert started. */ void - recordInsert(std::uint64_t depth, std::chrono::steady_clock::time_point begin) noexcept + recordInsert(std::chrono::steady_clock::time_point begin) noexcept { auto const elapsedUs = static_cast(std::chrono::duration_cast( diff --git a/src/libxrpl/telemetry/NullTelemetry.cpp b/src/libxrpl/telemetry/NullTelemetry.cpp index 92cc7303ef..baa418b19a 100644 --- a/src/libxrpl/telemetry/NullTelemetry.cpp +++ b/src/libxrpl/telemetry/NullTelemetry.cpp @@ -6,15 +6,26 @@ * unconditionally returns a NullTelemetry that does nothing. * * When XRPL_ENABLE_TELEMETRY IS defined, the OTel virtual methods - * (getTracer, startSpan) return noop tracers/spans. The makeTelemetry() - * factory in this file is not used in that case -- Telemetry.cpp provides - * its own factory that can return the real TelemetryImpl. + * (getTracer, startSpan, getMeter) return noop tracers, spans and meters, so + * the class stays concrete in both configurations. Every pure virtual the base + * declares behind that guard must be overridden here for that to hold. The + * makeTelemetry() factory in this file is not used in that case -- + * Telemetry.cpp provides its own factory that can return the real + * TelemetryImpl. */ #include +#ifndef XRPL_ENABLE_TELEMETRY +// beast::Journal is named only by the compiled-out makeTelemetry() below, so +// this include belongs to that configuration and would be unused in the other. +#include +#endif + #ifdef XRPL_ENABLE_TELEMETRY #include +#include +#include #include #include #include @@ -129,6 +140,14 @@ public: return opentelemetry::nostd::shared_ptr( new opentelemetry::trace::NoopSpan(nullptr)); } + + opentelemetry::nostd::shared_ptr + getMeter(std::string_view) override + { + static auto noopMeter = opentelemetry::nostd::shared_ptr( + new opentelemetry::metrics::NoopMeter()); + return noopMeter; + } #endif }; diff --git a/src/libxrpl/telemetry/SpanGuard.cpp b/src/libxrpl/telemetry/SpanGuard.cpp index 0782a78abc..de46821674 100644 --- a/src/libxrpl/telemetry/SpanGuard.cpp +++ b/src/libxrpl/telemetry/SpanGuard.cpp @@ -41,6 +41,7 @@ #include #include #include +#include #include #include #include @@ -474,6 +475,21 @@ SpanGuard::getTraceBytes() const return result; } +bool +SpanGuard::hasCurrentContext() noexcept +{ + // Read the active span straight out of the runtime context. GetSpan() + // would be shorter but heap-allocates a DefaultSpan whenever the context + // holds no span, which is the case this predicate exists to keep free. + // The two conditions below are the ones injectToProtobuf() returns early + // on, so a true result means that injector will write all three fields. + auto const ctx = opentelemetry::context::RuntimeContext::GetCurrent(); + auto const value = ctx.GetValue(otel_trace::kSpanKey); + auto const* const span = + opentelemetry::nostd::get_if>(&value); + return span != nullptr && *span && (*span)->GetContext().IsValid(); +} + void SpanGuard::injectCurrentContextToProtobuf(protocol::TraceContext& proto) { diff --git a/src/libxrpl/tx/Transactor.cpp b/src/libxrpl/tx/Transactor.cpp index 09aa4487da..8e97e730d9 100644 --- a/src/libxrpl/tx/Transactor.cpp +++ b/src/libxrpl/tx/Transactor.cpp @@ -1564,17 +1564,25 @@ Transactor::operator()() telemetry::tx_apply_span::transactor, txID.data(), txID.kBytes); - // "apply" — the third apply-pipeline stage, after preflight and preclaim. - span.setAttribute(telemetry::tx_apply_span::attr::stage, telemetry::tx_apply_span::val::apply); - if (auto const* fmt = TxFormats::getInstance().findByType(ctx_.tx.getTxnType())) - span.setAttribute(telemetry::tx_apply_span::attr::txType, fmt->getName().c_str()); - // The ledger being worked on (seq + parent hash) — correlates this apply - // stage to the ledger/consensus trace it is building into. - span.setAttribute( - telemetry::tx_apply_span::attr::currentLedgerSeq, static_cast(view().seq())); - span.setAttribute( - telemetry::tx_apply_span::attr::currentLedgerHash, - to_string(view().header().parentHash).c_str()); + // Guard the attribute work behind the active check, as preflight does in + // applySteps.cpp: this runs for every transaction applied, and the type + // lookup and the parent-hash string are not free. + if (span) + { + // "apply" — the third apply-pipeline stage, after preflight and preclaim. + span.setAttribute( + telemetry::tx_apply_span::attr::stage, telemetry::tx_apply_span::val::apply); + if (auto const* fmt = TxFormats::getInstance().findByType(ctx_.tx.getTxnType())) + span.setAttribute(telemetry::tx_apply_span::attr::txType, fmt->getName().c_str()); + // The ledger being worked on (seq + parent hash) — correlates this apply + // stage to the ledger/consensus trace it is building into. + span.setAttribute( + telemetry::tx_apply_span::attr::currentLedgerSeq, + static_cast(view().seq())); + span.setAttribute( + telemetry::tx_apply_span::attr::currentLedgerHash, + to_string(view().header().parentHash).c_str()); + } JLOG(j_.trace()) << "apply: " << ctx_.tx.getTransactionID(); @@ -1655,13 +1663,19 @@ Transactor::operator()() std::optional&& metadata = std::nullopt) -> ApplyResult { JLOG(j_.trace()) << (canApply ? "applied " : "not applied ") << transToken(result); - span.setAttribute(telemetry::tx_apply_span::attr::terResult, transToken(result).c_str()); - span.setAttribute(telemetry::tx_apply_span::attr::applied, canApply); - // Mark the span as errored when the transaction was not applied or the - // engine result is not a success, so failed applies surface in span-status - // error counts alongside preflight and preclaim. - if (!canApply || !isTesSuccess(result)) - span.setError(transToken(result)); + // Also guarded: transToken() is a lookup returning a string, and this + // funnel runs on every exit path. + if (span) + { + span.setAttribute( + telemetry::tx_apply_span::attr::terResult, transToken(result).c_str()); + span.setAttribute(telemetry::tx_apply_span::attr::applied, canApply); + // Mark the span as errored when the transaction was not applied or + // the engine result is not a success, so failed applies surface in + // span-status error counts alongside preflight and preclaim. + if (!canApply || !isTesSuccess(result)) + span.setError(transToken(result)); + } return {result, canApply, std::move(metadata)}; }; diff --git a/src/test/nodestore/DatabaseConfig_test.cpp b/src/test/nodestore/DatabaseConfig_test.cpp index 0f132081f0..04669bc6be 100644 --- a/src/test/nodestore/DatabaseConfig_test.cpp +++ b/src/test/nodestore/DatabaseConfig_test.cpp @@ -33,18 +33,23 @@ #include #include -#include #include #include #include #include #include -#include #include #include #include #include +#ifdef XRPL_ENABLE_TELEMETRY +// std::ranges::sort / std::ranges::find_if and std::optional / std::nullopt are +// used only by the gauge-helper suite below, so they are guarded like its uses. +#include +#include +#endif + namespace xrpl::node_store { class DatabaseConfig_test : public beast::unit_test::Suite diff --git a/src/tests/libxrpl/CMakeLists.txt b/src/tests/libxrpl/CMakeLists.txt index d624685529..9d9bc64691 100644 --- a/src/tests/libxrpl/CMakeLists.txt +++ b/src/tests/libxrpl/CMakeLists.txt @@ -116,14 +116,6 @@ if(telemetry) "${OTEL_IN_MEMORY_EXPORTER_LIB}" opentelemetry-cpp::opentelemetry-cpp ) - # ValidationTracker lives in src/xrpld/ (not libxrpl), so we compile its - # implementation directly into the test binary. The src/ include path its - # tests need is added unconditionally above. - target_sources( - xrpl_tests - PRIVATE - ${CMAKE_SOURCE_DIR}/src/xrpld/telemetry/detail/ValidationTracker.cpp - ) else() # MetricsRegistry lives in xrpld; compile its .cpp directly into the test # target so the no-op path can be tested without linking all of xrpld. @@ -135,4 +127,18 @@ else() ) endif() +# ValidationTracker lives in src/xrpld/ (not libxrpl), so we compile its +# implementation directly into the test binary and put src/ on the include path +# so its tests can reach headers. +# +# Both are unconditional: the class carries no telemetry guards, so its tests +# compile and run in every build. Gating them would leave the test file (which +# is likewise unguarded) without the header it includes and without the +# definitions it calls. +target_include_directories(xrpl_tests PRIVATE ${CMAKE_SOURCE_DIR}/src) +target_sources( + xrpl_tests + PRIVATE ${CMAKE_SOURCE_DIR}/src/xrpld/telemetry/detail/ValidationTracker.cpp +) + gtest_discover_tests(xrpl_tests DISCOVERY_TIMEOUT 60) diff --git a/src/tests/libxrpl/ledger/AcquireStats.cpp b/src/tests/libxrpl/ledger/AcquireStats.cpp index 4836c23fc5..54d09b33ba 100644 --- a/src/tests/libxrpl/ledger/AcquireStats.cpp +++ b/src/tests/libxrpl/ledger/AcquireStats.cpp @@ -10,10 +10,19 @@ * that a stalled run and a healthy run produce different, exact readings. * * AcquireStats is header-only, so nothing from xrpld needs to be linked here. + * + * The counts are held in telemetry::Counter, which carries no storage and + * records nothing when telemetry is compiled out. Every expectation on a + * non-zero count therefore has a mirror expectation of zero, so this file + * pins both configurations rather than leaving one of them unasserted. + * Expectations of zero that hold either way -- a deferral not disturbing the + * timeout counter, for instance -- are asserted unconditionally. */ #include +#include + #include #include @@ -23,10 +32,13 @@ namespace { using xrpl::AcquireStats; +namespace telemetry = xrpl::telemetry; /** * A fresh instance reads zero on every counter, so a later test can attribute - * every increment to its own call rather than to construction. + * every increment to its own call rather than to construction. The reading is + * the same in both configurations, which is what lets a caller report these + * accessors without knowing which build it is in. */ TEST(AcquireStatsTest, StartsAtZero) { @@ -54,27 +66,67 @@ TEST(AcquireStatsTest, CountersAdvanceIndependently) stats.recordDeferral(); stats.recordDeferral(); stats.recordDeferral(); - EXPECT_EQ(stats.getDeferrals(), 3u); + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(stats.getDeferrals(), 3u); + } + else + { + EXPECT_EQ(stats.getDeferrals(), 0u); + } + // Zero either way: a deferral must not reach the other counters. EXPECT_EQ(stats.getTimeouts(), 0u); EXPECT_EQ(stats.getGiveUps(), 0u); stats.recordTimeout(); - EXPECT_EQ(stats.getTimeouts(), 1u); - EXPECT_EQ(stats.getDeferrals(), 3u); + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(stats.getTimeouts(), 1u); + EXPECT_EQ(stats.getDeferrals(), 3u); + } + else + { + EXPECT_EQ(stats.getTimeouts(), 0u); + EXPECT_EQ(stats.getDeferrals(), 0u); + } stats.recordGiveUp(); - EXPECT_EQ(stats.getGiveUps(), 1u); - EXPECT_EQ(stats.getTimeouts(), 1u); - EXPECT_EQ(stats.getDeferrals(), 3u); + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(stats.getGiveUps(), 1u); + EXPECT_EQ(stats.getTimeouts(), 1u); + EXPECT_EQ(stats.getDeferrals(), 3u); + } + else + { + EXPECT_EQ(stats.getGiveUps(), 0u); + EXPECT_EQ(stats.getTimeouts(), 0u); + EXPECT_EQ(stats.getDeferrals(), 0u); + } stats.recordCompletion(); stats.recordCompletion(); - EXPECT_EQ(stats.getCompletions(), 2u); + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(stats.getCompletions(), 2u); + } + else + { + EXPECT_EQ(stats.getCompletions(), 0u); + } EXPECT_EQ(stats.getSweepEvictions(), 0u); stats.recordSweepEviction(); - EXPECT_EQ(stats.getSweepEvictions(), 1u); - EXPECT_EQ(stats.getCompletions(), 2u); + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(stats.getSweepEvictions(), 1u); + EXPECT_EQ(stats.getCompletions(), 2u); + } + else + { + EXPECT_EQ(stats.getSweepEvictions(), 0u); + EXPECT_EQ(stats.getCompletions(), 0u); + } // Nothing above records an abort, so both abort counters stay at zero. EXPECT_EQ(stats.getAborts(), 0u); @@ -92,12 +144,28 @@ TEST(AcquireStatsTest, AbortDistinguishesPartialWork) AcquireStats stats; stats.recordAbort(false); - EXPECT_EQ(stats.getAborts(), 1u); + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(stats.getAborts(), 1u); + } + else + { + EXPECT_EQ(stats.getAborts(), 0u); + } + // Zero either way: a cheap abort must not reach the partial-work counter. EXPECT_EQ(stats.getAbortsWithPartialWork(), 0u); stats.recordAbort(true); - EXPECT_EQ(stats.getAborts(), 2u); - EXPECT_EQ(stats.getAbortsWithPartialWork(), 1u); + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(stats.getAborts(), 2u); + EXPECT_EQ(stats.getAbortsWithPartialWork(), 1u); + } + else + { + EXPECT_EQ(stats.getAborts(), 0u); + EXPECT_EQ(stats.getAbortsWithPartialWork(), 0u); + } // An abort never counts as a completion or a give-up. EXPECT_EQ(stats.getCompletions(), 0u); @@ -120,13 +188,29 @@ TEST(AcquireStatsTest, StalledShapeIsDistinguishable) stalled.recordAbort(true); } - EXPECT_EQ(stalled.getDeferrals(), 1000u); + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(stalled.getDeferrals(), 1000u); + EXPECT_EQ(stalled.getSweepEvictions(), 5u); + EXPECT_EQ(stalled.getAborts(), 5u); + EXPECT_EQ(stalled.getAbortsWithPartialWork(), 5u); + } + else + { + // With telemetry compiled out every counter reads zero, so the stalled + // shape is not readable at all. That is the intended reading, not a + // healthy one. + EXPECT_EQ(stalled.getDeferrals(), 0u); + EXPECT_EQ(stalled.getSweepEvictions(), 0u); + EXPECT_EQ(stalled.getAborts(), 0u); + EXPECT_EQ(stalled.getAbortsWithPartialWork(), 0u); + } + + // Zero either way, and the point of the test: no timeout accrued, so the + // give-up path cannot fire. EXPECT_EQ(stalled.getTimeouts(), 0u); EXPECT_EQ(stalled.getGiveUps(), 0u); EXPECT_EQ(stalled.getCompletions(), 0u); - EXPECT_EQ(stalled.getSweepEvictions(), 5u); - EXPECT_EQ(stalled.getAborts(), 5u); - EXPECT_EQ(stalled.getAbortsWithPartialWork(), 5u); } /** @@ -145,10 +229,22 @@ TEST(AcquireStatsTest, HealthyShapeIsDistinguishable) for (int i = 0; i < 27; ++i) healthy.recordCompletion(); - EXPECT_EQ(healthy.getDeferrals(), 10u); - EXPECT_EQ(healthy.getTimeouts(), 7u); - EXPECT_EQ(healthy.getGiveUps(), 1u); - EXPECT_EQ(healthy.getCompletions(), 27u); + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(healthy.getDeferrals(), 10u); + EXPECT_EQ(healthy.getTimeouts(), 7u); + EXPECT_EQ(healthy.getGiveUps(), 1u); + EXPECT_EQ(healthy.getCompletions(), 27u); + } + else + { + EXPECT_EQ(healthy.getDeferrals(), 0u); + EXPECT_EQ(healthy.getTimeouts(), 0u); + EXPECT_EQ(healthy.getGiveUps(), 0u); + EXPECT_EQ(healthy.getCompletions(), 0u); + } + + // Zero either way: nothing above sweeps or aborts. EXPECT_EQ(healthy.getSweepEvictions(), 0u); EXPECT_EQ(healthy.getAborts(), 0u); EXPECT_EQ(healthy.getAbortsWithPartialWork(), 0u); @@ -184,9 +280,18 @@ TEST(AcquireStatsTest, ConcurrentRecordingLosesNothing) for (auto& th : threads) th.join(); - // 4 threads * 1000 iterations, stated independently of the loop bounds. - EXPECT_EQ(stats.getDeferrals(), std::uint64_t{4000}); - EXPECT_EQ(stats.getCompletions(), std::uint64_t{4000}); + if constexpr (telemetry::kEnabled) + { + // 4 threads * 1000 iterations, stated independently of the loop bounds. + EXPECT_EQ(stats.getDeferrals(), std::uint64_t{4000}); + EXPECT_EQ(stats.getCompletions(), std::uint64_t{4000}); + } + else + { + // Nothing was recorded, so concurrent calls also have nothing to lose. + EXPECT_EQ(stats.getDeferrals(), std::uint64_t{0}); + EXPECT_EQ(stats.getCompletions(), std::uint64_t{0}); + } // The threads record only those two events, so the rest stay at zero. EXPECT_EQ(stats.getTimeouts(), 0u); diff --git a/src/tests/libxrpl/nodestore/Backend.cpp b/src/tests/libxrpl/nodestore/Backend.cpp index f2648f3e12..0b37815843 100644 --- a/src/tests/libxrpl/nodestore/Backend.cpp +++ b/src/tests/libxrpl/nodestore/Backend.cpp @@ -182,7 +182,15 @@ TEST_P(BackendTypeTest, write_stats_reported_only_when_measured) auto backend = makeOpenBackend(); auto const stats = backend->getWriteStats(); - if (GetParam() == "nudb") + // NuDB records the write path only when telemetry is compiled in; without + // it there is no consumer, so it reports absence like every other backend. +#ifdef XRPL_ENABLE_TELEMETRY + bool const measures = GetParam() == "nudb"; +#else + bool const measures = false; +#endif + + if (measures) { if (!stats.has_value()) FAIL() << "nudb must report write stats"; diff --git a/src/tests/libxrpl/nodestore/Database.cpp b/src/tests/libxrpl/nodestore/Database.cpp index 2cf0e783e8..3c57397b53 100644 --- a/src/tests/libxrpl/nodestore/Database.cpp +++ b/src/tests/libxrpl/nodestore/Database.cpp @@ -225,8 +225,15 @@ TEST_P(NodeStoreDatabaseTest, write_stats_forwarded_from_backend) // Before any write, only a measuring backend answers at all. auto const initial = db->getWriteStats(); - ASSERT_EQ(initial.has_value(), GetParam() == "nudb") - << "only nudb measures its write path; backend=" << GetParam(); + // NuDB records the write path only when telemetry is compiled in; without + // it nothing measures, so the negative path below covers every backend. +#ifdef XRPL_ENABLE_TELEMETRY + bool const measures = GetParam() == "nudb"; +#else + bool const measures = false; +#endif + ASSERT_EQ(initial.has_value(), measures) + << "only nudb measures its write path, and only with telemetry; backend=" << GetParam(); if (!initial) { diff --git a/src/tests/libxrpl/nodestore/NuDBFactory.cpp b/src/tests/libxrpl/nodestore/NuDBFactory.cpp index a7d6a11dca..09d398fa14 100644 --- a/src/tests/libxrpl/nodestore/NuDBFactory.cpp +++ b/src/tests/libxrpl/nodestore/NuDBFactory.cpp @@ -354,6 +354,11 @@ TEST(NuDBFactory, configuration_parsing) // so these counters plus Little's Law are the only way to separate queuing // from service time from outside the library. +#ifdef XRPL_ENABLE_TELEMETRY +// The three tests below pin the cumulative write-path counters, which are +// recorded only when telemetry is compiled in. Without it the write-load gauge +// is their only consumer and it reads the depth atomic directly, so the depth +// samples and the two clock reads per insert are skipped on that hot path. TEST(NuDBFactory, write_stats_accumulate_per_insert) { TempDir const tempDir; @@ -637,6 +642,7 @@ TEST(NuDBFactory, write_load_reports_writer_depth) backend->close(); } +#endif // XRPL_ENABLE_TELEMETRY TEST(NuDBFactory, data_persistence) { diff --git a/src/tests/libxrpl/telemetry/MetricsRegistry.cpp b/src/tests/libxrpl/telemetry/MetricsRegistry.cpp index d5781a36a8..10f64aa188 100644 --- a/src/tests/libxrpl/telemetry/MetricsRegistry.cpp +++ b/src/tests/libxrpl/telemetry/MetricsRegistry.cpp @@ -464,10 +464,8 @@ TEST(MetricsRegistryScaledMean, default_scale_is_one) #include -#include #include #include -#include using namespace xrpl; @@ -735,8 +733,8 @@ public: [[nodiscard]] std::optional const& getTrapTxID() const override { - static std::optional const empty; - return empty; + static std::optional const kEmpty; + return kEmpty; } DatabaseCon& getWalletDB() override diff --git a/src/tests/libxrpl/telemetry/Recording.cpp b/src/tests/libxrpl/telemetry/Recording.cpp new file mode 100644 index 0000000000..77e970818f --- /dev/null +++ b/src/tests/libxrpl/telemetry/Recording.cpp @@ -0,0 +1,135 @@ +/** + * Tests for the telemetry recording utilities. + * + * Compiled in every build. Each test asserts the compiled-in behaviour when + * kEnabled and the compiled-out behaviour otherwise, so both configurations + * are pinned by the same file rather than one of them going unasserted. + */ + +#include + +#include + +#include +#include +#include +#include + +using namespace xrpl; + +// kEnabled must agree with the macro that drives every branch below. Asserted +// at compile time so a mismatch cannot reach the runtime assertions. +#ifdef XRPL_ENABLE_TELEMETRY +static_assert(telemetry::kEnabled, "kEnabled must be true when the macro is defined"); +#else +static_assert(!telemetry::kEnabled, "kEnabled must be false when the macro is absent"); +#endif + +// Counter's copy semantics must NOT depend on the configuration. When telemetry +// is compiled in, the std::atomic member deletes all four implicitly; when it is +// compiled out, Counter declares them deleted itself. Without that, an owning +// class would be non-copyable in one build and copyable in the other. Asserted +// unconditionally, because the whole point is that both builds agree. +static_assert(!std::is_copy_constructible_v>); +static_assert(!std::is_copy_assignable_v>); +static_assert(!std::is_move_constructible_v>); +static_assert(!std::is_move_assignable_v>); + +// Stopwatch and Mirror hold ordinary values, so they stay copyable in both +// builds; only Counter needed the explicit deletions above. +static_assert(std::is_copy_constructible_v); + +// The compiled-out forms must be empty types. That is what lets an owning class +// declare the member unconditionally: the storage collapses to padding, and the +// work disappears entirely. +TEST(Recording, compiled_out_types_are_empty) +{ + if constexpr (telemetry::kEnabled) + { + EXPECT_FALSE(std::is_empty_v>); + EXPECT_FALSE(std::is_empty_v); + } + else + { + EXPECT_TRUE(std::is_empty_v>); + EXPECT_TRUE(std::is_empty_v); + } +} + +// A fresh counter reads zero in both configurations -- the one value that must +// agree, since callers may report it unconditionally. +TEST(Recording, counter_starts_at_zero) +{ + telemetry::Counter<> counter; + EXPECT_EQ(counter.load(), 0U); +} + +// add() accumulates exactly when compiled in, and stays at zero when not. +TEST(Recording, counter_accumulates_only_when_compiled_in) +{ + telemetry::Counter<> counter; + counter.add(); + counter.add(4); + + if constexpr (telemetry::kEnabled) + { + EXPECT_EQ(counter.load(), 5U); + } + else + { + EXPECT_EQ(counter.load(), 0U); + } +} + +// The default template argument is std::uint64_t; an explicit type is honoured. +TEST(Recording, counter_honours_its_value_type) +{ + static_assert(std::is_same_v{}.load()), std::uint64_t>); + static_assert( + std::is_same_v{}.load()), std::uint32_t>); + + telemetry::Counter counter; + counter.add(7); + EXPECT_EQ(counter.load(), telemetry::kEnabled ? 7U : 0U); +} + +// A stopwatch measures a real interval when compiled in and reports exactly +// zero when not, so callers can report elapsedUs() unconditionally. +TEST(Recording, stopwatch_measures_only_when_compiled_in) +{ + telemetry::Stopwatch const timer; + std::this_thread::sleep_for(std::chrono::milliseconds{2}); + auto const elapsed = timer.elapsedUs(); + + if constexpr (telemetry::kEnabled) + { + EXPECT_GE(elapsed, std::chrono::microseconds{1000}); + } + else + { + EXPECT_EQ(elapsed, std::chrono::microseconds{0}); + } +} + +// restart() moves the origin forward, so the interval measured after it is +// shorter than the one before it. +TEST(Recording, stopwatch_restart_resets_the_origin) +{ + telemetry::Stopwatch timer; + std::this_thread::sleep_for(std::chrono::milliseconds{4}); + auto const beforeRestart = timer.elapsedUs(); + + timer.restart(); + auto const afterRestart = timer.elapsedUs(); + + if constexpr (telemetry::kEnabled) + { + EXPECT_GE(beforeRestart, std::chrono::microseconds{2000}); + EXPECT_LT(afterRestart, beforeRestart); + } + else + { + EXPECT_EQ(beforeRestart, std::chrono::microseconds{0}); + EXPECT_EQ(afterRestart, std::chrono::microseconds{0}); + } +} diff --git a/src/tests/libxrpl/telemetry/ValidationTracker.cpp b/src/tests/libxrpl/telemetry/ValidationTracker.cpp index 2542316e7a..9d176377a0 100644 --- a/src/tests/libxrpl/telemetry/ValidationTracker.cpp +++ b/src/tests/libxrpl/telemetry/ValidationTracker.cpp @@ -377,3 +377,39 @@ TEST_F(ValidationTrackerTest, GrossAgreementsCountInitialOnly) // Additive invariant: gross agree + gross miss == ledgers reconciled. EXPECT_EQ(tracker_.totalAgreementsEver() + tracker_.totalMissedEver(), 5u); } + +// --------------------------------------------------------------- +// 12. Pending map stays bounded when nothing ever reconciles +// reconcile() is the only pruning path, and it runs only from +// the observable-gauge callbacks -- which need telemetry both +// compiled in and enabled. A node with telemetry off, or with +// [telemetry] enabled=0, therefore never reconciles, so the +// record methods have to bound the map themselves. +// --------------------------------------------------------------- +TEST_F(ValidationTrackerTest, PendingStaysBoundedWithoutReconcile) +{ + constexpr std::uint64_t kFirstBatch = 4000; + constexpr std::uint64_t kSecondBatch = 8000; + + for (std::uint64_t i = 0; i < kFirstBatch; ++i) + tracker_.recordOurValidation(makeHash(i), static_cast(i)); + + auto const afterFirst = tracker_.pendingCount(); + + for (std::uint64_t i = kFirstBatch; i < kSecondBatch; ++i) + tracker_.recordOurValidation(makeHash(i), static_cast(i)); + + auto const afterSecond = tracker_.pendingCount(); + + // Bounded at all: far fewer entries retained than recorded. + EXPECT_LT(afterFirst, kFirstBatch); + + // Bounded at a fixed cap, not merely growing more slowly: doubling the + // input leaves the size unchanged. Asserted without naming the private + // constant, so the test survives a change to its value. + EXPECT_EQ(afterFirst, afterSecond); + + // Every recorded validation is still counted, so bounding the map does not + // cost us the lifetime totals the gauges report. + EXPECT_EQ(tracker_.totalValidationsSent(), kSecondBatch); +} diff --git a/src/xrpld/app/consensus/RCLConsensus.cpp b/src/xrpld/app/consensus/RCLConsensus.cpp index a7790fbdd6..83d0cc56f5 100644 --- a/src/xrpld/app/consensus/RCLConsensus.cpp +++ b/src/xrpld/app/consensus/RCLConsensus.cpp @@ -25,6 +25,7 @@ #include #endif #include +#include #include #include @@ -71,8 +72,6 @@ #include #include #include -#include -#include #include @@ -96,6 +95,13 @@ #include #include +#ifdef XRPL_ENABLE_TELEMETRY +// The Telemetry interface and the shared segment names are named only by +// startRoundTracing(), which is telemetry-enabled code. +#include +#include +#endif + namespace xrpl { RCLConsensus::RCLConsensus( @@ -284,7 +290,12 @@ RCLConsensus::Adaptor::propose(RCLCxPeerPos::Proposal const& proposal) // Inject the current thread's active span context (e.g. the consensus // round span) so receiving peers can link their proposal.receive span // as a child of this trace. - telemetry::SpanGuard::injectCurrentContextToProtobuf(*prop.mutable_trace_context()); + // + // The helper injects only when a span is actually active, so a node with + // telemetry compiled out, disabled by config, or simply not tracing this + // round sends no TraceContext at all rather than an empty one that makes + // every peer take its has_trace_context() branch for nothing. + telemetry::injectCurrentContext(prop); app_.getOverlay().broadcast(prop); } @@ -361,15 +372,23 @@ RCLConsensus::Adaptor::onClose( // Child of the round span via its captured context (roundSpan_ is a // thread-free SpanGuard, so parent explicitly via its context). auto span = telemetry::SpanGuard::childSpan(cs::ledgerClose, roundSpanContext_); - span.setAttribute(cs::attr::ledgerSeq, static_cast(ledger.ledger->header().seq) + 1); - span.setAttribute(cs::attr::mode, toDisplayString(mode).c_str()); - span.setAttribute( - cs::attr::txCountOpen, static_cast(app_.getOpenLedger().current()->txCount())); - span.setAttribute( - cs::attr::closeTimeResolutionMs, - static_cast( - std::chrono::duration_cast(ledger.closeTimeResolution()) - .count())); + // setAttribute is the only consumer of everything read here, so the block is + // guarded on the span being live. Unguarded, every round takes the open + // ledger's currentMutex_ and copies a shared_ptr just to read txCount, and + // builds a mode string, for attributes no one may be recording. + if (span) + { + span.setAttribute( + cs::attr::ledgerSeq, static_cast(ledger.ledger->header().seq) + 1); + span.setAttribute(cs::attr::mode, toDisplayString(mode).c_str()); + span.setAttribute( + cs::attr::txCountOpen, static_cast(app_.getOpenLedger().current()->txCount())); + span.setAttribute( + cs::attr::closeTimeResolutionMs, + static_cast( + std::chrono::duration_cast(ledger.closeTimeResolution()) + .count())); + } bool const wrongLCL = mode == ConsensusMode::WrongLedger; bool const proposing = mode == ConsensusMode::Proposing; @@ -512,17 +531,24 @@ RCLConsensus::Adaptor::onAccept( }); } +// Not static: the guarded body reads app_ and roundSpanContext_. With telemetry +// compiled out it returns an empty handle and touches no member, so clang-tidy +// sees a method that could be static. +// NOLINTBEGIN(readability-convert-member-functions-to-static) std::shared_ptr RCLConsensus::Adaptor::makeAcceptSpan(Result const& result) { + // The whole body is telemetry: the guard, its attributes and the captured + // context serve the accept span only. With telemetry compiled out the handle + // stays empty, so accepting a ledger does not allocate a control block for a + // span that can never record. doAccept only hands the handle to + // activateIfLive(), which tests it, so an empty handle is safe on both the + // sync (onForceAccept) and async (onAccept) paths. +#ifdef XRPL_ENABLE_TELEMETRY namespace cs = telemetry::consensus::span; auto span = std::make_shared( telemetry::SpanGuard::childSpan(cs::accept, roundSpanContext_)); - span->setAttribute(cs::attr::proposers, static_cast(result.proposers)); - span->setAttribute( - cs::attr::roundTimeMs, static_cast(result.roundTime.read().count())); - // Round duration as a native histogram, alongside the span attribute above. // The two answer different questions and neither replaces the other: the // attribute says how long THIS round took, readable inside the trace next @@ -541,29 +567,36 @@ RCLConsensus::Adaptor::makeAcceptSpan(Result const& result) telemetry::metric::consensusRoundDurationMs, "Wall-clock duration of a completed consensus round in milliseconds", result.roundTime.read().count()); - span->setAttribute(cs::attr::quorum, static_cast(app_.getValidators().quorum())); - span->setAttribute(cs::attr::disputesCount, static_cast(result.disputes.size())); - char const* stateStr = [&] { - switch (result.state) - { - case ConsensusState::Yes: - return "yes"; - case ConsensusState::MovedOn: - return "moved_on"; - case ConsensusState::Expired: - return "expired"; - default: - return "no"; - } - }(); - span->setAttribute(cs::attr::consensusState, stateStr); - // Capture the accept span's context so createValidationSpan() — which - // runs on the jtACCEPT worker thread — can link the validation.send - // span to the accept span (matching the design diagram and the - // "validation follows acceptance" causal model). + // Every attribute below exists only for the span, so the whole block — + // attributes and the context capture — is guarded on the span being live. if (*span) { + span->setAttribute(cs::attr::proposers, static_cast(result.proposers)); + span->setAttribute( + cs::attr::roundTimeMs, static_cast(result.roundTime.read().count())); + span->setAttribute(cs::attr::quorum, static_cast(app_.getValidators().quorum())); + span->setAttribute(cs::attr::disputesCount, static_cast(result.disputes.size())); + char const* stateStr = [&] { + switch (result.state) + { + case ConsensusState::Yes: + return "yes"; + case ConsensusState::MovedOn: + return "moved_on"; + case ConsensusState::Expired: + return "expired"; + default: + return "no"; + } + }(); + span->setAttribute(cs::attr::consensusState, stateStr); + + // Capture the accept span's context so createValidationSpan() — which + // runs on the jtACCEPT worker thread — can link the validation.send + // span to the accept span (matching the design diagram and the + // "validation follows acceptance" causal model). + // // span is a thread-free SpanGuard handed to the JtAccept worker // (onAccept), which ends it there. spanContext() captures the guard's // own span, so accept.apply parents via acceptSpanContext_ regardless @@ -572,7 +605,11 @@ RCLConsensus::Adaptor::makeAcceptSpan(Result const& result) acceptSpanContext_ = span->spanContext(); } return span; +#else + return {}; +#endif } +// NOLINTEND(readability-convert-member-functions-to-static) void RCLConsensus::Adaptor::doAccept( @@ -649,6 +686,10 @@ RCLConsensus::Adaptor::doAccept( cs::attr::closeTimeVoteBins, static_cast(rawCloseTimes.peers.size())); doAcceptSpan.setAttribute( cs::attr::disputesResolvedCount, static_cast(result.disputes.size())); + // prevRes and dir feed the resolution_direction attribute and nothing else, + // so both are guarded on the span being active. Unguarded, every accepted + // ledger builds a std::string that no one reads. + if (doAcceptSpan) { auto const prevRes = prevLedger.closeTimeResolution(); auto const dir = [&]() -> std::string { @@ -682,6 +723,10 @@ RCLConsensus::Adaptor::doAccept( JLOG(j_.debug()) << "Building canonical tx set: " << retriableTxs.key(); + // txCount and the per-transaction event feed the span and nothing else, so + // both are guarded on the span being active. Unguarded, every accepted + // ledger builds one 64-character hash string per transaction that no one + // reads. int64_t txCount = 0; for (auto const& item : *result.txns.map) { @@ -689,9 +734,12 @@ RCLConsensus::Adaptor::doAccept( { retriableTxs.insert(std::make_shared(SerialIter{item.slice()})); JLOG(j_.debug()) << " Tx: " << item.key(); - ++txCount; - auto const txHash = to_string(item.key()); - doAcceptSpan.addEvent(cs::event::txIncluded, {{cs::attr::txId, txHash}}); + if (doAcceptSpan) + { + ++txCount; + auto const txHash = to_string(item.key()); + doAcceptSpan.addEvent(cs::event::txIncluded, {{cs::attr::txId, txHash}}); + } } catch (std::exception const& ex) { @@ -699,7 +747,10 @@ RCLConsensus::Adaptor::doAccept( JLOG(j_.warn()) << " Tx: " << item.key() << " throws: " << ex.what(); } } - doAcceptSpan.setAttribute(cs::attr::txCount, txCount); + if (doAcceptSpan) + { + doAcceptSpan.setAttribute(cs::attr::txCount, txCount); + } auto built = buildLCL( prevLedger, @@ -986,7 +1037,10 @@ void RCLConsensus::Adaptor::validate(RCLCxLedger const& ledger, RCLTxSet const& txns, bool proposing) { auto valSpan = createValidationSpan(); - if (valSpan) + // Testing the guard as well as the optional matters: a guard that exists but + // is not live still evaluates its arguments, and the ledger_hash attribute + // below turns a 32-byte hash into a 64-character string. + if (valSpan && *valSpan) { namespace cs = telemetry::consensus::span; valSpan->setAttribute(cs::attr::ledgerSeq, static_cast(ledger.seq())); @@ -1005,7 +1059,7 @@ RCLConsensus::Adaptor::validate(RCLCxLedger const& ledger, RCLTxSet const& txns, validationTime = lastValidationTime_ + 1s; lastValidationTime_ = validationTime; - if (valSpan) + if (valSpan && *valSpan) { valSpan->setAttribute( telemetry::consensus::span::attr::validationSignTime, @@ -1090,7 +1144,11 @@ RCLConsensus::Adaptor::validate(RCLCxLedger const& ledger, RCLTxSet const& txns, // `serialized`, so it is not covered by validation authenticity. // Downstream consumers treat it as advisory only. A signature-covered // trace context is a possible future enhancement. - telemetry::SpanGuard::injectCurrentContextToProtobuf(*val.mutable_trace_context()); + // + // As on the proposal path, the helper injects only when a span is actually + // active, so a node that is not tracing sends no TraceContext at all + // rather than an empty one. + telemetry::injectCurrentContext(val); app_.getOverlay().broadcast(val); // Publish to all our subscribers: @@ -1100,9 +1158,16 @@ RCLConsensus::Adaptor::validate(RCLCxLedger const& ledger, RCLTxSet const& txns, if (auto* mr = app_.getMetricsRegistry()) { mr->incrementValidationsSent(); +#ifdef XRPL_ENABLE_TELEMETRY // Record our validation for the agreement tracker so it can // compare against network-validated ledgers. - mr->getValidationTracker().recordOurValidation(ledger.id(), ledger.seq()); + // + // Only when enabled: recording takes the tracker's lock and inserts an + // entry, and nothing reconciles or drains those entries unless the + // observable gauges are running. + if (mr->isEnabled()) + mr->getValidationTracker().recordOurValidation(ledger.id(), ledger.seq()); +#endif } } @@ -1295,9 +1360,18 @@ RCLConsensus::Adaptor::updateOperatingMode(std::size_t const positions) const app_.getOPs().setMode(OperatingMode::CONNECTED); } +// Neither is static: both guarded bodies read the span-context members. With +// telemetry compiled out one body is empty and the other returns std::nullopt, +// so clang-tidy sees two methods that could be static. +// NOLINTBEGIN(readability-convert-member-functions-to-static) void RCLConsensus::Adaptor::startRoundTracing(RCLCxLedger const& prevLgr) { + // The whole body is telemetry: every member it touches exists only to carry + // span state. It is compiled out rather than left to the early return below, + // because the work above that return — two virtual Telemetry calls and the + // strategy string compare — would otherwise run once per round for nothing. +#ifdef XRPL_ENABLE_TELEMETRY namespace cs = telemetry::consensus::span; // Capture the prior round's context BEFORE the new span overwrites @@ -1372,11 +1446,16 @@ RCLConsensus::Adaptor::startRoundTracing(RCLCxLedger const& prevLgr) // reset() on a different worker than it was emplaced on, so spanContext() // captures its own span and no scope work is needed. roundSpanContext_ = roundSpan_->spanContext(); +#endif } std::optional RCLConsensus::Adaptor::createValidationSpan() { + // The whole body is telemetry: it only builds a span from stored contexts. + // Compiled out, it yields std::nullopt, so validate() takes neither branch + // that reads the ledger hash into a string. +#ifdef XRPL_ENABLE_TELEMETRY namespace cs = telemetry::consensus::span; // Prefer linking to the accept span (matches the design diagram and @@ -1395,7 +1474,11 @@ RCLConsensus::Adaptor::createValidationSpan() } return telemetry::SpanGuard::linkedSpan(cs::validationSend, roundSpanContext_); +#else + return std::nullopt; +#endif } +// NOLINTEND(readability-convert-member-functions-to-static) void RCLConsensus::Adaptor::onPhaseEvent(std::string_view eventName, std::string_view phaseLabel) diff --git a/src/xrpld/app/ledger/AcquireStats.h b/src/xrpld/app/ledger/AcquireStats.h index 7ed93e85c7..27d5f219b5 100644 --- a/src/xrpld/app/ledger/AcquireStats.h +++ b/src/xrpld/app/ledger/AcquireStats.h @@ -1,6 +1,7 @@ #pragma once -#include +#include + #include namespace xrpl { @@ -34,8 +35,8 @@ namespace xrpl { * \--- recordTimeout() ------->| | * | | * InboundLedger ---- recordGiveUp() -------->| AcquireStats | - * |--- recordCompletion() ---->| (7 atomic | - * \--- recordAbort() --------->| counters) | + * |--- recordCompletion() ---->| (9 counters) | + * \--- recordAbort() --------->| | * | | * InboundLedgers ---- recordSweepEviction() ->+-----------------+ * | @@ -71,14 +72,18 @@ namespace xrpl { * // the same as healthy. Check completions before concluding anything. * @endcode * - * @note Thread-safe. Every counter is an independent relaxed atomic, so any - * number of threads may record concurrently without losing an - * increment. Because the counters are independent, no read across two - * of them is a consistent snapshot. Compare rates over an interval - * rather than instantaneous values. + * @note Thread-safe. Every counter is independent, so any number of threads + * may record concurrently without losing an increment. Because the + * counters are independent, no read across two of them is a consistent + * snapshot. Compare rates over an interval rather than instantaneous + * values. * @note All counters are monotonic for the life of the process and are never * reset, so a reader differences successive samples to get a rate. They * saturate only on 64-bit wraparound, which is unreachable in practice. + * @note The counts exist only to be reported, so they are held in + * telemetry::Counter. In a build with telemetry compiled out the + * counters carry no storage, recording is a no-op, and every accessor + * returns 0. The public API is the same in both builds. * @note Completions count acquisitions that ended successfully, whether the * data came from peers or was already present locally. They do not * count an acquisition that is still in flight. @@ -95,9 +100,9 @@ public: void recordDeferral(bool ledgerAcquisition = false) { - deferrals_.fetch_add(1, std::memory_order_relaxed); + deferrals_.add(); if (ledgerAcquisition) - ledgerDeferrals_.fetch_add(1, std::memory_order_relaxed); + ledgerDeferrals_.add(); } /** @@ -109,9 +114,9 @@ public: void recordTimeout(bool ledgerAcquisition = false) { - timeouts_.fetch_add(1, std::memory_order_relaxed); + timeouts_.add(); if (ledgerAcquisition) - ledgerTimeouts_.fetch_add(1, std::memory_order_relaxed); + ledgerTimeouts_.add(); } /** @@ -121,7 +126,7 @@ public: void recordGiveUp() { - giveUps_.fetch_add(1, std::memory_order_relaxed); + giveUps_.add(); } /** @@ -134,9 +139,9 @@ public: void recordAbort(bool hadPartialWork) { - aborts_.fetch_add(1, std::memory_order_relaxed); + aborts_.add(); if (hadPartialWork) - abortsWithPartialWork_.fetch_add(1, std::memory_order_relaxed); + abortsWithPartialWork_.add(); } /** @@ -145,7 +150,7 @@ public: void recordCompletion() { - completions_.fetch_add(1, std::memory_order_relaxed); + completions_.add(); } /** @@ -154,7 +159,7 @@ public: void recordSweepEviction() { - sweepEvictions_.fetch_add(1, std::memory_order_relaxed); + sweepEvictions_.add(); } /** @@ -163,7 +168,7 @@ public: [[nodiscard]] std::uint64_t getDeferrals() const { - return deferrals_.load(std::memory_order_relaxed); + return deferrals_.load(); } /** @@ -172,7 +177,7 @@ public: [[nodiscard]] std::uint64_t getTimeouts() const { - return timeouts_.load(std::memory_order_relaxed); + return timeouts_.load(); } /** @@ -185,7 +190,7 @@ public: [[nodiscard]] std::uint64_t getLedgerDeferrals() const { - return ledgerDeferrals_.load(std::memory_order_relaxed); + return ledgerDeferrals_.load(); } /** @@ -198,7 +203,7 @@ public: [[nodiscard]] std::uint64_t getLedgerTimeouts() const { - return ledgerTimeouts_.load(std::memory_order_relaxed); + return ledgerTimeouts_.load(); } /** @@ -207,7 +212,7 @@ public: [[nodiscard]] std::uint64_t getGiveUps() const { - return giveUps_.load(std::memory_order_relaxed); + return giveUps_.load(); } /** @@ -216,7 +221,7 @@ public: [[nodiscard]] std::uint64_t getAborts() const { - return aborts_.load(std::memory_order_relaxed); + return aborts_.load(); } /** @@ -227,7 +232,7 @@ public: [[nodiscard]] std::uint64_t getAbortsWithPartialWork() const { - return abortsWithPartialWork_.load(std::memory_order_relaxed); + return abortsWithPartialWork_.load(); } /** @@ -236,7 +241,7 @@ public: [[nodiscard]] std::uint64_t getCompletions() const { - return completions_.load(std::memory_order_relaxed); + return completions_.load(); } /** @@ -245,54 +250,54 @@ public: [[nodiscard]] std::uint64_t getSweepEvictions() const { - return sweepEvictions_.load(std::memory_order_relaxed); + return sweepEvictions_.load(); } private: /** * Timer jobs skipped because the job lane was at its limit. */ - std::atomic deferrals_{0}; + telemetry::Counter<> deferrals_; /** * Deferrals attributable to ledger acquisition alone. */ - std::atomic ledgerDeferrals_{0}; + telemetry::Counter<> ledgerDeferrals_; /** * Timeouts attributable to ledger acquisition alone. */ - std::atomic ledgerTimeouts_{0}; + telemetry::Counter<> ledgerTimeouts_; /** * Timer bodies that ran and advanced the retry count. */ - std::atomic timeouts_{0}; + telemetry::Counter<> timeouts_; /** * Acquisitions that exhausted their retry budget. */ - std::atomic giveUps_{0}; + telemetry::Counter<> giveUps_; /** * Acquisitions destroyed before finishing. */ - std::atomic aborts_{0}; + telemetry::Counter<> aborts_; /** * Aborts that discarded a partly built map. */ - std::atomic abortsWithPartialWork_{0}; + telemetry::Counter<> abortsWithPartialWork_; /** * Acquisitions that finished successfully. */ - std::atomic completions_{0}; + telemetry::Counter<> completions_; /** * Idle acquisitions evicted by the sweep. */ - std::atomic sweepEvictions_{0}; + telemetry::Counter<> sweepEvictions_; }; } // namespace xrpl diff --git a/src/xrpld/app/ledger/detail/LedgerMaster.cpp b/src/xrpld/app/ledger/detail/LedgerMaster.cpp index 715d605b45..746a13f9b8 100644 --- a/src/xrpld/app/ledger/detail/LedgerMaster.cpp +++ b/src/xrpld/app/ledger/detail/LedgerMaster.cpp @@ -345,10 +345,16 @@ LedgerMaster::setValidLedger(std::shared_ptr const& l) (void)maxLedgerDifference_; validLedgerSeq_ = l->header().seq; +#ifdef XRPL_ENABLE_TELEMETRY // Record the network-validated ledger for the agreement tracker so it // can compare against our own validations. - if (auto* mr = app_.getMetricsRegistry()) + // + // Only when enabled: recording takes the tracker's lock and inserts an + // entry, and nothing reconciles or drains those entries unless the + // observable gauges are running. + if (auto* mr = app_.getMetricsRegistry(); mr && mr->isEnabled()) mr->getValidationTracker().recordNetworkValidation(l->header().hash, l->header().seq); +#endif app_.getOPs().updateLocalTx(*l); app_.getSHAMapStore().onLedgerClosed(getValidatedLedger()); diff --git a/src/xrpld/app/ledger/detail/TimeoutCounter.cpp b/src/xrpld/app/ledger/detail/TimeoutCounter.cpp index 3e9961bfb1..5992bfe8f7 100644 --- a/src/xrpld/app/ledger/detail/TimeoutCounter.cpp +++ b/src/xrpld/app/ledger/detail/TimeoutCounter.cpp @@ -8,6 +8,7 @@ #include #include #include +#include #include #include @@ -67,7 +68,12 @@ TimeoutCounter::queueJob(ScopedLockType& sl) // Counted separately from timeouts: this path re-arms the timer // without running invokeOnTimer, so timeouts_ does not advance and the // give-up test that reads it cannot fire while the lane stays full. - app_.getAcquireStats().recordDeferral(isLedgerAcquisition()); + // + // Compiled out when telemetry is off: the argument compares the job + // name against a string literal on every deferred tick of every + // in-flight task, and a no-op count would still evaluate it. + if constexpr (telemetry::kEnabled) + app_.getAcquireStats().recordDeferral(isLedgerAcquisition()); JLOG(journal_.debug()) << "Deferring " << queueJobParameter_.jobName << " timer due to load"; setTimer(sl); @@ -92,7 +98,12 @@ TimeoutCounter::invokeOnTimer() if (!progress_) { ++timeouts_; - app_.getAcquireStats().recordTimeout(isLedgerAcquisition()); + // Same argument cost as the deferral above -- one job-name string + // comparison per no-progress tick -- so it is compiled out the same + // way. timeouts_ stays outside the block because the give-up test + // reads it. + if constexpr (telemetry::kEnabled) + app_.getAcquireStats().recordTimeout(isLedgerAcquisition()); JLOG(journal_.debug()) << "Timeout(" << timeouts_ << ") " << " acquiring " << hash_; onTimer(false, sl); diff --git a/src/xrpld/app/misc/NetworkOPs.cpp b/src/xrpld/app/misc/NetworkOPs.cpp index 1a71eb7fc3..4ac2d968c6 100644 --- a/src/xrpld/app/misc/NetworkOPs.cpp +++ b/src/xrpld/app/misc/NetworkOPs.cpp @@ -1599,23 +1599,40 @@ NetworkOPsImp::processTransaction( // SpanGuard is thread-free (holds no Scope), so it is safe to store here // and end on the batch worker thread that later applies this transaction — // no detach step is needed. - auto span = std::make_shared(txProcessSpan(transaction->getID())); - span->setAttribute(tx_span::attr::txHash, to_string(transaction->getID()).c_str()); - span->setAttribute(tx_span::attr::local, bLocal); - // The current (open) ledger index at submission/relay time — the ledger - // being worked on. Correlates this tx.process to the ledger trace; the tx - // has not yet been applied to a specific ledger here, so there is no hash. - span->setAttribute( - tx_span::attr::currentLedgerSeq, - static_cast(ledgerMaster_.getCurrentLedgerIndex())); - if (auto const& stx = transaction->getSTransaction()) + // Left null when telemetry is compiled out: there is no span to own, so + // nothing is allocated for one. The transaction pipeline already accepts a + // null span -- both doTransaction* overloads default it to nullptr -- and + // every use tests it. Without this the make_shared allocated once per + // submitted and relayed transaction to hold an empty object. + std::shared_ptr span; +#ifdef XRPL_ENABLE_TELEMETRY + span = std::make_shared(txProcessSpan(transaction->getID())); +#endif + // Guarded on the span being live because these values are not free and this + // runs for every submitted and relayed transaction: the hash string + // allocates, and the open-ledger index takes the ledger master's lock. With + // telemetry compiled out the span is null; with it compiled in the block is + // skipped when telemetry is disabled at runtime or the transaction category + // is off. + if (span && *span) { - if (auto const* fmt = TxFormats::getInstance().findByType(stx->getTxnType())) - span->setAttribute(tx_span::attr::txType, fmt->getName().c_str()); + span->setAttribute(tx_span::attr::txHash, to_string(transaction->getID()).c_str()); + span->setAttribute(tx_span::attr::local, bLocal); + // The current (open) ledger index at submission/relay time — the ledger + // being worked on. Correlates this tx.process to the ledger trace; the + // tx has not yet been applied to a specific ledger, so there is no hash. span->setAttribute( - tx_span::attr::fee, static_cast(stx->getFieldAmount(sfFee).xrp().drops())); - span->setAttribute( - tx_span::attr::sequence, static_cast(stx->getSeqProxy().value())); + tx_span::attr::currentLedgerSeq, + static_cast(ledgerMaster_.getCurrentLedgerIndex())); + if (auto const& stx = transaction->getSTransaction()) + { + if (auto const* fmt = TxFormats::getInstance().findByType(stx->getTxnType())) + span->setAttribute(tx_span::attr::txType, fmt->getName().c_str()); + span->setAttribute( + tx_span::attr::fee, static_cast(stx->getFieldAmount(sfFee).xrp().drops())); + span->setAttribute( + tx_span::attr::sequence, static_cast(stx->getSeqProxy().value())); + } } auto ev = jobQueue_.makeLoadEvent(JtTxnProc, "ProcessTXN"); @@ -1626,12 +1643,14 @@ NetworkOPsImp::processTransaction( if (bLocal) { - span->setAttribute(tx_span::attr::path, tx_span::val::sync); + if (span) + span->setAttribute(tx_span::attr::path, tx_span::val::sync); doTransactionSync(transaction, bUnlimited, failType, std::move(span)); } else { - span->setAttribute(tx_span::attr::path, tx_span::val::async); + if (span) + span->setAttribute(tx_span::attr::path, tx_span::val::async); doTransactionAsync(transaction, bUnlimited, failType, std::move(span)); } } @@ -2024,7 +2043,7 @@ NetworkOPsImp::apply(std::unique_lock& batchLock) // Inject the tx.process span's trace context so the // receiving node can link its tx.receive span as a child. if (e.span && *e.span) - telemetry::injectSpanContext(*e.span, *tx.mutable_trace_context()); + telemetry::injectSpanContext(*e.span, tx); // FIXME: This should be when we received it registry_.get().getOverlay().relay(e.transaction->getID(), tx, *toSkip); e.transaction->setBroadcast(); @@ -2922,8 +2941,14 @@ bool NetworkOPsImp::recvValidation(std::shared_ptr const& val, std::string const& source) { JLOG(journal_.trace()) << "recvValidation " << val->getLedgerHash() << " from " << source; +#ifdef XRPL_ENABLE_TELEMETRY + // One per validation received. The registry is always constructed, so the + // null test never short-circuits: without the guard every validation pays + // a virtual lookup and an out-of-line call whose body is empty. Nothing + // outside the validations_checked_total metric reads the counter. if (auto* mr = registry_.get().getMetricsRegistry()) mr->incrementValidationsChecked(); +#endif std::unique_lock lock(validationsMutex_); BypassAccept bypassAccept = BypassAccept::No; diff --git a/src/xrpld/app/misc/detail/TxQ.cpp b/src/xrpld/app/misc/detail/TxQ.cpp index fda5f264f4..cd415112a8 100644 --- a/src/xrpld/app/misc/detail/TxQ.cpp +++ b/src/xrpld/app/misc/detail/TxQ.cpp @@ -771,17 +771,25 @@ TxQ::apply( return ScopedSpanGuard( TraceCategory::Transactions, txq_span::prefix::txq, txq_span::op::enqueue); }(); - span.setAttribute(txq_span::attr::txHash, to_string(tx->getTransactionID()).c_str()); - if (auto const* fmt = TxFormats::getInstance().findByType(tx->getTxnType())) - span.setAttribute(txq_span::attr::txType, fmt->getName().c_str()); - // The ledger being worked on (open/tentative apply or in-flight consensus - // build) — correlates this enqueue to the ledger trace in every context. - span.setAttribute(txq_span::attr::currentLedgerSeq, static_cast(view.seq())); - span.setAttribute( - txq_span::attr::currentLedgerHash, to_string(view.header().parentHash).c_str()); - // Default outcome; overridden below on the direct-apply and queued paths. - // Every other early return leaves the tx rejected from the queue. - span.setAttribute(txq_span::attr::txqStatus, txq_span::val::rejected); + // Guarded on the span being recorded: this runs for every transaction and + // again for each one replayed on an open-ledger rebuild, and the two hash + // strings each allocate. The compiled-out guard's operator bool() is a + // literal false, so the block disappears in that build. + if (span) + { + span.setAttribute(txq_span::attr::txHash, to_string(tx->getTransactionID()).c_str()); + if (auto const* fmt = TxFormats::getInstance().findByType(tx->getTxnType())) + span.setAttribute(txq_span::attr::txType, fmt->getName().c_str()); + // The ledger being worked on (open/tentative apply or in-flight + // consensus build) — correlates this enqueue to the ledger trace in + // every context. + span.setAttribute(txq_span::attr::currentLedgerSeq, static_cast(view.seq())); + span.setAttribute( + txq_span::attr::currentLedgerHash, to_string(view.header().parentHash).c_str()); + // Default outcome; overridden below on the direct-apply and queued + // paths. Every other early return leaves the tx rejected from the queue. + span.setAttribute(txq_span::attr::txqStatus, txq_span::val::rejected); + } // See if the transaction is valid, properly formed, // etc. before doing potentially expensive queue @@ -1528,13 +1536,27 @@ TxQ::accept(Application& app, OpenView& view) ScopedSpanGuard txSpan( TraceCategory::Transactions, txq_span::prefix::txq, txq_span::op::acceptTx); - txSpan.setAttribute(txq_span::attr::txHash, to_string(candidateIter->txID).c_str()); - txSpan.setAttribute( - txq_span::attr::retriesRemaining, - static_cast(candidateIter->retriesRemaining)); + // Guarded on the span being recorded: this runs for every queued + // transaction the loop tries to apply to the open ledger, on every + // ledger close, and the hash string allocates 64 characters. The + // compiled-out guard's operator bool() is a literal false, so the + // block disappears in that build. + if (txSpan) + { + txSpan.setAttribute(txq_span::attr::txHash, to_string(candidateIter->txID).c_str()); + txSpan.setAttribute( + txq_span::attr::retriesRemaining, + static_cast(candidateIter->retriesRemaining)); + } auto const [txnResult, didApply, _metadata] = candidateIter->apply(app, view, j_); - txSpan.setAttribute(txq_span::attr::terCode, transToken(txnResult).c_str()); + // The apply above is the transaction itself and always runs; only + // the attribute reading its result is telemetry. transToken() is a + // lookup that builds a string, so it is guarded too. + if (txSpan) + { + txSpan.setAttribute(txq_span::attr::terCode, transToken(txnResult).c_str()); + } if (didApply) { diff --git a/src/xrpld/overlay/detail/PeerImp.cpp b/src/xrpld/overlay/detail/PeerImp.cpp index 20eeebb58e..5f0b327713 100644 --- a/src/xrpld/overlay/detail/PeerImp.cpp +++ b/src/xrpld/overlay/detail/PeerImp.cpp @@ -21,7 +21,6 @@ #include #include #include -#include #include #include #include @@ -75,7 +74,7 @@ #include #include #include -#include +#include #include #include #include @@ -116,6 +115,17 @@ #include #include +#ifdef XRPL_ENABLE_TELEMETRY +// The consensus receive-span factories are named only by the +// telemetry-enabled blocks in this file. +#include + +// The TMGetObjectByHash metric names and label keys. Every use of them sits in +// an XRPL_METRIC_* argument list, and those macros expand to nothing when +// telemetry is off, so the include is guarded like its uses. +#include +#endif + using namespace std::chrono_literals; namespace xrpl { @@ -1398,19 +1408,37 @@ PeerImp::handleTransaction( using namespace telemetry; // SpanGuard is thread-free (holds no Scope), so it is safe to hand to // a job-queue worker and end on that thread — no detach step is needed. - auto span = std::make_shared(txReceiveSpan(txID, *m)); - span->setAttribute(tx_span::attr::txHash, to_string(txID).c_str()); - span->setAttribute(tx_span::attr::peerId, static_cast(id_)); - // The current (open) ledger index when the relayed tx was received — - // the ledger being worked on. Correlates this tx.receive to the ledger - // trace; not yet applied to a specific ledger here, so no hash. - span->setAttribute( - tx_span::attr::currentLedgerSeq, - static_cast(app_.getLedgerMaster().getCurrentLedgerIndex())); - if (auto const* fmt = TxFormats::getInstance().findByType(stx->getTxnType())) - span->setAttribute(tx_span::attr::txType, fmt->getName().c_str()); - if (auto const version = getVersion(); !version.empty()) - span->setAttribute(tx_span::attr::peerVersion, version.c_str()); + // Left null when telemetry is compiled out: there is no span to own, so + // nothing is allocated for one. Every use below tests it, the job + // capture and activateIfLive() accept a null handle, and the transaction + // pipeline already takes a null span by default. Without this the + // make_shared allocated once per inbound transaction, duplicates + // included, to hold an empty object. + std::shared_ptr span; +#ifdef XRPL_ENABLE_TELEMETRY + span = std::make_shared(txReceiveSpan(txID, *m)); +#endif + // Guarded on the span being live because these values are not free and + // this runs for every inbound transaction, including duplicates: the + // hash string allocates, and the open-ledger index takes the ledger + // master's lock. With telemetry compiled out the span is null; with it + // compiled in the block is skipped when telemetry is disabled at runtime + // or the transaction category is off. + if (span && *span) + { + span->setAttribute(tx_span::attr::txHash, to_string(txID).c_str()); + span->setAttribute(tx_span::attr::peerId, static_cast(id_)); + // The current (open) ledger index when the relayed tx was received + // — the ledger being worked on. Correlates this tx.receive to the + // ledger trace; not yet applied to a specific ledger, so no hash. + span->setAttribute( + tx_span::attr::currentLedgerSeq, + static_cast(app_.getLedgerMaster().getCurrentLedgerIndex())); + if (auto const* fmt = TxFormats::getInstance().findByType(stx->getTxnType())) + span->setAttribute(tx_span::attr::txType, fmt->getName().c_str()); + if (auto const version = getVersion(); !version.empty()) + span->setAttribute(tx_span::attr::peerVersion, version.c_str()); + } // Note: suppressed and txStatus are set once at each exit path // (not as defaults here) to avoid OTel SDK attribute duplication. @@ -1434,7 +1462,8 @@ PeerImp::handleTransaction( */ if (stx->isFlag(tfInnerBatchTxn)) { - span->setAttribute(tx_span::attr::txStatus, tx_span::val::rejectedInnerBatch); + if (span) + span->setAttribute(tx_span::attr::txStatus, tx_span::val::rejectedInnerBatch); JLOG(pJournal_.warn()) << "Ignoring Network relayed Tx containing " "tfInnerBatchTxn (handleTransaction)."; fee_.update(resource::kFeeModerateBurdenPeer, "inner batch txn"); @@ -1447,11 +1476,13 @@ PeerImp::handleTransaction( if (!app_.getHashRouter().shouldProcess(txID, id_, flags, kTxInterval)) { - span->setAttribute(tx_span::attr::suppressed, true); + if (span) + span->setAttribute(tx_span::attr::suppressed, true); // we have seen this transaction recently if (any(flags & HashRouterFlags::BAD)) { - span->setAttribute(tx_span::attr::txStatus, tx_span::val::knownBad); + if (span) + span->setAttribute(tx_span::attr::txStatus, tx_span::val::knownBad); fee_.update(resource::kFeeUselessData, "known bad"); JLOG(pJournal_.debug()) << "Ignoring known bad tx " << txID; } @@ -1460,7 +1491,8 @@ PeerImp::handleTransaction( // Recently-seen but not flagged bad — this is the plain // duplicate-suppression path. Mark it explicitly so the // span never exits as "new". - span->setAttribute(tx_span::attr::txStatus, tx_span::val::suppressed); + if (span) + span->setAttribute(tx_span::attr::txStatus, tx_span::val::suppressed); // Erase only if the server has seen this tx. If the server // has not seen this tx then the tx could not have been @@ -1477,7 +1509,8 @@ PeerImp::handleTransaction( return; } - span->setAttribute(tx_span::attr::suppressed, false); + if (span) + span->setAttribute(tx_span::attr::suppressed, false); JLOG(pJournal_.debug()) << "Got tx " << txID; bool checkSignature = true; @@ -1502,12 +1535,14 @@ PeerImp::handleTransaction( if (app_.getLedgerMaster().getValidatedLedgerAge() > 4min) { - span->setAttribute(tx_span::attr::txStatus, tx_span::val::droppedNoSync); + if (span) + span->setAttribute(tx_span::attr::txStatus, tx_span::val::droppedNoSync); JLOG(pJournal_.trace()) << "No new transactions until synchronized"; } else if (app_.getJobQueue().getJobCount(JtTransaction) > app_.config().maxTransactions) { - span->setAttribute(tx_span::attr::txStatus, tx_span::val::droppedQueueFull); + if (span) + span->setAttribute(tx_span::attr::txStatus, tx_span::val::droppedQueueFull); overlay_.incJqTransOverflow(); JLOG(pJournal_.info()) << "Transaction queue is full"; } @@ -2107,26 +2142,39 @@ PeerImp::onMessage(std::shared_ptr const& m) // Create a receive span that links to the sender's trace context // (if propagated). shared_ptr keeps it alive across the job boundary. // The receive span is a thread-free SpanGuard handed to the job worker; - // no scope to strip. - auto consSpan = std::make_shared(telemetry::proposalReceiveSpan(set)); - consSpan->setAttribute(telemetry::consensus::span::attr::proposalTrusted, isTrusted); - consSpan->setAttribute( - telemetry::consensus::span::attr::round, static_cast(set.proposeseq())); - // First 16 hex chars (8 bytes) of each hash — enough to disambiguate - // peer positions and prior ledgers without exporting full 32-byte - // hashes on every receive event. - consSpan->setAttribute( - telemetry::consensus::span::attr::prevLedgerPrefix, - to_string(prevLedger).substr(0, 16).c_str()); - consSpan->setAttribute( - telemetry::consensus::span::attr::positionHashPrefix, - to_string(proposeHash).substr(0, 16).c_str()); + // no scope to strip. The handle stays empty when telemetry is compiled + // out, so nothing is allocated on a path every inbound proposal takes. + // The job body only carries the handle to hold the span alive, so an + // empty handle is safe there. + std::shared_ptr span; +#ifdef XRPL_ENABLE_TELEMETRY + span = std::make_shared(telemetry::proposalReceiveSpan(set)); +#endif + // Every attribute below exists only for the span, so the block is guarded + // on the span being live. Unguarded, each inbound proposal — trusted or + // not — builds two full hex strings and a substring of each, four string + // allocations no one reads. + if (span && *span) + { + span->setAttribute(telemetry::consensus::span::attr::proposalTrusted, isTrusted); + span->setAttribute( + telemetry::consensus::span::attr::round, static_cast(set.proposeseq())); + // First 16 hex chars (8 bytes) of each hash — enough to disambiguate + // peer positions and prior ledgers without exporting full 32-byte + // hashes on every receive event. + span->setAttribute( + telemetry::consensus::span::attr::prevLedgerPrefix, + to_string(prevLedger).substr(0, 16).c_str()); + span->setAttribute( + telemetry::consensus::span::attr::positionHashPrefix, + to_string(proposeHash).substr(0, 16).c_str()); + } std::weak_ptr const weak = shared_from_this(); app_.getJobQueue().addJob( isTrusted ? JtProposalT : JtProposalUt, "checkPropose", - [weak, isTrusted, m, proposal, sp = std::move(consSpan)]() { + [weak, isTrusted, m, proposal, sp = std::move(span)]() { if (auto peer = weak.lock()) peer->checkPropose(isTrusted, m, proposal); }); @@ -2633,8 +2681,19 @@ PeerImp::onMessage(std::shared_ptr const& m) } val->setSeen(closeTime); } - valSpan.setAttribute(peer_span::attr::ledgerHash, to_string(val->getLedgerHash()).c_str()); - valSpan.setAttribute(peer_span::attr::fullValidation, val->isFull()); + // setAttribute evaluates its arguments even when telemetry is compiled + // out, and to_string() heap-allocates a 64-character hex string. This + // runs before the duplicate check below, so without the guard every + // peer's copy of every validation pays for that string. The guard is + // false when telemetry is compiled out, switched off in the config, or + // the Peer trace category is disabled; a span that exists but was + // sampled out still pays. + if (valSpan) + { + valSpan.setAttribute( + peer_span::attr::ledgerHash, to_string(val->getLedgerHash()).c_str()); + valSpan.setAttribute(peer_span::attr::fullValidation, val->isFull()); + } if (!isCurrent( app_.getValidations().parms(), @@ -2692,20 +2751,33 @@ PeerImp::onMessage(std::shared_ptr const& m) // Create a receive span that links to the sender's trace context // (if propagated). shared_ptr keeps it alive across the job boundary. // The receive span is a thread-free SpanGuard handed to the job worker; - // no scope to strip. - auto consSpan = - std::make_shared(telemetry::validationReceiveSpan(*m)); - consSpan->setAttribute(telemetry::consensus::span::attr::validationTrusted, isTrusted); - if (val->isFieldPresent(sfLedgerSequence)) + // no scope to strip. The handle stays empty when telemetry is compiled + // out, so nothing is allocated on a path every inbound validation + // takes. The job body only carries the handle to hold the span alive, + // so an empty handle is safe there. + std::shared_ptr span; +#ifdef XRPL_ENABLE_TELEMETRY + span = std::make_shared(telemetry::validationReceiveSpan(*m)); +#endif + // Every attribute below exists only for the span, so the block is + // guarded on the span being live. Unguarded, each inbound validation + // pays the field lookups and time conversions here. The span is built + // before the drop decision below on purpose, so a dropped validation + // is still traced; the guard removes the cost, not the span. + if (span && *span) { - consSpan->setAttribute( - telemetry::consensus::span::attr::ledgerSeq, - static_cast(val->getFieldU32(sfLedgerSequence))); + span->setAttribute(telemetry::consensus::span::attr::validationTrusted, isTrusted); + if (val->isFieldPresent(sfLedgerSequence)) + { + span->setAttribute( + telemetry::consensus::span::attr::ledgerSeq, + static_cast(val->getFieldU32(sfLedgerSequence))); + } + span->setAttribute(telemetry::consensus::span::attr::fullValidation, val->isFull()); + span->setAttribute( + telemetry::consensus::span::attr::validationSignTime, + static_cast(val->getSignTime().time_since_epoch().count())); } - consSpan->setAttribute(telemetry::consensus::span::attr::fullValidation, val->isFull()); - consSpan->setAttribute( - telemetry::consensus::span::attr::validationSignTime, - static_cast(val->getSignTime().time_since_epoch().count())); if (!isTrusted && (tracking_.load() == Tracking::Diverged)) { @@ -2719,7 +2791,7 @@ PeerImp::onMessage(std::shared_ptr const& m) app_.getJobQueue().addJob( isTrusted ? JtValidationT : JtValidationUt, name, - [weak, val, m, key, sp = std::move(consSpan)]() { + [weak, val, m, key, sp = std::move(span)]() { if (auto peer = weak.lock()) peer->checkValidation(val, key, m); }); @@ -2929,8 +3001,10 @@ PeerImp::processGetObjectByHash(std::shared_ptr con // Time the whole loop once, not each iteration: the loop can run up to // kHardMaxReplyNodes times, so per-iteration clock reads would cost more - // than the lookups they measure. - auto const lookupStart = std::chrono::steady_clock::now(); + // than the lookups they measure. Both clock reads serve only the metric + // recorded below, so neither happens when telemetry is compiled out -- + // Stopwatch holds no state in that build. + telemetry::Stopwatch const lookupTimer; for (int i = 0; i < iterLimit; ++i) { @@ -2956,8 +3030,9 @@ PeerImp::processGetObjectByHash(std::shared_ptr con newObj.set_ledgerseq(obj.ledgerseq()); } - auto const lookupElapsed = std::chrono::duration_cast( - std::chrono::steady_clock::now() - lookupStart); + // Measured here rather than at the call below, which would fold the fee + // computation and charge() into the reported lookup latency. + auto const lookupElapsed = lookupTimer.elapsedUs(); // Apply work-proportional charge. `charge()` posts the disconnect // step (if any) back to strand_, so it is safe to call from this @@ -2972,12 +3047,20 @@ PeerImp::processGetObjectByHash(std::shared_ptr con resource::Charge const fee = computeGetObjectByHashFee(requested, reply.objects_size()); charge(fee, "processed get object by hash request"); + // Called unconditionally: every statement in the body is an XRPL_METRIC_* + // argument, and those macros discard their arguments when telemetry is + // compiled out. All four values here are already computed for the request + // itself, so passing them costs nothing. recordGetObjectMetrics(requested, reply.objects_size(), lookupElapsed, fee); JLOG(pJournal_.trace()) << "GetObj: " << reply.objects_size() << " of " << requested; send(std::make_shared(reply, protocol::mtGET_OBJECTS)); } +// Reads app_ through the metric macros when telemetry is compiled in and +// touches no member when it is not, so clang-tidy asks for it to be static. +// Making it static would give the two builds different signatures. +// NOLINTBEGIN(readability-convert-member-functions-to-static) void PeerImp::recordGetObjectMetrics( int const requested, @@ -3006,9 +3089,9 @@ PeerImp::recordGetObjectMetrics( // negative -- the counter takes an unsigned amount, where a wrap would // read as ~1.8e19 rather than as an error. // - // Written as two calls rather than a loop over a {hit, miss} pair: the - // macros expand to empty statements in a telemetry-off build, which would - // leave a loop's induction variable unused and fail the -Werror build. + // Written as two calls rather than a loop over a {hit, miss} pair: the two + // amounts come from different expressions, so there is no single value to + // iterate over. XRPL_METRIC_COUNTER_ADD_LABELED( app_, kGetObjectLookupsTotal, @@ -3024,6 +3107,8 @@ PeerImp::recordGetObjectMetrics( {{kLabelResult, std::string(kResultMiss)}}); } +// NOLINTEND(readability-convert-member-functions-to-static) + void PeerImp::onMessage(std::shared_ptr const& m) { diff --git a/src/xrpld/overlay/detail/PeerImp.h b/src/xrpld/overlay/detail/PeerImp.h index 942cd2a21d..02a58dc7c2 100644 --- a/src/xrpld/overlay/detail/PeerImp.h +++ b/src/xrpld/overlay/detail/PeerImp.h @@ -821,8 +821,9 @@ private: * * Records `getobject_request_objects`, `getobject_lookup_us`, * `getobject_charge`, and both label values of - * `getobject_lookups_total`. Compiles to nothing when telemetry is - * disabled, because the `XRPL_METRIC_*` macros do. + * `getobject_lookups_total`. Every statement is an `XRPL_METRIC_*` record, + * and those macros discard their arguments when telemetry is disabled, so + * the body costs nothing in that build and the call site needs no guard. * * @param requested Objects the peer asked for (`objects_size()`). * @param found Objects returned, i.e. the reply's object count. diff --git a/src/xrpld/perflog/detail/PerfLogImp.cpp b/src/xrpld/perflog/detail/PerfLogImp.cpp index 5573522687..4173ea1a7c 100644 --- a/src/xrpld/perflog/detail/PerfLogImp.cpp +++ b/src/xrpld/perflog/detail/PerfLogImp.cpp @@ -2,7 +2,12 @@ #include #include + +#ifdef XRPL_ENABLE_TELEMETRY +// Only the recording calls below and the metric macros' expansion name the +// registry, and neither survives with telemetry compiled out. #include +#endif #include #include @@ -339,8 +344,10 @@ PerfLogImp::rpcStart(std::string const& method, std::uint64_t const requestId) // above are released: the OTel call path allocates and takes locks // inside the SDK, so holding methodsMutex across it would widen a // process-wide critical section for no reason. Mirrors rpcEnd(). +#ifdef XRPL_ENABLE_TELEMETRY if (auto* mr = app_.getMetricsRegistry()) mr->recordRpcStarted(method); +#endif // A value that must be able to decrease (UpDownCounter), added at its // call site with no MetricsRegistry member/init-line/method. Paired with @@ -398,6 +405,7 @@ PerfLogImp::rpcEnd(std::string const& method, std::uint64_t const requestId, boo // Record RPC completion in OTel metrics pipeline. Mirrors the // rpcStart() instrumentation so the finished/errored counters and // duration histogram advance with every call. +#ifdef XRPL_ENABLE_TELEMETRY if (auto* mr = app_.getMetricsRegistry()) { if (finish) @@ -409,6 +417,7 @@ PerfLogImp::rpcEnd(std::string const& method, std::uint64_t const requestId, boo mr->recordRpcErrored(method, durationUs.count()); } } +#endif // Matching -1 for the +1 recorded in rpcStart(). Placed after the early // returns above so it runs only when this request's methods-map entry was @@ -432,10 +441,17 @@ PerfLogImp::jobQueue(JobType const type, std::string const& name) ++counter->second.value.queued; } +#ifdef XRPL_ENABLE_TELEMETRY // Record job enqueue in OTel metrics pipeline, after the lock above is // released so the SDK's work stays outside the critical section. + // + // Guarded because this runs three times per job, counting jobStart and + // jobFinish below. recordJobQueued's body is compiled out, but the call is + // not: the virtual getMetricsRegistry() and the JobTypes::name() map lookup + // that builds its argument both still happen. if (auto* mr = app_.getMetricsRegistry()) mr->recordJobQueued(JobTypes::name(type), name); +#endif } void @@ -470,8 +486,10 @@ PerfLogImp::jobStart( // released. jobsMutex is process-wide and taken by every worker thread // on every job, so the SDK's allocation and internal locking must not // run inside it. +#ifdef XRPL_ENABLE_TELEMETRY if (auto* mr = app_.getMetricsRegistry()) mr->recordJobStarted(JobTypes::name(type), name, dur.count()); +#endif } void @@ -499,8 +517,10 @@ PerfLogImp::jobFinish(JobType const type, std::string const& name, microseconds // Record job finish in OTel metrics pipeline, after the locks above // are released, for the same reason as jobStart(). +#ifdef XRPL_ENABLE_TELEMETRY if (auto* mr = app_.getMetricsRegistry()) mr->recordJobFinished(JobTypes::name(type), name, dur.count()); +#endif } void diff --git a/src/xrpld/rpc/detail/PathRequest.cpp b/src/xrpld/rpc/detail/PathRequest.cpp index 2ac8ed9515..9d9c1e290c 100644 --- a/src/xrpld/rpc/detail/PathRequest.cpp +++ b/src/xrpld/rpc/detail/PathRequest.cpp @@ -602,7 +602,14 @@ PathRequest::findPaths( span.setAttribute( pathfind_span::attr::numSourceAssets, static_cast(sourceAssets.size())); +#ifdef XRPL_ENABLE_TELEMETRY + // Only the numPaths attribute at the end of this function reads this, so it + // is not maintained at all when telemetry is compiled out. An #ifdef rather + // than `if (span)`, because the attribute cannot read a variable that does + // not exist, and a counter kept up to date but never read is an unused + // variable, which fails the build. std::int64_t totalPaths = 0; +#endif for (auto const& asset : sourceAssets) { if (continueCallback && !continueCallback()) @@ -622,7 +629,9 @@ PathRequest::findPaths( auto ps = pathfinder->getBestPaths( kMaxPaths, fullLiquidityPath, context_[asset], asset.getIssuer(), continueCallback); context_[asset] = ps; +#ifdef XRPL_ENABLE_TELEMETRY totalPaths += static_cast(ps.size()); +#endif auto const& sourceAccount = [&] { if (!isXRP(asset.getIssuer())) @@ -725,7 +734,9 @@ PathRequest::findPaths( } } +#ifdef XRPL_ENABLE_TELEMETRY span.setAttribute(pathfind_span::attr::numPaths, totalPaths); +#endif /* The resource fee is based on the number of source currencies used. The minimum cost is 50 and the maximum is 400. The cost increases @@ -748,21 +759,34 @@ PathRequest::doUpdate( // nests under it. doUpdate does not yield, so scoping is safe. auto span = ScopedSpanGuard( TraceCategory::Rpc, pathfind_span::prefix::pathfind, pathfind_span::op::compute); - span.setAttribute(pathfind_span::attr::fast, fast); - // to_string(Issue) renders a non-XRP asset as "/" with the - // issuer as a plaintext Base58 address, so it cannot be emitted as-is: every - // account reaching a span is hashed first. Redact just the issuer and keep - // the currency, which is what this attribute is for. An MPT asset renders as - // its issuance ID and carries no address, so it needs no redaction. - span.setAttribute( - pathfind_span::attr::destCurrency, - saDstAmount_.asset().visit( - [](Issue const& issue) { - return isXRP(issue.account) - ? to_string(issue.currency) - : redactAccount(toBase58(issue.account)) + "/" + to_string(issue.currency); - }, - [](MPTIssue const& mpt) { return to_string(mpt.getMptID()); })); + // Guarded on the span being live because setAttribute's arguments are + // evaluated whatever the build, and doUpdate is hot: PathRequestManager + // calls it once per active path_find subscription on every ledger close, so + // a node with N subscriptions pays this N times a close. The destCurrency + // value costs a base58check encode of the issuer (two SHA-256 rounds), a + // SHA-512Half over the result and three string allocations. The compiled-out + // guard's operator bool() is a literal false, so the block disappears + // entirely in that build; with telemetry compiled in it is skipped when + // telemetry is disabled at runtime or the pathfind category is off. + if (span) + { + span.setAttribute(pathfind_span::attr::fast, fast); + // to_string(Issue) renders a non-XRP asset as "/" with + // the issuer as a plaintext Base58 address, so it cannot be emitted + // as-is: every account reaching a span is hashed first. Redact just the + // issuer and keep the currency, which is what this attribute is for. An + // MPT asset renders as its issuance ID and carries no address, so it + // needs no redaction. + span.setAttribute( + pathfind_span::attr::destCurrency, + saDstAmount_.asset().visit( + [](Issue const& issue) { + return isXRP(issue.account) + ? to_string(issue.currency) + : redactAccount(toBase58(issue.account)) + "/" + to_string(issue.currency); + }, + [](MPTIssue const& mpt) { return to_string(mpt.getMptID()); })); + } JLOG(journal_.debug()) << iIdentifier_ << " update " << (fast ? "fast" : "normal"); diff --git a/src/xrpld/rpc/detail/PathRequestManager.cpp b/src/xrpld/rpc/detail/PathRequestManager.cpp index 794b03124e..07345e325a 100644 --- a/src/xrpld/rpc/detail/PathRequestManager.cpp +++ b/src/xrpld/rpc/detail/PathRequestManager.cpp @@ -3,7 +3,6 @@ #include #include #include -#include #include #include @@ -16,17 +15,26 @@ #include #include #include -#include #include #include #include #include #include -#include #include #include +// Needed only by the update_all span in updateAll(), which is compiled out when +// telemetry is off. Without the same guard here they would be unused includes in +// that build, which clang-tidy's misc-include-cleaner rejects. +#ifdef XRPL_ENABLE_TELEMETRY +#include + +#include + +#include +#endif // XRPL_ENABLE_TELEMETRY + namespace xrpl { /** @@ -75,12 +83,16 @@ PathRequestManager::updateAll(std::shared_ptr const& inLedger) cache = getAssetCache(inLedger, true); } +#ifdef XRPL_ENABLE_TELEMETRY using namespace telemetry; - // updateAll runs on every ledger close. Skip span emission when there are - // no active path subscriptions, to avoid a steady stream of empty spans at - // mainnet close cadence. All other work still runs unchanged (notably the - // isNewPathRequest() flag reset below), so behaviour matches the pre-span - // code path. + // Nothing outside telemetry reads this block, and updateAll runs on every + // ledger close, so it is compiled out entirely when telemetry is off rather + // than left to construct a stub guard and discard two attributes per close. + // No span object exists to test here, so the guard has to be an #ifdef. + // + // Skip span emission when there are no active path subscriptions, to avoid + // a steady stream of empty spans at mainnet close cadence. All other work + // still runs unchanged (notably the isNewPathRequest() flag reset below). // // Scoped, so the pathfind.compute spans that doUpdate() creates below on // this thread nest under it. std::optional because ScopedSpanGuard is @@ -93,6 +105,7 @@ PathRequestManager::updateAll(std::shared_ptr const& inLedger) span->setAttribute(pathfind_span::attr::ledgerIndex, static_cast(inLedger->seq())); span->setAttribute(pathfind_span::attr::numRequests, static_cast(requests.size())); } +#endif // XRPL_ENABLE_TELEMETRY bool newRequests = app_.getLedgerMaster().isNewPathRequest(); bool mustBreak = false; diff --git a/src/xrpld/rpc/detail/RPCHandler.cpp b/src/xrpld/rpc/detail/RPCHandler.cpp index df8c795c3f..e3400b746f 100644 --- a/src/xrpld/rpc/detail/RPCHandler.cpp +++ b/src/xrpld/rpc/detail/RPCHandler.cpp @@ -223,6 +223,14 @@ callMethod(JsonContext& context, Method method, std::string const& name, Object& } } +// Telemetry-only helper, so it is compiled out with telemetry off. Left +// running it would cost several json lookups, up to two string copies and a +// handler-table lookup on every failed request, for a name nobody records. +// Its result IS the span name, so `if (span)` cannot guard it: at that point +// no span exists to test. The single call site is gated the same way, which +// also keeps this file-static function referenced in both configurations. +#ifdef XRPL_ENABLE_TELEMETRY + // Resolve the span suffix / command attribute for a request that failed in // fillHandler. Returns the canonical handler name for a recognized command // (a finite, bounded set) or the literal "unknown" for a request that omits @@ -256,6 +264,8 @@ resolveCommandSpanName(JsonContext const& context) : std::string_view{rpc_span::val::unknownCommand}; } +#endif // XRPL_ENABLE_TELEMETRY + } // namespace Status @@ -264,6 +274,11 @@ doCommand(rpc::JsonContext& context, json::Value& result) Handler const* handler = nullptr; if (auto error = fillHandler(context, handler)) { + // Every statement below only feeds the error span, and the span name + // itself comes from resolveCommandSpanName(), so there is no span + // object to test with `if (span)`. With telemetry off the whole block + // is compiled out and a storm of malformed requests pays nothing. +#ifdef XRPL_ENABLE_TELEMETRY // Bound the span name and command attribute to the finite set of // registered handler names (plus "unknown") — see the helper for why // raw request input must not reach the telemetry pipeline. @@ -279,6 +294,7 @@ doCommand(rpc::JsonContext& context, json::Value& result) : std::string_view(rpc_span::val::user)); span.setAttribute(rpc_span::attr::rpcStatus, rpc_span::val::error); span.setError(getErrorInfo(error).token.cStr()); +#endif // XRPL_ENABLE_TELEMETRY injectError(error, result); return error; diff --git a/src/xrpld/rpc/detail/ServerHandler.cpp b/src/xrpld/rpc/detail/ServerHandler.cpp index ebf318b247..3c797d3fcb 100644 --- a/src/xrpld/rpc/detail/ServerHandler.cpp +++ b/src/xrpld/rpc/detail/ServerHandler.cpp @@ -483,7 +483,15 @@ ServerHandler::processSession( // else collapses to "unknown". Emitting the raw string would let request // input drive unbounded label cardinality. Mirrors the HTTP path's // resolveCommandSpanName(). - span.setAttribute(rpc_span::attr::command, resolveWsCommandSpanName(jv, app_.config())); + // + // The guard is required because the resolver is a call argument: it runs + // even when setAttribute itself is an empty no-op. Without it, every + // WebSocket message pays for the JSON member lookups, a string copy and a + // handler-registry lookup that nothing reads. + if (span) + { + span.setAttribute(rpc_span::attr::command, resolveWsCommandSpanName(jv, app_.config())); + } auto is = std::static_pointer_cast(session->appDefined); if (is->getConsumer().disconnect(journal_)) diff --git a/src/xrpld/rpc/handlers/orderbook/PathFind.cpp b/src/xrpld/rpc/handlers/orderbook/PathFind.cpp index 6da4f2cdca..44830619e4 100644 --- a/src/xrpld/rpc/handlers/orderbook/PathFind.cpp +++ b/src/xrpld/rpc/handlers/orderbook/PathFind.cpp @@ -25,16 +25,27 @@ doPathFind(rpc::JsonContext& context) // thread) nest under it. doPathFind does not yield, so scoping is safe. auto span = ScopedSpanGuard( TraceCategory::Rpc, pathfind_span::prefix::pathfind, pathfind_span::op::request); - // Addresses are hashed before emission for privacy. Read through a const - // reference: the non-const json::Value::operator[] inserts a null for a - // missing key, which would make PathRequest::parseJson's isMember() checks - // see an absent field as present and return Malformed instead of Missing. - // Reading for telemetry must not alter what the request looks like. - auto const& params = std::as_const(context.params); - if (auto const& src = params[jss::source_account]; src.isString()) - span.setAttribute(pathfind_span::attr::sourceAccount, redactAccount(src.asString())); - if (auto const& dst = params[jss::destination_account]; dst.isString()) - span.setAttribute(pathfind_span::attr::destAccount, redactAccount(dst.asString())); + // Guarded on the span being live because setAttribute's arguments are + // evaluated whatever the build, and neither is free: asString() copies the + // address out of the JSON and redactAccount() takes a SHA-512Half over it. + // That is two copies and two hashes on every path_find call. The + // compiled-out guard's operator bool() is a literal false, so the block + // disappears entirely in that build; with telemetry compiled in it is + // skipped when telemetry is disabled at runtime or the category is off. + if (span) + { + // Addresses are hashed before emission for privacy. Read through a + // const reference: the non-const json::Value::operator[] inserts a null + // for a missing key, which would make PathRequest::parseJson's + // isMember() checks see an absent field as present and return Malformed + // instead of Missing. Reading for telemetry must not alter what the + // request looks like. + auto const& params = std::as_const(context.params); + if (auto const& src = params[jss::source_account]; src.isString()) + span.setAttribute(pathfind_span::attr::sourceAccount, redactAccount(src.asString())); + if (auto const& dst = params[jss::destination_account]; dst.isString()) + span.setAttribute(pathfind_span::attr::destAccount, redactAccount(dst.asString())); + } if (context.app.config().pathSearchMax == 0) return rpcError(RpcNotSupported); diff --git a/src/xrpld/rpc/handlers/orderbook/RipplePathFind.cpp b/src/xrpld/rpc/handlers/orderbook/RipplePathFind.cpp index 1b5e881d2b..49b90a6e4c 100644 --- a/src/xrpld/rpc/handlers/orderbook/RipplePathFind.cpp +++ b/src/xrpld/rpc/handlers/orderbook/RipplePathFind.cpp @@ -34,16 +34,27 @@ doRipplePathFind(rpc::JsonContext& context) // span's log lines stay trace-correlated. auto span = ScopedSpanGuard( TraceCategory::Rpc, pathfind_span::prefix::pathfind, pathfind_span::op::request); - // Addresses are hashed before emission for privacy. Read through a const - // reference: the non-const json::Value::operator[] inserts a null for a - // missing key, which would make PathRequest::parseJson's isMember() checks - // see an absent field as present and return Malformed instead of Missing. - // Reading for telemetry must not alter what the request looks like. - auto const& params = std::as_const(context.params); - if (auto const& src = params[jss::source_account]; src.isString()) - span.setAttribute(pathfind_span::attr::sourceAccount, redactAccount(src.asString())); - if (auto const& dst = params[jss::destination_account]; dst.isString()) - span.setAttribute(pathfind_span::attr::destAccount, redactAccount(dst.asString())); + // Guarded on the span being live because setAttribute's arguments are + // evaluated whatever the build, and neither is free: asString() copies the + // address out of the JSON and redactAccount() takes a SHA-512Half over it. + // That is two copies and two hashes on every ripple_path_find call. The + // compiled-out guard's operator bool() is a literal false, so the block + // disappears entirely in that build; with telemetry compiled in it is + // skipped when telemetry is disabled at runtime or the category is off. + if (span) + { + // Addresses are hashed before emission for privacy. Read through a + // const reference: the non-const json::Value::operator[] inserts a null + // for a missing key, which would make PathRequest::parseJson's + // isMember() checks see an absent field as present and return Malformed + // instead of Missing. Reading for telemetry must not alter what the + // request looks like. + auto const& params = std::as_const(context.params); + if (auto const& src = params[jss::source_account]; src.isString()) + span.setAttribute(pathfind_span::attr::sourceAccount, redactAccount(src.asString())); + if (auto const& dst = params[jss::destination_account]; dst.isString()) + span.setAttribute(pathfind_span::attr::destAccount, redactAccount(dst.asString())); + } if (context.app.config().pathSearchMax == 0) return rpcError(RpcNotSupported); diff --git a/src/xrpld/telemetry/ConsensusReceiveTracing.h b/src/xrpld/telemetry/ConsensusReceiveTracing.h index 0a1a4458dc..498eaa6fca 100644 --- a/src/xrpld/telemetry/ConsensusReceiveTracing.h +++ b/src/xrpld/telemetry/ConsensusReceiveTracing.h @@ -42,9 +42,14 @@ #include #include #include + +#ifdef XRPL_ENABLE_TELEMETRY +// The trace-context validator and std::uint8_t are named only by the +// telemetry-enabled branches below. #include #include +#endif namespace xrpl::telemetry { diff --git a/src/xrpld/telemetry/MetricsRegistry.cpp b/src/xrpld/telemetry/MetricsRegistry.cpp index 8642d31b23..7373218c6d 100644 --- a/src/xrpld/telemetry/MetricsRegistry.cpp +++ b/src/xrpld/telemetry/MetricsRegistry.cpp @@ -22,6 +22,10 @@ #include +// Unguarded because the constructor's `beast::Journal journal` parameter is +// declared in both configurations; only the member it initialises is guarded. +#include + #ifdef XRPL_ENABLE_TELEMETRY // The app and overlay includes below are why @@ -55,7 +59,6 @@ #include #include #include -#include #include #include #include diff --git a/src/xrpld/telemetry/MetricsRegistry.h b/src/xrpld/telemetry/MetricsRegistry.h index 5609a58982..acc9a4110a 100644 --- a/src/xrpld/telemetry/MetricsRegistry.h +++ b/src/xrpld/telemetry/MetricsRegistry.h @@ -148,16 +148,17 @@ * instrumentation site. */ +#ifdef XRPL_ENABLE_TELEMETRY +// The tracker is held and exposed only in this configuration, where the gauge +// callbacks that drain it exist. #include +#endif #include #include -#include #include -#include #include -#include #include #include #include @@ -168,6 +169,13 @@ #include #include #include + +// These three serve only the telemetry-only members below, so they are guarded +// like their uses: std::atomic by callbacksDetached_, std::function by the +// ObserveFn sink, std::shared_ptr by provider_. +#include +#include +#include #endif namespace xrpl { @@ -653,10 +661,15 @@ public: void incrementTxqDropped(std::string_view reason); +#ifdef XRPL_ENABLE_TELEMETRY /** * Access the validation agreement tracker. * Used by consensus and ledger hooks to record our validations and * network validations so the tracker can compute agreement percentages. + * + * Guarded, along with the tracker itself, because only the observable-gauge + * callbacks read it and those exist only in this configuration. Recording + * into it is not free: each call takes its lock and inserts an entry. * @return Reference to the internal ValidationTracker instance. */ ValidationTracker& @@ -665,7 +678,6 @@ public: return validationTracker_; } -#ifdef XRPL_ENABLE_TELEMETRY /** * Access the shared OTel Meter for call-site instrument creation. * Used by the XRPL_METRIC_* macros (MetricMacros.h) so new synchronous @@ -765,15 +777,17 @@ private: */ bool const enabled_; +#ifdef XRPL_ENABLE_TELEMETRY /** * Tracks validation agreement between this node and the network. - * Lives outside the XRPL_ENABLE_TELEMETRY guard because it is - * always safe to record events; the gauge callback simply won't - * fire when telemetry is disabled. + * + * Guarded because reconcile() -- which resolves and then prunes recorded + * events -- runs only from the observable-gauge callbacks. Recording + * without it accumulates one entry per validated ledger, so the tracker + * exists only where something drains it. */ ValidationTracker validationTracker_; -#ifdef XRPL_ENABLE_TELEMETRY /** * Reference to Application services for gauge callbacks. * Only needed when OTel is compiled in, since observable gauge diff --git a/src/xrpld/telemetry/PropagationHelpers.h b/src/xrpld/telemetry/PropagationHelpers.h index ce64733d82..0497ef0826 100644 --- a/src/xrpld/telemetry/PropagationHelpers.h +++ b/src/xrpld/telemetry/PropagationHelpers.h @@ -13,16 +13,29 @@ * +--- TraceBytes -----+ * | | * injectSpanContext(span, proto) + * ^ + * | delegates, once there is something to write + * injectSpanContext(span, message) <-- preferred entry point * - * @note When XRPL_ENABLE_TELEMETRY is disabled, getTraceBytes() returns - * {.valid=false}, so injectSpanContext becomes a no-op with zero overhead. + * SpanGuard::hasCurrentContext() + * | gates + * v + * injectCurrentContext(message) <-- no span handle needed + * | + * +--> SpanGuard::injectCurrentContextToProtobuf(proto) + * + * @note Prefer the helpers that take the whole message. They are true + * no-ops when nothing is recorded, because they decide whether to create + * the TraceContext submessage at all. The overload taking a + * protocol::TraceContext& cannot be: its caller has already created the + * submessage and set its has-bit before this code runs. * * Usage: * @code * // Send side — inject from a SpanGuard reference: * protocol::TMTransaction tx; * // ... populate tx fields ... - * injectSpanContext(mySpanGuard, *tx.mutable_trace_context()); + * injectSpanContext(mySpanGuard, tx); * overlay.relay(txID, tx, toSkip); * @endcode * @@ -60,4 +73,55 @@ injectSpanContext(SpanGuard const& span, protocol::TraceContext& proto) proto.set_trace_flags(bytes.traceFlags); } +/** + * Inject an active span's trace context into a message that carries an + * optional TraceContext submessage. + * + * Takes the parent message rather than the submessage so the decision to + * create the submessage stays here. `mutable_trace_context()` on a protobuf + * optional field allocates the submessage and sets its has-bit, so a caller + * that passes `*msg.mutable_trace_context()` puts an empty TraceContext on + * the wire whenever nothing is recorded, and makes receiving peers take + * their has_trace_context() branch for nothing. + * + * @param span The span whose context to propagate; may be inactive. + * @param msg The message to populate. Untouched when nothing is recorded. + */ +template +void +injectSpanContext(SpanGuard const& span, Message& msg) +{ + auto const bytes = span.getTraceBytes(); + if (!bytes.valid) + return; + + injectSpanContext(span, *msg.mutable_trace_context()); +} + +/** + * Inject this thread's currently active span context into a message that + * carries an optional TraceContext submessage. + * + * For senders that have no SpanGuard in hand and want whatever span is + * ambient on the calling thread, such as the consensus round span. + * + * Tests for a live context before touching the message, so a build with + * telemetry compiled out, disabled by config, or simply not tracing this + * round sends no TraceContext submessage at all. Calling + * mutable_trace_context() unconditionally would create it and set its + * has-bit, putting an empty TraceContext on the wire and making every + * receiving peer take its has_trace_context() branch for nothing. + * + * @param msg The message to populate. Untouched when no span is active. + */ +template +void +injectCurrentContext(Message& msg) +{ + if (!SpanGuard::hasCurrentContext()) + return; + + SpanGuard::injectCurrentContextToProtobuf(*msg.mutable_trace_context()); +} + } // namespace xrpl::telemetry diff --git a/src/xrpld/telemetry/TxTracing.h b/src/xrpld/telemetry/TxTracing.h index 02900fa9b4..5baf01df2d 100644 --- a/src/xrpld/telemetry/TxTracing.h +++ b/src/xrpld/telemetry/TxTracing.h @@ -16,9 +16,14 @@ #include #include #include + +#ifdef XRPL_ENABLE_TELEMETRY +// The span-id validator and std::uint8_t are named only by the +// telemetry-enabled branches below. #include #include +#endif namespace xrpl::telemetry { diff --git a/src/xrpld/telemetry/ValidationTracker.h b/src/xrpld/telemetry/ValidationTracker.h index ac80f5cdf5..61f711602f 100644 --- a/src/xrpld/telemetry/ValidationTracker.h +++ b/src/xrpld/telemetry/ValidationTracker.h @@ -253,6 +253,17 @@ public: uint64_t totalValidationsChecked() const; + /** + * Number of ledgers currently held awaiting reconciliation. + * + * Never exceeds kMaxPendingEvents: the record methods enforce that bound + * as they insert, so the map stays bounded whether or not anything ever + * reconciles or reads it. + * @return Size of the pending map. + */ + [[nodiscard]] std::size_t + pendingCount() const; + /** @} */ private: @@ -397,6 +408,26 @@ private: void evictOldPending(TimePoint now); + /** + * Hold pending_ at kMaxPendingEvents by dropping its oldest entry. + * + * Called on the insert path, because that is the only place the bound can + * be guaranteed. reconcile() also prunes, but it runs only while the gauge + * callbacks are registered, which needs telemetry both compiled in and + * enabled -- so a node with telemetry off, or with [telemetry] enabled=0, + * would otherwise grow this map by one entry per validated ledger forever. + * + * Drops the oldest entry rather than the least useful one: the map is + * unordered, so this is a linear scan, but it runs at most once per + * recorded validation and only once the map is already full. + * + * @param justRecorded Hash inserted by the caller, kept even if the scan + * finds it oldest (equal timestamps make that possible). + * @note Caller must hold mutex_. + */ + void + boundPending(uint256 const& justRecorded); + /** * Scan a window deque and flip the first non-agreed entry matching * the given ledger hash to agreed. diff --git a/src/xrpld/telemetry/detail/ValidationTracker.cpp b/src/xrpld/telemetry/detail/ValidationTracker.cpp index c7f9c599bc..89e74c8f56 100644 --- a/src/xrpld/telemetry/detail/ValidationTracker.cpp +++ b/src/xrpld/telemetry/detail/ValidationTracker.cpp @@ -31,6 +31,7 @@ ValidationTracker::recordOurValidation(uint256 const& ledgerHash, LedgerIndex se } evt.weValidated = true; totalValidationsSent_.fetch_add(1, std::memory_order_relaxed); + boundPending(ledgerHash); } void @@ -46,6 +47,33 @@ ValidationTracker::recordNetworkValidation(uint256 const& ledgerHash, LedgerInde } evt.networkValidated = true; totalValidationsChecked_.fetch_add(1, std::memory_order_relaxed); + boundPending(ledgerHash); +} + +void +ValidationTracker::boundPending(uint256 const& justRecorded) +{ + if (pending_.size() <= kMaxPendingEvents) + return; + + auto oldest = pending_.end(); + for (auto it = pending_.begin(); it != pending_.end(); ++it) + { + if (it->first == justRecorded) + continue; + if (oldest == pending_.end() || it->second.recordTime < oldest->second.recordTime) + oldest = it; + } + + if (oldest != pending_.end()) + pending_.erase(oldest); +} + +std::size_t +ValidationTracker::pendingCount() const +{ + std::scoped_lock const lock(mutex_); + return pending_.size(); } void