From be0d3d1c89adb95b8dd0f4a30648f32bd9f4393a Mon Sep 17 00:00:00 2001 From: aticie Date: Mon, 24 Aug 2026 18:00:39 +0200 Subject: [PATCH] fix opentelemetry context detach issues --- .../middlewares/opentelemetry_middleware.py | 57 ++++++++----------- 1 file changed, 23 insertions(+), 34 deletions(-) diff --git a/taskiq/middlewares/opentelemetry_middleware.py b/taskiq/middlewares/opentelemetry_middleware.py index 6d2d8bfc..06f55833 100644 --- a/taskiq/middlewares/opentelemetry_middleware.py +++ b/taskiq/middlewares/opentelemetry_middleware.py @@ -312,40 +312,6 @@ def pre_execute(self, message: TaskiqMessage) -> TaskiqMessage: ) return message - def post_save( # pylint: disable=R6301 - self, - message: TaskiqMessage, - result: TaskiqResult[T], - ) -> None: - """ - This function closes span from `pre_execute`. - - :param message: received message. - :param result: result of the execution. - """ - logger.debug("post_execute task_id=%s", message.task_id) - - # retrieve and finish the Span - ctx = retrieve_context(message) - - if ctx is None: - logger.warning("no existing span found for task_id=%s", message.task_id) - return - - span, activation, token = ctx - - if span.is_recording(): - span.set_attribute(_TASK_TAG_KEY, _TASK_EXECUTE) - set_attributes_from_context(span, message.labels) - span.set_attribute(_TASK_NAME_KEY, message.task_name) - - activation.__exit__(None, None, None) - detach_context(message) - # if the process sending the task is not instrumented - # there's no incoming context and no token to detach - if token is not None: - context_api.detach(token) # type: ignore[arg-type] - def on_error( self, message: TaskiqMessage, @@ -399,6 +365,8 @@ def post_execute( :param message: received message. :param result: result of the execution. """ + logger.debug("post_execute task_id=%s", message.task_id) + if result.is_err: retry_on_error = message.labels.get("retry_on_error") if isinstance(retry_on_error, str): @@ -441,3 +409,24 @@ def post_execute( -1, attributes={"task_name": message.task_name}, ) + + # retrieve and finish the Span + ctx = retrieve_context(message) + + if ctx is None: + logger.warning("no existing span found for task_id=%s", message.task_id) + return + + span, activation, token = ctx + + if span.is_recording(): + span.set_attribute(_TASK_TAG_KEY, _TASK_EXECUTE) + set_attributes_from_context(span, message.labels) + span.set_attribute(_TASK_NAME_KEY, message.task_name) + + activation.__exit__(None, None, None) + detach_context(message) + # if the process sending the task is not instrumented + # there's no incoming context and no token to detach + if token is not None: + context_api.detach(token) # type: ignore[arg-type] \ No newline at end of file