Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
Prev Previous commit
fix(scheduler): preserve cancellation cleanup with send timeouts
  • Loading branch information
Kuang-xianxin committed Sep 19, 2026
commit ad691239faab66a90ceaf4f07f0caacbbcac81b8
11 changes: 10 additions & 1 deletion taskiq/cli/scheduler/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -395,8 +395,17 @@ async def run( # noqa: C901
is_ready_to_send
and task.schedule_id not in running_schedules
):
if send_timeout is not None:
send_coro = send_with_timeout(
self.scheduler,
source,
task,
timeout=send_timeout,
)
else:
send_coro = send(self.scheduler, source, task)
send_task = self._event_loop.create_task(
send(self.scheduler, source, task),
send_coro,
# We need to set the name of the task
# to be able to discard its reference
# after it is done.
Expand Down
12 changes: 10 additions & 2 deletions tests/cli/scheduler/test_shutdown.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,9 +33,13 @@ async def run(self) -> None:
self.finished.set()


@pytest.mark.parametrize("entrypoint", ["loop", "api", "cli"])
@pytest.mark.parametrize(
("entrypoint", "send_timeout"),
[("loop", None), ("api", None), ("cli", None), ("loop", 60), ("cli", 60)],
)
async def test_shutdown_drains_sends_and_updates(
entrypoint: str,
send_timeout: float | None,
monkeypatch: pytest.MonkeyPatch,
caplog: pytest.LogCaptureFixture,
) -> None:
Expand Down Expand Up @@ -88,12 +92,16 @@ async def shutdown(self) -> None:
modules=[],
update_interval=0,
configure_logging=False,
send_timeout=send_timeout,
),
)
elif entrypoint == "api":
coroutine = run_scheduler_task(scheduler, interval=timedelta(0))
else:
coroutine = SchedulerLoop(scheduler).run(update_interval=timedelta(0))
coroutine = SchedulerLoop(scheduler).run(
update_interval=timedelta(0),
send_timeout=send_timeout,
)

scheduler_task = asyncio.create_task(coroutine)
cleanup_started = asyncio.gather(
Expand Down
You are viewing a condensed version of this merge commit. You can view the full changes here.