Skip to content
Open
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
7 changes: 6 additions & 1 deletion python/src/smooai_observability/bootstrap/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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"),
Expand Down Expand Up @@ -139,13 +141,15 @@ 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,
environment=env.environment,
release=env.release,
traces_endpoint=traces_endpoint,
metrics_endpoint=metrics_endpoint,
logs_endpoint=logs_endpoint,
static_headers=static_headers,
token_provider=token_provider,
)
Expand Down Expand Up @@ -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]
)
Expand Down
23 changes: 23 additions & 0 deletions python/src/smooai_observability/otel/auth_injecting_exporter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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."""

Expand Down
48 changes: 46 additions & 2 deletions python/src/smooai_observability/otel/setup_otel_sdk.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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:
Expand All @@ -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:
Expand Down Expand Up @@ -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
Expand All @@ -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,
)
Expand Down Expand Up @@ -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
Expand Down
100 changes: 100 additions & 0 deletions python/tests/test_otel_logs.py
Original file line number Diff line number Diff line change
@@ -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
Loading