diff --git a/python/src/smooai_observability/bootstrap/__init__.py b/python/src/smooai_observability/bootstrap/__init__.py index 8bf1d21..078c0b1 100644 --- a/python/src/smooai_observability/bootstrap/__init__.py +++ b/python/src/smooai_observability/bootstrap/__init__.py @@ -11,7 +11,7 @@ ## Required env vars SMOOAI_OBSERVABILITY_ENDPOINT — base ingest URL (e.g. "https://api.smoo.ai"). - ``/v1/traces`` + ``/v1/metrics`` are appended. + ``/v1/traces`` + ``/v1/metrics`` + ``/v1/logs`` are appended. ## Auth (pick ONE; pre-minted token wins if both present) SMOOAI_OBSERVABILITY_TOKEN — pre-minted Bearer JWT (not refreshed). @@ -44,6 +44,7 @@ class BootstrapEnv: endpoint: str | None = None traces_endpoint: str | None = None metrics_endpoint: str | None = None + logs_endpoint: str | None = None token: str | None = None auth_url: str | None = None client_id: str | None = None @@ -98,6 +99,7 @@ def bootstrap_observability( endpoint=o.endpoint or os.environ.get("SMOOAI_OBSERVABILITY_ENDPOINT"), traces_endpoint=o.traces_endpoint or os.environ.get("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT"), metrics_endpoint=o.metrics_endpoint or os.environ.get("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT"), + logs_endpoint=o.logs_endpoint or os.environ.get("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT"), token=o.token or os.environ.get("SMOOAI_OBSERVABILITY_TOKEN"), auth_url=o.auth_url or os.environ.get("SMOOAI_OBSERVABILITY_AUTH_URL"), client_id=o.client_id or os.environ.get("SMOOAI_OBSERVABILITY_CLIENT_ID"), @@ -139,6 +141,7 @@ def bootstrap_observability( traces_endpoint = env.traces_endpoint or (f"{_strip_trailing_slash(env.endpoint)}/v1/traces" if env.endpoint else None) metrics_endpoint = env.metrics_endpoint or (f"{_strip_trailing_slash(env.endpoint)}/v1/metrics" if env.endpoint else None) + logs_endpoint = env.logs_endpoint or (f"{_strip_trailing_slash(env.endpoint)}/v1/logs" if env.endpoint else None) otel_handle = _maybe_setup_otel( service_name=env.service_name, @@ -146,6 +149,7 @@ def bootstrap_observability( release=env.release, traces_endpoint=traces_endpoint, metrics_endpoint=metrics_endpoint, + logs_endpoint=logs_endpoint, static_headers=static_headers, token_provider=token_provider, ) @@ -192,6 +196,7 @@ def _maybe_setup_otel(**kwargs: object) -> object: release=kwargs.get("release"), # type: ignore[arg-type] otlp_endpoint=kwargs.get("traces_endpoint"), # type: ignore[arg-type] otlp_metrics_endpoint=kwargs.get("metrics_endpoint"), # type: ignore[arg-type] + otlp_logs_endpoint=kwargs.get("logs_endpoint"), # type: ignore[arg-type] otlp_headers=dict(kwargs.get("static_headers") or {}), # type: ignore[arg-type] token_provider=kwargs.get("token_provider"), # type: ignore[arg-type] ) diff --git a/python/src/smooai_observability/otel/auth_injecting_exporter.py b/python/src/smooai_observability/otel/auth_injecting_exporter.py index bb6a14b..64bac9c 100644 --- a/python/src/smooai_observability/otel/auth_injecting_exporter.py +++ b/python/src/smooai_observability/otel/auth_injecting_exporter.py @@ -18,6 +18,7 @@ from __future__ import annotations +from opentelemetry.exporter.otlp.proto.http._log_exporter import OTLPLogExporter from opentelemetry.exporter.otlp.proto.http.metric_exporter import OTLPMetricExporter from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter @@ -56,6 +57,28 @@ def _export(self, serialized_data, timeout_sec=None): # type: ignore[override] return resp +class AuthInjectingLogExporter(OTLPLogExporter): + """OTLP/HTTP log exporter that injects a fresh Bearer per export.""" + + def __init__( + self, + *, + endpoint: str, + token_provider: TokenProvider, + headers: dict[str, str] | None = None, + timeout: float | None = None, + ) -> None: + super().__init__(endpoint=endpoint, headers=headers, timeout=timeout) + self._token_provider = token_provider + + def _export(self, serialized_data, timeout_sec=None): # type: ignore[override] + _inject(self._session.headers, self._token_provider) + resp = super()._export(serialized_data, timeout_sec) + if getattr(resp, "status_code", None) == 401: + self._token_provider.invalidate() + return resp + + class AuthInjectingMetricExporter(OTLPMetricExporter): """OTLP/HTTP metric exporter that injects a fresh Bearer per export.""" diff --git a/python/src/smooai_observability/otel/setup_otel_sdk.py b/python/src/smooai_observability/otel/setup_otel_sdk.py index 5591502..c092866 100644 --- a/python/src/smooai_observability/otel/setup_otel_sdk.py +++ b/python/src/smooai_observability/otel/setup_otel_sdk.py @@ -36,6 +36,7 @@ class SetupOtelOptions: service_name: str = DEFAULT_SERVICE_NAME otlp_endpoint: str | None = None otlp_metrics_endpoint: str | None = None + otlp_logs_endpoint: str | None = None otlp_headers: dict[str, str] = field(default_factory=dict) environment: str | None = None release: str | None = None @@ -49,10 +50,14 @@ class SetupOtelOptions: class OtelSdkHandle: tracer_provider: Any = None meter_provider: Any = None + logger_provider: Any = None + # The stdlib ``logging.Handler`` bridging root logging → OTLP logs. Stored so + # ``shutdown`` can detach it from the root logger (undo the global mutation). + log_handler: Any = None enabled: bool = False def flush(self, timeout_millis: int = 2_000) -> None: - for provider in (self.tracer_provider, self.meter_provider): + for provider in (self.tracer_provider, self.meter_provider, self.logger_provider): if provider is None: continue try: @@ -62,7 +67,14 @@ def flush(self, timeout_millis: int = 2_000) -> None: def shutdown(self) -> None: global _installed - for provider in (self.tracer_provider, self.meter_provider): + if self.log_handler is not None: + try: + import logging + + logging.getLogger().removeHandler(self.log_handler) + except Exception: + pass + for provider in (self.tracer_provider, self.meter_provider, self.logger_provider): if provider is None: continue try: @@ -91,8 +103,13 @@ def setup_otel_sdk(options: SetupOtelOptions | None = None, **kwargs: Any) -> Ot opts = options or SetupOtelOptions(**kwargs) try: + import logging + from opentelemetry import metrics as otel_metrics from opentelemetry import trace as otel_trace + from opentelemetry._logs import set_logger_provider + from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler + from opentelemetry.sdk._logs.export import BatchLogRecordProcessor from opentelemetry.sdk.metrics import MeterProvider from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader from opentelemetry.sdk.resources import Resource @@ -104,10 +121,12 @@ def setup_otel_sdk(options: SetupOtelOptions | None = None, **kwargs: Any) -> Ot return _installed try: + from opentelemetry.exporter.otlp.proto.http._log_exporter import OTLPLogExporter from opentelemetry.exporter.otlp.proto.http.metric_exporter import OTLPMetricExporter from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter from .auth_injecting_exporter import ( + AuthInjectingLogExporter, AuthInjectingMetricExporter, AuthInjectingTraceExporter, ) @@ -160,13 +179,38 @@ def setup_otel_sdk(options: SetupOtelOptions | None = None, **kwargs: Any) -> Ot ) meter_provider = MeterProvider(resource=resource, metric_readers=[metric_reader]) + # --- logs ----------------------------------------------------------- + log_endpoint = opts.otlp_logs_endpoint or os.environ.get("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT") or os.environ.get("OTEL_EXPORTER_OTLP_ENDPOINT") + if opts.token_provider is not None and log_endpoint: + log_exporter: Any = AuthInjectingLogExporter( + endpoint=log_endpoint, + token_provider=opts.token_provider, + headers=opts.otlp_headers or None, + ) + elif log_endpoint: + log_exporter = OTLPLogExporter(endpoint=log_endpoint, headers=opts.otlp_headers or None) + else: + log_exporter = OTLPLogExporter(headers=opts.otlp_headers or None) + + logger_provider = LoggerProvider(resource=resource) + logger_provider.add_log_record_processor(BatchLogRecordProcessor(log_exporter)) + # Bridge stdlib logging → OTel log records. The handler reads the ACTIVE + # span context per record, so logs emitted inside a span carry its real + # W3C trace_id/span_id (verified in tests). Attached to the root logger + # so any stdlib logger (including @smooai/logger's bridge) is captured. + log_handler = LoggingHandler(level=logging.NOTSET, logger_provider=logger_provider) + if not opts.skip_start: otel_trace.set_tracer_provider(tracer_provider) otel_metrics.set_meter_provider(meter_provider) + set_logger_provider(logger_provider) + logging.getLogger().addHandler(log_handler) handle = OtelSdkHandle( tracer_provider=tracer_provider, meter_provider=meter_provider, + logger_provider=logger_provider, + log_handler=log_handler, enabled=True, ) _installed = handle diff --git a/python/tests/test_otel_logs.py b/python/tests/test_otel_logs.py new file mode 100644 index 0000000..fea2f97 --- /dev/null +++ b/python/tests/test_otel_logs.py @@ -0,0 +1,100 @@ +"""OTLP logs signal + trace correlation tests. + +The logs signal reuses the same endpoint/auth/enable path as traces + metrics +and bridges stdlib ``logging`` → OTel log records via ``LoggingHandler``. The +load-bearing behavior is that a log emitted inside an active span carries that +span's real W3C trace_id/span_id — that's what correlates logs to traces in the +product. We verify it against the exact ``LoggingHandler`` class ``setup_otel_sdk`` +attaches, using an in-memory exporter (no network). +""" + +from __future__ import annotations + +import logging + +import pytest +from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler +from opentelemetry.sdk._logs.export import InMemoryLogExporter, SimpleLogRecordProcessor +from opentelemetry.sdk.trace import TracerProvider + +from smooai_observability.otel import ( + SetupOtelOptions, + reset_otel_sdk_for_tests, + setup_otel_sdk, +) + + +@pytest.fixture(autouse=True) +def _reset(): + reset_otel_sdk_for_tests() + yield + reset_otel_sdk_for_tests() + + +def test_setup_wires_logger_provider(): + handle = setup_otel_sdk( + SetupOtelOptions( + service_name="svc", + otlp_logs_endpoint="https://example.test/v1/logs", + skip_start=True, + ) + ) + assert handle.enabled is True + assert handle.logger_provider is not None + assert handle.log_handler is not None + + +def test_skip_start_does_not_touch_root_logger(): + before = list(logging.getLogger().handlers) + setup_otel_sdk( + SetupOtelOptions( + service_name="svc", + otlp_logs_endpoint="https://example.test/v1/logs", + skip_start=True, + ) + ) + assert list(logging.getLogger().handlers) == before + + +def test_log_within_span_carries_trace_context(): + exporter = InMemoryLogExporter() + provider = LoggerProvider() + provider.add_log_record_processor(SimpleLogRecordProcessor(exporter)) + handler = LoggingHandler(level=logging.NOTSET, logger_provider=provider) + + log = logging.getLogger("bridge_correlation_test") + log.setLevel(logging.DEBUG) + log.propagate = False + log.addHandler(handler) + + tracer = TracerProvider().get_tracer("test") + with tracer.start_as_current_span("work") as span: + ctx = span.get_span_context() + log.warning("inside the span") + + records = exporter.get_finished_logs() + assert records, "expected a log record to be exported" + lr = records[-1].log_record + # Real W3C ids from the active span — not a fabricated/zero id. + assert lr.trace_id == ctx.trace_id + assert lr.span_id == ctx.span_id + assert lr.trace_id != 0 and lr.span_id != 0 + + +def test_log_outside_span_has_no_trace_context(): + exporter = InMemoryLogExporter() + provider = LoggerProvider() + provider.add_log_record_processor(SimpleLogRecordProcessor(exporter)) + handler = LoggingHandler(level=logging.NOTSET, logger_provider=provider) + + log = logging.getLogger("bridge_no_span_test") + log.setLevel(logging.DEBUG) + log.propagate = False + log.addHandler(handler) + + log.warning("no active span") + + records = exporter.get_finished_logs() + assert records + # No span active → zero ids (OTel's "invalid" sentinel), not a random uuid. + assert records[-1].log_record.trace_id == 0