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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Fixed

- `tilebox-workflows`: Fix logging initialization and workflow startup with OpenTelemetry 1.45 and newer.

## [0.63.2] - 2026-10-04

### Fixed
Expand Down
2 changes: 1 addition & 1 deletion prek.toml
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ hooks = [

[[repos]]
repo = "https://github.com/charliermarsh/ruff-pre-commit"
rev = "v0.16.8"
rev = "v0.16.10"
hooks = [
{
id = "ruff-check",
Expand Down
9 changes: 8 additions & 1 deletion tilebox-datasets/tests/test_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,10 +59,17 @@ def test_client_url_environment(
open_channel_mock.assert_called_once_with(expected_url, "runner-key", rpc_method_prefix=None)


def test_heavy_imports_are_lazy() -> None:
@pytest.mark.parametrize("grpc_verbosity", [None, "DEBUG"])
def test_heavy_imports_are_lazy(monkeypatch: pytest.MonkeyPatch, grpc_verbosity: str | None) -> None:
monkeypatch.delenv("GRPC_VERBOSITY", raising=False)
if grpc_verbosity is not None:
monkeypatch.setenv("GRPC_VERBOSITY", grpc_verbosity)
code = (
"import os\n"
"import sys\n"
"import tilebox.datasets as datasets\n"
f"assert os.environ['GRPC_VERBOSITY'] == {grpc_verbosity or 'ERROR'!r}\n"
"assert 'grpc' not in sys.modules\n"
"assert not {'pandas', 'xarray'} & sys.modules.keys()\n"
"assert set(datasets.__all__) <= set(dir(datasets))\n"
"from tilebox.datasets.query.time_interval import TimeInterval\n"
Expand Down
3 changes: 3 additions & 0 deletions tilebox-datasets/tilebox/datasets/__init__.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
from typing import TYPE_CHECKING, Any

# Apply native logging defaults before any submodule can initialize gRPC.
import _tilebox.grpc # noqa: F401

if TYPE_CHECKING:
from tilebox.datasets.aio.timeseries import TimeseriesCollection, TimeseriesDataset
from tilebox.datasets.datapoints import iter_datapoints
Expand Down
6 changes: 6 additions & 0 deletions tilebox-grpc/_tilebox/grpc/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
import os

# Suppress native gRPC fork-log spam while preserving explicit user settings.
# Remove once our minimum grpcio version resolves https://github.com/grpc/grpc/issues/42293.
# This must run before gRPC initializes; it also affects other Abseil-based native logging.
os.environ.setdefault("GRPC_VERBOSITY", "ERROR")
16 changes: 12 additions & 4 deletions tilebox-grpc/tests/test_fork_logging.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,10 @@

Run on macOS to check RPCs and report the known fork-log issue:

uv run --isolated --no-project --with grpcio==1.84.0 --with pytest pytest -c /dev/null -p no:cacheprovider -rx tilebox-grpc/tests/test_fork_logging.py
uv run --isolated --no-project --with-editable ./tilebox-grpc --with grpcio==1.84.0 --with pytest pytest -c /dev/null -p no:cacheprovider -rx tilebox-grpc/tests/test_fork_logging.py

RPC failures fail the test. Known fork diagnostics produce XFAIL after the RPC checks pass.
RPC failures or diagnostics with Tilebox's default verbosity fail the test.
With explicit INFO verbosity, known fork diagnostics produce XFAIL after the RPC checks pass.
Linux may use vfork and not reproduce the diagnostics. This check uses no cloud credentials.
"""

Expand All @@ -18,7 +19,11 @@

# RPCs must survive CLI subprocess launches; known fork-log noise is reported separately.
@pytest.mark.skipif(sys.platform == "win32", reason="Windows does not use POSIX fork handlers")
def test_cli_spawn_during_rpcs() -> None:
@pytest.mark.parametrize("grpc_verbosity", [None, "INFO"])
def test_cli_spawn_during_rpcs(grpc_verbosity: str | None) -> None:
env = {key: value for key, value in os.environ.items() if key != "GRPC_VERBOSITY"}
if grpc_verbosity is not None:
env["GRPC_VERBOSITY"] = grpc_verbosity
result = subprocess.run( # noqa: S603
[
sys.executable,
Expand All @@ -29,6 +34,7 @@ def test_cli_spawn_during_rpcs() -> None:
from concurrent.futures import ThreadPoolExecutor
from threading import Event

import _tilebox.grpc
import grpc

with ThreadPoolExecutor(max_workers=4) as executor:
Expand Down Expand Up @@ -65,9 +71,11 @@ def poll():
capture_output=True,
text=True,
timeout=60,
env={**os.environ, "GRPC_VERBOSITY": "INFO"},
env=env,
check=False,
)
assert result.returncode == 0, result.stderr
if grpc_verbosity is None:
assert "FD from fork parent still in poll list" not in result.stderr, result.stderr
if "FD from fork parent still in poll list" in result.stderr:
pytest.xfail("Known gRPC fork diagnostics: https://github.com/grpc/grpc/issues/42293")
10 changes: 5 additions & 5 deletions tilebox-workflows/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -27,14 +27,14 @@ dependencies = [
"grpcio>=1.84.0",
"google-auth[requests]>=2.29",
"azure-identity>=1.23",
"opentelemetry-api>=1.43.0",
"opentelemetry-exporter-otlp-proto-http>=1.43.0",
"opentelemetry-sdk>=1.43.0",
"opentelemetry-api>=1.44.0",
"opentelemetry-exporter-otlp-proto-http>=1.44.0",
"opentelemetry-sdk>=1.44.0",
"tenacity>=8",
"boto3>=1.40.2",
"obstore>=0.8.2",
"opentelemetry-proto>=1.43.0",
"opentelemetry-instrumentation-logging>=0.64b0",
"opentelemetry-proto>=1.44.0",
"opentelemetry-instrumentation-logging>=0.65b0",
"msgspec>=0.19",
]

Expand Down
65 changes: 65 additions & 0 deletions tilebox-workflows/tests/observability/test_logging.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
import msgspec
import pytest
from affine import Affine
from opentelemetry._logs import SeverityNumber
from opentelemetry.sdk._logs.export import InMemoryLogRecordExporter, SimpleLogRecordProcessor
from opentelemetry.sdk.trace import Span, TracerProvider

Expand All @@ -37,6 +38,70 @@ def isolated_logging(monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureF
caplog.set_level(logging.ERROR, logger=internal_logger.name)


@pytest.mark.parametrize("from_environment", [False, True])
@pytest.mark.parametrize(
("endpoint", "expected"),
[
("https://collector.example", "https://collector.example/v1/logs"),
("https://collector.example/tenant", "https://collector.example/tenant/v1/logs"),
("https://collector.example/tenant/", "https://collector.example/tenant/v1/logs"),
("https://collector.example/tenant/v1/logs", "https://collector.example/tenant/v1/logs"),
],
)
def test_otel_log_exporter_endpoint(
monkeypatch: pytest.MonkeyPatch, endpoint: str, expected: str, from_environment: bool
) -> None:
monkeypatch.setenv("OTEL_LOGS_ENDPOINT", endpoint if from_environment else "https://ignored.example")
headers = {"Authorization": "Bearer test-key"}
with patch.object(observability, "OTLPLogExporter", return_value=InMemoryLogRecordExporter()) as factory:
processor = observability._otel_log_exporter(None if from_environment else endpoint, headers=headers)
try:
factory.assert_called_once_with(endpoint=expected, headers=headers)
finally:
processor.shutdown()


@pytest.mark.parametrize("formatted", [False, True])
@pytest.mark.parametrize(
("level", "severity", "text"),
[(logging.WARNING, SeverityNumber.WARN, "WARN"), (logging.CRITICAL, SeverityNumber.FATAL, "FATAL")],
)
def test_otel_handler_preserves_record_fields(formatted: bool, level: int, severity: SeverityNumber, text: str) -> None:
exporter = InMemoryLogRecordExporter()
provider = observability.LoggerProvider()
provider.add_log_record_processor(SimpleLogRecordProcessor(exporter))
handler = observability.OTELLoggingHandler(logger_provider=provider)
if formatted:
handler.setFormatter(logging.Formatter("%(name)s: %(message)s"))
record = logging.LogRecord("task.example", level, "workflow.py", 42, "task %s", ("started",), None)
record.created = 1234.5
try:
handler.handle(record)
(exported,) = exporter.get_finished_logs()
assert exported.instrumentation_scope is not None
assert exported.instrumentation_scope.name == "task.example"
assert exported.log_record.timestamp == 1_234_500_000_000
assert exported.log_record.observed_timestamp is not None
assert exported.log_record.severity_number == severity
assert exported.log_record.severity_text == text
assert exported.log_record.body == ("task.example: task started" if formatted else "task started")
assert not exported.log_record.attributes
finally:
provider.shutdown()


def test_otel_handler_suppresses_recursive_exports() -> None:
handler = observability.OTELLoggingHandler(logger_provider=observability.LoggerProvider())
record = logging.LogRecord("task.example", logging.INFO, __file__, 0, "outer", (), None)
recursive_record = logging.LogRecord("task.example", logging.WARNING, __file__, 0, "exporter diagnostic", (), None)
with patch.object(observability, "get_otel_logger") as get_logger:
emit = get_logger.return_value.emit
emit.side_effect = lambda _: handler.emit(recursive_record)
handler.emit(record)
handler.emit(record)
assert [call.args[0].body for call in emit.call_args_list] == ["outer", "outer"]


@pytest.mark.usefixtures("isolated_logging")
@pytest.mark.parametrize(
("task_level", "internal_level"), [(logging.DEBUG, logging.ERROR), (logging.ERROR, logging.DEBUG)]
Expand Down
25 changes: 25 additions & 0 deletions tilebox-workflows/tests/observability/test_tracing.py
Original file line number Diff line number Diff line change
@@ -1,13 +1,38 @@
from collections.abc import Iterator
from unittest.mock import patch

import pytest
from opentelemetry.context import Context
from opentelemetry.sdk.trace import ReadableSpan, Span, TracerProvider
from opentelemetry.sdk.trace.export import SpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter

from tilebox.workflows.observability import tracing


@pytest.mark.parametrize("from_environment", [False, True])
@pytest.mark.parametrize(
("endpoint", "expected"),
[
("https://collector.example", "https://collector.example/v1/traces"),
("https://collector.example/tenant", "https://collector.example/tenant/v1/traces"),
("https://collector.example/tenant/", "https://collector.example/tenant/v1/traces"),
("https://collector.example/tenant/v1/traces", "https://collector.example/tenant/v1/traces"),
],
)
def test_otel_span_exporter_endpoint(
monkeypatch: pytest.MonkeyPatch, endpoint: str, expected: str, from_environment: bool
) -> None:
monkeypatch.setenv("OTEL_TRACES_ENDPOINT", endpoint if from_environment else "https://ignored.example")
headers = {"Authorization": "Bearer test-key"}
with patch.object(tracing, "OTLPSpanExporter", return_value=InMemorySpanExporter()) as factory:
processor = tracing._otel_span_exporter(None if from_environment else endpoint, headers=headers)
try:
factory.assert_called_once_with(endpoint=expected, headers=headers)
finally:
processor.shutdown()


class RecordingSpanProcessor(SpanProcessor):
def __init__(self) -> None:
self.span_names: list[str] = []
Expand Down
9 changes: 8 additions & 1 deletion tilebox-workflows/tests/runner/test_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,10 +42,17 @@ def execute(self, context: ExecutionContext) -> None:
assert cache["result"] == b"processed"


def test_task_authoring_imports_are_lazy() -> None:
@pytest.mark.parametrize("grpc_verbosity", [None, "DEBUG"])
def test_task_authoring_imports_are_lazy(monkeypatch: pytest.MonkeyPatch, grpc_verbosity: str | None) -> None:
monkeypatch.delenv("GRPC_VERBOSITY", raising=False)
if grpc_verbosity is not None:
monkeypatch.setenv("GRPC_VERBOSITY", grpc_verbosity)
code = (
"import os\n"
"import sys\n"
"import tilebox.workflows as workflows\n"
f"assert os.environ['GRPC_VERBOSITY'] == {grpc_verbosity or 'ERROR'!r}\n"
"assert 'grpc' not in sys.modules\n"
"heavy = {'pandas', 'xarray', 'boto3', 'google.cloud.storage', 'google.auth', 'azure.identity', 'obstore', 'ipywidgets', 'opentelemetry.sdk'}\n"
"assert not heavy & sys.modules.keys()\n"
"assert set(workflows.__all__) <= set(dir(workflows))\n"
Expand Down
3 changes: 3 additions & 0 deletions tilebox-workflows/tilebox/workflows/__init__.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
import os
from typing import TYPE_CHECKING, Any

# Apply native logging defaults before any submodule can initialize gRPC.
import _tilebox.grpc # noqa: F401

if TYPE_CHECKING:
from tilebox.workflows.client import Client
from tilebox.workflows.data import Job
Expand Down
44 changes: 36 additions & 8 deletions tilebox-workflows/tilebox/workflows/observability/logging.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,18 +8,20 @@
import sys
import threading
import traceback
from contextvars import ContextVar
from datetime import timedelta
from importlib.metadata import PackageNotFoundError, version
from typing import Any, ClassVar, TextIO
from uuid import UUID, uuid4

import msgspec
from opentelemetry._logs import LogRecord, get_logger_provider
from opentelemetry._logs import get_logger as get_otel_logger
from opentelemetry.exporter.otlp.proto.http._log_exporter import (
DEFAULT_LOGS_EXPORT_PATH,
OTLPLogExporter,
_append_logs_path,
)
from opentelemetry.instrumentation.logging.handler import LoggingHandler
from opentelemetry.instrumentation.log_utils import std_to_otel
from opentelemetry.sdk._logs import LoggerProvider
from opentelemetry.sdk._logs.export import BatchLogRecordProcessor
from opentelemetry.sdk.resources import (
Expand All @@ -34,7 +36,6 @@
Resource,
)
from opentelemetry.semconv.attributes import exception_attributes
from opentelemetry.util.types import _ExtendedAttributes

from tilebox.workflows._serialization import normalize_log_value
from tilebox.workflows.observability._log_pipe import _PipeWriter, _StructuredHandler
Expand Down Expand Up @@ -110,15 +111,42 @@ def _sanitize_otel_attributes(attributes: dict[str, Any]) -> dict[str, Any]:
return {str(key): _sanitize_otel_attribute_value(value) for key, value in attributes.items()}


class OTELLoggingHandler(LoggingHandler):
def _get_attributes(self, record: logging.LogRecord) -> _ExtendedAttributes:
class OTELLoggingHandler(logging.Handler):
def __init__(self, level: int = logging.NOTSET, logger_provider: LoggerProvider | None = None) -> None:
super().__init__(level)
self._logger_provider = logger_provider or get_logger_provider()
self._emitting = ContextVar("tilebox_otel_emitting", default=False)

def emit(self, record: logging.LogRecord) -> None:
# Export processors can log themselves; do not recursively export those records.
if self._emitting.get():
return
token = self._emitting.set(True)
try:
get_otel_logger(record.name, logger_provider=self._logger_provider).emit(
LogRecord(
timestamp=int(record.created * 1e9),
severity_number=std_to_otel(record.levelno),
severity_text={"WARNING": "WARN", "CRITICAL": "FATAL"}.get(record.levelname, record.levelname),
body=self.format(record) if self.formatter else record.getMessage(),
attributes=self._get_attributes(record),
)
)
finally:
self._emitting.reset(token)

def flush(self) -> None:
if callable(force_flush := getattr(self._logger_provider, "force_flush", None)):
# Match OTEL's handler: flushing under the logging lock can deadlock.
threading.Thread(target=force_flush).start()

def _get_attributes(self, record: logging.LogRecord) -> dict[str, Any]:
cached = getattr(record, "_tilebox_otel_attributes", None)
if cached is not None:
return cached
attributes = _sanitize_otel_attributes(_record_attributes(record))

# the default implementation returns attributes for the filepath, lineno and function of the log record
# we don't want that by default, so we override it to return an empty dict
# Export only task context and exception details, not stdlib logging metadata.
if record.exc_info:
exctype, value, tb = record.exc_info
if exctype is not None:
Expand Down Expand Up @@ -237,7 +265,7 @@ def _otel_log_exporter(
)

if not endpoint.endswith(DEFAULT_LOGS_EXPORT_PATH):
endpoint = _append_logs_path(endpoint)
endpoint = f"{endpoint.rstrip('/')}/{DEFAULT_LOGS_EXPORT_PATH}"

if export_interval is None:
export_interval_env = os.environ.get(_OTEL_EXPORT_INTERVAL_ENV_VAR, None)
Expand Down
3 changes: 1 addition & 2 deletions tilebox-workflows/tilebox/workflows/observability/tracing.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@
from opentelemetry.exporter.otlp.proto.http.trace_exporter import (
DEFAULT_TRACES_EXPORT_PATH,
OTLPSpanExporter,
_append_trace_path,
)
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace import TracerProvider
Expand Down Expand Up @@ -159,7 +158,7 @@ def _otel_span_exporter(
)

if not endpoint.endswith(DEFAULT_TRACES_EXPORT_PATH):
endpoint = _append_trace_path(endpoint)
endpoint = f"{endpoint.rstrip('/')}/{DEFAULT_TRACES_EXPORT_PATH}"

if export_interval is None:
export_interval_env = os.environ.get(_OTEL_EXPORT_INTERVAL_ENV_VAR, None)
Expand Down
Loading
Loading