Skip to content

Commit 6eaf254

Browse files
Propagate external cancellation out of Job.close() (#612)
1 parent b456399 commit 6eaf254

3 files changed

Lines changed: 117 additions & 2 deletions

File tree

‎CHANGES/606.bugfix‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
Fixed ``Job.close()`` and ``Job.wait()`` swallowing an external cancellation of the calling task (exact on Python 3.11+ via ``Task.cancelling()``; on older versions the cancellation propagates whenever the job is still running).

‎aiojobs/_job.py‎

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -132,9 +132,28 @@ async def _close(self, timeout: Optional[float]) -> None:
132132
scheduler = self._scheduler
133133
try:
134134
async with asyncio_timeout(timeout):
135-
await self._task
135+
if sys.version_info >= (3, 11):
136+
await self._task
137+
else:
138+
# Cancelling the current task would be forwarded to
139+
# self._task through _fut_waiter, making the two
140+
# cancellation sources indistinguishable; shield the
141+
# job so an external cancellation leaves it running
142+
# and detectable below.
143+
await asyncio.shield(self._task)
136144
except asyncio.CancelledError:
137-
pass
145+
# Either the cancelled job finished or the task running
146+
# close() was itself cancelled; re-raise in the second case
147+
# so callers stay cancellable.
148+
if sys.version_info >= (3, 11):
149+
ctask = asyncio.current_task()
150+
assert ctask is not None
151+
if ctask.cancelling() > 0:
152+
raise
153+
elif not self._task.done():
154+
# The job is still running, so the CancelledError cannot
155+
# have come from awaiting it.
156+
raise
138157
except asyncio.TimeoutError as exc:
139158
if self._explicit:
140159
raise

‎tests/test_job.py‎

Lines changed: 95 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import asyncio
2+
import sys
23
from collections.abc import Awaitable
34
from contextlib import suppress
45
from typing import Callable, NoReturn
@@ -239,6 +240,100 @@ async def coro2() -> None:
239240
await job.close()
240241

241242

243+
async def test_job_close_cancelled_from_outside(scheduler: Scheduler) -> None:
244+
"""An external cancellation of the task running close() must propagate.
245+
246+
The CancelledError raised by awaiting the job's cancelled task is
247+
expected and swallowed, but a cancellation aimed at the caller of
248+
close() itself is not ours to swallow.
249+
"""
250+
cancel_seen = asyncio.Event()
251+
unblock = asyncio.Event()
252+
253+
async def coro() -> None:
254+
try:
255+
await asyncio.sleep(60)
256+
except asyncio.CancelledError:
257+
cancel_seen.set()
258+
await unblock.wait()
259+
raise
260+
261+
job = await scheduler.spawn(coro())
262+
closer = asyncio.ensure_future(job.close())
263+
# Once the job saw the cancellation, close() is suspended awaiting
264+
# the job task, which is blocked until unblock is set.
265+
await cancel_seen.wait()
266+
closer.cancel()
267+
await asyncio.wait({closer})
268+
assert closer.cancelled()
269+
270+
unblock.set()
271+
with suppress(asyncio.CancelledError):
272+
await job.wait()
273+
274+
275+
@pytest.mark.skipif(
276+
sys.version_info < (3, 11),
277+
reason="distinguishing this race needs Task.cancelling()",
278+
)
279+
async def test_job_close_cancelled_from_outside_completion_race(
280+
scheduler: Scheduler,
281+
) -> None:
282+
"""The job task may finish in the same loop iteration the caller of
283+
close() is cancelled in; the external cancellation must still
284+
propagate."""
285+
cancel_seen = asyncio.Event()
286+
unblock = asyncio.Event()
287+
288+
async def coro() -> None:
289+
try:
290+
await asyncio.sleep(60)
291+
except asyncio.CancelledError:
292+
cancel_seen.set()
293+
await unblock.wait()
294+
raise
295+
296+
job = await scheduler.spawn(coro())
297+
closer = asyncio.ensure_future(job.close())
298+
await cancel_seen.wait()
299+
# Unblock the job before cancelling the closer: the job task then
300+
# completes and schedules its done callbacks before the closer wakes
301+
# up and detaches them.
302+
unblock.set()
303+
closer.cancel()
304+
await asyncio.wait({closer})
305+
assert closer.cancelled()
306+
with suppress(asyncio.CancelledError):
307+
await job.wait()
308+
309+
310+
async def test_job_close_timeout_source_traceback(
311+
make_scheduler: _MakeScheduler,
312+
) -> None:
313+
"""In debug mode the close timeout report carries the source traceback."""
314+
loop = asyncio.get_running_loop()
315+
loop.set_debug(True)
316+
try:
317+
handler = mock.Mock()
318+
scheduler = await make_scheduler(close_timeout=0.01, exception_handler=handler)
319+
320+
async def coro() -> None:
321+
try:
322+
await asyncio.sleep(60)
323+
except asyncio.CancelledError:
324+
await asyncio.sleep(60)
325+
326+
job = await scheduler.spawn(coro())
327+
await scheduler.close()
328+
assert job.closed
329+
assert handler.called
330+
context = handler.call_args[0][1]
331+
assert context["message"] == "Job closing timed out"
332+
assert "source_traceback" in context
333+
finally:
334+
loop.set_debug(False)
335+
336+
242337
async def test_job_await_closed(scheduler: Scheduler) -> None:
243338
async def coro() -> int:
244339
return 5

0 commit comments

Comments
 (0)