#!/usr/bin/env python3 """Telemetry Validation Suite for xrpld. Validates that the full telemetry stack is emitting expected data after a workload run. Queries Tempo (spans), Prometheus (metrics), Loki (logs), and Grafana (dashboards) APIs to produce a pass/fail report. Validation categories: 1. Span validation — Every required span type in expected_spans.json, each carrying its required attributes 2. Metric validation — SpanMetrics, StatsD, and MetricsRegistry OTLP metrics are non-zero, and each group's required_labels reach Prometheus with non-empty values 3. Log-trace correlation — Loki logs contain trace_id/span_id fields 4. Dashboard validation — Every dashboard uid in expected_metrics.json provisions and loads (panel count only, not panel data) 5. External parity — Span attrs, metric existence, and value sanity for external dashboard parity (validator-health, peer-quality, node-health) Usage: python3 validate_telemetry.py --report /tmp/validation-report.json # Custom API endpoints: python3 validate_telemetry.py \\ --tempo http://localhost:3200 \\ --prometheus http://localhost:9090 \\ --loki http://localhost:3100 \\ --grafana http://localhost:3000 """ import argparse import asyncio import fnmatch import json import logging import re import sys import time from dataclasses import dataclass, field from pathlib import Path from typing import Any import aiohttp # Loki's default query window is the last hour. A validation run finishes in # minutes, but bounding the range explicitly keeps the query reproducible when # someone re-runs it later to investigate a result. LOG_QUERY_WINDOW_SECONDS = 4 * 60 * 60 logger = logging.getLogger("validate_telemetry") # --------------------------------------------------------------------------- # Configuration defaults # --------------------------------------------------------------------------- DEFAULT_TEMPO = "http://localhost:3200" DEFAULT_PROMETHEUS = "http://localhost:9090" DEFAULT_LOKI = "http://localhost:3100" DEFAULT_GRAFANA = "http://localhost:3000" SCRIPT_DIR = Path(__file__).parent EXPECTED_SPANS_FILE = SCRIPT_DIR / "expected_spans.json" EXPECTED_METRICS_FILE = SCRIPT_DIR / "expected_metrics.json" # Some beast::insight gauges/counters (ledger-age, peer-finder, overlay # traffic) only populate after the node validates ledgers and sustains peer # traffic, then travel a 1s periodic OTLP export + a 15s Prometheus scrape # before they are queryable. On a slow CI runner the fixed post-workload wait # can end before that pipeline settles, so a single query races and reports 0 # series. Poll each missing metric for up to this long (covering two scrape # cycles) before failing, so the check is robust to runner speed. METRIC_POLL_TIMEOUT_SEC = 45.0 METRIC_POLL_INTERVAL_SEC = 5.0 # All metrics are polled concurrently against ONE shared deadline, so the # metric phase costs a single poll window instead of one per metric. This caps # how many /api/v1/series requests are in flight at a time, so the fan-out does # not hammer the single-container Prometheus the harness runs. METRIC_POLL_CONCURRENCY = 8 # --------------------------------------------------------------------------- # Data classes # --------------------------------------------------------------------------- @dataclass class CheckResult: """Result of a single validation check. Attributes: name: Check identifier (e.g., "span.rpc.ws_message"). category: Validation category (span, metric, log, dashboard). passed: Whether the check passed. message: Human-readable description of the result. details: Optional additional data (counts, values, etc.). """ name: str category: str passed: bool message: str details: dict[str, Any] = field(default_factory=dict) def to_dict(self) -> dict[str, Any]: """Serialize to a JSON-compatible dict.""" return { "name": self.name, "category": self.category, "passed": self.passed, "message": self.message, "details": self.details, } @dataclass class ValidationReport: """Aggregated validation report. Attributes: checks: List of all individual check results. start_time: ISO timestamp when validation started. end_time: ISO timestamp when validation completed. """ checks: list[CheckResult] = field(default_factory=list) start_time: str = "" end_time: str = "" @property def total_checks(self) -> int: """Total number of checks executed.""" return len(self.checks) @property def passed(self) -> int: """Number of checks that passed.""" return sum(1 for c in self.checks if c.passed) @property def failed(self) -> int: """Number of checks that failed.""" return sum(1 for c in self.checks if not c.passed) @property def all_passed(self) -> bool: """Whether all checks passed.""" return self.failed == 0 def add(self, check: CheckResult) -> None: """Add a check result to the report.""" self.checks.append(check) status = "PASS" if check.passed else "FAIL" logger.info("[%s] %s: %s", status, check.name, check.message) def to_dict(self) -> dict[str, Any]: """Serialize to a JSON-compatible dict.""" return { "summary": { "total": self.total_checks, "passed": self.passed, "failed": self.failed, "all_passed": self.all_passed, }, "start_time": self.start_time, "end_time": self.end_time, "checks": [c.to_dict() for c in self.checks], } # --------------------------------------------------------------------------- # Tempo API helpers # --------------------------------------------------------------------------- def _log_query_window() -> dict[str, str]: """Loki query_range bounds covering a validation run. Returns: start/end parameters in nanoseconds since the epoch. """ now = time.time() return { "start": str(int((now - LOG_QUERY_WINDOW_SECONDS) * 1_000_000_000)), "end": str(int(now * 1_000_000_000)), } async def _tempo_search( session: aiohttp.ClientSession, tempo_url: str, query: str, limit: int = 20, ) -> list[dict[str, Any]]: """Search traces in Tempo using TraceQL. Args: session: aiohttp client session. tempo_url: Base URL for Tempo API (e.g., http://localhost:3200). query: TraceQL query string. limit: Maximum number of traces to return. Returns: List of trace summary dicts from Tempo search results. """ params = {"q": query, "limit": str(limit)} async with session.get(f"{tempo_url}/api/search", params=params) as resp: data = await resp.json() return data.get("traces", []) async def _tempo_get_trace( session: aiohttp.ClientSession, tempo_url: str, trace_id: str, ) -> list[dict[str, Any]]: """Fetch a full trace from Tempo by trace ID. Returns the list of spans extracted from the OTLP-format response. Args: session: aiohttp client session. tempo_url: Tempo API base URL. trace_id: Hex trace ID string. Returns: Flat list of span dicts with 'name' and 'attributes' keys. """ async with session.get(f"{tempo_url}/api/traces/{trace_id}") as resp: data = await resp.json() spans: list[dict[str, Any]] = [] for batch in data.get("batches", []): for scope_spans in batch.get("scopeSpans", []): spans.extend(scope_spans.get("spans", [])) return spans def _otlp_span_attr_keys(span: dict[str, Any]) -> set[str]: """Extract all attribute key names from an OTLP span. Args: span: OTLP span dict with an 'attributes' list. Returns: Set of attribute key strings. """ return {a["key"] for a in span.get("attributes", []) if "key" in a} def _span_name_matches(emitted_name: str, expected_name: str) -> bool: """Test an emitted span name against a name from expected_spans.json. Contract names are either literals or globs containing "*" (for example "rpc.command.*"). Literals are compared for exact equality so a longer emitted name cannot satisfy a shorter contract: "consensus.accept.apply" must not stand in for "consensus.accept". Args: emitted_name: Span name as reported by Tempo. expected_name: Span name or glob pattern from expected_spans.json. Returns: True when the emitted name satisfies the expected name. """ if "*" in expected_name: return fnmatch.fnmatchcase(emitted_name, expected_name) return emitted_name == expected_name # --------------------------------------------------------------------------- # Span Validation (Tempo API) # --------------------------------------------------------------------------- async def validate_spans( session: aiohttp.ClientSession, tempo_url: str, report: ValidationReport, ) -> None: """Validate that all expected spans appear in Tempo. Queries the Tempo TraceQL API for each expected span name and checks that traces exist. Also validates required attributes on spans and parent-child relationships. Args: session: aiohttp client session. tempo_url: Base URL for Tempo API (e.g., http://localhost:3200). report: ValidationReport to accumulate results. """ logger.info("--- Span Validation (Tempo) ---") # Load expected spans. with open(EXPECTED_SPANS_FILE) as f: expected = json.load(f) # Check service registration. try: async with session.get( f"{tempo_url}/api/v2/search/tag/resource.service.name/values" ) as resp: data = await resp.json() tag_values = data.get("tagValues", []) services = [tv.get("value", "") for tv in tag_values] has_xrpld = "xrpld" in services report.add( CheckResult( name="span.service_registration", category="span", passed=has_xrpld, message=( f"Service 'xrpld' registered (found: {services})" if has_xrpld else f"Service 'xrpld' NOT found (found: {services})" ), ) ) except Exception as exc: report.add( CheckResult( name="span.service_registration", category="span", passed=False, message=f"Tempo API unreachable: {exc}", ) ) return # Diagnostic: list all available operations (span names) for the xrpld # service. This output appears in CI logs and helps debug missing-span # failures without needing to reproduce the full stack locally. try: async with session.get( f"{tempo_url}/api/v2/search/tag/span.name/values" ) as resp: ops_data = await resp.json() tag_values = ops_data.get("tagValues", []) operations = [tv.get("value", "") for tv in tag_values] logger.info( "Tempo operations (%d total): %s", len(operations), operations, ) except Exception as exc: logger.warning("Failed to fetch Tempo operations: %s", exc) # Concrete probe names for wildcard span entries. Exact-match TraceQL can't # match a literal "*", so a representative operation name is substituted. # Wildcards without a known concrete example (e.g. grpc. when no # gRPC client runs) are skipped when marked optional. wildcard_probes = {"rpc.command.*": "rpc.command.server_info"} # Check each expected span. for span_def in expected["spans"]: span_name = span_def["name"] is_optional = span_def.get("optional", False) check_name = f"span.{span_name}" if "*" in span_name: operation = wildcard_probes.get(span_name) if operation is None: # No concrete probe. Optional wildcards (e.g. grpc.*) are skipped; # a required one would be a config error worth surfacing. if is_optional: logger.info( "[SKIP] %s: optional wildcard span with no concrete " "probe (not exercised by the workload)", check_name, ) continue report.add( CheckResult( name=check_name, category="span", passed=False, message=f"{span_name}: required wildcard has no probe name", ) ) continue else: operation = span_name try: query = '{resource.service.name="xrpld" && name="' + operation + '"}' traces = await _tempo_search(session, tempo_url, query, limit=5) count = len(traces) # Optional spans only fire under specific traffic (mode changes, # missing-ledger fetch, fee escalation). Absence is not a failure — # mirror the parent-child "skip" handling so CI stays green. if count == 0 and is_optional: logger.info( "[SKIP] %s: optional span not emitted under this workload", check_name, ) report.add( CheckResult( name=check_name, category="span", passed=True, message=f"{span_name}: optional, not emitted (skipped)", details={"trace_count": 0, "optional": True}, ) ) continue report.add( CheckResult( name=check_name, category="span", passed=count > 0, message=( f"{span_name}: {count} traces found" if count > 0 else f"{span_name}: 0 traces (expected > 0)" ), details={"trace_count": count}, ) ) # Validate required attributes on first trace. if count > 0 and span_def.get("required_attributes"): await _check_attributes_on_first_trace( session, tempo_url, traces, span_def, report ) except Exception as exc: report.add( CheckResult( name=check_name, category="span", passed=False, message=f"{span_name}: query failed ({exc})", ) ) # Validate parent-child relationships. for rel in expected.get("parent_child_relationships", []): # Skip relationships marked with "skip: true" (e.g., cross-thread # parent-child that requires a C++ fix to propagate span context). if rel.get("skip", False): reason = rel.get("skip_reason", "marked skip in expected_spans.json") logger.info( "[SKIP] span.hierarchy.%s->%s: %s", rel["parent"], rel["child"], reason, ) continue await _validate_parent_child(session, tempo_url, rel, report) async def _check_attributes_on_first_trace( session: aiohttp.ClientSession, tempo_url: str, traces: list[dict[str, Any]], span_def: dict[str, Any], report: ValidationReport, ) -> None: """Fetch the first trace and check the span's required attributes. Fetching the trace is a second network call, so it carries its own error handling. Letting it fall through to the caller's handler would add a second result under the span's own check name, which has already recorded the trace as found -- one entry passing and one failing for the same name, inflating the check total and blaming the trace-existence check for an attribute-fetch failure. Args: session: aiohttp client session. tempo_url: Base URL for the Tempo API. traces: Traces returned for this span, most recent first. span_def: The span's entry from expected_spans.json. report: ValidationReport to accumulate results. """ span_name = span_def["name"] try: trace_id = traces[0].get("traceID", "") if not trace_id: return spans = await _tempo_get_trace(session, tempo_url, trace_id) await _validate_span_attributes_otlp(spans, span_def, report) except Exception as exc: report.add( CheckResult( name=f"span.attrs.{span_name}", category="span", passed=False, message=f"{span_name}: attribute check failed ({exc})", ) ) async def _validate_span_attributes_otlp( spans: list[dict[str, Any]], span_def: dict[str, Any], report: ValidationReport, ) -> None: """Check that the contract's own span carries its required attributes. Only spans whose name matches ``span_def["name"]`` are inspected. Attributes are never borrowed from siblings: many span types share keys such as ledger_seq or tx_hash, so a trace-wide scan would satisfy every one of those contracts from a single carrier span and make the per-span contract unenforceable. A span type passes when at least one instance of it carries every required attribute. When none does, the closest instance's missing keys are reported. Args: spans: Every OTLP span dict in the fetched trace. span_def: Span definition from expected_spans.json. report: ValidationReport to accumulate results. """ required_attrs = span_def.get("required_attributes", []) if not required_attrs: return span_name = span_def["name"] check_name = f"span.attrs.{span_name}" matching = [s for s in spans if _span_name_matches(s.get("name", ""), span_name)] if not matching: report.add( CheckResult( name=check_name, category="span", passed=False, message=( f"{span_name}: no span named '{span_name}' in the fetched " "trace, cannot verify its attributes" ), details={"required": required_attrs, "instances": 0}, ) ) return # Keep the instance that is missing the fewest required attributes, so the # failure message names the closest witness rather than an arbitrary one. best_found: set[str] = set() best_missing: list[str] = list(required_attrs) for span in matching: found = _otlp_span_attr_keys(span) missing = [a for a in required_attrs if a not in found] if len(missing) < len(best_missing): best_found, best_missing = found, missing if not best_missing: break report.add( CheckResult( name=check_name, category="span", passed=not best_missing, message=( f"{span_name}: all {len(required_attrs)} attributes present" if not best_missing else f"{span_name}: no '{span_name}' span carried all " f"{len(required_attrs)} required attributes; closest of " f"{len(matching)} instance(s) missing {best_missing}" ), details={ "required": required_attrs, "found": sorted(best_found), "missing": best_missing, "instances": len(matching), }, ) ) async def _validate_parent_child( session: aiohttp.ClientSession, tempo_url: str, relationship: dict[str, Any], report: ValidationReport, ) -> None: """Validate a parent-child span relationship in Tempo traces. Args: session: aiohttp client session. tempo_url: Base URL for Tempo API. relationship: Dict with 'parent' and 'child' span names. report: ValidationReport to accumulate results. """ parent_name = relationship["parent"] child_name = relationship["child"] try: # Query traces for the parent span. query = '{resource.service.name="xrpld" && name="' + parent_name + '"}' traces = await _tempo_search(session, tempo_url, query, limit=3) if not traces: report.add( CheckResult( name=f"span.hierarchy.{parent_name}->{child_name}", category="span", passed=False, message=f"No {parent_name} traces to check hierarchy", ) ) return # Check if child spans exist within parent traces. Names are matched # exactly (globs for wildcard contracts) — a substring test let a # longer emitted name satisfy a shorter contract, so # consensus.round -> consensus.accept passed on a # consensus.accept.apply span alone. found_child = False for trace_summary in traces: trace_id = trace_summary.get("traceID", "") if not trace_id: continue spans = await _tempo_get_trace(session, tempo_url, trace_id) if any( _span_name_matches(span.get("name", ""), child_name) for span in spans ): found_child = True break report.add( CheckResult( name=f"span.hierarchy.{parent_name}->{child_name}", category="span", passed=found_child, message=( f"Found {child_name} as child of {parent_name}" if found_child else f"{child_name} not found in {parent_name} traces" ), ) ) except Exception as exc: report.add( CheckResult( name=f"span.hierarchy.{parent_name}->{child_name}", category="span", passed=False, message=f"Hierarchy check failed: {exc}", ) ) # --------------------------------------------------------------------------- # Metric Validation (Prometheus API) # --------------------------------------------------------------------------- async def _log_prometheus_metric_names( session: aiohttp.ClientSession, prometheus_url: str ) -> None: """Log the harness-relevant metric names Prometheus currently knows. Diagnostic only — this output appears in CI logs and helps debug name mismatches between expected_metrics.json and actual emissions. Failures are warnings, never check failures. Args: session: aiohttp client session. prometheus_url: Prometheus base URL. """ try: async with session.get( f"{prometheus_url}/api/v1/label/__name__/values" ) as resp: label_data = await resp.json() all_metrics = label_data.get("data", []) relevant = [ m for m in all_metrics if m.startswith( ( "span_", "rpc_method", "cache_", "txq_", "object_count", "load_factor", "nodestore", "ledgermaster", "peer_finder", "jobq_", "total_bytes", "total_messages", "validation_agreement", "validator_health", "peer_quality", "ledger_economy", "state_tracking", "storage_detail", ) ) ] logger.info( "Prometheus metrics (relevant, %d of %d total): %s", len(relevant), len(all_metrics), relevant, ) except Exception as exc: logger.warning("Failed to fetch Prometheus metric names: %s", exc) def _metric_check_targets( expected: dict[str, Any], ) -> tuple[list[tuple[str, str]], list[tuple[str, str, list[str]]]]: """Flatten expected_metrics.json into the two lists of check targets. Args: expected: The parsed expected_metrics.json contract. Returns: A (metric targets, label targets) pair. Metric targets are (group, metric selector) for every name under a group's "metrics". Label targets are (group, label, that group's metric selectors) for every name under a group's "required_labels" — read for every group that declares it, not just spanmetrics. Nothing read the key at all until this was added, so the four labels the spanmetrics group documented as required had never actually been checked. """ groups = [ (category_key, category_data) for category_key, category_data in expected.items() if category_key not in ("description", "grafana_dashboards") ] targets = [ (category_key, metric_name) for category_key, category_data in groups for metric_name in category_data.get("metrics", []) ] label_targets = [ (category_key, label, category_data.get("metrics", [])) for category_key, category_data in groups for label in category_data.get("required_labels", []) ] return targets, label_targets async def validate_metrics( session: aiohttp.ClientSession, prometheus_url: str, report: ValidationReport, ) -> None: """Validate that expected metrics appear in Prometheus with non-zero values. Two kinds of check come out of expected_metrics.json: every name under a group's "metrics" must have at least one series, and every label under a group's "required_labels" must reach at least one of that group's series with a non-empty value. Args: session: aiohttp client session. prometheus_url: Base URL for Prometheus API (e.g., http://localhost:9090). report: ValidationReport to accumulate results. """ logger.info("--- Metric Validation (Prometheus) ---") await _log_prometheus_metric_names(session, prometheus_url) with open(EXPECTED_METRICS_FILE) as f: expected = json.load(f) # Flatten the contract, then poll every target concurrently against ONE # shared deadline. Polling them serially made each metric own its own # timeout, so the waits were additive: 58 metrics x 45 s = 43.5 min, which # overran the CI job budget and lost the artifact-upload and summary # diagnostics. Sharing the deadline bounds the whole phase to a single # poll window. targets, label_targets = _metric_check_targets(expected) deadline = time.monotonic() + METRIC_POLL_TIMEOUT_SEC sem = asyncio.Semaphore(METRIC_POLL_CONCURRENCY) # Both kinds of check share the one deadline and the one concurrency bound, # so the label checks cost no extra poll window and add no extra load. metric_checks, label_checks = await asyncio.gather( asyncio.gather( *( _check_prometheus_metric( session, prometheus_url, metric_name, category, deadline, sem ) for category, metric_name in targets ) ), asyncio.gather( *( _check_metric_label( session, prometheus_url, category, label, metric_selectors, deadline, sem, ) for category, label, metric_selectors in label_targets ) ), ) # Add in contract order, not completion order, so the report and its log # lines stay deterministic across runs. Existence checks keep their place # ahead of the label checks, so no existing check's position moves. for check in [*metric_checks, *label_checks]: report.add(check) def _selector_with_label(metric_selector: str, label: str) -> str: """Add a "label is present and non-empty" matcher to a metric selector. ``