From 666f0d90196f3f810cf5679b0a0bd8dd85e0bf46 Mon Sep 17 00:00:00 2001 From: Morgan Wowk Date: Mon, 17 Aug 2026 15:47:06 -0700 Subject: [PATCH] feat(orchestrator): trace per-execution processing time by container_execution.id --- .../instrumentation/orchestrator_tracing.py | 24 +++++++++++++++++++ cloud_pipelines_backend/orchestrator_sql.py | 15 ++++++++---- 2 files changed, 35 insertions(+), 4 deletions(-) create mode 100644 cloud_pipelines_backend/instrumentation/orchestrator_tracing.py diff --git a/cloud_pipelines_backend/instrumentation/orchestrator_tracing.py b/cloud_pipelines_backend/instrumentation/orchestrator_tracing.py new file mode 100644 index 0000000..a0d2808 --- /dev/null +++ b/cloud_pipelines_backend/instrumentation/orchestrator_tracing.py @@ -0,0 +1,24 @@ +"""OTel spans for the orchestrator processing loop so per-execution processing time is measurable.""" + +from __future__ import annotations + +import collections.abc +import contextlib + +from opentelemetry import trace +from opentelemetry.trace import StatusCode + +_tracer = trace.get_tracer("tangle.orchestrator") + + +@contextlib.contextmanager +def operation_span( + name: str, *, attributes: dict[str, object] | None = None +) -> collections.abc.Iterator[trace.Span]: + with _tracer.start_as_current_span(name, attributes=attributes) as span: + try: + yield span + except Exception as exception: + span.set_status(StatusCode.ERROR) + span.record_exception(exception) + raise diff --git a/cloud_pipelines_backend/orchestrator_sql.py b/cloud_pipelines_backend/orchestrator_sql.py index d363b12..56ab7b8 100644 --- a/cloud_pipelines_backend/orchestrator_sql.py +++ b/cloud_pipelines_backend/orchestrator_sql.py @@ -25,6 +25,7 @@ from .instrumentation import bugsnag_instrumentation from .instrumentation import contextual_logging from .instrumentation import metrics as app_metrics +from .instrumentation import orchestrator_tracing _logger = logging.getLogger(__name__) @@ -229,10 +230,16 @@ def internal_process_running_executions_queue(self, session: orm.Session): f"Before processing running container execution. Queries duration {queries_duration_ms}ms." ) try: - self.internal_process_one_running_execution( - session=session, - container_execution=running_container_execution, - ) + with orchestrator_tracing.operation_span( + "orchestrator.process_running_execution", + attributes={ + "container_execution.id": running_container_execution.id + }, + ): + self.internal_process_one_running_execution( + session=session, + container_execution=running_container_execution, + ) except Exception as ex: _logger.exception("Error processing running container execution") session.rollback()