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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,11 @@ to include examples, links to docs, or any other relevant information.

### Fixed

- `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 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
Expand Down
51 changes: 51 additions & 0 deletions temporalio/contrib/opentelemetry/_context.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
"""Attach OpenTelemetry contexts so that the matching detach is always safe."""

from __future__ import annotations

from collections.abc import Iterator
from contextlib import contextmanager
from contextvars import Token

import opentelemetry.context
from opentelemetry.context import 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
59 changes: 20 additions & 39 deletions temporalio/contrib/opentelemetry/_interceptor.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import temporalio.nexus.system.workflow_service.models
import temporalio.worker
import temporalio.workflow
from temporalio.contrib.opentelemetry._context import attached_context
from temporalio.exceptions import ApplicationError, ApplicationErrorCategory

# OpenTelemetry dynamically, lazily chooses its context implementation at
Expand Down Expand Up @@ -183,8 +184,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
try:
with attached_context(context):
with self.tracer.start_as_current_span(
name,
attributes=attributes,
Expand Down Expand Up @@ -219,9 +219,6 @@ def _start_as_current_span(
)
)
raise
finally:
if token and context is opentelemetry.context.get_current():
opentelemetry.context.detach(token)

def _completed_workflow_span(
self, params: _CompletedWorkflowSpanParams
Expand Down Expand Up @@ -556,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
token = opentelemetry.context.attach(context)
try:
with attached_context(context):
# This won't be created if there was no context header
self._completed_span(
f"HandleQuery:{input.query}",
Expand All @@ -567,13 +563,6 @@ async def handle_query(self, input: temporalio.worker.HandleQueryInput) -> Any:
kind=opentelemetry.trace.SpanKind.SERVER,
)
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)

def handle_update_validator(
self, input: temporalio.worker.HandleUpdateInput
Expand Down Expand Up @@ -644,31 +633,23 @@ def _top_level_workflow_context(
success = False
exception: Exception | None = None
# Run under this context
token = opentelemetry.context.attach(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,
)

# 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)
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]
Expand Down
25 changes: 5 additions & 20 deletions temporalio/contrib/opentelemetry/_otel_interceptor.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import temporalio.nexus.system.workflow_service.models
import temporalio.worker
import temporalio.workflow
from temporalio.contrib.opentelemetry._context import attached_context
from temporalio.contrib.opentelemetry._tracer_provider import (
ReplaySafeTracerProvider,
)
Expand Down Expand Up @@ -125,8 +126,7 @@ def _maybe_span(
yield
return

token = opentelemetry.context.attach(context) if context else None
try:
with attached_context(context):
with tracer.start_as_current_span(
name,
attributes=attributes,
Expand All @@ -148,9 +148,6 @@ def _maybe_span(
)
)
raise
finally:
if token and context is opentelemetry.context.get_current():
opentelemetry.context.detach(token)


class OpenTelemetryInterceptor(
Expand Down Expand Up @@ -334,8 +331,7 @@ async def execute_activity(
self, input: temporalio.worker.ExecuteActivityInput
) -> Any:
context = _headers_to_context(input.headers)
token = opentelemetry.context.attach(context)
try:
with attached_context(context):
info = temporalio.activity.info()
with _maybe_span(
get_tracer(__name__),
Expand All @@ -349,9 +345,6 @@ async def execute_activity(
kind=opentelemetry.trace.SpanKind.SERVER,
):
return await super().execute_activity(input)
finally:
if context is opentelemetry.context.get_current():
opentelemetry.context.detach(token)


class _TracingNexusOperationInboundInterceptor(
Expand All @@ -368,12 +361,8 @@ 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)
try:
with attached_context(context):
yield
finally:
if context is opentelemetry.context.get_current():
opentelemetry.context.detach(token)

async def execute_nexus_operation_start(
self, input: temporalio.worker.ExecuteNexusOperationStartInput
Expand Down Expand Up @@ -501,12 +490,8 @@ 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)
try:
with attached_context(context):
yield
finally:
if context is opentelemetry.context.get_current():
opentelemetry.context.detach(token)


class _TracingWorkflowOutboundInterceptor(
Expand Down
Loading
Loading