diff --git a/src/splunk_ao/exporter/__init__.py b/src/splunk_ao/exporter/__init__.py index 063136b..bc4c50a 100644 --- a/src/splunk_ao/exporter/__init__.py +++ b/src/splunk_ao/exporter/__init__.py @@ -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", diff --git a/src/splunk_ao/exporter/config.py b/src/splunk_ao/exporter/config.py index 5887b81..44287d0 100644 --- a/src/splunk_ao/exporter/config.py +++ b/src/splunk_ao/exporter/config.py @@ -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, @@ -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)) diff --git a/src/splunk_ao/exporter/diagnostics.py b/src/splunk_ao/exporter/diagnostics.py new file mode 100644 index 0000000..459a2ea --- /dev/null +++ b/src/splunk_ao/exporter/diagnostics.py @@ -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 + + 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 diff --git a/src/splunk_ao/exporter/o11y.py b/src/splunk_ao/exporter/o11y.py index 69b00ec..962b83f 100644 --- a/src/splunk_ao/exporter/o11y.py +++ b/src/splunk_ao/exporter/o11y.py @@ -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, @@ -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, ) diff --git a/src/splunk_ao/exporter/sink.py b/src/splunk_ao/exporter/sink.py index cb8f3d8..b32a608 100644 --- a/src/splunk_ao/exporter/sink.py +++ b/src/splunk_ao/exporter/sink.py @@ -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: @@ -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: @@ -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) diff --git a/src/splunk_ao/exporter/span_transform.py b/src/splunk_ao/exporter/span_transform.py index 60daab2..75ede8c 100644 --- a/src/splunk_ao/exporter/span_transform.py +++ b/src/splunk_ao/exporter/span_transform.py @@ -10,6 +10,7 @@ from opentelemetry.sdk.trace.export import SpanExporter, SpanExportResult from splunk_ao.converter.attribute_mapping import normalize_attributes_for_export +from splunk_ao.exporter.diagnostics import ExportHealth, get_export_health ROUTING_ATTRIBUTE_KEYS = frozenset( { @@ -81,6 +82,11 @@ def delegate(self) -> SpanExporter: """Return the wrapped exporter for SDK composition and diagnostics.""" return self._delegate + @property + def export_health(self) -> ExportHealth: + """Return the delegate's receiver-acknowledgement health snapshot.""" + return get_export_health(self._delegate) + def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: """Export normalized immutable copies of the supplied spans.""" return self._delegate.export( diff --git a/src/splunk_ao/exporter/standalone.py b/src/splunk_ao/exporter/standalone.py index 0aab409..9fb80e2 100644 --- a/src/splunk_ao/exporter/standalone.py +++ b/src/splunk_ao/exporter/standalone.py @@ -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 StandaloneConfig +from splunk_ao.deployment import DeploymentMode, StandaloneConfig from splunk_ao.exporter.config import ( ExporterConfig, ExporterFactory, @@ -32,5 +32,10 @@ def build_standalone_exporter( ) -> SpanExporter: """Build an OTLP exporter authenticated for standalone Splunk AO.""" return build_exporter( - config.otlp_endpoint, _standalone_auth_header(config), routing, _exporter_factory, **exporter_kwargs + config.otlp_endpoint, + _standalone_auth_header(config), + routing, + deployment=DeploymentMode.STANDALONE, + _exporter_factory=_exporter_factory, + **exporter_kwargs, ) diff --git a/src/splunk_ao/logger/logger.py b/src/splunk_ao/logger/logger.py index 91deaf7..8ebd893 100644 --- a/src/splunk_ao/logger/logger.py +++ b/src/splunk_ao/logger/logger.py @@ -46,6 +46,7 @@ from splunk_ao.deployment import DeploymentMode, O11yConfig, StandaloneConfig, resolve_deployment from splunk_ao.exceptions import SplunkAOLoggerException from splunk_ao.exporter import ( + ExportHealth, SpanSink, build_o11y_exporter, build_span_sink, @@ -53,6 +54,7 @@ create_otel_resource, resolve_routing, ) +from splunk_ao.exporter.diagnostics import get_export_health from splunk_ao.logger.control import ControlAppliesTo, ControlCheckStage, ControlResult from splunk_ao.logger.task_handler import ThreadPoolTaskHandler from splunk_ao.projects import Projects @@ -401,6 +403,11 @@ def __init__( atexit.register(self.terminate) self._auto_enable_agent_control_if_available() + @property + def export_health(self) -> ExportHealth: + """Return the current receiver-acknowledgement health snapshot.""" + return get_export_health(getattr(self, "_sink", None)) + def _set_current_parent(self, parent: StepWithChildSpans | None) -> None: """Set the proprietary parent and mirror its open chain in OTel context.""" super()._set_current_parent(parent) diff --git a/src/splunk_ao/otel.py b/src/splunk_ao/otel.py index 2386c16..2b1f5ea 100644 --- a/src/splunk_ao/otel.py +++ b/src/splunk_ao/otel.py @@ -27,7 +27,14 @@ _session_id_context, ) from splunk_ao.deployment import DeploymentMode, O11yConfig, StandaloneConfig -from splunk_ao.exporter import RoutingAttrs, build_o11y_exporter, build_standalone_exporter, resolve_routing +from splunk_ao.exporter import ( + ExportHealth, + RoutingAttrs, + build_o11y_exporter, + build_standalone_exporter, + resolve_routing, +) +from splunk_ao.exporter.diagnostics import get_export_health logger = logging.getLogger(__name__) @@ -143,6 +150,11 @@ def export(self, spans: Sequence[ReadableSpan]) -> SpanExportResult: """Export through the shared immutable normalization pipeline.""" return self._delegate.export(spans) + @property + def export_health(self) -> ExportHealth: + """Return the current receiver-acknowledgement health snapshot.""" + return get_export_health(self._delegate) + def force_flush(self, timeout_millis: int = 30000) -> bool: """Flush the delegate exporter.""" return self._delegate.force_flush(timeout_millis) @@ -261,6 +273,11 @@ def force_flush(self, timeout_millis: int = 40000) -> bool: """Force immediate export of all pending spans with specified timeout.""" return self._processor.force_flush(timeout_millis) + @property + def export_health(self) -> ExportHealth: + """Return the current receiver-acknowledgement health snapshot.""" + return get_export_health(self._exporter) + @property def exporter(self) -> SpanExporter: """Access to the underlying Splunk AO OTLP exporter instance.""" diff --git a/tests/test_exporter_config.py b/tests/test_exporter_config.py index 2a5cb03..9d33565 100644 --- a/tests/test_exporter_config.py +++ b/tests/test_exporter_config.py @@ -3,7 +3,7 @@ from typing import Any from unittest.mock import patch -from splunk_ao.deployment import StandaloneConfig +from splunk_ao.deployment import DeploymentMode, StandaloneConfig from splunk_ao.exporter.config import RoutingAttrs, create_otel_resource, routing_resource_attributes from splunk_ao.exporter.span_transform import NormalizingSpanExporter from splunk_ao.exporter.standalone import build_standalone_exporter, resolve_standalone_exporter_config @@ -171,6 +171,14 @@ def exporter_factory(**kwargs: Any) -> object: } +def test_standalone_exporter_passes_deployment_explicitly_to_diagnostics() -> None: + with patch("splunk_ao.exporter.config.DiagnosticOTLPSpanExporter") as diagnostic_exporter: + exporter = build_standalone_exporter(make_standalone_cfg(), make_routing()) + + assert exporter.delegate is diagnostic_exporter.return_value + assert diagnostic_exporter.call_args.kwargs["deployment"] == DeploymentMode.STANDALONE + + def test_healthz_probe_no_longer_called_on_exporter_construction() -> None: with patch("httpx.get") as httpx_get: build_standalone_exporter( diff --git a/tests/test_exporter_diagnostics.py b/tests/test_exporter_diagnostics.py new file mode 100644 index 0000000..ef8e64d --- /dev/null +++ b/tests/test_exporter_diagnostics.py @@ -0,0 +1,242 @@ +"""Tests for successful-response OTLP rejection diagnostics.""" + +from __future__ import annotations + +import json +from concurrent.futures import ThreadPoolExecutor +from dataclasses import FrozenInstanceError +from typing import Any +from unittest.mock import MagicMock + +import pytest +from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ExportTraceServiceResponse +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import SimpleSpanProcessor, SpanExportResult +from requests import Response +from requests.exceptions import Timeout + +from splunk_ao.deployment import DeploymentMode +from splunk_ao.exporter.diagnostics import ( + DiagnosticOTLPSpanExporter, + ExportHealth, + _ExportHealthTracker, + _RejectionDetail, +) +from splunk_ao.exporter.sink import BatchConfig, build_span_sink +from splunk_ao.exporter.span_transform import NormalizingSpanExporter +from splunk_ao.logger.logger import SplunkAOLogger +from splunk_ao.otel import SplunkAOSpanProcessor + + +class FakeSession: + def __init__(self, *responses: Response | Exception) -> None: + self.headers: dict[str, str] = {} + self.responses = list(responses) + self.post_calls = 0 + + def post(self, **_: Any) -> Response: + outcome = self.responses[min(self.post_calls, len(self.responses) - 1)] + self.post_calls += 1 + if isinstance(outcome, Exception): + raise outcome + return outcome + + def close(self) -> None: + pass + + +def response(status_code: int = 200, body: bytes = b"", content_type: str = "application/json") -> Response: + result = Response() + result.status_code = status_code + result._content = body + result.headers["Content-Type"] = content_type + result.reason = "test response" + return result + + +def diagnostic_exporter(session: FakeSession) -> DiagnosticOTLPSpanExporter: + return DiagnosticOTLPSpanExporter( + endpoint="https://ingest.example.test/v2/trace/otlp", + headers={"X-SF-Token": "placeholder-secret"}, + session=session, + deployment=DeploymentMode.O11Y, + timeout=0.01, + ) + + +def test_splunk_json_rejection_returns_failure_without_retry_or_sensitive_logging() -> None: + session = FakeSession( + response( + body=( + b'{"valid":0,"invalid":{"nilServiceName":["private payload"],"bad.attribute":{"nested":"not logged"}}}' + ) + ) + ) + exporter = diagnostic_exporter(session) + logger = MagicMock() + exporter._health_tracker._logger = logger + + result = exporter.export(()) + + assert result == SpanExportResult.FAILURE + assert session.post_calls == 1 + assert exporter.export_health.last_failure is not None + assert exporter.export_health.last_failure.category == "rejected" + log_message = logger.error.call_args.args[1] + assert "nilServiceName=1" in log_message + assert "bad.attribute=1" in log_message + assert "private payload" not in log_message + assert "not logged" not in log_message + assert exporter._attempt_local.response is None + + +def test_positive_otlp_partial_success_returns_failure_without_retry() -> None: + acknowledgement = ExportTraceServiceResponse() + acknowledgement.partial_success.rejected_spans = 3 + acknowledgement.partial_success.error_message = "do not retain this backend detail" + session = FakeSession(response(body=acknowledgement.SerializeToString(), content_type="application/x-protobuf")) + exporter = diagnostic_exporter(session) + + result = exporter.export(()) + + assert result == SpanExportResult.FAILURE + assert session.post_calls == 1 + assert exporter.export_health.last_failure is not None + assert "Rejected spans: 3" in exporter.export_health.last_failure.message + assert "backend detail" not in exporter.export_health.last_failure.message + assert exporter._attempt_local.response is None + + +@pytest.mark.parametrize("rejected_spans", [3, "3"]) +def test_positive_json_otlp_partial_success_returns_failure_without_backend_detail(rejected_spans: int | str) -> None: + body = json.dumps( + {"partialSuccess": {"rejectedSpans": rejected_spans, "errorMessage": "do not retain this backend detail"}} + ).encode() + exporter = diagnostic_exporter(FakeSession(response(body=body))) + + assert exporter.export(()) == SpanExportResult.FAILURE + assert exporter.export_health.last_failure is not None + assert "Rejected spans: 3" in exporter.export_health.last_failure.message + assert "backend detail" not in exporter.export_health.last_failure.message + + +@pytest.mark.parametrize( + ("body", "content_type"), + [ + (b"", "application/x-protobuf"), + (b'{"valid":1}', "application/json"), + (b'{"other":"value"}', "application/json"), + (b'{"partialSuccess":{"rejectedSpans":"0","errorMessage":"advisory"}}', "application/json"), + (b'{"partialSuccess":{"rejectedSpans":"invalid"}}', "application/json"), + (b"not-json", "application/json"), + (b"\x80", "application/x-protobuf"), + ], +) +def test_successful_or_unrecognized_acknowledgements_remain_successful(body: bytes, content_type: str) -> None: + exporter = diagnostic_exporter(FakeSession(response(body=body, content_type=content_type))) + + assert exporter.export(()) == SpanExportResult.SUCCESS + assert exporter.export_health == ExportHealth(healthy=True, consecutive_failures=0, last_failure=None) + assert exporter._attempt_local.response is None + + +def test_zero_rejected_spans_with_advisory_message_is_successful() -> None: + acknowledgement = ExportTraceServiceResponse() + acknowledgement.partial_success.error_message = "advisory only" + exporter = diagnostic_exporter( + FakeSession(response(body=acknowledgement.SerializeToString(), content_type="application/protobuf")) + ) + + assert exporter.export(()) == SpanExportResult.SUCCESS + assert exporter.export_health.healthy is True + + +def test_ordinary_http_failure_is_left_to_the_standard_exporter() -> None: + exporter = diagnostic_exporter(FakeSession(response(body=b'{"valid":1}'), response(401))) + logger = MagicMock() + exporter._health_tracker._logger = logger + + assert exporter.export(()) == SpanExportResult.SUCCESS + assert exporter.export_health.healthy is True + assert exporter.export(()) == SpanExportResult.FAILURE + assert exporter.export_health.healthy is None + logger.error.assert_not_called() + assert exporter._attempt_local.response is None + + +def test_transport_exception_is_unchanged_and_clears_response_snapshot() -> None: + exporter = diagnostic_exporter(FakeSession(Timeout("standard exporter failure"))) + exporter._attempt_local.response = object() + + with pytest.raises(Timeout, match="standard exporter failure"): + exporter.export(()) + + assert exporter.export_health.healthy is None + assert exporter._attempt_local.response is None + + +def test_repeated_rejections_are_rate_limited_and_recovery_logs_once() -> None: + clock = MagicMock(return_value=10.0) + logger = MagicMock() + tracker = _ExportHealthTracker( + DeploymentMode.O11Y, "https://ingest.example.test/v2/trace/otlp", clock=clock, logger=logger + ) + + rejection = _RejectionDetail(200, "Rejection categories: nilServiceName=1.") + tracker.record_rejection(rejection) + tracker.record_rejection(rejection) + assert logger.error.call_count == 1 + + clock.return_value = 71.0 + tracker.record_rejection(rejection) + assert logger.error.call_count == 2 + assert "1 repeated rejections suppressed" in logger.error.call_args.args[1] + + tracker.record_success() + tracker.record_success() + assert logger.info.call_count == 1 + assert tracker.health == ExportHealth(healthy=True, consecutive_failures=0, last_failure=None) + + +def test_export_health_snapshots_are_immutable_and_thread_safe() -> None: + tracker = _ExportHealthTracker(DeploymentMode.O11Y, "https://ingest.example.test/v2/trace/otlp", logger=MagicMock()) + rejection = _RejectionDetail(200) + + with ThreadPoolExecutor(max_workers=8) as executor: + list(executor.map(lambda _: tracker.record_rejection(rejection), range(100))) + + health = tracker.health + assert health.consecutive_failures == 100 + assert health.last_failure is not None + assert health.last_failure.consecutive_failures == 100 + with pytest.raises(FrozenInstanceError): + health.consecutive_failures = 0 + + +def test_logger_path_reports_rejection_without_failing_the_logged_operation() -> None: + exporter = diagnostic_exporter(FakeSession(response(body=b'{"valid":0}'))) + sink = build_span_sink(exporter, BatchConfig(schedule_delay_millis=60_000)) + logger = SplunkAOLogger(project_id="project-id", agent_stream_id="stream-id", _sink=sink) + try: + logger.add_single_llm_span_trace(input="question", output="answer", model="test-model") + + assert logger.flush() is None + assert sink.export_health.healthy is False + assert logger.export_health == sink.export_health + finally: + sink.shutdown() + + +def test_user_wired_processor_reports_rejection_without_failing_span_completion() -> None: + exporter = diagnostic_exporter(FakeSession(response(body=b'{"valid":0}'))) + normalizing_exporter = NormalizingSpanExporter(exporter, Resource.create()) + processor = SplunkAOSpanProcessor(SpanProcessor=SimpleSpanProcessor, _exporter=normalizing_exporter) + provider = TracerProvider() + provider.add_span_processor(processor) + + with provider.get_tracer("test").start_as_current_span("operation"): + pass + + assert processor.export_health.healthy is False + provider.shutdown() diff --git a/tests/test_exporter_o11y.py b/tests/test_exporter_o11y.py index e6eb374..bdc329c 100644 --- a/tests/test_exporter_o11y.py +++ b/tests/test_exporter_o11y.py @@ -1,10 +1,11 @@ """Tests for Splunk Observability Cloud OTLP exporter configuration.""" from typing import Any +from unittest.mock import patch import pytest -from splunk_ao.deployment import O11yConfig +from splunk_ao.deployment import DeploymentMode, O11yConfig from splunk_ao.exporter.config import RoutingAttrs from splunk_ao.exporter.o11y import build_o11y_exporter, resolve_o11y_exporter_config from splunk_ao.exporter.span_transform import NormalizingSpanExporter @@ -105,3 +106,11 @@ def exporter_factory(**kwargs: Any) -> object: "endpoint": "https://ingest.us1.observability.splunkcloud.com/v2/trace/otlp", "headers": {"X-SF-Token": "tok", "projectid": "pid", "logstreamid": "lsid"}, } + + +def test_o11y_exporter_passes_deployment_explicitly_to_diagnostics() -> None: + with patch("splunk_ao.exporter.config.DiagnosticOTLPSpanExporter") as diagnostic_exporter: + exporter = build_o11y_exporter(O11yConfig(realm="us1", sf_token="tok"), make_routing()) + + assert exporter.delegate is diagnostic_exporter.return_value + assert diagnostic_exporter.call_args.kwargs["deployment"] == DeploymentMode.O11Y diff --git a/tests/test_otel_native_paths.py b/tests/test_otel_native_paths.py index f655718..88f0bc6 100644 --- a/tests/test_otel_native_paths.py +++ b/tests/test_otel_native_paths.py @@ -320,6 +320,9 @@ def test_exporter_delegates_force_flush_and_shutdown() -> None: factory = RecordingExporterFactory() exporter = build_exporter(DeploymentMode.O11Y, factory) + assert exporter.export_health.healthy is None + assert exporter.export((make_span(),)) == SpanExportResult.SUCCESS + assert exporter.export_health.healthy is None assert exporter.force_flush(1234) is True exporter.shutdown() @@ -327,6 +330,17 @@ def test_exporter_delegates_force_flush_and_shutdown() -> None: assert factory.exporter.shutdown_calls == 1 +def test_processor_forwards_export_health() -> None: + factory = RecordingExporterFactory() + exporter = build_exporter(DeploymentMode.O11Y, factory) + processor = SplunkAOSpanProcessor(SpanProcessor=RecordingSpanProcessor, _exporter=exporter) + + assert processor.export_health.healthy is None + exporter.export((make_span(),)) + assert processor.export_health.healthy is None + processor.shutdown() + + def test_processor_does_not_put_routing_on_span_attributes() -> None: exporter = RecordingExporter() processor = SplunkAOSpanProcessor(SpanProcessor=RecordingSpanProcessor, _exporter=exporter)