diff --git a/.github/scripts/otel-naming/test_check_otel_naming.py b/.github/scripts/otel-naming/test_check_otel_naming.py index c24bc063b4..4fd225ad6e 100644 --- a/.github/scripts/otel-naming/test_check_otel_naming.py +++ b/.github/scripts/otel-naming/test_check_otel_naming.py @@ -738,28 +738,54 @@ class RuleDDashboards(unittest.TestCase): ["bogus_label"], ) + # Rule D returns early on an empty L1 set, so a test that passes one asserts + # nothing. Each test below passes a nonempty L1 set and puts `bogus_label` in + # the same expression as the labels under test: the assertion then pins both + # halves at once — the accepted labels are absent from the result, and Rule D + # demonstrably ran because it flagged the bad one. + def test_builtin_labels_not_flagged(self): self.assertEqual( - self._run('"expr": "sum by (le, span_name, exported_instance) (x)"', set()), - [], + self._run( + '"expr": "sum by (le, span_name, exported_instance, bogus_label) (x)"', + {"command"}, + ), + ["bogus_label"], ) def test_external_infra_labels_not_flagged(self): # EXTERNAL_INFRA_LABELS (perf-iac identity labels with no in-tree - # source) must be recognized as valid, distinct from `builtins`. - expr = "sum by (" + ", ".join(sorted(chk.EXTERNAL_INFRA_LABELS)) + ") (x)" - self.assertEqual(self._run(f'"expr": "{expr}"', set()), []) + # source) must be recognized as valid, distinct from `builtins`. The + # names are spelled out rather than joined from chk.EXTERNAL_INFRA_LABELS + # because building the query from the set that validates it passes for + # whatever that set happens to hold — including an empty one. + self.assertEqual( + self._run( + '"expr": "sum by (xrpl_branch, xrpl_node_role, bogus_label) (x)"', + {"command"}, + ), + ["bogus_label"], + ) def test_prometheus_name_label_not_flagged(self): # `__name__` is the Prometheus reserved metric-name label; the renamed # system-*.json dashboards use `sum by (le, __name__)`. self.assertEqual( - self._run('"expr": "sum by (le, __name__) (rate(x[5m]))"', set()), - [], + self._run( + '"expr": "sum by (le, __name__, bogus_label) (rate(x[5m]))"', + {"command"}, + ), + ["bogus_label"], ) def test_l1_label_passes(self): - self.assertEqual(self._run('"q": "{command=\\"x\\"}"', {"command"}), []) + # `by (...)` form, not a `{command="x"}` selector: a dashboard stores the + # query inside a JSON string, so its quotes are escaped on disk and the + # selector branch extracts nothing from them here. + self.assertEqual( + self._run('"expr": "sum by (command, bogus_label) (x)"', {"command"}), + ["bogus_label"], + ) def test_traceql_span_prefix_stripped(self): # `span.establish_count` must validate against the bare L1 key. @@ -772,7 +798,15 @@ class RuleDDashboards(unittest.TestCase): ) def test_traceql_resource_prefix_stripped(self): - self.assertEqual(self._run('"q": "{resource.service_name=\\"x\\"}"', set()), []) + # `resource.service_name` must validate against the bare builtin, same as + # the `span.` case above. + self.assertEqual( + self._run( + '"expr": "count_over_time(x) by (resource.service_name, bogus_label)"', + {"command"}, + ), + ["bogus_label"], + ) def test_native_metric_label_passes(self): # `job_type` / `reason` are emitted by MetricsRegistry, not span attrs. diff --git a/OpenTelemetryPlan/00-tracing-fundamentals.md b/OpenTelemetryPlan/00-tracing-fundamentals.md index 9c6f96d7af..1c7675243a 100644 --- a/OpenTelemetryPlan/00-tracing-fundamentals.md +++ b/OpenTelemetryPlan/00-tracing-fundamentals.md @@ -297,9 +297,9 @@ XRPL has a unique advantage: its core workflows produce **globally unique 256-bi Transaction: STTx::getTransactionID() → uint256 tid_ TMTransaction::rawTransaction → recompute hash from bytes -Consensus: ConsensusProposal::prevLedger_ → uint256 (previous ledger hash) - ConsensusProposal::position_ → uint256 (TxSet hash) - LedgerHeader::seq → uint32_t (ledger sequence) +Consensus: ConsensusProposal::previousLedger_ → uint256 (previous ledger hash) + ConsensusProposal::position_ → uint256 (TxSet hash) + LedgerHeader::seq → uint32_t (ledger sequence) Validation: STValidation::getLedgerHash() → uint256 STValidation::getNodeID() → NodeID (160-bit) diff --git a/OpenTelemetryPlan/02-design-decisions.md b/OpenTelemetryPlan/02-design-decisions.md index b3caffa54c..70c29093ad 100644 --- a/OpenTelemetryPlan/02-design-decisions.md +++ b/OpenTelemetryPlan/02-design-decisions.md @@ -236,17 +236,19 @@ keys (the dotted form is reserved for resource scope per §2.3.3). #### Transaction Attributes -| Key | Type | Description | -| -------------- | ------ | ------------------------------------- | -| `tx_hash` | string | Transaction hash (hex) | -| `tx_type` | string | `"Payment"`, `"OfferCreate"`, etc. | -| `tx_account` | string | Source account (redacted in prod) | -| `tx_sequence` | int64 | Account sequence number | -| `tx_fee` | int64 | Fee in drops | -| `tx_result` | string | `"tesSUCCESS"`, `"tecPATH_DRY"`, etc. | -| `ledger_index` | int64 | Ledger containing transaction | -| `relay_count` | int64 | Peers the transaction was relayed to | -| `suppressed` | bool | `true` when HashRouter dropped a dup | +| Key | Type | Description | +| -------------------- | ------ | ------------------------------------- | +| `tx_hash` | string | Transaction hash (hex) | +| `tx_type` | string | `"Payment"`, `"OfferCreate"`, etc. | +| `tx_account` | string | Source account (redacted in prod) | +| `tx_sequence` | int64 | Account sequence number | +| `tx_fee` | int64 | Fee in drops | +| `tx_result` | string | `"tesSUCCESS"`, `"tecPATH_DRY"`, etc. | +| `current_ledger_seq` | int64 | Open ledger the transaction targeted | +| `relay_count` | int64 | Peers the transaction was relayed to | +| `suppressed` | bool | `true` when HashRouter dropped a dup | + +> **Note:** `current_ledger_seq` and `ledger_seq` are the same concept — a ledger's sequence number — but they name different ledgers, so the design keeps two keys rather than one. `current_ledger_seq` is the open or in-flight ledger a transaction's work was applied into; it is named after the RPC field `ledger_current_index`. `ledger_seq` (see [Ledger & Job Attributes](#ledger--job-attributes)) is a closed or validated ledger, set by the ledger and consensus spans. Neither is spelled `ledger_index`: per rule 2 of [Telemetry span attribute naming](../CONTRIBUTING.md#telemetry-span-attribute-naming), one concept gets one key reused verbatim, and a different referent is disambiguated with a prefix rather than a synonym. #### Consensus Attributes @@ -289,7 +291,7 @@ keys (the dotted form is reserved for resource scope per §2.3.3). | Key | Type | Description | | --------------------------- | ------- | --------------------------------- | | `ledger_hash` | string | Ledger hash | -| `ledger_index` | int64 | Ledger sequence/index | +| `ledger_seq` | int64 | Closed/validated ledger sequence | | `close_time_ripple_epoch_s` | int64 | Close time (Ripple epoch seconds) | | `ledger_tx_count` | int64 | Transaction count | | `job_type` | string | Job type name | @@ -348,11 +350,11 @@ The following table summarizes what data is collected by category: | Category | Attributes Collected | Purpose | | --------------- | ---------------------------------------------------------------------------------------------------------------- | ---------------------------- | -| **Transaction** | `tx_hash`, `tx_type`, `tx_result`, `tx_fee`, `ledger_index` | Trace transaction lifecycle | +| **Transaction** | `tx_hash`, `tx_type`, `tx_result`, `tx_fee`, `current_ledger_seq` | Trace transaction lifecycle | | **Consensus** | `consensus_round`, `consensus_phase`, `consensus_mode`, `proposers`, `round_time_ms` | Analyze consensus timing | | **RPC** | `command`, `version`, `rpc_status`, `duration_ms` | Monitor RPC performance | | **Peer** | `peer_id` (public key), `peer_latency_ms`, `message_type`, `message_size_bytes` | Network topology analysis | -| **Ledger** | `ledger_hash`, `ledger_index`, `close_time`, `ledger_tx_count` | Ledger progression tracking | +| **Ledger** | `ledger_hash`, `ledger_seq`, `close_time`, `ledger_tx_count` | Ledger progression tracking | | **Job** | `job_type`, `job_queue_ms`, `job_worker` | JobQueue performance | | **PathFinding** | `pathfind_fast`, `pathfind_search_level`, `pathfind_num_paths`, `pathfind_ledger_index`, `pathfind_num_requests` | Payment path analysis | | **TxQ** | `txq_queue_depth`, `txq_fee_level`, `txq_eviction_reason` | Queue depth and fee tracking | @@ -589,9 +591,12 @@ A PerfLog entry is a JSON object with fields such as `time`, `method`, - No request-level detail - No causal relationships - Single-node perspective + - Aggregation happens on the StatsD server, not in the process In xrpld, Beast Insight is used through `increment` (counters), `gauge` -(point-in-time values), and `timing` (durations) calls. +(point-in-time values), and `timing` (durations) calls. A `timing` call sends +each measured value as its own raw `|ms` sample, so the histogram a dashboard +reads is built by the StatsD server from that stream of values. #### OpenTelemetry (NEW) @@ -600,6 +605,7 @@ In xrpld, Beast Insight is used through `increment` (counters), `gauge` - **Cross-node correlation** via `trace_id` - Parent-child span relationships - Rich attributes per span + - A `Histogram` instrument that aggregates **at the point of measure** - Industry standard (CNCF) - **Limitations**: - Requires collector infrastructure @@ -609,6 +615,13 @@ A span is created via `startSpan` (e.g. `"tx.relay"`), annotated with attributes such as `tx_hash` and `peer_id`, and is automatically linked to its parent through the active context. +OpenTelemetry is not only spans. The same SDK offers a `Histogram` instrument, +and a `Record()` call folds the value straight into bucket counts inside the +process — no per-event record is shipped and no server-side aggregation step is +needed. That is what makes it affordable in a hot loop where one span per event +would not be, and it is the one thing Beast Insight cannot do, because its +`timing` path ships raw values and aggregates them on the StatsD server. + ### 2.6.3 When to Use Each | Scenario | PerfLog | StatsD | OpenTelemetry | @@ -619,6 +632,15 @@ parent through the active context. | "Which node delayed consensus?" | ❌ | ❌ | ✅ | | "What happened on node X at time T?" | ✅ | ❌ | ✅ | | "Show me the TX journey across 5 nodes" | ❌ | ❌ | ✅ | +| "p99 NodeStore fetch latency?" | ❌ | ❌ | ✅ | + +The last row is the case a span cannot answer. One `TMGetObjectByHash` message +requests up to `tuning::kHardMaxReplyNodes` objects, so a span per NodeStore +fetch is not affordable in that loop. Instead the fetch loop's wall time is +recorded once per message into an OpenTelemetry `Histogram` +(`getobject_lookup_us`), and the quantile is read off its buckets. StatsD is +marked ❌ because that instrument is recorded on the native OpenTelemetry metrics +path, not through Beast Insight. ### 2.6.4 Coexistence Strategy diff --git a/OpenTelemetryPlan/05-configuration-reference.md b/OpenTelemetryPlan/05-configuration-reference.md index c649582502..0edf8f20f6 100644 --- a/OpenTelemetryPlan/05-configuration-reference.md +++ b/OpenTelemetryPlan/05-configuration-reference.md @@ -77,7 +77,7 @@ The parser `makeTelemetrySetup()` in `src/libxrpl/telemetry/TelemetryConfig.cpp` > available on both. Components that hold a `ServiceRegistry&` (e.g. > `NetworkOPsImp`) call `registry_.get().getTelemetry()`. Components that > still hold an `Application&` (e.g. `ServerHandler`, `PeerImp`, -> `RCLConsensusAdaptor`) call `app_.getTelemetry()` directly. +> `RCLConsensus::Adaptor`) call `app_.getTelemetry()` directly. --- @@ -200,7 +200,7 @@ An example `xrpld RPC Performance` dashboard (uid `xrpld-rpc-performance`) sourc ### 5.8.4 Example Dashboard: Transaction Tracing -An example `xrpld Transaction Tracing` dashboard (uid `xrpld-tx-tracing`) over Tempo provides three panels: transaction throughput (`tx.receive` rate, stat), cross-node relay count (average `span.relay_count` on `tx.relay`, timeseries), and a table of transaction validation errors (`tx.validate` with `status.code=error`). +An example `xrpld Transaction Tracing` dashboard (uid `xrpld-tx-tracing`) over Tempo provides three panels: transaction throughput (`tx.receive` rate, stat), cross-node relay count (average `span.relay_count` on `tx.relay`, timeseries), and a table of transaction validation errors (`tx.validate` with `status = error`). ### 5.8.5 TraceQL Query Examples @@ -211,19 +211,19 @@ Common queries for xrpld traces: {resource.service.name="xrpld" && span.tx_hash="ABC123..."} # Find slow RPC commands (>100ms) -{resource.service.name="xrpld" && name=~"rpc.command.*"} | duration > 100ms +{resource.service.name="xrpld" && name=~"rpc.command.*"} | { duration > 100ms } # Find consensus rounds taking >5 seconds -{resource.service.name="xrpld" && name="consensus.round"} | duration > 5s +{resource.service.name="xrpld" && name="consensus.round"} | { duration > 5s } # Find failed transactions with error details -{resource.service.name="xrpld" && name="tx.validate" && status.code=error} +{resource.service.name="xrpld" && name="tx.validate" && status = error} # Find transactions relayed to many peers -{resource.service.name="xrpld" && name="tx.relay"} | span.relay_count > 10 +{resource.service.name="xrpld" && name="tx.relay"} | { span.relay_count > 10 } # Compare latency across nodes -{resource.service.name="xrpld" && name="rpc.command.account_info"} | avg(duration) by (resource.service.instance.id) +{resource.service.name="xrpld" && name="rpc.command.account_info"} | avg_over_time(duration) by (resource.service.instance.id) ``` ### 5.8.6 Correlation with PerfLog diff --git a/OpenTelemetryPlan/06-implementation-phases.md b/OpenTelemetryPlan/06-implementation-phases.md index f5f22818df..3027088d07 100644 --- a/OpenTelemetryPlan/06-implementation-phases.md +++ b/OpenTelemetryPlan/06-implementation-phases.md @@ -97,7 +97,7 @@ gantt | ---- | -------------------------------------------------------------------------- | | 2.1 | Implement W3C Trace Context HTTP header extraction | | 2.2 | Instrument `ServerHandler::onRequest()` | -| 2.3 | Instrument `RPCHandler::doCommand()` | +| 2.3 | Instrument `xrpl::rpc::doCommand()` | | 2.4 | Add RPC-specific attributes | | 2.5 | Instrument WebSocket handler | | 2.6 | PathFinding instrumentation (`pathfind.request`, `pathfind.compute` spans) | @@ -162,19 +162,19 @@ and [Phase3_taskList.md Task 3.9](./Phase3_taskList.md) for the full implementat ### Tasks -| Task | Description | -| ---- | ---------------------------------------------- | -| 4.1 | Instrument `RCLConsensusAdaptor::startRound()` | -| 4.2 | Instrument phase transitions | -| 4.3 | Instrument proposal handling | -| 4.4 | Instrument validation handling | -| 4.5 | Add consensus-specific attributes | -| 4.6 | Correlate with transaction traces | -| 4.7 | Validator list and manifest tracing | -| 4.8 | Amendment voting tracing | -| 4.9 | SHAMap sync tracing | -| 4.10 | Multi-validator integration tests | -| 4.11 | Performance validation | +| Task | Description | +| ---- | --------------------------------------- | +| 4.1 | Instrument `RCLConsensus::startRound()` | +| 4.2 | Instrument phase transitions | +| 4.3 | Instrument proposal handling | +| 4.4 | Instrument validation handling | +| 4.5 | Add consensus-specific attributes | +| 4.6 | Correlate with transaction traces | +| 4.7 | Validator list and manifest tracing | +| 4.8 | Amendment voting tracing | +| 4.9 | SHAMap sync tracing | +| 4.10 | Multi-validator integration tests | +| 4.11 | Performance validation | ### Exit Criteria diff --git a/OpenTelemetryPlan/07-observability-backends.md b/OpenTelemetryPlan/07-observability-backends.md index 4ebb6028fd..ecca8320d3 100644 --- a/OpenTelemetryPlan/07-observability-backends.md +++ b/OpenTelemetryPlan/07-observability-backends.md @@ -245,7 +245,7 @@ A Tempo-backed dashboard (uid `xrpld-node-overview`) with four panels: - **Active Nodes** (stat): count of distinct `resource.service.instance.id` values seen for the `xrpld` service. - **Total Transactions (1h)** (stat): count of `tx.receive` spans. -- **Error Rate** (gauge, percent): ratio of `status.code=error` spans to all spans, with yellow/red thresholds at 1%/5%. +- **Error Rate** (gauge, percent): ratio of `status = error` spans to all spans, with yellow/red thresholds at 1%/5%. - **Service Map** (nodeGraph): Tempo-generated service dependency graph. ### 7.6.3 Alert Rules diff --git a/cfg/xrpld-example.cfg b/cfg/xrpld-example.cfg index 70f6611368..52242c9cb6 100644 --- a/cfg/xrpld-example.cfg +++ b/cfg/xrpld-example.cfg @@ -1779,15 +1779,16 @@ validators.txt # # batch_size=512 # -# Maximum number of spans exported in a single batch. Default: 512. +# Maximum number of spans exported in a single batch. Must be at least 1 +# and must not exceed max_queue_size. Default: 512. # # batch_delay_ms=5000 # # Maximum delay (milliseconds) before a partial batch is flushed. -# Default: 5000 (5 seconds). +# Must be at least 1. Default: 5000 (5 seconds). # # max_queue_size=2048 # -# Maximum number of spans queued in memory before drops occur. -# Default: 2048. +# Maximum number of spans queued in memory before drops occur. Must be +# at least 1 and at least as large as batch_size. Default: 2048. # diff --git a/include/xrpl/telemetry/DeterministicIdGenerator.h b/include/xrpl/telemetry/DeterministicIdGenerator.h index 60d4aa5c4a..390b5b6f84 100644 --- a/include/xrpl/telemetry/DeterministicIdGenerator.h +++ b/include/xrpl/telemetry/DeterministicIdGenerator.h @@ -75,7 +75,7 @@ namespace xrpl::telemetry { * @code * // With an active parent span, startSpan() inherits the parent's trace_id * // and the SDK does NOT call GenerateTraceId(), so no PendingTraceId is used. - * // auto child = parentGuard.childSpan(rpc_span::op::process); // random/parent trace_id + * // auto child = parentGuard.childSpan(rpc_span::prefix::command); // random/parent trace_id * @endcode */ class DeterministicIdGenerator final : public opentelemetry::sdk::trace::IdGenerator diff --git a/include/xrpl/telemetry/SpanGuard.h b/include/xrpl/telemetry/SpanGuard.h index 6a31ecd83a..e71d00f950 100644 --- a/include/xrpl/telemetry/SpanGuard.h +++ b/include/xrpl/telemetry/SpanGuard.h @@ -99,7 +99,7 @@ * auto ctx = span.spanContext(); * * // Thread B: create child with captured context - * auto child = SpanGuard::childSpan(rpc_span::op::process, ctx); + * auto child = SpanGuard::childSpan(rpc_span::prefix::command, ctx); * @endcode * * 4. Conditional check (rarely needed — methods are no-ops on null): @@ -155,6 +155,22 @@ * }); * @endcode * + * 8. Internal work inside a category whose default role is Server: + * @code + * #include + * using namespace xrpl::telemetry; + * + * // Only the inbound handler is the server side of a remote call. + * // Work below it is internal, so pass the role explicitly: the + * // category default (Server) would read as a second inbound + * // request and leave an unpaired edge in a service graph. + * auto span = SpanGuard::span( + * TraceCategory::Rpc, + * rpc_span::prefix::rpc, + * rpc_span::op::process, + * SpanRole::Internal); + * @endcode + * * @note Thread safety: SpanGuard is thread-free. It holds only the * span (no Scope), so it never binds to a thread-local context stack * and may be moved to and destroyed on any thread. To make a span the @@ -198,6 +214,25 @@ namespace xrpl::telemetry { */ enum class TraceCategory { Rpc, Transactions, Consensus, Peer, Ledger }; +/** + * Role a span plays in a call relationship. Each value maps to the OTel + * span kind of the same name; see Telemetry::startSpan() for what those + * mean. + * + * Orthogonal to TraceCategory. The category names the subsystem and gates + * the span on config (`trace_rpc=1`); the role says whether the span + * handles a remote call or is internal work. An Rpc-category span can be + * either: the inbound request handler is Server, everything it calls into + * is Internal. + * + * FromCategory takes the category's own role, so a call site that does not + * care passes nothing. Pick a role explicitly where the category default + * is wrong: trace backends pair Server with Client and Consumer with + * Producer, so internal work left as Server becomes an unpaired edge in a + * service graph. + */ +enum class SpanRole { FromCategory, Internal, Server, Client, Producer, Consumer }; + /** * Raw trace context bytes for cross-node propagation. * @@ -310,9 +345,16 @@ public: * @param cat Trace subsystem category. * @param prefix Span name prefix (e.g. "rpc.command"). * @param name Span name suffix (e.g. "submit"). + * @param role Call-relationship role; defaults to the category's own + * role. Pass Internal for work the category maps to Server or Consumer + * but that handles no remote call. */ [[nodiscard]] static SpanGuard - span(TraceCategory cat, std::string_view prefix, std::string_view name) noexcept; + span( + TraceCategory cat, + std::string_view prefix, + std::string_view name, + SpanRole role = SpanRole::FromCategory) noexcept; /** * Create a span that always starts a fresh trace root. @@ -327,10 +369,16 @@ public: * @param cat Trace subsystem category. * @param prefix Span name prefix (e.g. "peer"). * @param name Span name suffix (e.g. "validation.receive"). + * @param role Call-relationship role; defaults to the category's own + * role. See span(). * @return An active root-span guard, or a null guard if disabled. */ [[nodiscard]] static SpanGuard - freshRoot(TraceCategory cat, std::string_view prefix, std::string_view name) noexcept; + freshRoot( + TraceCategory cat, + std::string_view prefix, + std::string_view name, + SpanRole role = SpanRole::FromCategory) noexcept; // --- Child / linked span creation ---------------------------------- @@ -607,10 +655,12 @@ public: * using namespace xrpl::telemetry; * * ScopedSpanGuard span( - * TraceCategory::Rpc, rpc_span::prefix::command, commandName); - * span.setAttribute(rpc_span::attr::command, commandName); - * // childSpan parents to `span` because it is active on this thread - * auto child = span.childSpan(rpc_span::op::process); + * TraceCategory::Rpc, rpc_span::prefix::rpc, rpc_span::op::process); + * // childSpan takes the name verbatim, so pass a full dotted constant, + * // never a bare op:: suffix. The child parents to `span` because + * // `span` is active on this thread. + * auto child = span.childSpan(rpc_span::prefix::command); + * child.setAttribute(rpc_span::attr::command, commandName); * @endcode * * 2. Capture on this thread, hand off to another (edge case): @@ -657,8 +707,14 @@ public: * @param cat Trace subsystem category. * @param prefix Span name prefix (e.g. "rpc.command"). * @param name Span name suffix (e.g. "submit"). + * @param role Call-relationship role; defaults to the category's own + * role. See SpanGuard::span(). */ - ScopedSpanGuard(TraceCategory cat, std::string_view prefix, std::string_view name) noexcept; + ScopedSpanGuard( + TraceCategory cat, + std::string_view prefix, + std::string_view name, + SpanRole role = SpanRole::FromCategory) noexcept; ~ScopedSpanGuard(); @@ -677,10 +733,16 @@ public: * @param cat Trace subsystem category. * @param prefix Span name prefix. * @param name Span name suffix. + * @param role Call-relationship role; defaults to the category's own + * role. See SpanGuard::span(). * @return An active scoped root-span guard, or a null one if disabled. */ [[nodiscard]] static ScopedSpanGuard - freshRoot(TraceCategory cat, std::string_view prefix, std::string_view name) noexcept; + freshRoot( + TraceCategory cat, + std::string_view prefix, + std::string_view name, + SpanRole role = SpanRole::FromCategory) noexcept; // --- Child / linked span creation ---------------------------------- @@ -972,13 +1034,21 @@ public: operator=(SpanGuard const&) = delete; [[nodiscard]] static SpanGuard - span(TraceCategory, std::string_view, std::string_view) noexcept + span( + TraceCategory, + std::string_view, + std::string_view, + SpanRole = SpanRole::FromCategory) noexcept { return {}; } [[nodiscard]] static SpanGuard - freshRoot(TraceCategory, std::string_view, std::string_view) noexcept + freshRoot( + TraceCategory, + std::string_view, + std::string_view, + SpanRole = SpanRole::FromCategory) noexcept { return {}; } @@ -1106,7 +1176,11 @@ class ScopedSpanGuard ScopedSpanGuard() = default; public: - ScopedSpanGuard(TraceCategory, std::string_view, std::string_view) noexcept + ScopedSpanGuard( + TraceCategory, + std::string_view, + std::string_view, + SpanRole = SpanRole::FromCategory) noexcept { } /** @@ -1127,7 +1201,11 @@ public: operator=(ScopedSpanGuard const&) = delete; [[nodiscard]] static ScopedSpanGuard - freshRoot(TraceCategory, std::string_view, std::string_view) noexcept + freshRoot( + TraceCategory, + std::string_view, + std::string_view, + SpanRole = SpanRole::FromCategory) noexcept { return {}; } diff --git a/include/xrpl/telemetry/Telemetry.h b/include/xrpl/telemetry/Telemetry.h index 68efe4965c..ab694abdae 100644 --- a/include/xrpl/telemetry/Telemetry.h +++ b/include/xrpl/telemetry/Telemetry.h @@ -53,14 +53,18 @@ * * 2. Child span for a sub-operation (scoped child): * @code - * auto parent = SpanGuard::span( + * auto parent = ScopedSpanGuard( * TraceCategory::Rpc, rpc_span::prefix::rpc, rpc_span::op::process); * { - * auto child = parent.childSpan(rpc_span::op::process); - * child.setAttribute(rpc_span::attr::version, apiVersion); + * auto child = parent.childSpan(rpc_span::prefix::command); + * child.setAttribute(rpc_span::attr::version, static_cast(apiVersion)); * // child ends here * } * @endcode + * childSpan() parents to the ambient scope, so the parent must be a + * ScopedSpanGuard. A plain SpanGuard is not ambient: pass its spanContext() + * to childSpan(name, ctx) instead. childSpan() takes the name verbatim, so + * pass a full dotted constant, never a bare op:: suffix. * * 3. Unrelated span (cross-scope, same thread): * @code @@ -78,7 +82,7 @@ * auto ctx = parentGuard.spanContext(); * * // Thread B: create child span with explicit parent - * auto child = SpanGuard::childSpan(rpc_span::op::process, ctx); + * auto child = SpanGuard::childSpan(rpc_span::prefix::command, ctx); * @endcode * * @note Thread safety: The Telemetry interface is safe for concurrent reads diff --git a/src/libxrpl/telemetry/SpanGuard.cpp b/src/libxrpl/telemetry/SpanGuard.cpp index 52be7ad560..574c51a358 100644 --- a/src/libxrpl/telemetry/SpanGuard.cpp +++ b/src/libxrpl/telemetry/SpanGuard.cpp @@ -86,7 +86,10 @@ SpanContext::SpanContext(std::shared_ptr impl) : impl_(std::move(impl)) bool SpanContext::isValid() const noexcept { - return impl_ != nullptr; + // Holding a Context is not proof of holding a span. GetCurrent() hands back + // an empty Context on a thread with no active span, and threadLocalContext() + // wraps that too. Ask the Context for its span instead of trusting impl_. + return impl_ != nullptr && otel_trace::GetSpan(impl_->ctx)->GetContext().IsValid(); } // ===== SpanGuard::Impl ==================================================== @@ -174,11 +177,13 @@ namespace { constexpr char const* kLinkTypeKey = "link_type"; constexpr char const* kLinkTypeFollowsFrom = "follows_from"; -// Map a TraceCategory to an OTel SpanKind so Tempo's service-graph / -// RED metrics see the correct direction. RPC spans are emitted at the -// server entry point (handler dispatch), Peer spans at inbound-message -// receipt. Transactions / Consensus / Ledger are internal processing -// and keep the default kInternal. +// Per-category default OTel SpanKind, used when a call site passes no +// SpanRole. A category cannot tell an inbound entry point from the +// internal work under it, so RPC and Peer default to the entry-point +// kind and any call site below the entry point passes SpanRole::Internal +// instead. Transactions / Consensus / Ledger are internal throughout. +// The kind drives direction in Tempo's service-graph / RED metrics, +// which pair kServer with kClient and kConsumer with kProducer. otel_trace::SpanKind categoryToSpanKind(TraceCategory cat) { @@ -196,6 +201,38 @@ categoryToSpanKind(TraceCategory cat) return otel_trace::SpanKind::kInternal; // unreachable } +/** + * Resolve the span kind to start a span with. + * + * An explicit SpanRole wins; SpanRole::FromCategory falls back to the + * category default above. Role and category are separate axes, so a single + * category can emit both an inbound handler and the internal work under it. + * + * @param cat Trace subsystem category. Read only for SpanRole::FromCategory. + * @param role Role the caller asked for. + * @return The OTel span kind for this span. + */ +[[nodiscard]] otel_trace::SpanKind +resolveSpanKind(TraceCategory cat, SpanRole role) +{ + switch (role) + { + case SpanRole::FromCategory: + return categoryToSpanKind(cat); + case SpanRole::Internal: + return otel_trace::SpanKind::kInternal; + case SpanRole::Server: + return otel_trace::SpanKind::kServer; + case SpanRole::Client: + return otel_trace::SpanKind::kClient; + case SpanRole::Producer: + return otel_trace::SpanKind::kProducer; + case SpanRole::Consumer: + return otel_trace::SpanKind::kConsumer; + } + return categoryToSpanKind(cat); // unreachable +} + /** * Join a span-name prefix and suffix into the dotted full name. * @@ -227,7 +264,11 @@ joinSpanName(std::string_view prefix, std::string_view name) noexcept } // namespace SpanGuard -SpanGuard::span(TraceCategory cat, std::string_view prefix, std::string_view name) noexcept +SpanGuard::span( + TraceCategory cat, + std::string_view prefix, + std::string_view name, + SpanRole role) noexcept { auto* tel = Telemetry::getInstance(); if ((tel == nullptr) || !tel->isEnabled() || !isCategoryEnabled(*tel, cat)) @@ -235,11 +276,15 @@ SpanGuard::span(TraceCategory cat, std::string_view prefix, std::string_view nam auto const fullName = joinSpanName(prefix, name); if (!fullName) return {}; - return SpanGuard(std::make_unique(tel->startSpan(*fullName, categoryToSpanKind(cat)))); + return SpanGuard(std::make_unique(tel->startSpan(*fullName, resolveSpanKind(cat, role)))); } SpanGuard -SpanGuard::freshRoot(TraceCategory cat, std::string_view prefix, std::string_view name) noexcept +SpanGuard::freshRoot( + TraceCategory cat, + std::string_view prefix, + std::string_view name, + SpanRole role) noexcept { auto* tel = Telemetry::getInstance(); if ((tel == nullptr) || !tel->isEnabled() || !isCategoryEnabled(*tel, cat)) @@ -250,7 +295,7 @@ SpanGuard::freshRoot(TraceCategory cat, std::string_view prefix, std::string_vie // Force a fresh trace root: do NOT inherit this thread's active span. auto rootCtx = opentelemetry::context::Context{otel_trace::kIsRootSpanKey, true}; return SpanGuard( - std::make_unique(tel->startSpan(*fullName, rootCtx, categoryToSpanKind(cat)))); + std::make_unique(tel->startSpan(*fullName, rootCtx, resolveSpanKind(cat, role)))); } // ===== Child / linked span creation ======================================== @@ -630,8 +675,9 @@ ScopedSpanGuard::~ScopedSpanGuard() ScopedSpanGuard::ScopedSpanGuard( TraceCategory cat, std::string_view prefix, - std::string_view name) noexcept - : ScopedSpanGuard(SpanGuard::span(cat, prefix, name)) + std::string_view name, + SpanRole role) noexcept + : ScopedSpanGuard(SpanGuard::span(cat, prefix, name, role)) { } @@ -639,9 +685,10 @@ ScopedSpanGuard ScopedSpanGuard::freshRoot( TraceCategory cat, std::string_view prefix, - std::string_view name) noexcept + std::string_view name, + SpanRole role) noexcept { - return ScopedSpanGuard(SpanGuard::freshRoot(cat, prefix, name)); + return ScopedSpanGuard(SpanGuard::freshRoot(cat, prefix, name, role)); } ScopedSpanGuard diff --git a/src/libxrpl/telemetry/TelemetryConfig.cpp b/src/libxrpl/telemetry/TelemetryConfig.cpp index 4cbbbf2a98..4f3698383c 100644 --- a/src/libxrpl/telemetry/TelemetryConfig.cpp +++ b/src/libxrpl/telemetry/TelemetryConfig.cpp @@ -8,11 +8,15 @@ * See cfg/xrpld-example.cfg for the full list of available options. */ +#include #include #include #include #include +#include +#include +#include #include namespace xrpl::telemetry { @@ -60,6 +64,71 @@ constexpr std::uint32_t batchDelayMs = 5000u; constexpr std::uint32_t maxQueueSize = 2048u; } // namespace dflt +/** + * Smallest accepted value for the three batch settings. + * + * All three size a queue or a timer, so zero is meaningless for every one of + * them. The OTel BatchSpanProcessor takes them as given and does not validate, + * so the config parser is the only place a nonsense value can be rejected. + */ +constexpr std::uint32_t kMinBatchSetting = 1u; + +/** + * Section name used in error messages, so the operator knows where to look. + */ +constexpr char const* kSectionLabel = "[telemetry]"; + +/** + * Read a config value and reject anything outside minValue..UINT32_MAX. + * + * Section::get() lets boost::bad_lexical_cast escape. That derives from + * std::bad_cast, not std::runtime_error, so a mistyped value gives the operator + * a bare "bad cast" naming no key. Wrap it and rethrow with the key name. + * + * @param section The [telemetry] section to read from. + * @param name Key to read, as documented in cfg/xrpld-example.cfg. + * @param absentValue Value returned when the key is absent. + * @param minValue Smallest accepted value. + * @return The configured value, or absentValue if the key is absent. + * @note Throws std::runtime_error for a value that is not a whole number, and + * for one out of range, with a different message for each. + */ +[[nodiscard]] std::uint32_t +readBounded( + Section const& section, + char const* name, + std::uint32_t absentValue, + std::uint32_t minValue) +{ + // Read as signed. boost::lexical_cast to an unsigned type wraps a leading + // minus instead of failing ("-1" yields 4294967295), so reading signed is + // the only way to see a negative value and reject it below. + std::optional parsed; + try + { + parsed = section.get(name); + } + catch (...) + { + Throw( + std::string("Invalid value '") + name + "' in " + kSectionLabel + + ": must be a whole number."); + } + + if (!parsed) + return absentValue; + + constexpr auto maxValue = static_cast(std::numeric_limits::max()); + if (*parsed < static_cast(minValue) || *parsed > maxValue) + { + Throw( + std::string("Invalid value '") + name + "' in " + kSectionLabel + ": must be between " + + std::to_string(minValue) + " and " + std::to_string(maxValue) + "."); + } + + return static_cast(*parsed); +} + /** * Derive a human-readable network type label from the numeric network ID. * @param networkId The network identifier from [network_id] config. @@ -108,10 +177,22 @@ makeTelemetrySetup( // traces; volume reduction is delegated to the collector's tail sampling. // setup.samplingRatio is a const member fixed at 1.0; nothing to parse. - setup.batchSize = section.valueOr(key::batchSize, dflt::batchSize); + setup.batchSize = readBounded(section, key::batchSize, dflt::batchSize, kMinBatchSetting); setup.batchDelay = std::chrono::milliseconds{ - section.valueOr(key::batchDelayMs, dflt::batchDelayMs)}; - setup.maxQueueSize = section.valueOr(key::maxQueueSize, dflt::maxQueueSize); + readBounded(section, key::batchDelayMs, dflt::batchDelayMs, kMinBatchSetting)}; + setup.maxQueueSize = + readBounded(section, key::maxQueueSize, dflt::maxQueueSize, kMinBatchSetting); + + // The OTel SDK documents max_export_batch_size <= max_queue_size as a + // precondition of BatchSpanProcessorOptions and does not enforce it, so + // reject the pair here rather than hand the SDK a state it forbids. + if (setup.batchSize > setup.maxQueueSize) + { + Throw( + std::string("Invalid value '") + key::batchSize + "' in " + kSectionLabel + + ": must not exceed '" + key::maxQueueSize + "' (" + std::to_string(setup.maxQueueSize) + + ")."); + } setup.networkId = networkId; setup.networkType = networkTypeFromId(networkId); diff --git a/src/tests/libxrpl/telemetry/TelemetryConfig.cpp b/src/tests/libxrpl/telemetry/TelemetryConfig.cpp index 152511e0a5..bef61c3d25 100644 --- a/src/tests/libxrpl/telemetry/TelemetryConfig.cpp +++ b/src/tests/libxrpl/telemetry/TelemetryConfig.cpp @@ -4,8 +4,85 @@ #include +#include +#include +#include +#include +#include +#include +#include + using namespace xrpl; +namespace { + +/** + * Batch-setting keys of the [telemetry] section. + * + * Spelled once so every case below matches what the parser reads. A + * misspelling cannot hide: the accepting cases would see the default instead + * of the value they wrote, and the rejecting cases would stop rejecting. + */ +namespace key { +constexpr char const* batchSize = "batch_size"; +constexpr char const* batchDelayMs = "batch_delay_ms"; +constexpr char const* maxQueueSize = "max_queue_size"; +} // namespace key + +/** + * The upper bound quoted in the expected messages below. + * + * makeTelemetrySetup() derives it from std::uint32_t, so pin the literal to + * that type here rather than repeating an unanchored number in 5 messages. + */ +static_assert(std::numeric_limits::max() == 4294967295u); + +using KeyValue = std::pair; + +/** + * Parse a [telemetry] section holding only the given keys. + * + * A key that is not listed stays absent, so its default applies. + * + * @param values Key/value pairs to write into the section. + * @return The populated Setup struct. + */ +telemetry::Telemetry::Setup +parseBatch(std::initializer_list values) +{ + Section section; + for (auto const& [name, value] : values) + section.set(name, value); + return telemetry::makeTelemetrySetup(section, "nHUtest123", "2.0.0", 0); +} + +/** + * Parse and return the rejection message. + * + * Only std::runtime_error is caught. A boost::bad_lexical_cast escaping the + * parser derives from std::bad_cast, so it propagates and fails the test + * instead of being mistaken for a clean rejection. That is the point of the + * not-a-number cases. + * + * @param values Key/value pairs to write into the section. + * @return The exception message, or "" if the parse succeeded. + */ +std::string +batchRejection(std::initializer_list values) +{ + try + { + static_cast(parseBatch(values)); + return {}; + } + catch (std::runtime_error const& e) + { + return e.what(); + } +} + +} // namespace + TEST(TelemetryConfig, setup_defaults) { telemetry::Telemetry::Setup const s; @@ -39,6 +116,11 @@ TEST(TelemetryConfig, parse_empty_section) EXPECT_EQ(setup.serviceVersion, "2.0.0"); EXPECT_EQ(setup.serviceInstanceId, "nHUtest123"); EXPECT_DOUBLE_EQ(setup.samplingRatio, 1.0); + // An absent key takes the documented default. setup_defaults covers the + // struct's own initializers; these three cover the parser applying them. + EXPECT_EQ(setup.batchSize, 512u); + EXPECT_EQ(setup.batchDelay, std::chrono::milliseconds{5000}); + EXPECT_EQ(setup.maxQueueSize, 2048u); EXPECT_TRUE(setup.traceRpc); EXPECT_TRUE(setup.traceTransactions); EXPECT_TRUE(setup.traceConsensus); @@ -83,6 +165,127 @@ TEST(TelemetryConfig, parse_full_section) EXPECT_FALSE(setup.traceLedger); } +TEST(TelemetryConfig, batch_settings_accept_the_lower_bound_exactly) +{ + auto const setup = + parseBatch({{key::batchSize, "1"}, {key::batchDelayMs, "1"}, {key::maxQueueSize, "1"}}); + EXPECT_EQ(setup.batchSize, 1u); + EXPECT_EQ(setup.batchDelay, std::chrono::milliseconds{1}); + EXPECT_EQ(setup.maxQueueSize, 1u); +} + +TEST(TelemetryConfig, batch_settings_accept_the_upper_bound_exactly) +{ + auto const setup = parseBatch( + {{key::batchSize, "4294967295"}, + {key::batchDelayMs, "4294967295"}, + {key::maxQueueSize, "4294967295"}}); + EXPECT_EQ(setup.batchSize, 4294967295u); + EXPECT_EQ(setup.batchDelay, std::chrono::milliseconds{4294967295}); + EXPECT_EQ(setup.maxQueueSize, 4294967295u); +} + +TEST(TelemetryConfig, batch_size_zero_is_rejected) +{ + EXPECT_EQ( + batchRejection({{key::batchSize, "0"}}), + "Invalid value 'batch_size' in [telemetry]: must be between 1 and 4294967295."); +} + +TEST(TelemetryConfig, batch_delay_ms_zero_is_rejected) +{ + EXPECT_EQ( + batchRejection({{key::batchDelayMs, "0"}}), + "Invalid value 'batch_delay_ms' in [telemetry]: must be between 1 and 4294967295."); +} + +TEST(TelemetryConfig, max_queue_size_zero_is_rejected) +{ + EXPECT_EQ( + batchRejection({{key::maxQueueSize, "0"}}), + "Invalid value 'max_queue_size' in [telemetry]: must be between 1 and 4294967295."); +} + +TEST(TelemetryConfig, batch_size_not_a_number_is_rejected_as_runtime_error) +{ + // Section::get() reaches boost::lexical_cast, which throws a std::bad_cast. + // Catching only std::runtime_error is the point: this fails unless the + // parser turned that into a message naming the key. + EXPECT_EQ( + batchRejection({{key::batchSize, "abc"}}), + "Invalid value 'batch_size' in [telemetry]: must be a whole number."); +} + +TEST(TelemetryConfig, batch_delay_ms_not_a_number_is_rejected_as_runtime_error) +{ + EXPECT_EQ( + batchRejection({{key::batchDelayMs, "abc"}}), + "Invalid value 'batch_delay_ms' in [telemetry]: must be a whole number."); +} + +TEST(TelemetryConfig, max_queue_size_not_a_number_is_rejected_as_runtime_error) +{ + EXPECT_EQ( + batchRejection({{key::maxQueueSize, "abc"}}), + "Invalid value 'max_queue_size' in [telemetry]: must be a whole number."); +} + +TEST(TelemetryConfig, batch_size_fractional_is_rejected) +{ + // A batch counts spans, so "512.5" must not silently truncate to 512. + EXPECT_EQ( + batchRejection({{key::batchSize, "512.5"}}), + "Invalid value 'batch_size' in [telemetry]: must be a whole number."); +} + +TEST(TelemetryConfig, batch_settings_reject_negative_rather_than_wrapping) +{ + // boost::lexical_cast to an unsigned type turns "-1" into 4294967295 + // instead of failing, so a negative must land on the range check. + EXPECT_EQ( + batchRejection({{key::batchSize, "-1"}}), + "Invalid value 'batch_size' in [telemetry]: must be between 1 and 4294967295."); + EXPECT_EQ( + batchRejection({{key::batchDelayMs, "-1"}}), + "Invalid value 'batch_delay_ms' in [telemetry]: must be between 1 and 4294967295."); + EXPECT_EQ( + batchRejection({{key::maxQueueSize, "-1"}}), + "Invalid value 'max_queue_size' in [telemetry]: must be between 1 and 4294967295."); +} + +TEST(TelemetryConfig, max_queue_size_above_the_upper_bound_is_rejected) +{ + EXPECT_EQ( + batchRejection({{key::maxQueueSize, "4294967296"}}), + "Invalid value 'max_queue_size' in [telemetry]: must be between 1 and 4294967295."); +} + +TEST(TelemetryConfig, batch_size_above_max_queue_size_is_rejected) +{ + // The OTel SDK documents max_export_batch_size <= max_queue_size as a + // precondition and does not enforce it, so the parser must. + EXPECT_EQ( + batchRejection({{key::batchSize, "600"}, {key::maxQueueSize, "512"}}), + "Invalid value 'batch_size' in [telemetry]: must not exceed 'max_queue_size' (512)."); +} + +TEST(TelemetryConfig, batch_size_above_a_lowered_max_queue_size_is_rejected) +{ + // The likely operator mistake: lowering only max_queue_size and leaving + // batch_size at its 512 default. + EXPECT_EQ( + batchRejection({{key::maxQueueSize, "256"}}), + "Invalid value 'batch_size' in [telemetry]: must not exceed 'max_queue_size' (256)."); +} + +TEST(TelemetryConfig, batch_size_equal_to_max_queue_size_is_accepted) +{ + // The cross-check rejects only batchSize > maxQueueSize, so equal passes. + auto const setup = parseBatch({{key::batchSize, "512"}, {key::maxQueueSize, "512"}}); + EXPECT_EQ(setup.batchSize, 512u); + EXPECT_EQ(setup.maxQueueSize, 512u); +} + TEST(TelemetryConfig, null_telemetry_factory) { telemetry::Telemetry::Setup setup; diff --git a/src/xrpld/rpc/detail/RPCHandler.cpp b/src/xrpld/rpc/detail/RPCHandler.cpp index e3400b746f..317adee22a 100644 --- a/src/xrpld/rpc/detail/RPCHandler.cpp +++ b/src/xrpld/rpc/detail/RPCHandler.cpp @@ -164,8 +164,10 @@ callMethod(JsonContext& context, Method method, std::string const& name, Object& { // Scoped so this command nests under rpc.process and becomes the ambient // parent of any command-internal spans (e.g. pathfind.request). Coro-aware - // storage keeps the scope correct across doRipplePathFind's yield. - auto span = ScopedSpanGuard(TraceCategory::Rpc, rpc_span::prefix::command, name); + // storage keeps the scope correct across doRipplePathFind's yield. Internal + // rather than Server: the inbound boundary is above rpc.process. + auto span = + ScopedSpanGuard(TraceCategory::Rpc, rpc_span::prefix::command, name, SpanRole::Internal); span.setAttribute(rpc_span::attr::command, name.c_str()); span.setAttribute(rpc_span::attr::version, static_cast(context.apiVersion)); span.setAttribute( @@ -283,7 +285,9 @@ doCommand(rpc::JsonContext& context, json::Value& result) // registered handler names (plus "unknown") — see the helper for why // raw request input must not reach the telemetry pipeline. auto const cmdName = resolveCommandSpanName(context); - auto span = ScopedSpanGuard(TraceCategory::Rpc, rpc_span::prefix::command, cmdName); + // Internal for the same reason as the success path above. + auto span = ScopedSpanGuard( + TraceCategory::Rpc, rpc_span::prefix::command, cmdName, SpanRole::Internal); span.setAttribute(rpc_span::attr::command, cmdName); // Mirror the attribute set callMethod() puts on a successful command // span, so error spans stay filterable by API version and role. diff --git a/src/xrpld/rpc/detail/ServerHandler.cpp b/src/xrpld/rpc/detail/ServerHandler.cpp index 16184df92e..18f8ee77a4 100644 --- a/src/xrpld/rpc/detail/ServerHandler.cpp +++ b/src/xrpld/rpc/detail/ServerHandler.cpp @@ -710,7 +710,9 @@ ServerHandler::processRequest( // yield in doRipplePathFind: the coro-aware context storage moves this // scope with the coroutine on resume (it is never stranded on a worker's // thread-local stack), so nesting and log-trace correlation both hold. - auto span = ScopedSpanGuard(TraceCategory::Rpc, rpc_span::prefix::rpc, rpc_span::op::process); + // Internal, not Server: the inbound boundary is rpc.http_request above. + auto span = ScopedSpanGuard( + TraceCategory::Rpc, rpc_span::prefix::rpc, rpc_span::op::process, SpanRole::Internal); auto rpcJ = app_.getJournal("RPC"); // Tracks whether any failure occurred. Set on every error path (early