diff --git a/packages/traceloop-sdk/tests/test_fortifyroot_auth_warnings.py b/packages/traceloop-sdk/tests/test_fortifyroot_auth_warnings.py new file mode 100644 index 0000000000..baeeeec58c --- /dev/null +++ b/packages/traceloop-sdk/tests/test_fortifyroot_auth_warnings.py @@ -0,0 +1,280 @@ +from __future__ import annotations + +import logging +import os +import threading +from http.server import BaseHTTPRequestHandler, HTTPServer + +import pytest +from grpc import RpcError, StatusCode +from opentelemetry.sdk._logs.export import LogExportResult, LogExporter +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter + +from traceloop.sdk.exporters.auth_warnings import ( + _AuthWarningClientProxy, + reset_auth_warning_state_for_tests, +) + +pytestmark = pytest.mark.fr + + +class _FakeHTTPResponse: + ok = False + reason = "Unauthorized" + text = "unauthorized" + + def __init__(self, status_code): + self.status_code = status_code + + +class _FakeAuthRpcError(RpcError): + def __init__(self, code): + super().__init__() + self._code = code + + def code(self): + return self._code + + +class _CapturingLogExporter(LogExporter): + def __init__(self): + self.bodies = [] + + def export(self, batch): + for record in batch: + self.bodies.append(str(record.log_record.body)) + return LogExportResult.SUCCESS + + def shutdown(self): + return None + + +def _make_traces_http_exporter(endpoint="http://localhost:4318"): + from traceloop.sdk.tracing.tracing import init_spans_exporter + + return init_spans_exporter(endpoint, {}) + + +def _make_metrics_http_exporter(endpoint="http://localhost:4318"): + from traceloop.sdk.metrics.metrics import init_metrics_exporter + + return init_metrics_exporter(endpoint, {}) + + +def _make_logs_http_exporter(endpoint="http://localhost:4318"): + from traceloop.sdk.logging.logging import init_logging_exporter + + return init_logging_exporter(endpoint, {}) + + +def _make_traces_grpc_exporter(endpoint="grpc://localhost:4317"): + from traceloop.sdk.tracing.tracing import init_spans_exporter + + return init_spans_exporter(endpoint, {}) + + +def _make_metrics_grpc_exporter(endpoint="localhost:4317"): + from traceloop.sdk.metrics.metrics import init_metrics_exporter + + return init_metrics_exporter(endpoint, {}) + + +def _make_logs_grpc_exporter(endpoint="localhost:4317"): + from traceloop.sdk.logging.logging import init_logging_exporter + + return init_logging_exporter(endpoint, {}) + + +def _rejecting_post(status_code): + def post(*args, **kwargs): + return _FakeHTTPResponse(status_code) + + return post + + +def _auth_warnings(caplog): + return [ + record + for record in caplog.records + if "FortifyRoot SDK auth warning" in record.getMessage() + ] + + +@pytest.mark.parametrize( + "factory,signal,status_code", + [ + (_make_traces_http_exporter, "traces", 401), + (_make_traces_http_exporter, "traces", 403), + (_make_metrics_http_exporter, "metrics", 401), + (_make_metrics_http_exporter, "metrics", 403), + (_make_logs_http_exporter, "logs", 401), + (_make_logs_http_exporter, "logs", 403), + ], +) +def test_http_exporters_warn_on_auth_failure_once(factory, signal, status_code, caplog): + reset_auth_warning_state_for_tests() + exporter = factory() + assert hasattr(exporter, "_session") + exporter._session.post = _rejecting_post(status_code) + + with caplog.at_level("WARNING"): + exporter._export(b"payload") + exporter._export(b"payload") + + auth_warnings = _auth_warnings(caplog) + assert len(auth_warnings) == 1 + warning = auth_warnings[0].getMessage() + assert signal in warning + assert f"HTTP {status_code}" in warning + assert "invalid, revoked, deleted, or missing permissions" in warning + assert "Telemetry will not reach the OTLP endpoint" in warning + assert "Telemetry will not reach FortifyRoot" not in warning + + +def test_http_exporter_dedupes_by_endpoint(caplog): + reset_auth_warning_state_for_tests() + exporter_one = _make_traces_http_exporter("http://collector-one:4318") + exporter_two = _make_traces_http_exporter("http://collector-two:4318") + assert hasattr(exporter_one, "_session") + assert hasattr(exporter_two, "_session") + exporter_one._session.post = _rejecting_post(401) + exporter_two._session.post = _rejecting_post(401) + + with caplog.at_level("WARNING"): + exporter_one._export(b"payload") + exporter_two._export(b"payload") + + auth_warnings = _auth_warnings(caplog) + assert len(auth_warnings) == 2 + messages = [record.getMessage() for record in auth_warnings] + assert any("collector-one:4318" in message for message in messages) + assert any("collector-two:4318" in message for message in messages) + + +@pytest.mark.parametrize( + "factory,signal,code", + [ + (_make_traces_grpc_exporter, "traces", StatusCode.UNAUTHENTICATED), + (_make_traces_grpc_exporter, "traces", StatusCode.PERMISSION_DENIED), + (_make_metrics_grpc_exporter, "metrics", StatusCode.UNAUTHENTICATED), + (_make_metrics_grpc_exporter, "metrics", StatusCode.PERMISSION_DENIED), + (_make_logs_grpc_exporter, "logs", StatusCode.UNAUTHENTICATED), + (_make_logs_grpc_exporter, "logs", StatusCode.PERMISSION_DENIED), + ], +) +def test_grpc_exporter_client_proxy_warns_on_auth_failure( + factory, + signal, + code, + caplog, +): + reset_auth_warning_state_for_tests() + + class FakeClient: + def Export(self, *args, **kwargs): # noqa: N802 + raise _FakeAuthRpcError(code) + + exporter = factory() + assert isinstance(exporter._client, _AuthWarningClientProxy) + exporter._client._client = FakeClient() + + with caplog.at_level("WARNING"), pytest.raises(_FakeAuthRpcError): + exporter._client.Export() + + assert "FortifyRoot SDK auth warning" in caplog.text + assert signal in caplog.text + assert f"gRPC {code.name}" in caplog.text + + +def test_trace_http_exporter_warns_end_to_end_with_mock_collector(caplog): + reset_auth_warning_state_for_tests() + + class RejectingCollector(BaseHTTPRequestHandler): + def do_POST(self): # noqa: N802 + self.send_response(401) + self.end_headers() + self.wfile.write(b"unauthorized") + + def log_message(self, format, *args): + return None + + server = HTTPServer(("127.0.0.1", 0), RejectingCollector) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + endpoint = f"http://127.0.0.1:{server.server_port}" + exporter = _make_traces_http_exporter(endpoint) + with caplog.at_level("WARNING"): + exporter.export([]) + finally: + server.shutdown() + server.server_close() + thread.join(timeout=2) + + assert "FortifyRoot SDK auth warning" in caplog.text + assert "traces" in caplog.text + assert "HTTP 401" in caplog.text + + +def test_auth_warning_log_is_not_reexported_through_otel_logging_handler(caplog): + reset_auth_warning_state_for_tests() + + from traceloop.sdk import Traceloop + from traceloop.sdk.logging.logging import ( + LoggerWrapper, + is_fortifyroot_logging_handler, + ) + from traceloop.sdk.tracing.tracing import TracerWrapper + + saved_tracer_instance = getattr(TracerWrapper, "instance", None) + saved_logger_instance = getattr(LoggerWrapper, "instance", None) + saved_logging_enabled = os.environ.get("TRACELOOP_LOGGING_ENABLED") + root_logger = logging.getLogger() + log_exporter = _CapturingLogExporter() + + if hasattr(TracerWrapper, "instance"): + del TracerWrapper.instance + if hasattr(LoggerWrapper, "instance"): + del LoggerWrapper.instance + + try: + os.environ["TRACELOOP_LOGGING_ENABLED"] = "true" + Traceloop.init( + app_name="test-auth-warning-log-filter", + api_endpoint="http://localhost:4318", + api_key="fr-test", + disable_batch=True, + exporter=InMemorySpanExporter(), + logging_exporter=log_exporter, + ) + exporter = _make_traces_http_exporter("http://localhost:4318") + assert hasattr(exporter, "_session") + exporter._session.post = _rejecting_post(401) + + with caplog.at_level("WARNING"): + exporter._export(b"payload") + + provider = LoggerWrapper.get_logging_provider() + assert provider is not None + provider.force_flush() + + assert "FortifyRoot SDK auth warning" in caplog.text + assert not any( + "FortifyRoot SDK auth warning" in body for body in log_exporter.bodies + ) + finally: + for handler in list(root_logger.handlers): + if is_fortifyroot_logging_handler(handler): + root_logger.removeHandler(handler) + if hasattr(LoggerWrapper, "instance"): + del LoggerWrapper.instance + if saved_logger_instance is not None: + LoggerWrapper.instance = saved_logger_instance + if hasattr(TracerWrapper, "instance"): + del TracerWrapper.instance + if saved_tracer_instance is not None: + TracerWrapper.instance = saved_tracer_instance + if saved_logging_enabled is None: + os.environ.pop("TRACELOOP_LOGGING_ENABLED", None) + else: + os.environ["TRACELOOP_LOGGING_ENABLED"] = saved_logging_enabled diff --git a/packages/traceloop-sdk/traceloop/sdk/exporters/auth_warnings.py b/packages/traceloop-sdk/traceloop/sdk/exporters/auth_warnings.py new file mode 100644 index 0000000000..0168a3a4cf --- /dev/null +++ b/packages/traceloop-sdk/traceloop/sdk/exporters/auth_warnings.py @@ -0,0 +1,170 @@ +# NOTE: +# This file has been added by FortifyRoot. +# Original project: https://github.com/traceloop/openllmetry + +import logging +import threading +import time +from typing import Any, Protocol, cast +from urllib.parse import urlparse + +from grpc import RpcError, StatusCode +from opentelemetry.exporter.otlp.proto.grpc._log_exporter import ( + OTLPLogExporter as BaseGRPCLogExporter, +) +from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import ( + OTLPMetricExporter as BaseGRPCMetricExporter, +) +from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import ( + OTLPSpanExporter as BaseGRPCSpanExporter, +) +from opentelemetry.exporter.otlp.proto.http._log_exporter import ( + OTLPLogExporter as BaseHTTPLogExporter, +) +from opentelemetry.exporter.otlp.proto.http.metric_exporter import ( + OTLPMetricExporter as BaseHTTPMetricExporter, +) +from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( + OTLPSpanExporter as BaseHTTPSpanExporter, +) + +AUTH_WARNING_LOGGER_NAME = "fortifyroot.sdk.exporters.auth_warnings" + +logger = logging.getLogger(AUTH_WARNING_LOGGER_NAME) + +_HTTP_AUTH_STATUS_CODES = {401, 403} +_GRPC_AUTH_STATUS_CODES = { + StatusCode.UNAUTHENTICATED, + StatusCode.PERMISSION_DENIED, +} +_AUTH_WARNING_TTL_SECONDS = 60 * 60 +_WARNED_AUTH_FAILURES: dict[tuple[str, str, str], float] = {} +_WARNED_AUTH_FAILURES_LOCK = threading.Lock() + + +class _HTTPExporterProtocol(Protocol): + def _export(self, *args: Any, **kwargs: Any) -> Any: ... + + +def _endpoint_label(endpoint: Any) -> str: + raw_endpoint = str(endpoint or "").strip() + if not raw_endpoint: + return "unknown endpoint" + parsed = urlparse(raw_endpoint) + if parsed.netloc: + return parsed.netloc + return raw_endpoint + + +def _warn_once(signal: str, status: str, endpoint: Any) -> None: + endpoint_label = _endpoint_label(endpoint) + key = (signal, status, endpoint_label) + now = time.monotonic() + with _WARNED_AUTH_FAILURES_LOCK: + last_warned_at = _WARNED_AUTH_FAILURES.get(key) + if ( + last_warned_at is not None + and now - last_warned_at < _AUTH_WARNING_TTL_SECONDS + ): + return + _WARNED_AUTH_FAILURES[key] = now + + logger.warning( + "FortifyRoot SDK auth warning: %s export was rejected by the configured OTLP endpoint (%s) with %s. " + "If this endpoint is FortifyRoot, the SDK API key may be invalid, revoked, deleted, or missing permissions. " + "Telemetry will not reach the OTLP endpoint until a valid credential is configured.", + signal, + endpoint_label, + status, + ) + + +def _warn_for_http_auth_failure(signal: str, endpoint: Any, response: Any) -> None: + status_code = getattr(response, "status_code", None) + if status_code in _HTTP_AUTH_STATUS_CODES: + _warn_once(signal, f"HTTP {status_code}", endpoint) + + +def _call_next_http_export(exporter: Any, *args: Any, **kwargs: Any) -> Any: + next_exporter = cast( + _HTTPExporterProtocol, + super(_HTTPAuthWarningMixin, exporter), + ) + return next_exporter._export(*args, **kwargs) + + +class _AuthWarningClientProxy: + def __init__(self, client: Any, signal: str, endpoint: Any) -> None: + self._client = client + self._signal = signal + self._endpoint = endpoint + + def Export(self, *args: Any, **kwargs: Any) -> Any: # noqa: N802 + try: + return self._client.Export(*args, **kwargs) + except RpcError as exc: + code = exc.code() + if code in _GRPC_AUTH_STATUS_CODES: + _warn_once(self._signal, f"gRPC {code.name}", self._endpoint) + raise + + def __getattr__(self, name: str) -> Any: + return getattr(self._client, name) + + +class _HTTPAuthWarningMixin: + _fortifyroot_export_signal = "telemetry" + + def _export(self, *args: Any, **kwargs: Any) -> Any: + response = _call_next_http_export(self, *args, **kwargs) + _warn_for_http_auth_failure( + self._fortifyroot_export_signal, + getattr(self, "_endpoint", ""), + response, + ) + return response + + +class _GRPCAuthWarningMixin: + _fortifyroot_export_signal = "telemetry" + + def __init__(self, *args: Any, **kwargs: Any) -> None: + super().__init__(*args, **kwargs) + client = getattr(self, "_client", None) + if client is not None: + # Upstream OTel currently only invokes `.Export` on `_client`. + # This proxy intentionally does not subclass the generated stub. + self._client = _AuthWarningClientProxy( + client, + self._fortifyroot_export_signal, + getattr(self, "_endpoint", ""), + ) + + +class FortifyRootHTTPSpanExporter(_HTTPAuthWarningMixin, BaseHTTPSpanExporter): + _fortifyroot_export_signal = "traces" + + +class FortifyRootGRPCSpanExporter(_GRPCAuthWarningMixin, BaseGRPCSpanExporter): + _fortifyroot_export_signal = "traces" + + +class FortifyRootHTTPMetricExporter(_HTTPAuthWarningMixin, BaseHTTPMetricExporter): + _fortifyroot_export_signal = "metrics" + + +class FortifyRootGRPCMetricExporter(_GRPCAuthWarningMixin, BaseGRPCMetricExporter): + _fortifyroot_export_signal = "metrics" + + +class FortifyRootHTTPLogExporter(_HTTPAuthWarningMixin, BaseHTTPLogExporter): + _fortifyroot_export_signal = "logs" + + +class FortifyRootGRPCLogExporter(_GRPCAuthWarningMixin, BaseGRPCLogExporter): + _fortifyroot_export_signal = "logs" + + +def reset_auth_warning_state_for_tests() -> None: + with _WARNED_AUTH_FAILURES_LOCK: + _WARNED_AUTH_FAILURES.clear() diff --git a/packages/traceloop-sdk/traceloop/sdk/logging/logging.py b/packages/traceloop-sdk/traceloop/sdk/logging/logging.py index 4af81e9d9e..3458f37d2d 100644 --- a/packages/traceloop-sdk/traceloop/sdk/logging/logging.py +++ b/packages/traceloop-sdk/traceloop/sdk/logging/logging.py @@ -1,20 +1,28 @@ import logging from typing import Dict, Optional, Any, cast +from urllib.parse import urlparse -from opentelemetry.exporter.otlp.proto.grpc._log_exporter import ( - OTLPLogExporter as GRPCExporter, -) -from opentelemetry.exporter.otlp.proto.http._log_exporter import ( - OTLPLogExporter as HTTPExporter, -) from opentelemetry.sdk.resources import Resource from opentelemetry.sdk._logs.export import LogExporter, BatchLogRecordProcessor from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler from opentelemetry.instrumentation.logging import LoggingInstrumentor +from traceloop.sdk.exporters.auth_warnings import ( + AUTH_WARNING_LOGGER_NAME, + FortifyRootGRPCLogExporter as GRPCExporter, + FortifyRootHTTPLogExporter as HTTPExporter, +) _FORTIFYROOT_LOGGING_HANDLER_MARKER = "_fortifyroot_logging_handler" +_FORTIFYROOT_INTERNAL_EXPORTER_LOGGER_PREFIX = "fortifyroot.sdk.exporters." + + +class _FortifyRootInternalLogFilter(logging.Filter): + def filter(self, record: logging.LogRecord) -> bool: + if record.name == AUTH_WARNING_LOGGER_NAME: + return False + return not record.name.startswith(_FORTIFYROOT_INTERNAL_EXPORTER_LOGGER_PREFIX) class LoggerWrapper(object): @@ -62,10 +70,14 @@ def get_logging_provider(cls) -> Optional[LoggerProvider]: def init_logging_exporter(endpoint: str, headers: Dict[str, str]) -> LogExporter: - if "http" in endpoint.lower() or "https" in endpoint.lower(): - return cast(LogExporter, HTTPExporter(endpoint=f"{endpoint}/v1/logs", headers=headers)) + trimmed_endpoint = endpoint.strip() + if urlparse(trimmed_endpoint).scheme.lower() in {"http", "https"}: + base_url = trimmed_endpoint.rstrip("/") + if not base_url.endswith("/v1/logs"): + base_url = f"{base_url}/v1/logs" + return cast(LogExporter, HTTPExporter(endpoint=base_url, headers=headers)) else: - return cast(LogExporter, GRPCExporter(endpoint=endpoint, headers=headers)) + return cast(LogExporter, GRPCExporter(endpoint=trimmed_endpoint, headers=headers)) def init_logging_provider( @@ -101,6 +113,7 @@ def _attach_root_logging_handler(logging_handler: LoggingHandler) -> None: root_logger.removeHandler(handler) setattr(logging_handler, _FORTIFYROOT_LOGGING_HANDLER_MARKER, True) + logging_handler.addFilter(_FortifyRootInternalLogFilter()) root_logger.addHandler(logging_handler) # Preserve existing app logging levels/formatters when they are already @@ -108,4 +121,3 @@ def _attach_root_logging_handler(logging_handler: LoggingHandler) -> None: # default of exporting INFO and above. if not had_non_fortifyroot_handlers and root_logger.level > logging.INFO: root_logger.setLevel(logging.INFO) - diff --git a/packages/traceloop-sdk/traceloop/sdk/metrics/metrics.py b/packages/traceloop-sdk/traceloop/sdk/metrics/metrics.py index 02d74dda94..b365f1457f 100644 --- a/packages/traceloop-sdk/traceloop/sdk/metrics/metrics.py +++ b/packages/traceloop-sdk/traceloop/sdk/metrics/metrics.py @@ -1,12 +1,7 @@ from collections.abc import Sequence from typing import Dict, Optional, Any +from urllib.parse import urlparse -from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import ( - OTLPMetricExporter as GRPCExporter, -) -from opentelemetry.exporter.otlp.proto.http.metric_exporter import ( - OTLPMetricExporter as HTTPExporter, -) from opentelemetry.semconv_ai import Meters from opentelemetry.sdk.metrics import MeterProvider from opentelemetry.sdk.metrics.export import ( @@ -17,6 +12,10 @@ from opentelemetry.sdk.resources import Resource from opentelemetry import metrics +from traceloop.sdk.exporters.auth_warnings import ( + FortifyRootGRPCMetricExporter as GRPCExporter, + FortifyRootHTTPMetricExporter as HTTPExporter, +) class MetricsWrapper(object): @@ -59,10 +58,14 @@ def set_static_params( def init_metrics_exporter(endpoint: str, headers: Dict[str, str]) -> MetricExporter: - if "http" in endpoint.lower() or "https" in endpoint.lower(): - return HTTPExporter(endpoint=f"{endpoint}/v1/metrics", headers=headers) + trimmed_endpoint = endpoint.strip() + if urlparse(trimmed_endpoint).scheme.lower() in {"http", "https"}: + base_url = trimmed_endpoint.rstrip("/") + if not base_url.endswith("/v1/metrics"): + base_url = f"{base_url}/v1/metrics" + return HTTPExporter(endpoint=base_url, headers=headers) else: - return GRPCExporter(endpoint=endpoint, headers=headers) + return GRPCExporter(endpoint=trimmed_endpoint, headers=headers) def init_metrics_provider( diff --git a/packages/traceloop-sdk/traceloop/sdk/tracing/tracing.py b/packages/traceloop-sdk/traceloop/sdk/tracing/tracing.py index 4473518cc5..127c509c3d 100644 --- a/packages/traceloop-sdk/traceloop/sdk/tracing/tracing.py +++ b/packages/traceloop-sdk/traceloop/sdk/tracing/tracing.py @@ -10,12 +10,6 @@ from colorama import Fore from opentelemetry import trace -from opentelemetry.exporter.otlp.proto.http.trace_exporter import ( - OTLPSpanExporter as HTTPExporter, -) -from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import ( - OTLPSpanExporter as GRPCExporter, -) from opentelemetry.sdk.resources import Resource from opentelemetry.sdk.trace import TracerProvider, SpanProcessor, ReadableSpan from opentelemetry.sdk.trace.sampling import Sampler @@ -35,6 +29,10 @@ from opentelemetry.semconv_ai import SpanAttributes from traceloop.sdk.images.image_uploader import ImageUploader from traceloop.sdk.instruments import Instruments +from traceloop.sdk.exporters.auth_warnings import ( + FortifyRootGRPCSpanExporter as GRPCExporter, + FortifyRootHTTPSpanExporter as HTTPExporter, +) from traceloop.sdk.tracing.content_allow_list import ContentAllowList from traceloop.sdk.utils import is_notebook from traceloop.sdk.utils.package_check import is_package_installed