Skip to content

Bug: Retrying task produces the error as the pipeline result #26

Description

@kai-nashi

Description

  1. Schedule a pipeline and wait for its result.
  2. The first task in the pipeline raises an exception and goes to retry.
  3. Retrieve the pipeline result as an error.
  4. The pipeline then successfully completes with the correct result.

Expected

Await the pipeline’s result while a task is retrying.

Actual

The first exception is returned as the result regardless of retries.

Code

import asyncio
import math

from taskiq import Context
from taskiq import SimpleRetryMiddleware
from taskiq import TaskiqDepends
from taskiq import InMemoryBroker
from taskiq_pipelines import Pipeline
from taskiq_pipelines import PipelineMiddleware
from taskiq_redis import RedisAsyncResultBackend
from taskiq_redis import RedisStreamBroker

result_backend = RedisAsyncResultBackend(
    redis_url="redis://localhost:6379",
)

broker = (
    RedisStreamBroker(
        url="redis://localhost:6379",
    )
    .with_middlewares(
        PipelineMiddleware(),
        SimpleRetryMiddleware(default_retry_count=3),
    )
    .with_result_backend(result_backend)
)

check_interval = 0.2


@broker.task("power", retry_on_error=True)
async def power(x: int, y: int = 2, context: "Context" = TaskiqDepends()) -> int:
    if context.message.labels.get("_retries", 0) == 0:
        raise ValueError()
    await asyncio.sleep(check_interval * 10)
    result = x**y
    print(f"{x} ** {y} = {result}")
    return result


@broker.task("sqrt", retry_on_error=True)
async def sqrt(value: int) -> float:
    result = math.sqrt(value)
    print(f"sqrt({value}) = {result}")
    return result  #


pipeline = Pipeline(broker).call_next(power).call_next(sqrt)


async def main():
    task = await pipeline.kiq(2)
    result = await task.wait_result(check_interval=check_interval)
    print("result:", result)


if __name__ == "__main__":
    asyncio.run(main())

Activity

  1. harlanrs commented on Dec 10, 2025

    @harlanrs

    Potentially related details but let me know if I should raise this as its own issue.

    I was running into a serialization issue when using both PipelineMiddleware and SmartRetryMiddleware with RedisStreamBroker and ListRedisScheduleSource.

    The pipeline would work fine until a task triggered a delayed retry which would be funneled through the scheduler instance.

    Here's an example of the pipeline labels logged in the task before retry:

    '_pipe_data': type=bytes, value=b'[{"step_type":"sequential","step_data":{"task_name":"<task_name>", ... }]'
    '_pipe_current_step': type=int, value=4
    

    And here's how these labels were logged after the task was triggered by the scheduler:

    '_pipe_data': type=str, value='W3sic3RlcF90eXBlIjoic2VxdW ...'
    '_pipe_current_step': type=str, value='4'
    

    Note that the _pipe_data label was converted from a bytes json array to a base64 string and _pipe_current_step was converted from an int to a string.

    This would lead to a JSONDecodeError in the post_save (line 44 of taskiq_pipeline/middleware.py) due to the pipeline label being a base64 string.

    I was able to bandaid the serialization error by adding a custom middleware first ahead of the other middlewares that specifically checked _pipe_data and base64 decoded it if needed. I am using the ORJSONSerializer so I could simply check if the label was a string but this might require an alternative check when using the standard json serializer.

    I also went ahead and converted _pipe_current_step back to an int from a string to match the label format pre-scheduler but haven't specifically verified if this is required.

    With this custom middleware in place and 'correcting' the labels post-scheduler, the pipeline appears to be fully functional again.

    Hope this provides some guidance or clues for others experiencing this issue.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions