From 1a2277bf14a169725b65543c6fc2e146d75b44fe Mon Sep 17 00:00:00 2001 From: DABH Date: Thu, 1 Oct 2026 00:58:43 -0500 Subject: [PATCH 1/6] Detach OpenTelemetry context only on the thread that attached it The tracing interceptor's best-effort detach checked that the attached Context was still current, so a generator finalized by GC on another thread would skip the detach instead of using a token that is invalid there. OpenTelemetry's threading instrumentation, which strands enables whenever an Agent is created, copies the same Context object into new threads, so that check passes on the wrong thread and detach logs "Failed to detach context". Record the attaching thread with the token and skip the detach anywhere else. The regression test runs the existing safe-detach scenario with the instrumentation active; it fails without this change. --- CHANGELOG.md | 4 ++ .../contrib/opentelemetry/_interceptor.py | 49 ++++++++++++------- .../opentelemetry/test_opentelemetry.py | 25 +++++++++- 3 files changed, 60 insertions(+), 18 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5a8995696..09685e8e8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -28,6 +28,10 @@ to include examples, links to docs, or any other relevant information. ### Fixed +- `temporalio.contrib.opentelemetry`: the tracing interceptor no longer logs `Failed to detach + context` when a workflow context is torn down on a different thread while OpenTelemetry's + threading instrumentation (enabled by strands, among others) is active; the context is now + detached only on the thread that attached it. ### Security ## [1.34.0] - 2026-09-30 diff --git a/temporalio/contrib/opentelemetry/_interceptor.py b/temporalio/contrib/opentelemetry/_interceptor.py index 6dca4596e..62080922b 100644 --- a/temporalio/contrib/opentelemetry/_interceptor.py +++ b/temporalio/contrib/opentelemetry/_interceptor.py @@ -3,8 +3,10 @@ from __future__ import annotations import dataclasses +import threading from collections.abc import Callable, Iterator, Mapping, Sequence from contextlib import contextmanager +from contextvars import Token from dataclasses import dataclass from typing import ( Any, @@ -59,6 +61,29 @@ _ContextT = TypeVar("_ContextT", bound=nexusrpc.handler.OperationContext) +def _attach_context(context: Context) -> tuple[Token[Context], int]: + """Attach ``context`` and remember the attaching thread for a safe detach.""" + return opentelemetry.context.attach(context), threading.get_ident() + + +def _detach_context(context: Context, token: Token[Context], thread_ident: int) -> None: + """Detach ``token`` only where it is valid. + + Generator finalization and GC can run a ``finally`` on a different thread + or ``contextvars.Context`` than the one that attached, where the token is + invalid and ``detach`` logs "Failed to detach context". Comparing the + active context is not enough on its own: OpenTelemetry's threading + instrumentation (enabled by strands, among others) propagates the same + ``Context`` object into new threads, so the attaching thread is required + too. + """ + if ( + threading.get_ident() == thread_ident + and context is opentelemetry.context.get_current() + ): + opentelemetry.context.detach(token) + + class TracingInterceptor(temporalio.client.Interceptor, temporalio.worker.Interceptor): """Interceptor that supports client and worker OpenTelemetry span creation and propagation. @@ -183,7 +208,7 @@ def _start_as_current_span( kind: opentelemetry.trace.SpanKind, context: Context | None = None, ) -> Iterator[None]: - token = opentelemetry.context.attach(context) if context else None + attached = _attach_context(context) if context else None try: with self.tracer.start_as_current_span( name, @@ -220,8 +245,8 @@ def _start_as_current_span( ) raise finally: - if token and context is opentelemetry.context.get_current(): - opentelemetry.context.detach(token) + if attached and context: + _detach_context(context, *attached) def _completed_workflow_span( self, params: _CompletedWorkflowSpanParams @@ -556,7 +581,7 @@ async def handle_query(self, input: temporalio.worker.HandleQueryInput) -> Any: # We need to put this interceptor on the context too context = self._set_on_context(context) # Run under context with new span - token = opentelemetry.context.attach(context) + attached = _attach_context(context) try: # This won't be created if there was no context header self._completed_span( @@ -568,12 +593,7 @@ async def handle_query(self, input: temporalio.worker.HandleQueryInput) -> Any: ) return await super().handle_query(input) finally: - # In some exceptional cases this finally is executed with a - # different contextvars.Context than the one the token was created - # on. As such we do a best effort detach to avoid using a mismatched - # token. - if context is opentelemetry.context.get_current(): - opentelemetry.context.detach(token) + _detach_context(context, *attached) def handle_update_validator( self, input: temporalio.worker.HandleUpdateInput @@ -644,7 +664,7 @@ def _top_level_workflow_context( success = False exception: Exception | None = None # Run under this context - token = opentelemetry.context.attach(context) + attached = _attach_context(context) try: yield None @@ -663,12 +683,7 @@ def _top_level_workflow_context( kind=opentelemetry.trace.SpanKind.INTERNAL, ) - # In some exceptional cases this finally is executed with a - # different contextvars.Context than the one the token was created - # on. As such we do a best effort detach to avoid using a mismatched - # token. - if context is opentelemetry.context.get_current(): - opentelemetry.context.detach(token) + _detach_context(context, *attached) def _context_to_headers( self, headers: Mapping[str, temporalio.api.common.v1.Payload] diff --git a/tests/contrib/opentelemetry/test_opentelemetry.py b/tests/contrib/opentelemetry/test_opentelemetry.py index 0ea9530e6..6dc296f08 100644 --- a/tests/contrib/opentelemetry/test_opentelemetry.py +++ b/tests/contrib/opentelemetry/test_opentelemetry.py @@ -1052,7 +1052,7 @@ async def test_opentelemetry_standalone_activity_tracing( assert start_activity_span.attributes["temporalActivityType"] == "tracing_activity" -def test_opentelemetry_safe_detach(): +def _assert_context_detach_is_safe() -> None: class _fake_self: def _load_workflow_context_carrier(*_args): return None @@ -1097,3 +1097,26 @@ def otel_context_error(record: logging.LogRecord) -> bool: assert capturer.find(otel_context_error) is None, ( "Detach from context message should not be logged" ) + + +def test_opentelemetry_safe_detach(): + _assert_context_detach_is_safe() + + +def test_opentelemetry_safe_detach_with_threading_instrumentation(): + # OpenTelemetry's threading instrumentation (strands turns it on when an + # Agent is created) propagates the current Context object into new + # threads, so a context-identity check alone would detach a token minted + # on another thread. + threading_instrumentation = pytest.importorskip( + "opentelemetry.instrumentation.threading" + ) + instrumentor = threading_instrumentation.ThreadingInstrumentor() + already_instrumented = instrumentor.is_instrumented_by_opentelemetry + if not already_instrumented: + instrumentor.instrument() + try: + _assert_context_detach_is_safe() + finally: + if not already_instrumented: + instrumentor.uninstrument() From 1c4b282c0cf1600c881fe131af1e610669250876 Mon Sep 17 00:00:00 2001 From: DABH Date: Thu, 1 Oct 2026 01:24:56 -0500 Subject: [PATCH 2/6] Cover OpenTelemetryInterceptor and compare the attaching Thread object OpenTelemetryPlugin installs OpenTelemetryInterceptor, whose four attach/detach sites had the same context-identity guard. Move the helpers into _context.py and use them from both interceptors. Compare the attaching threading.Thread object rather than its id, which the OS can reuse. The regression test is now a matrix over interceptor and threading instrumentation; the OpenTelemetryInterceptor case fails without this. --- CHANGELOG.md | 8 ++-- temporalio/contrib/opentelemetry/_context.py | 44 +++++++++++++++++++ .../contrib/opentelemetry/_interceptor.py | 40 ++++------------- .../opentelemetry/_otel_interceptor.py | 22 +++++----- .../opentelemetry/test_opentelemetry.py | 44 +++++++++++++++---- 5 files changed, 101 insertions(+), 57 deletions(-) create mode 100644 temporalio/contrib/opentelemetry/_context.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 09685e8e8..6729b0747 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -28,10 +28,10 @@ to include examples, links to docs, or any other relevant information. ### Fixed -- `temporalio.contrib.opentelemetry`: the tracing interceptor no longer logs `Failed to detach - context` when a workflow context is torn down on a different thread while OpenTelemetry's - threading instrumentation (enabled by strands, among others) is active; the context is now - detached only on the thread that attached it. +- `temporalio.contrib.opentelemetry`: `TracingInterceptor` and `OpenTelemetryInterceptor` no longer + log `Failed to detach context` when a context is torn down on a different thread while + OpenTelemetry's threading instrumentation (enabled by strands, among others) is active; a + context is now detached only on the thread that attached it. ### Security ## [1.34.0] - 2026-09-30 diff --git a/temporalio/contrib/opentelemetry/_context.py b/temporalio/contrib/opentelemetry/_context.py new file mode 100644 index 000000000..432773227 --- /dev/null +++ b/temporalio/contrib/opentelemetry/_context.py @@ -0,0 +1,44 @@ +"""Attach OpenTelemetry contexts so that the matching detach is always safe.""" + +from __future__ import annotations + +import threading +from contextvars import Token +from dataclasses import dataclass + +import opentelemetry.context +from opentelemetry.context import Context + + +@dataclass(frozen=True) +class AttachedContext: + """A context attached by :func:`attach_context`, with what its detach needs.""" + + context: Context + token: Token[Context] + thread: threading.Thread + + def detach(self) -> None: + """Detach the context only where the token is valid. + + Generator finalization and GC can run a ``finally`` on a different + thread or ``contextvars.Context`` than the one that attached, where the + token is invalid and ``opentelemetry.context.detach`` logs "Failed to + detach context". Checking that the attached context is still current is + not enough on its own: OpenTelemetry's threading instrumentation + (enabled by strands, among others) propagates the same ``Context`` + object into new threads. Requiring the attaching thread as well (the + ``Thread`` object, so a recycled thread id cannot match) closes that gap. + """ + if ( + threading.current_thread() is self.thread + and self.context is opentelemetry.context.get_current() + ): + opentelemetry.context.detach(self.token) + + +def attach_context(context: Context) -> AttachedContext: + """Attach ``context`` and remember what a safe detach needs.""" + return AttachedContext( + context, opentelemetry.context.attach(context), threading.current_thread() + ) diff --git a/temporalio/contrib/opentelemetry/_interceptor.py b/temporalio/contrib/opentelemetry/_interceptor.py index 62080922b..e698f22bd 100644 --- a/temporalio/contrib/opentelemetry/_interceptor.py +++ b/temporalio/contrib/opentelemetry/_interceptor.py @@ -3,10 +3,8 @@ from __future__ import annotations import dataclasses -import threading from collections.abc import Callable, Iterator, Mapping, Sequence from contextlib import contextmanager -from contextvars import Token from dataclasses import dataclass from typing import ( Any, @@ -38,6 +36,7 @@ import temporalio.nexus.system.workflow_service.models import temporalio.worker import temporalio.workflow +from temporalio.contrib.opentelemetry._context import attach_context from temporalio.exceptions import ApplicationError, ApplicationErrorCategory # OpenTelemetry dynamically, lazily chooses its context implementation at @@ -61,29 +60,6 @@ _ContextT = TypeVar("_ContextT", bound=nexusrpc.handler.OperationContext) -def _attach_context(context: Context) -> tuple[Token[Context], int]: - """Attach ``context`` and remember the attaching thread for a safe detach.""" - return opentelemetry.context.attach(context), threading.get_ident() - - -def _detach_context(context: Context, token: Token[Context], thread_ident: int) -> None: - """Detach ``token`` only where it is valid. - - Generator finalization and GC can run a ``finally`` on a different thread - or ``contextvars.Context`` than the one that attached, where the token is - invalid and ``detach`` logs "Failed to detach context". Comparing the - active context is not enough on its own: OpenTelemetry's threading - instrumentation (enabled by strands, among others) propagates the same - ``Context`` object into new threads, so the attaching thread is required - too. - """ - if ( - threading.get_ident() == thread_ident - and context is opentelemetry.context.get_current() - ): - opentelemetry.context.detach(token) - - class TracingInterceptor(temporalio.client.Interceptor, temporalio.worker.Interceptor): """Interceptor that supports client and worker OpenTelemetry span creation and propagation. @@ -208,7 +184,7 @@ def _start_as_current_span( kind: opentelemetry.trace.SpanKind, context: Context | None = None, ) -> Iterator[None]: - attached = _attach_context(context) if context else None + attached = attach_context(context) if context else None try: with self.tracer.start_as_current_span( name, @@ -245,8 +221,8 @@ def _start_as_current_span( ) raise finally: - if attached and context: - _detach_context(context, *attached) + if attached: + attached.detach() def _completed_workflow_span( self, params: _CompletedWorkflowSpanParams @@ -581,7 +557,7 @@ async def handle_query(self, input: temporalio.worker.HandleQueryInput) -> Any: # We need to put this interceptor on the context too context = self._set_on_context(context) # Run under context with new span - attached = _attach_context(context) + attached = attach_context(context) try: # This won't be created if there was no context header self._completed_span( @@ -593,7 +569,7 @@ async def handle_query(self, input: temporalio.worker.HandleQueryInput) -> Any: ) return await super().handle_query(input) finally: - _detach_context(context, *attached) + attached.detach() def handle_update_validator( self, input: temporalio.worker.HandleUpdateInput @@ -664,7 +640,7 @@ def _top_level_workflow_context( success = False exception: Exception | None = None # Run under this context - attached = _attach_context(context) + attached = attach_context(context) try: yield None @@ -683,7 +659,7 @@ def _top_level_workflow_context( kind=opentelemetry.trace.SpanKind.INTERNAL, ) - _detach_context(context, *attached) + attached.detach() def _context_to_headers( self, headers: Mapping[str, temporalio.api.common.v1.Payload] diff --git a/temporalio/contrib/opentelemetry/_otel_interceptor.py b/temporalio/contrib/opentelemetry/_otel_interceptor.py index ff07f0f1d..1727970f6 100644 --- a/temporalio/contrib/opentelemetry/_otel_interceptor.py +++ b/temporalio/contrib/opentelemetry/_otel_interceptor.py @@ -35,6 +35,7 @@ import temporalio.nexus.system.workflow_service.models import temporalio.worker import temporalio.workflow +from temporalio.contrib.opentelemetry._context import attach_context from temporalio.contrib.opentelemetry._tracer_provider import ( ReplaySafeTracerProvider, ) @@ -125,7 +126,7 @@ def _maybe_span( yield return - token = opentelemetry.context.attach(context) if context else None + attached = attach_context(context) if context else None try: with tracer.start_as_current_span( name, @@ -149,8 +150,8 @@ def _maybe_span( ) raise finally: - if token and context is opentelemetry.context.get_current(): - opentelemetry.context.detach(token) + if attached: + attached.detach() class OpenTelemetryInterceptor( @@ -334,7 +335,7 @@ async def execute_activity( self, input: temporalio.worker.ExecuteActivityInput ) -> Any: context = _headers_to_context(input.headers) - token = opentelemetry.context.attach(context) + attached = attach_context(context) try: info = temporalio.activity.info() with _maybe_span( @@ -350,8 +351,7 @@ async def execute_activity( ): return await super().execute_activity(input) finally: - if context is opentelemetry.context.get_current(): - opentelemetry.context.detach(token) + attached.detach() class _TracingNexusOperationInboundInterceptor( @@ -368,12 +368,11 @@ def __init__( @contextmanager def _top_level_context(self, headers: Mapping[str, str]) -> Iterator[None]: context = _nexus_headers_to_context(headers) - token = opentelemetry.context.attach(context) + attached = attach_context(context) try: yield finally: - if context is opentelemetry.context.get_current(): - opentelemetry.context.detach(token) + attached.detach() async def execute_nexus_operation_start( self, input: temporalio.worker.ExecuteNexusOperationStartInput @@ -501,12 +500,11 @@ async def handle_update_handler( @contextmanager def _top_level_workflow_context(self, input: _InputWithHeaders) -> Iterator[None]: context = _headers_to_context(input.headers) - token = opentelemetry.context.attach(context) + attached = attach_context(context) try: yield finally: - if context is opentelemetry.context.get_current(): - opentelemetry.context.detach(token) + attached.detach() class _TracingWorkflowOutboundInterceptor( diff --git a/tests/contrib/opentelemetry/test_opentelemetry.py b/tests/contrib/opentelemetry/test_opentelemetry.py index 6dc296f08..1b664339a 100644 --- a/tests/contrib/opentelemetry/test_opentelemetry.py +++ b/tests/contrib/opentelemetry/test_opentelemetry.py @@ -30,6 +30,9 @@ TracingWorkflowInboundInterceptor, ) from temporalio.contrib.opentelemetry import workflow as otel_workflow +from temporalio.contrib.opentelemetry._otel_interceptor import ( + _TracingWorkflowInboundInterceptor as _OtelTracingWorkflowInboundInterceptor, +) from temporalio.exceptions import ( ApplicationError, ApplicationErrorCategory, @@ -1052,7 +1055,7 @@ async def test_opentelemetry_standalone_activity_tracing( assert start_activity_span.attributes["temporalActivityType"] == "tracing_activity" -def _assert_context_detach_is_safe() -> None: +def _v1_workflow_context() -> Any: class _fake_self: def _load_workflow_context_carrier(*_args): return None @@ -1063,11 +1066,25 @@ def _set_on_context(self, ctx: Any): def _completed_span(*args: Any, **_kwargs: Any): pass - # create a context manager and force enter to happen on this thread - context_manager = TracingWorkflowInboundInterceptor._top_level_workflow_context( + return TracingWorkflowInboundInterceptor._top_level_workflow_context( _fake_self(), # type: ignore success_is_complete=True, ) + + +def _v2_workflow_context() -> Any: + class _fake_input: + headers: dict[str, Any] = {} + + return _OtelTracingWorkflowInboundInterceptor._top_level_workflow_context( + object(), # type: ignore + _fake_input(), # type: ignore + ) + + +def _assert_context_detach_is_safe(make_context_manager: Callable[[], Any]) -> None: + # create a context manager and force enter to happen on this thread + context_manager = make_context_manager() context_manager.__enter__() # move reference to context manager into queue @@ -1099,11 +1116,20 @@ def otel_context_error(record: logging.LogRecord) -> bool: ) -def test_opentelemetry_safe_detach(): - _assert_context_detach_is_safe() - - -def test_opentelemetry_safe_detach_with_threading_instrumentation(): +@pytest.mark.parametrize( + "make_context_manager", + [_v1_workflow_context, _v2_workflow_context], + ids=["TracingInterceptor", "OpenTelemetryInterceptor"], +) +@pytest.mark.parametrize( + "threading_instrumented", [False, True], ids=["plain", "threading-instrumented"] +) +def test_opentelemetry_safe_detach( + make_context_manager: Callable[[], Any], threading_instrumented: bool +): + if not threading_instrumented: + _assert_context_detach_is_safe(make_context_manager) + return # OpenTelemetry's threading instrumentation (strands turns it on when an # Agent is created) propagates the current Context object into new # threads, so a context-identity check alone would detach a token minted @@ -1116,7 +1142,7 @@ def test_opentelemetry_safe_detach_with_threading_instrumentation(): if not already_instrumented: instrumentor.instrument() try: - _assert_context_detach_is_safe() + _assert_context_detach_is_safe(make_context_manager) finally: if not already_instrumented: instrumentor.uninstrument() From bd614534a733e4863b797f8d1afe0259a39c55ae Mon Sep 17 00:00:00 2001 From: DABH Date: Thu, 1 Oct 2026 11:51:36 -0500 Subject: [PATCH 3/6] Detach where the token is valid, not where the thread matches Workflow activations run on a thread pool while the asyncio task keeps its contextvars.Context, so a context attached in one activation is legitimately detached in a later one on another thread. The thread check skipped those detaches and left the inner context attached for outer interceptors. Let contextvars decide instead: perform the reset that opentelemetry.context.detach performs and treat its ValueError for a foreign Context as nothing to detach. Covers the thread change with a unit test on both interceptors and a worker test whose executor alternates activations between two threads. --- CHANGELOG.md | 3 +- temporalio/contrib/opentelemetry/_context.py | 49 +++--- .../opentelemetry/test_opentelemetry.py | 166 +++++++++++++++++- 3 files changed, 190 insertions(+), 28 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6729b0747..3134fa923 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -31,7 +31,8 @@ to include examples, links to docs, or any other relevant information. - `temporalio.contrib.opentelemetry`: `TracingInterceptor` and `OpenTelemetryInterceptor` no longer log `Failed to detach context` when a context is torn down on a different thread while OpenTelemetry's threading instrumentation (enabled by strands, among others) is active; a - context is now detached only on the thread that attached it. + context is now detached exactly when its token is still valid in the current + `contextvars.Context`, which it stays when a workflow resumes on another pool thread. ### Security ## [1.34.0] - 2026-09-30 diff --git a/temporalio/contrib/opentelemetry/_context.py b/temporalio/contrib/opentelemetry/_context.py index 432773227..0038b7e36 100644 --- a/temporalio/contrib/opentelemetry/_context.py +++ b/temporalio/contrib/opentelemetry/_context.py @@ -2,7 +2,6 @@ from __future__ import annotations -import threading from contextvars import Token from dataclasses import dataclass @@ -16,29 +15,35 @@ class AttachedContext: context: Context token: Token[Context] - thread: threading.Thread - - def detach(self) -> None: - """Detach the context only where the token is valid. - - Generator finalization and GC can run a ``finally`` on a different - thread or ``contextvars.Context`` than the one that attached, where the - token is invalid and ``opentelemetry.context.detach`` logs "Failed to - detach context". Checking that the attached context is still current is - not enough on its own: OpenTelemetry's threading instrumentation - (enabled by strands, among others) propagates the same ``Context`` - object into new threads. Requiring the attaching thread as well (the - ``Thread`` object, so a recycled thread id cannot match) closes that gap. + + def detach(self) -> bool: + """Detach the context where its token is valid; return whether it was. + + The attach and the detach of one interceptor call can run on different + threads: workflow activations run on a thread pool while the asyncio + task keeps its ``contextvars.Context`` across them, so the token is + still valid there and the detach must happen. Generator finalization + and GC, on the other hand, can run a ``finally`` in a different + ``contextvars.Context``, where the token is invalid and + ``opentelemetry.context.detach`` logs "Failed to detach context". + Checking that the attached context is still current cannot tell those + apart, because OpenTelemetry's threading instrumentation (enabled by + strands, among others) propagates the same ``Context`` object into new + threads. Only ``contextvars`` knows which ``Context`` a token belongs + to, so this performs the reset that ``opentelemetry.context.detach`` + performs and treats its ``ValueError`` for a foreign ``Context`` as + "nothing to detach here". """ - if ( - threading.current_thread() is self.thread - and self.context is opentelemetry.context.get_current() - ): - opentelemetry.context.detach(self.token) + if self.context is not opentelemetry.context.get_current(): + return False + try: + self.token.var.reset(self.token) + except ValueError: + # The token was created in a different contextvars.Context. + return False + return True def attach_context(context: Context) -> AttachedContext: """Attach ``context`` and remember what a safe detach needs.""" - return AttachedContext( - context, opentelemetry.context.attach(context), threading.current_thread() - ) + return AttachedContext(context, opentelemetry.context.attach(context)) diff --git a/tests/contrib/opentelemetry/test_opentelemetry.py b/tests/contrib/opentelemetry/test_opentelemetry.py index 1b664339a..cedbf3423 100644 --- a/tests/contrib/opentelemetry/test_opentelemetry.py +++ b/tests/contrib/opentelemetry/test_opentelemetry.py @@ -1,20 +1,22 @@ from __future__ import annotations import asyncio +import contextvars import gc import logging import queue import threading import uuid from collections.abc import Callable, Generator, Iterable -from concurrent.futures import ThreadPoolExecutor +from concurrent.futures import Future, ThreadPoolExecutor from contextlib import contextmanager from dataclasses import dataclass from datetime import timedelta -from typing import Any, cast +from typing import Any, ParamSpec, TypeVar, cast import nexusrpc import opentelemetry.context +import opentelemetry.trace import pytest from opentelemetry import baggage, context from opentelemetry.sdk.trace import ReadableSpan, TracerProvider @@ -26,10 +28,13 @@ from temporalio.client import Client, WithStartWorkflowOperation, WorkflowUpdateStage from temporalio.common import RetryPolicy, WorkflowIDConflictPolicy from temporalio.contrib.opentelemetry import ( + OpenTelemetryInterceptor, TracingInterceptor, TracingWorkflowInboundInterceptor, + create_tracer_provider, ) from temporalio.contrib.opentelemetry import workflow as otel_workflow +from temporalio.contrib.opentelemetry._context import AttachedContext from temporalio.contrib.opentelemetry._otel_interceptor import ( _TracingWorkflowInboundInterceptor as _OtelTracingWorkflowInboundInterceptor, ) @@ -39,7 +44,14 @@ NexusOperationError, ) from temporalio.testing import WorkflowEnvironment -from temporalio.worker import UnsandboxedWorkflowRunner, Worker +from temporalio.worker import ( + ExecuteWorkflowInput, + Interceptor, + UnsandboxedWorkflowRunner, + Worker, + WorkflowInboundInterceptor, + WorkflowInterceptorClassInput, +) from tests.helpers import LogCapturer from tests.helpers.nexus import make_nexus_endpoint_name @@ -924,6 +936,7 @@ async def test_opentelemetry_context_restored_after_activity( detach_count = 0 original_attach = context.attach original_detach = context.detach + original_attached_detach = AttachedContext.detach def tracked_attach(ctx): # type:ignore[reportMissingParameterType] nonlocal attach_count @@ -935,8 +948,18 @@ def tracked_detach(token): # type:ignore[reportMissingParameterType] detach_count += 1 return original_detach(token) + # Spans detach through context.detach; the interceptors reset their own + # tokens directly, so count those detaches where they happen. + def tracked_attached_detach(self: AttachedContext) -> bool: + nonlocal detach_count + detached = original_attached_detach(self) + if detached: + detach_count += 1 + return detached + context.attach = tracked_attach context.detach = tracked_detach + AttachedContext.detach = tracked_attached_detach try: task_queue = f"task_queue_{uuid.uuid4()}" @@ -967,6 +990,7 @@ def tracked_detach(token): # type:ignore[reportMissingParameterType] finally: context.attach = original_attach context.detach = original_detach + AttachedContext.detach = original_attached_detach @activity.defn @@ -1055,7 +1079,7 @@ async def test_opentelemetry_standalone_activity_tracing( assert start_activity_span.attributes["temporalActivityType"] == "tracing_activity" -def _v1_workflow_context() -> Any: +def _v1_workflow_context(success_is_complete: bool = True) -> Any: class _fake_self: def _load_workflow_context_carrier(*_args): return None @@ -1068,7 +1092,7 @@ def _completed_span(*args: Any, **_kwargs: Any): return TracingWorkflowInboundInterceptor._top_level_workflow_context( _fake_self(), # type: ignore - success_is_complete=True, + success_is_complete=success_is_complete, ) @@ -1146,3 +1170,135 @@ def test_opentelemetry_safe_detach( finally: if not already_instrumented: instrumentor.uninstrument() + + +def _run_in_context_on_new_thread( + ctx: contextvars.Context, fn: Callable[[], Any] +) -> Any: + with ThreadPoolExecutor(max_workers=1) as pool: + return pool.submit(ctx.run, fn).result(timeout=5) + + +@pytest.mark.parametrize( + "make_context_manager", + [lambda: _v1_workflow_context(success_is_complete=False), _v2_workflow_context], + ids=["TracingInterceptor", "OpenTelemetryInterceptor"], +) +def test_opentelemetry_detach_after_thread_change( + make_context_manager: Callable[[], Any], +): + # Workflow activations run on a thread pool, so the asyncio task that + # entered one of these context managers can leave it on another thread + # while keeping the same contextvars.Context. The token is valid there and + # the detach must restore the outer context. + task_context = contextvars.copy_context() + outer = opentelemetry.context.set_value("outer", True) + task_context.run(opentelemetry.context.attach, outer) + context_manager = make_context_manager() + + _run_in_context_on_new_thread(task_context, context_manager.__enter__) + assert task_context.run(opentelemetry.context.get_current) is not outer + _run_in_context_on_new_thread( + task_context, lambda: context_manager.__exit__(None, None, None) + ) + assert task_context.run(opentelemetry.context.get_current) is outer + + +_P = ParamSpec("_P") +_T = TypeVar("_T") + + +class _AlternatingThreadExecutor(ThreadPoolExecutor): + """Runs each submission on a different thread than the one before it, the + way successive activations of a workflow can land on different pool + threads under load.""" + + def __init__(self) -> None: + super().__init__(max_workers=1) + self._pools = [ThreadPoolExecutor(max_workers=1) for _ in range(2)] + self._submissions = 0 + + def submit( + self, fn: Callable[_P, _T], /, *args: _P.args, **kwargs: _P.kwargs + ) -> Future[_T]: + pool = self._pools[self._submissions % len(self._pools)] + self._submissions += 1 + return pool.submit(fn, *args, **kwargs) + + def shutdown(self, wait: bool = True, *, cancel_futures: bool = False) -> None: + for pool in self._pools: + pool.shutdown(wait, cancel_futures=cancel_futures) + super().shutdown(wait, cancel_futures=cancel_futures) + + +@workflow.defn +class TimerWorkflow: + @workflow.run + async def run(self) -> None: + # Two activations: the start, and the timer firing. + await asyncio.sleep(0.01) + + +@pytest.mark.parametrize( + "make_interceptor", + [lambda: TracingInterceptor(get_tracer(__name__)), OpenTelemetryInterceptor], + ids=["TracingInterceptor", "OpenTelemetryInterceptor"], +) +async def test_opentelemetry_context_restored_after_activation_thread_change( + client: Client, + make_interceptor: Callable[[], Interceptor], + reset_otel_tracer_provider: Any, # type: ignore[reportUnusedParameter] +): + # OpenTelemetryInterceptor insists on a replay-safe global provider. + opentelemetry.trace.set_tracer_provider(create_tracer_provider()) + records: list[tuple[threading.Thread, threading.Thread, bool]] = [] + + class ContextRestoredWorkflowInboundInterceptor(WorkflowInboundInterceptor): + async def execute_workflow(self, input: ExecuteWorkflowInput) -> Any: + entered_on = threading.current_thread() + before = opentelemetry.context.get_current() + result = await super().execute_workflow(input) + records.append( + ( + entered_on, + threading.current_thread(), + opentelemetry.context.get_current() is before, + ) + ) + return result + + class ContextRestoredInterceptor(Interceptor): + def workflow_interceptor_class( + self, input: WorkflowInterceptorClassInput + ) -> type[WorkflowInboundInterceptor] | None: + return ContextRestoredWorkflowInboundInterceptor + + executor = _AlternatingThreadExecutor() + try: + async with Worker( + client, + task_queue=f"task_queue_{uuid.uuid4()}", + workflows=[TimerWorkflow], + # The first interceptor is the outermost, so it sees whatever the + # tracing interceptor leaves attached. + interceptors=[ContextRestoredInterceptor(), make_interceptor()], + workflow_task_executor=executor, + workflow_runner=UnsandboxedWorkflowRunner(), + ) as worker: + await client.execute_workflow( + TimerWorkflow.run, + id=f"workflow_{uuid.uuid4()}", + task_queue=worker.task_queue, + ) + finally: + executor.shutdown() + + assert len(records) == 1 + entered_on, left_on, context_restored = records[0] + assert entered_on is not left_on, ( + "The workflow was expected to change threads between its activations" + ) + assert context_restored, ( + "The tracing interceptor must detach its context even though the " + "workflow left it on a different thread than it entered on" + ) From 123605cb08255f46d8f7e8cf8e356d2f34f9fe6e Mon Sep 17 00:00:00 2001 From: DABH Date: Thu, 1 Oct 2026 12:00:03 -0500 Subject: [PATCH 4/6] Patch AttachedContext.detach through monkeypatch in the leak test mypy rejects assigning to a method (method-assign); pytest's monkeypatch does the same thing and restores it on teardown. --- tests/contrib/opentelemetry/test_opentelemetry.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/contrib/opentelemetry/test_opentelemetry.py b/tests/contrib/opentelemetry/test_opentelemetry.py index cedbf3423..53bcbe938 100644 --- a/tests/contrib/opentelemetry/test_opentelemetry.py +++ b/tests/contrib/opentelemetry/test_opentelemetry.py @@ -931,6 +931,7 @@ async def test_opentelemetry_context_restored_after_activity( client_with_tracing: Client, activity: Callable[[], None], expect_failure: bool, + monkeypatch: pytest.MonkeyPatch, ) -> None: attach_count = 0 detach_count = 0 @@ -959,7 +960,7 @@ def tracked_attached_detach(self: AttachedContext) -> bool: context.attach = tracked_attach context.detach = tracked_detach - AttachedContext.detach = tracked_attached_detach + monkeypatch.setattr(AttachedContext, "detach", tracked_attached_detach) try: task_queue = f"task_queue_{uuid.uuid4()}" @@ -990,7 +991,6 @@ def tracked_attached_detach(self: AttachedContext) -> bool: finally: context.attach = original_attach context.detach = original_detach - AttachedContext.detach = original_attached_detach @activity.defn From 8ab48b0f7b0371d7a5d71955f8471464858b5d10 Mon Sep 17 00:00:00 2001 From: DABH Date: Thu, 1 Oct 2026 12:49:02 -0500 Subject: [PATCH 5/6] Make the attach/detach helper a context manager attached_context(context) replaces the seven attach/try/finally/detach sites; the guarded detach moves to _detach, which the leak test counts. --- temporalio/contrib/opentelemetry/_context.py | 80 ++++++++++--------- .../contrib/opentelemetry/_interceptor.py | 50 +++++------- .../opentelemetry/_otel_interceptor.py | 23 ++---- .../opentelemetry/test_opentelemetry.py | 10 +-- 4 files changed, 71 insertions(+), 92 deletions(-) diff --git a/temporalio/contrib/opentelemetry/_context.py b/temporalio/contrib/opentelemetry/_context.py index 0038b7e36..a57f8d48d 100644 --- a/temporalio/contrib/opentelemetry/_context.py +++ b/temporalio/contrib/opentelemetry/_context.py @@ -2,48 +2,50 @@ from __future__ import annotations +from collections.abc import Iterator +from contextlib import contextmanager from contextvars import Token -from dataclasses import dataclass import opentelemetry.context from opentelemetry.context import Context -@dataclass(frozen=True) -class AttachedContext: - """A context attached by :func:`attach_context`, with what its detach needs.""" - - context: Context - token: Token[Context] - - def detach(self) -> bool: - """Detach the context where its token is valid; return whether it was. - - The attach and the detach of one interceptor call can run on different - threads: workflow activations run on a thread pool while the asyncio - task keeps its ``contextvars.Context`` across them, so the token is - still valid there and the detach must happen. Generator finalization - and GC, on the other hand, can run a ``finally`` in a different - ``contextvars.Context``, where the token is invalid and - ``opentelemetry.context.detach`` logs "Failed to detach context". - Checking that the attached context is still current cannot tell those - apart, because OpenTelemetry's threading instrumentation (enabled by - strands, among others) propagates the same ``Context`` object into new - threads. Only ``contextvars`` knows which ``Context`` a token belongs - to, so this performs the reset that ``opentelemetry.context.detach`` - performs and treats its ``ValueError`` for a foreign ``Context`` as - "nothing to detach here". - """ - if self.context is not opentelemetry.context.get_current(): - return False - try: - self.token.var.reset(self.token) - except ValueError: - # The token was created in a different contextvars.Context. - return False - return True - - -def attach_context(context: Context) -> AttachedContext: - """Attach ``context`` and remember what a safe detach needs.""" - return AttachedContext(context, opentelemetry.context.attach(context)) +@contextmanager +def attached_context(context: Context | None) -> Iterator[None]: + """Attach ``context`` for the block and detach it afterwards where possible. + + ``None`` attaches nothing. The block's ``finally`` can run in a different + ``contextvars.Context`` than the one that attached: a context manager + abandoned by an evicted workflow is finalized wherever garbage collection + happens to run. The token is not valid there, and + ``opentelemetry.context.detach`` would log "Failed to detach context" even + though there is nothing to detach. Checking that the attached context is + still current does not catch every such case, because OpenTelemetry's + threading instrumentation (enabled by strands, among others) propagates + the same ``Context`` object into new threads. The thread is no test + either: workflow activations move between pool threads while the asyncio + task keeps its ``contextvars.Context``, and those detaches must happen. + Only ``contextvars`` knows which ``Context`` a token belongs to, so + :func:`_detach` performs the reset ``detach`` performs and ignores the + ``ValueError`` raised for a token from another ``Context``. + """ + if context is None: + yield + return + token = opentelemetry.context.attach(context) + try: + yield + finally: + _detach(context, token) + + +def _detach(context: Context, token: Token[Context]) -> bool: + """Detach ``context`` if it is current and ``token`` is valid here.""" + if context is not opentelemetry.context.get_current(): + return False + try: + token.var.reset(token) + except ValueError: + # The token was created in a different contextvars.Context. + return False + return True diff --git a/temporalio/contrib/opentelemetry/_interceptor.py b/temporalio/contrib/opentelemetry/_interceptor.py index e698f22bd..f503d5ac5 100644 --- a/temporalio/contrib/opentelemetry/_interceptor.py +++ b/temporalio/contrib/opentelemetry/_interceptor.py @@ -36,7 +36,7 @@ import temporalio.nexus.system.workflow_service.models import temporalio.worker import temporalio.workflow -from temporalio.contrib.opentelemetry._context import attach_context +from temporalio.contrib.opentelemetry._context import attached_context from temporalio.exceptions import ApplicationError, ApplicationErrorCategory # OpenTelemetry dynamically, lazily chooses its context implementation at @@ -184,8 +184,7 @@ def _start_as_current_span( kind: opentelemetry.trace.SpanKind, context: Context | None = None, ) -> Iterator[None]: - attached = attach_context(context) if context else None - try: + with attached_context(context): with self.tracer.start_as_current_span( name, attributes=attributes, @@ -220,9 +219,6 @@ def _start_as_current_span( ) ) raise - finally: - if attached: - attached.detach() def _completed_workflow_span( self, params: _CompletedWorkflowSpanParams @@ -557,8 +553,7 @@ async def handle_query(self, input: temporalio.worker.HandleQueryInput) -> Any: # We need to put this interceptor on the context too context = self._set_on_context(context) # Run under context with new span - attached = attach_context(context) - try: + with attached_context(context): # This won't be created if there was no context header self._completed_span( f"HandleQuery:{input.query}", @@ -568,8 +563,6 @@ async def handle_query(self, input: temporalio.worker.HandleQueryInput) -> Any: kind=opentelemetry.trace.SpanKind.SERVER, ) return await super().handle_query(input) - finally: - attached.detach() def handle_update_validator( self, input: temporalio.worker.HandleUpdateInput @@ -640,26 +633,23 @@ def _top_level_workflow_context( success = False exception: Exception | None = None # Run under this context - attached = attach_context(context) - - try: - yield None - success = True - except temporalio.exceptions.FailureError as err: - # We only record the failure errors since those are the only ones - # that lead to workflow completions - exception = err - raise - finally: - # Create a completed span before detaching context - if exception or (success and success_is_complete): - self._completed_span( - f"CompleteWorkflow:{temporalio.workflow.info().workflow_type}", - exception=exception, - kind=opentelemetry.trace.SpanKind.INTERNAL, - ) - - attached.detach() + with attached_context(context): + try: + yield None + success = True + except temporalio.exceptions.FailureError as err: + # We only record the failure errors since those are the only ones + # that lead to workflow completions + exception = err + raise + finally: + # Create a completed span before detaching context + if exception or (success and success_is_complete): + self._completed_span( + f"CompleteWorkflow:{temporalio.workflow.info().workflow_type}", + exception=exception, + kind=opentelemetry.trace.SpanKind.INTERNAL, + ) def _context_to_headers( self, headers: Mapping[str, temporalio.api.common.v1.Payload] diff --git a/temporalio/contrib/opentelemetry/_otel_interceptor.py b/temporalio/contrib/opentelemetry/_otel_interceptor.py index 1727970f6..5207095dc 100644 --- a/temporalio/contrib/opentelemetry/_otel_interceptor.py +++ b/temporalio/contrib/opentelemetry/_otel_interceptor.py @@ -35,7 +35,7 @@ import temporalio.nexus.system.workflow_service.models import temporalio.worker import temporalio.workflow -from temporalio.contrib.opentelemetry._context import attach_context +from temporalio.contrib.opentelemetry._context import attached_context from temporalio.contrib.opentelemetry._tracer_provider import ( ReplaySafeTracerProvider, ) @@ -126,8 +126,7 @@ def _maybe_span( yield return - attached = attach_context(context) if context else None - try: + with attached_context(context): with tracer.start_as_current_span( name, attributes=attributes, @@ -149,9 +148,6 @@ def _maybe_span( ) ) raise - finally: - if attached: - attached.detach() class OpenTelemetryInterceptor( @@ -335,8 +331,7 @@ async def execute_activity( self, input: temporalio.worker.ExecuteActivityInput ) -> Any: context = _headers_to_context(input.headers) - attached = attach_context(context) - try: + with attached_context(context): info = temporalio.activity.info() with _maybe_span( get_tracer(__name__), @@ -350,8 +345,6 @@ async def execute_activity( kind=opentelemetry.trace.SpanKind.SERVER, ): return await super().execute_activity(input) - finally: - attached.detach() class _TracingNexusOperationInboundInterceptor( @@ -368,11 +361,8 @@ def __init__( @contextmanager def _top_level_context(self, headers: Mapping[str, str]) -> Iterator[None]: context = _nexus_headers_to_context(headers) - attached = attach_context(context) - try: + with attached_context(context): yield - finally: - attached.detach() async def execute_nexus_operation_start( self, input: temporalio.worker.ExecuteNexusOperationStartInput @@ -500,11 +490,8 @@ async def handle_update_handler( @contextmanager def _top_level_workflow_context(self, input: _InputWithHeaders) -> Iterator[None]: context = _headers_to_context(input.headers) - attached = attach_context(context) - try: + with attached_context(context): yield - finally: - attached.detach() class _TracingWorkflowOutboundInterceptor( diff --git a/tests/contrib/opentelemetry/test_opentelemetry.py b/tests/contrib/opentelemetry/test_opentelemetry.py index 53bcbe938..c42a84bec 100644 --- a/tests/contrib/opentelemetry/test_opentelemetry.py +++ b/tests/contrib/opentelemetry/test_opentelemetry.py @@ -33,8 +33,8 @@ TracingWorkflowInboundInterceptor, create_tracer_provider, ) +from temporalio.contrib.opentelemetry import _context as otel_context from temporalio.contrib.opentelemetry import workflow as otel_workflow -from temporalio.contrib.opentelemetry._context import AttachedContext from temporalio.contrib.opentelemetry._otel_interceptor import ( _TracingWorkflowInboundInterceptor as _OtelTracingWorkflowInboundInterceptor, ) @@ -937,7 +937,7 @@ async def test_opentelemetry_context_restored_after_activity( detach_count = 0 original_attach = context.attach original_detach = context.detach - original_attached_detach = AttachedContext.detach + original_context_detach = otel_context._detach def tracked_attach(ctx): # type:ignore[reportMissingParameterType] nonlocal attach_count @@ -951,16 +951,16 @@ def tracked_detach(token): # type:ignore[reportMissingParameterType] # Spans detach through context.detach; the interceptors reset their own # tokens directly, so count those detaches where they happen. - def tracked_attached_detach(self: AttachedContext) -> bool: + def tracked_context_detach(ctx: Any, token: Any) -> bool: nonlocal detach_count - detached = original_attached_detach(self) + detached = original_context_detach(ctx, token) if detached: detach_count += 1 return detached context.attach = tracked_attach context.detach = tracked_detach - monkeypatch.setattr(AttachedContext, "detach", tracked_attached_detach) + monkeypatch.setattr(otel_context, "_detach", tracked_context_detach) try: task_queue = f"task_queue_{uuid.uuid4()}" From 77a029f749040f5dc5cc93586a25a5c15cd23bb6 Mon Sep 17 00:00:00 2001 From: DABH Date: Thu, 1 Oct 2026 13:11:44 -0500 Subject: [PATCH 6/6] Count the interceptors' own detaches in the activity leak test The test patched opentelemetry.context.attach/detach process-wide, which also counts the threading instrumentation's attach/detach around every thread's run(). A pool thread started before the patch and exiting during the test adds a detach with no counted attach (9 attaches vs 10 detaches on 3.10 macos-arm). Record the interceptors' own detach results instead. --- .../opentelemetry/test_opentelemetry.py | 89 +++++++------------ 1 file changed, 34 insertions(+), 55 deletions(-) diff --git a/tests/contrib/opentelemetry/test_opentelemetry.py b/tests/contrib/opentelemetry/test_opentelemetry.py index c42a84bec..4c11e95f8 100644 --- a/tests/contrib/opentelemetry/test_opentelemetry.py +++ b/tests/contrib/opentelemetry/test_opentelemetry.py @@ -933,64 +933,43 @@ async def test_opentelemetry_context_restored_after_activity( expect_failure: bool, monkeypatch: pytest.MonkeyPatch, ) -> None: - attach_count = 0 - detach_count = 0 - original_attach = context.attach - original_detach = context.detach - original_context_detach = otel_context._detach - - def tracked_attach(ctx): # type:ignore[reportMissingParameterType] - nonlocal attach_count - attach_count += 1 - return original_attach(ctx) - - def tracked_detach(token): # type:ignore[reportMissingParameterType] - nonlocal detach_count - detach_count += 1 - return original_detach(token) - - # Spans detach through context.detach; the interceptors reset their own - # tokens directly, so count those detaches where they happen. - def tracked_context_detach(ctx: Any, token: Any) -> bool: - nonlocal detach_count - detached = original_context_detach(ctx, token) - if detached: - detach_count += 1 - return detached - - context.attach = tracked_attach - context.detach = tracked_detach - monkeypatch.setattr(otel_context, "_detach", tracked_context_detach) + # Every context the interceptors attach must be detached again, also when + # the activity raises. Counting through the interceptors' own detach keeps + # this independent of other users of opentelemetry.context, such as the + # threading instrumentation's attach/detach around every thread's run(). + detached: list[bool] = [] + original_detach = otel_context._detach - try: - task_queue = f"task_queue_{uuid.uuid4()}" - async with Worker( - client_with_tracing, - task_queue=task_queue, - workflows=[ContextClearWorkflow], - activities=[activity], - ): - with baggage_values({"user.id": "test-123"}): - try: - await client_with_tracing.execute_workflow( - ContextClearWorkflow.run, - id=f"workflow_{uuid.uuid4()}", - task_queue=task_queue, - ) - assert not expect_failure, ( - "This test should have raised an exception" - ) - except Exception: - assert expect_failure, "This test is not expeced to raise" + def tracked_detach(ctx: Any, token: Any) -> bool: + result = original_detach(ctx, token) + detached.append(result) + return result - assert attach_count == detach_count, ( - f"Context leak detected: {attach_count} attaches vs {detach_count} detaches. " - ) - assert attach_count > 0, "Expected at least one context attach/detach" + monkeypatch.setattr(otel_context, "_detach", tracked_detach) - finally: - context.attach = original_attach - context.detach = original_detach + task_queue = f"task_queue_{uuid.uuid4()}" + async with Worker( + client_with_tracing, + task_queue=task_queue, + workflows=[ContextClearWorkflow], + activities=[activity], + ): + with baggage_values({"user.id": "test-123"}): + try: + await client_with_tracing.execute_workflow( + ContextClearWorkflow.run, + id=f"workflow_{uuid.uuid4()}", + task_queue=task_queue, + ) + assert not expect_failure, "This test should have raised an exception" + except Exception: + assert expect_failure, "This test is not expeced to raise" + + assert detached, "Expected at least one context attach/detach" + assert all(detached), ( + f"Context leak detected: {detached.count(False)} of {len(detached)} " + "attached contexts were not detached" + ) @activity.defn