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
29 changes: 29 additions & 0 deletions .github/workflows/main.yml
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,35 @@ jobs:
paths: "test-report.xml"
if: always()

windows-runtime:
name: Test Windows worker runtime
strategy:
matrix:
python-version: ["3.11", "3.14"]
runs-on: windows-latest
steps:
- uses: actions/checkout@v7
with:
lfs: true
- name: Set up uv
uses: astral-sh/setup-uv@v10.1.0
with:
python-version: ${{ matrix.python-version }}
enable-cache: true
cache-python: true
cache-dependency-glob: "uv.lock"
- name: Sync
run: uv sync --package tilebox-workflows --frozen
- name: Test worker lifecycle and logging
run: >-
uv run --frozen --package tilebox-workflows pytest
tilebox-workflows/tests/runner/test_worker_server.py
tilebox-workflows/tests/runner/test_worker_concurrency.py
tilebox-workflows/tests/runner/test_runtime_logging.py
tilebox-workflows/tests/runner/test_log_stream.py
tilebox-workflows/tests/observability/test_logging.py
tilebox-workflows/tests/observability/test_tracing.py

minimum-dependencies:
name: Test minimum dependencies
runs-on: ubuntu-latest
Expand Down
16 changes: 15 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,9 +7,22 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

## [0.64.0] - 2026-10-06

### Added

- `tilebox-workflows`: Enable live workflow logs in the Windows CLI and improve local worker startup on Windows and Unix.

### Changed

- `tilebox-workflows`: Remove the legacy CLI log pipe. Older Unix CLIs can still execute workflows, but require an upgrade to receive structured live logs.

### Fixed

- `tilebox-workflows`: Fix logging initialization and workflow startup with OpenTelemetry 1.45 and newer.
- `tilebox-workflows`: Always apply Tilebox service metadata to logs and traces, with explicit resource attributes overriding those defaults. Initialize tracing with the same runtime identity as logging instead of relying on OpenTelemetry's defaults.
- `tilebox-workflows`: Bound the final API log flush so stalled log exports do not delay CLI worker shutdown.
- `tilebox-workflows`: Preserve large integer attributes in live CLI logs as strings instead of dropping the log record.

## [0.63.2] - 2026-10-04

Expand Down Expand Up @@ -588,7 +601,8 @@ the first client that does not cache data (since it's already on the local file
- Released under the [MIT](https://opensource.org/license/mit) license.
- Released packages: `tilebox-datasets`, `tilebox-workflows`, `tilebox-storage`, `tilebox-grpc`

[Unreleased]: https://github.com/tilebox/tilebox-python/compare/v0.63.2...HEAD
[Unreleased]: https://github.com/tilebox/tilebox-python/compare/v0.64.0...HEAD
[0.64.0]: https://github.com/tilebox/tilebox-python/compare/v0.63.2...v0.64.0
[0.63.2]: https://github.com/tilebox/tilebox-python/compare/v0.63.1...v0.63.2
[0.63.1]: https://github.com/tilebox/tilebox-python/compare/v0.63.0...v0.63.1
[0.63.0]: https://github.com/tilebox/tilebox-python/compare/v0.62.0...v0.63.0
Expand Down
154 changes: 76 additions & 78 deletions tilebox-workflows/tests/observability/test_logging.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
from tilebox.workflows.observability import _logging as structured_logging
from tilebox.workflows.observability import logging as observability
from tilebox.workflows.observability import tracing
from tilebox.workflows.observability._log_stream import _LogQueue, _OTLPQueueHandler
from tilebox.workflows.observability._logging import StructuredLogger, internal_logger, logger, root_logger, task_logger


Expand All @@ -32,8 +33,7 @@ def isolated_logging(monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureF
monkeypatch.setattr(observability, "_api_handler", None)
monkeypatch.setattr(observability, "_console_handlers", [])
monkeypatch.setattr(observability, "_console_configured", False)
monkeypatch.setattr(observability, "_writer", None)
monkeypatch.delenv("TILEBOX_LOG_FD", raising=False)
monkeypatch.setattr(observability, "managed_log_queue", None)
caplog.set_level(logging.INFO, logger=task_logger.name)
caplog.set_level(logging.ERROR, logger=internal_logger.name)

Expand Down Expand Up @@ -190,20 +190,20 @@ def test_export_attributes_capture_context_once_without_mutating_record(structur


@pytest.mark.parametrize("stdlib", [False, True])
@pytest.mark.parametrize("pipe_first", [False, True])
def test_codecs_run_once_per_record_across_handlers(stdlib: bool, pipe_first: bool) -> None:
@pytest.mark.parametrize("stream_first", [False, True])
def test_codecs_run_once_per_record_across_handlers(stdlib: bool, stream_first: bool) -> None:
codec = registry.find(Affine)
assert codec is not None
encode = MagicMock(wraps=codec.encode)
target = logging.Logger("normalized-records", logging.INFO) # noqa: LOG001 -- isolated from process-wide loggers
writer = MagicMock(spec=observability._PipeWriter)
queue = _LogQueue()
exporters = [InMemoryLogRecordExporter(), InMemoryLogRecordExporter()]
providers = [observability.LoggerProvider() for _ in exporters]
handlers: list[logging.Handler] = [observability._StructuredHandler(writer), tracing.SpanEventLoggingHandler()]
handlers: list[logging.Handler] = [_OTLPQueueHandler(queue), tracing.SpanEventLoggingHandler()]
for provider, exporter in zip(providers, exporters, strict=True):
provider.add_log_record_processor(SimpleLogRecordProcessor(exporter))
handlers.append(observability.OTELLoggingHandler(logger_provider=provider))
for handler in handlers if pipe_first else reversed(handlers):
for handler in handlers if stream_first else reversed(handlers):
target.addHandler(handler)
tracer_provider = TracerProvider()
value = {"transform": Affine(2, 3, 5, 7, 11, 13)}
Expand All @@ -226,9 +226,14 @@ def test_codecs_run_once_per_record_across_handlers(stdlib: bool, pipe_first: bo
assert event.attributes["input"] == '{"transform":[2.0,3.0,5.0,7.0,11.0,13.0]}'
# Records hold snapshots, not references to task-owned mutable containers.
value.clear()
for call in writer.submit.call_args_list:
assert call.args[0]["attributes"]["input"] == {"transform": [2, 3, 5, 7, 11, 13]}
assert writer.submit.call_count == 2
queue.seal()
stream_records = list(queue.subscribe())
assert [record.body.string_value for record in stream_records] == ["record-0", "record-1"]
for record in stream_records:
attributes = {item.key: item.value for item in record.attributes}
transform = attributes["input"].kvlist_value.values[0]
assert transform.key == "transform"
assert [item.double_value for item in transform.value.array_value.values] == [2, 3, 5, 7, 11, 13]
for exporter in exporters:
records = [item.log_record for item in exporter.get_finished_logs()]
assert [record.body for record in records] == ["record-0", "record-1"]
Expand All @@ -243,43 +248,34 @@ def test_codecs_run_once_per_record_across_handlers(stdlib: bool, pipe_first: bo

@pytest.mark.usefixtures("isolated_logging")
@pytest.mark.parametrize("external_first", [False, True])
def test_external_exports_and_console_do_not_replace_api_or_pipe(
monkeypatch: pytest.MonkeyPatch, external_first: bool
) -> None:
read_fd, write_fd = os.pipe()
monkeypatch.setenv("TILEBOX_LOG_FD", str(write_fd))
def test_external_exports_and_console_do_not_replace_api_or_stream(external_first: bool) -> None:
queue = _LogQueue()
root_logger.addHandler(_OTLPQueueHandler(queue))
api, external = InMemoryLogRecordExporter(), InMemoryLogRecordExporter()
first, second = StringIO(), StringIO()
try:
with patch.object(observability, "_otel_log_exporter", return_value=SimpleLogRecordProcessor(external)):
if external_first:
observability.configure_otel_logging(endpoint="https://external.example")
observability.configure_console_logging(stream=first)
with patch.object(observability, "_otel_log_exporter", return_value=SimpleLogRecordProcessor(api)):
observability.initialize_logging("https://api.tilebox.com", "test-key")
with patch.object(observability, "_otel_log_exporter", return_value=SimpleLogRecordProcessor(external)):
if not external_first:
observability.configure_otel_logging(endpoint="https://external.example")
observability.configure_console_logging(stream=first)
observability.configure_console_logging(stream=second, reconfigure=False)
observability.configure_log_level(logging.DEBUG)
assert internal_logger.level == logging.ERROR
task_logger.info("all outputs")
observability.configure_console_logging(enabled=False)
observability.initialize_logging("https://ignored.example", "ignored-key")
task_logger.debug("exports only")
assert first.getvalue().count("all outputs") == second.getvalue().count("all outputs") == 1
assert "exports only" not in first.getvalue() + second.getvalue()
for exporter in (api, external):
assert [item.log_record.body for item in exporter.get_finished_logs()] == ["all outputs", "exports only"]
assert observability._writer is not None
observability._writer.close()
records = [msgspec.json.decode(line) for line in os.read(read_fd, 65536).splitlines()]
assert [record["message"] for record in records] == ["all outputs", "exports only"]
finally:
if observability._writer is not None:
observability._writer.close()
os.close(read_fd)
with patch.object(observability, "_otel_log_exporter", return_value=SimpleLogRecordProcessor(external)):
if external_first:
observability.configure_otel_logging(endpoint="https://external.example")
observability.configure_console_logging(stream=first)
with patch.object(observability, "_otel_log_exporter", return_value=SimpleLogRecordProcessor(api)):
observability.initialize_logging("https://api.tilebox.com", "test-key")
with patch.object(observability, "_otel_log_exporter", return_value=SimpleLogRecordProcessor(external)):
if not external_first:
observability.configure_otel_logging(endpoint="https://external.example")
observability.configure_console_logging(stream=first)
observability.configure_console_logging(stream=second, reconfigure=False)
observability.configure_log_level(logging.DEBUG)
assert internal_logger.level == logging.ERROR
task_logger.info("all outputs")
observability.configure_console_logging(enabled=False)
observability.initialize_logging("https://ignored.example", "ignored-key")
task_logger.debug("exports only")
assert first.getvalue().count("all outputs") == second.getvalue().count("all outputs") == 1
assert "exports only" not in first.getvalue() + second.getvalue()
for exporter in (api, external):
assert [item.log_record.body for item in exporter.get_finished_logs()] == ["all outputs", "exports only"]
queue.seal()
assert [record.body.string_value for record in queue.subscribe()] == ["all outputs", "exports only"]


@pytest.mark.usefixtures("isolated_logging")
Expand Down Expand Up @@ -320,40 +316,42 @@ def __str__(self) -> str:

cyclic: list[Any] = []
cyclic.append(cyclic)
read_fd, write_fd = os.pipe()
monkeypatch.setenv("TILEBOX_LOG_FD", str(write_fd))
queue = _LogQueue()
monkeypatch.setattr(observability, "managed_log_queue", queue)
root_logger.addHandler(_OTLPQueueHandler(queue))
exporter = InMemoryLogRecordExporter()
try:
with patch.object(observability, "_otel_log_exporter", return_value=SimpleLogRecordProcessor(exporter)):
observability.initialize_logging("https://api.tilebox.com", "test-key")
StructuredLogger(task_logger).info(
"input",
input=Input(Affine(2, 3, 5, 7, 11, 13), Path("image.tif")),
counts={date(2026, 9, 16): 3},
cyclic=cyclic,
unsupported=Broken(),
)
assert observability._writer is not None
observability._writer.close()
record = msgspec.json.decode(os.read(read_fd, 65536))
assert record["attributes"] == {
"input": {"transform": [2, 3, 5, 7, 11, 13], "path": "image.tif"},
"counts": {"2026-09-16": 3},
"cyclic": "[[...]]",
"unsupported": "<unserializable Broken>",
}
attributes = exporter.get_finished_logs()[0].log_record.attributes
assert attributes is not None
assert isinstance(attributes["input"], str)
assert isinstance(attributes["counts"], str)
assert msgspec.json.decode(attributes["input"]) == record["attributes"]["input"]
assert msgspec.json.decode(attributes["counts"]) == record["attributes"]["counts"]
assert attributes["cyclic"] == "[[...]]"
assert attributes["unsupported"] == "<unserializable Broken>"
finally:
if observability._writer is not None:
observability._writer.close()
os.close(read_fd)
with patch.object(observability, "_otel_log_exporter", return_value=SimpleLogRecordProcessor(exporter)):
observability.initialize_logging("https://api.tilebox.com", "test-key")
StructuredLogger(task_logger).info(
"input",
input=Input(Affine(2, 3, 5, 7, 11, 13), Path("image.tif")),
counts={date(2026, 9, 16): 3},
cyclic=cyclic,
unsupported=Broken(),
)
queue.seal()
[record] = queue.subscribe()
stream_attributes = {item.key: item.value for item in record.attributes}
input_attributes = {item.key: item.value for item in stream_attributes["input"].kvlist_value.values}
assert [item.double_value for item in input_attributes["transform"].array_value.values] == [2, 3, 5, 7, 11, 13]
assert input_attributes["path"].string_value == "image.tif"
counts = stream_attributes["counts"].kvlist_value.values
assert len(counts) == 1
assert counts[0].key == "2026-09-16"
assert counts[0].value.int_value == 3
assert stream_attributes["cyclic"].string_value == "[[...]]"
assert stream_attributes["unsupported"].string_value == "<unserializable Broken>"
attributes = exporter.get_finished_logs()[0].log_record.attributes
assert attributes is not None
assert isinstance(attributes["input"], str)
assert isinstance(attributes["counts"], str)
assert msgspec.json.decode(attributes["input"]) == {
"transform": [2, 3, 5, 7, 11, 13],
"path": "image.tif",
}
assert msgspec.json.decode(attributes["counts"]) == {"2026-09-16": 3}
assert attributes["cyclic"] == "[[...]]"
assert attributes["unsupported"] == "<unserializable Broken>"


@pytest.mark.parametrize("without_otel", [False, True])
Expand Down
88 changes: 85 additions & 3 deletions tilebox-workflows/tests/observability/test_tracing.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,19 @@
import json
import os
import subprocess
import sys
from collections.abc import Iterator
from importlib.metadata import version
from unittest.mock import patch

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

from tilebox.workflows.observability import logging as observability_logging
from tilebox.workflows.observability import tracing


Expand Down Expand Up @@ -54,11 +61,12 @@ def force_flush(self, timeout_millis: int = 30000) -> bool: # noqa: ARG002

@pytest.fixture(autouse=True)
def reset_tilebox_tracing() -> Iterator[None]:
tracing._set_tilebox_tracer_provider(TracerProvider())
original_provider = tracing._get_tilebox_tracer_provider()
tracing._workflow_tracers.clear()
tracing._set_tilebox_tracer_provider(TracerProvider(resource=observability_logging._get_default_resource()))
yield
tracing._set_tilebox_tracer_provider(TracerProvider())
tracing._workflow_tracers.clear()
tracing._set_tilebox_tracer_provider(original_provider)


@pytest.fixture
Expand All @@ -74,6 +82,80 @@ def create_processor(*args: object, **kwargs: object) -> RecordingSpanProcessor:
return processors


@pytest.mark.parametrize("otel_service_name", ["unknown_service:python", "unknown_service:python.exe", "other-default"])
def test_exported_traces_use_tilebox_resource_defaults(otel_service_name: str) -> None:
# Exercise actual module initialization, independent of OTEL's platform-specific defaults.
instance_id = "66d615f3-7d53-4a31-bc94-94cbb9d9ffa2"
environment = {key: value for key, value in os.environ.items() if not key.startswith(("TILEBOX_", "OTEL_"))}
environment.update(
TILEBOX_RUNTIME_ID=instance_id,
OTEL_SERVICE_NAME=otel_service_name,
OTEL_RESOURCE_ATTRIBUTES="service.instance.id=otel-generated",
)
result = subprocess.run(
[
sys.executable,
"-c",
"""
import json
from unittest.mock import patch
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
from tilebox.workflows.observability import tracing
from tilebox.workflows.observability.logging import _get_default_resource

exporter = InMemorySpanExporter()
with patch.object(tracing, "_otel_span_exporter", return_value=SimpleSpanProcessor(exporter)):
tracer = tracing.WorkflowTracer(service=None, url="https://api.tilebox.com", token=None)
with tracer.span("task"):
pass
[span] = exporter.get_finished_spans()
print(json.dumps([dict(span.resource.attributes), dict(_get_default_resource().attributes)]))
""",
],
env=environment,
check=True,
capture_output=True,
text=True,
timeout=10,
)
expected = {
"service.name": "tilebox-python",
"service.namespace": "tilebox.workflows",
"service.version": version("tilebox-workflows"),
"service.instance.id": instance_id,
}
trace_resource, log_resource = json.loads(result.stdout)
assert {key: trace_resource.get(key) for key in expected} == expected
assert {key: log_resource.get(key) for key in expected} == expected


@pytest.mark.parametrize("service_name", [None, "custom-worker", "unknown_service:python.exe"])
@pytest.mark.parametrize("explicit_resource", [False, True])
def test_exported_traces_preserve_configured_resources(service_name: str | None, explicit_resource: bool) -> None:
attributes = {"service.instance.id": "custom-instance", "deployment": "test"}
if service_name is not None:
attributes["service.name"] = service_name
resource = Resource(attributes, schema_url="https://example.com/schema")
if not explicit_resource:
tracing._set_tilebox_tracer_provider(TracerProvider(resource=resource))
exporter = InMemorySpanExporter()
with patch.object(tracing, "_otel_span_exporter", return_value=SimpleSpanProcessor(exporter)):
tracer = tracing.WorkflowTracer(
service=resource if explicit_resource else None, url="https://api.tilebox.com", token=None
)
with tracer.span("task"):
pass

[span] = exporter.get_finished_spans()
for result in (span.resource, observability_logging._get_default_resource(resource)):
assert {key: result.attributes.get(key) for key in attributes} == attributes
assert result.attributes["service.name"] == (service_name or "tilebox-python")
assert result.attributes["service.namespace"] == "tilebox.workflows"
assert result.attributes["service.version"] == version("tilebox-workflows")
assert result.schema_url == resource.schema_url


def test_workflow_tracers_do_not_share_client_span_processors(
span_processors: list[RecordingSpanProcessor],
) -> None:
Expand Down
Loading
Loading