Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions src/splunk_ao/exporter/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,15 @@
resolve_routing,
routing_resource_attributes,
)
from splunk_ao.exporter.diagnostics import ExportFailure, ExportHealth
from splunk_ao.exporter.o11y import build_o11y_exporter, resolve_o11y_exporter_config
from splunk_ao.exporter.sink import BatchConfig, SpanSink, build_batch_processor, build_span_sink
from splunk_ao.exporter.standalone import build_standalone_exporter, resolve_standalone_exporter_config

__all__ = [
"BatchConfig",
"ExportFailure",
"ExportHealth",
"ExporterConfig",
"RoutingAttrs",
"SpanSink",
Expand Down
9 changes: 8 additions & 1 deletion src/splunk_ao/exporter/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

from splunk_ao.constants import DEFAULT_AGENT_STREAM_NAME, DEFAULT_PROJECT_NAME
from splunk_ao.deployment import DeploymentMode
from splunk_ao.exporter.diagnostics import DiagnosticOTLPSpanExporter
from splunk_ao.exporter.span_transform import NormalizingSpanExporter
from splunk_ao.utils.env_helpers import (
_get_agent_stream_from_env,
Expand Down Expand Up @@ -156,10 +157,16 @@ def build_exporter(
endpoint: str,
auth_header: tuple[str, str],
routing: RoutingAttrs,
deployment: DeploymentMode,
_exporter_factory: ExporterFactory = OTLPSpanExporter,
**exporter_kwargs: Any,
) -> SpanExporter:
"""Build an OTLP HTTP exporter from shared resolved configuration."""
config = resolve_exporter_config(endpoint, auth_header, routing)
delegate = _exporter_factory(endpoint=config.endpoint, headers=config.headers, **exporter_kwargs)
if _exporter_factory is OTLPSpanExporter:
delegate: SpanExporter = DiagnosticOTLPSpanExporter(
endpoint=config.endpoint, headers=config.headers, deployment=deployment, **exporter_kwargs
)
else:
delegate = _exporter_factory(endpoint=config.endpoint, headers=config.headers, **exporter_kwargs)
return NormalizingSpanExporter(delegate, create_otel_resource(routing))
284 changes: 284 additions & 0 deletions src/splunk_ao/exporter/diagnostics.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,284 @@
"""Diagnostics for successful OTLP responses that reject spans."""

from __future__ import annotations

import json
import logging
import re
import threading
import time
from collections.abc import Callable, Sequence
from dataclasses import dataclass
from typing import Any, Literal
from urllib.parse import urlsplit

from google.protobuf.message import DecodeError
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ExportTraceServiceResponse
from opentelemetry.sdk.trace import ReadableSpan
from opentelemetry.sdk.trace.export import SpanExportResult
from requests import Response

from splunk_ao.deployment import DeploymentMode

_logger = logging.getLogger("splunk_ao.exporter")

ExportFailureCategory = Literal["rejected"]

_MAX_RESPONSE_BYTES = 64 * 1024
_MAX_MESSAGE_LENGTH = 512
_MAX_REJECTION_KEYS = 8
_DEFAULT_LOG_INTERVAL_SECONDS = 60.0
_SAFE_KEY = re.compile(r"[^a-zA-Z0-9_.-]+")
_PROTOBUF_CONTENT_TYPES = frozenset(
{"application/protobuf", "application/vnd.google.protobuf", "application/x-protobuf"}
)


@dataclass(frozen=True)
class ExportFailure:
"""Sanitized details for the most recent acknowledged rejection."""

category: ExportFailureCategory
message: str
status_code: int
consecutive_failures: int


@dataclass(frozen=True)
class ExportHealth:
"""Immutable acknowledgement health.

``healthy`` is ``None`` before export and after ordinary transport or
non-2xx failures, ``True`` after an accepted 2xx response, and ``False``
after a 2xx response that rejects telemetry.
"""

healthy: bool | None
consecutive_failures: int
last_failure: ExportFailure | None


UNKNOWN_EXPORT_HEALTH = ExportHealth(healthy=None, consecutive_failures=0, last_failure=None)


@dataclass(frozen=True)
class _RejectionDetail:
status_code: int
detail: str | None = None


@dataclass(frozen=True)
class _ResponseSnapshot:
status_code: int
content_type: str
body: bytes


class _ExportHealthTracker:
def __init__(
self,
deployment: DeploymentMode,
endpoint: str,
*,
clock: Callable[[], float] = time.monotonic,
log_interval_seconds: float = _DEFAULT_LOG_INTERVAL_SECONDS,
logger: logging.Logger = _logger,
) -> None:
self._deployment = deployment
self._endpoint_host = urlsplit(endpoint).hostname or "configured endpoint"
self._clock = clock
self._log_interval_seconds = log_interval_seconds
self._logger = logger
self._lock = threading.Lock()
self._health = UNKNOWN_EXPORT_HEALTH
self._last_logged_rejection: tuple[int, str | None] | None = None
self._last_logged_at = 0.0
self._suppressed_rejections = 0

@property
def health(self) -> ExportHealth:
with self._lock:
return self._health

def record_rejection(self, detail: _RejectionDetail) -> None:
message = self._rejection_message(detail)
now = self._clock()
log_message: str | None = None
with self._lock:
consecutive_failures = self._health.consecutive_failures + 1
failure = ExportFailure(
category="rejected",
message=message,
status_code=detail.status_code,
consecutive_failures=consecutive_failures,
)
fingerprint = (detail.status_code, detail.detail)
interval_elapsed = now - self._last_logged_at >= self._log_interval_seconds
if self._health.healthy is not False or fingerprint != self._last_logged_rejection or interval_elapsed:
log_message = message
if self._suppressed_rejections:
log_message += f" ({self._suppressed_rejections} repeated rejections suppressed)"
self._last_logged_rejection = fingerprint
self._last_logged_at = now
self._suppressed_rejections = 0
else:
self._suppressed_rejections += 1
self._health = ExportHealth(healthy=False, consecutive_failures=consecutive_failures, last_failure=failure)
if log_message is not None:
self._logger.error("%s", log_message)

def record_success(self) -> None:
with self._lock:
should_log_recovery = self._health.healthy is False
self._health = ExportHealth(healthy=True, consecutive_failures=0, last_failure=None)
self._last_logged_rejection = None
self._last_logged_at = 0.0
self._suppressed_rejections = 0
if should_log_recovery:
self._logger.info(
"Splunk AO OTLP export recovered for %s endpoint %s.", self._deployment.value, self._endpoint_host
)

def record_unknown(self) -> None:
"""Clear acknowledgement health when the standard exporter reports another failure."""
with self._lock:
self._health = UNKNOWN_EXPORT_HEALTH
self._last_logged_rejection = None
self._last_logged_at = 0.0
self._suppressed_rejections = 0

def _rejection_message(self, rejection: _RejectionDetail) -> str:
detail = f" {rejection.detail}" if rejection.detail else ""
return (
f"Splunk AO OTLP ingest at {self._deployment.value} endpoint "
f"{self._endpoint_host} returned HTTP {rejection.status_code} but rejected some or all spans."
f"{detail}"
)[:_MAX_MESSAGE_LENGTH]


class DiagnosticOTLPSpanExporter(OTLPSpanExporter):
"""OTLP HTTP exporter that detects rejection acknowledgements returned with HTTP 2xx."""

def __init__(self, *args: Any, deployment: DeploymentMode, **kwargs: Any) -> None:
configured_endpoint = str(kwargs.get("endpoint", ""))
super().__init__(*args, **kwargs)
self._attempt_local = threading.local()
self._health_tracker = _ExportHealthTracker(deployment, getattr(self, "_endpoint", configured_endpoint))

@property
def export_health(self) -> ExportHealth:
return self._health_tracker.health

def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult:
self._attempt_local.response = None
try:
try:
result = super().export(spans)
except Exception:
self._health_tracker.record_unknown()
raise
if result != SpanExportResult.SUCCESS:
self._health_tracker.record_unknown()
return result
Comment on lines +181 to +183

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 minor (design): On any non-SUCCESS result (e.g. HTTP 401/500 or transport error) export_health.healthy is reset to None (UNKNOWN) rather than False. This means a consumer polling export_health cannot distinguish "never exported / transport failing" from "healthy" — only 2xx-acknowledged span rejections flip it to False. This appears intentional (the PR preserves standard OTel handling for non-2xx), but the healthy tri-state semantics (None vs True vs False) are subtle and undocumented on the public ExportHealth dataclass. Worth a docstring note so downstream consumers don't treat healthy is not False as "all good".

🤖 Generated by the Astra agent


response = getattr(self._attempt_local, "response", None)
rejection = _classify_successful_response(response) if response is not None else None
if rejection is not None:
self._health_tracker.record_rejection(rejection)
return SpanExportResult.FAILURE

self._health_tracker.record_success()
return SpanExportResult.SUCCESS
finally:
self._attempt_local.response = None

def _export(self, serialized_data: bytes, timeout_sec: float | None = None) -> Response:
response = super()._export(serialized_data, timeout_sec)
self._attempt_local.response = _snapshot_response(response) if 200 <= response.status_code < 300 else None
return response


def get_export_health(owner: object) -> ExportHealth:
"""Read acknowledgement health from an SDK exporter ownership surface."""
health = getattr(owner, "export_health", UNKNOWN_EXPORT_HEALTH)
return health if isinstance(health, ExportHealth) else UNKNOWN_EXPORT_HEALTH


def _snapshot_response(response: Response) -> _ResponseSnapshot:
content_type = response.headers.get("Content-Type", "").partition(";")[0].strip().lower()
return _ResponseSnapshot(
status_code=response.status_code, content_type=content_type, body=bytes(response.content[:_MAX_RESPONSE_BYTES])
)


def _classify_successful_response(response: _ResponseSnapshot) -> _RejectionDetail | None:
if not response.body:
return None
if response.content_type == "application/json" or response.content_type.endswith("+json"):
return _classify_json_acknowledgement(response.body, response.status_code)
if response.content_type in _PROTOBUF_CONTENT_TYPES:
return _classify_protobuf_acknowledgement(response.body, response.status_code)
return None


def _classify_json_acknowledgement(body: bytes, status_code: int) -> _RejectionDetail | None:
try:
payload = json.loads(body)
except (UnicodeDecodeError, json.JSONDecodeError):
return None
if not isinstance(payload, dict):
return None

partial_success = payload.get("partialSuccess")
if isinstance(partial_success, dict):
rejected_spans = _positive_json_integer(partial_success.get("rejectedSpans"))
if rejected_spans is not None:
return _RejectionDetail(status_code, f"Rejected spans: {rejected_spans}.")

invalid = payload.get("invalid")
if payload.get("valid") != 0 and not invalid:
return None
return _RejectionDetail(status_code, _summarize_rejections(invalid))


def _positive_json_integer(value: object) -> int | None:
if isinstance(value, bool):
return None
if isinstance(value, int):
return value if value > 0 else None
if isinstance(value, str) and value.isdecimal():
parsed = int(value)
return parsed if parsed > 0 else None
return None


def _classify_protobuf_acknowledgement(body: bytes, status_code: int) -> _RejectionDetail | None:
response = ExportTraceServiceResponse()
try:
response.ParseFromString(body)
except DecodeError:
return None
rejected_spans = response.partial_success.rejected_spans
if rejected_spans <= 0:
return None
return _RejectionDetail(status_code, f"Rejected spans: {rejected_spans}.")


def _summarize_rejections(invalid: object) -> str | None:
if not invalid:
return None
if not isinstance(invalid, dict):
return "The response contained an explicit rejection."

summaries: list[str] = []
for key, value in sorted(invalid.items(), key=lambda item: str(item[0]))[:_MAX_REJECTION_KEYS]:
safe_key = _SAFE_KEY.sub("_", str(key))[:48] or "unknown"
if isinstance(value, (list, tuple, dict, set)):
count = len(value)
elif isinstance(value, int) and not isinstance(value, bool):
count = max(0, value)
else:
count = 1
summaries.append(f"{safe_key}={count}")
return f"Rejection categories: {', '.join(summaries)}." if summaries else None
9 changes: 7 additions & 2 deletions src/splunk_ao/exporter/o11y.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.trace.export import SpanExporter

from splunk_ao.deployment import O11yConfig
from splunk_ao.deployment import DeploymentMode, O11yConfig
from splunk_ao.exporter.config import (
ExporterConfig,
ExporterFactory,
Expand All @@ -32,5 +32,10 @@ def build_o11y_exporter(
) -> SpanExporter:
"""Build an OTLP exporter authenticated for Splunk Observability Cloud."""
return build_exporter(
config.otlp_endpoint, _o11y_auth_header(config), routing, _exporter_factory, **exporter_kwargs
config.otlp_endpoint,
_o11y_auth_header(config),
routing,
deployment=DeploymentMode.O11Y,
_exporter_factory=_exporter_factory,
**exporter_kwargs,
)
14 changes: 12 additions & 2 deletions src/splunk_ao/exporter/sink.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@
from opentelemetry.sdk.trace import ReadableSpan, TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor, SpanExporter

from splunk_ao.exporter.diagnostics import ExportHealth, get_export_health


@dataclass
class BatchConfig:
Expand All @@ -19,11 +21,19 @@ class BatchConfig:
class SpanSink:
"""SDK-owned abstraction over a batch processor and tracer provider."""

def __init__(self, processor: BatchSpanProcessor, provider: TracerProvider) -> None:
def __init__(
self, processor: BatchSpanProcessor, provider: TracerProvider, exporter: SpanExporter | None = None
) -> None:
self._processor = processor
self._provider = provider
self._exporter = exporter
self._shutdown = False

@property
def export_health(self) -> ExportHealth:
"""Return the current receiver-acknowledgement health snapshot."""
return get_export_health(self._exporter)

def emit(self, span: ReadableSpan) -> None:
"""Enqueue a completed span without flushing the batch."""
if self._shutdown:
Expand Down Expand Up @@ -60,4 +70,4 @@ def build_span_sink(exporter: SpanExporter, batch_config: BatchConfig | None = N
processor = build_batch_processor(exporter, batch_config)
provider = TracerProvider()
provider.add_span_processor(processor)
return SpanSink(processor, provider)
return SpanSink(processor, provider, exporter)
Loading
Loading