fix(tasks): clean up the singleflight cache on cancellation

`_cached` deduplicates the scheduler's outbound fetches — Miniflux unread
counts and MCP tool snapshots — by parking every concurrent caller on one
shared Future. Two cancellation paths left that Future stranded.

A waiter awaited the shared Future directly, so cancelling the waiter
cancelled the Future the owner and every other waiter were using. It now
awaits through `asyncio.shield`.

The owner removed its pending entry inside the success and `except Exception`
branches. `CancelledError` is a `BaseException`, so it took neither: the key
stayed in `_shared_cache_pending` pointing at a Future nobody would ever
resolve, and every later caller for that key waited forever. Cleanup moves to
a `finally` that is synchronous on purpose, and the owner cancels its own
Future so current waiters wake while a later caller can still retry.

Ported from public `dev` (`ce04dc1d`, #6174 upstream), with its test.
This commit is contained in:
Léo
2026-09-30 16:05:09 +02:00
parent 6105702901
commit 6e4b3aa5bd
2 changed files with 101 additions and 4 deletions
+15 -4
View File
@@ -97,19 +97,30 @@ async def _cached(key: Tuple, ttl: float, fetch: Callable[[], Awaitable[Any]]) -
pending = fut
owner = True
if not owner:
return await pending
# A cancelled waiter must not cancel the shared Future for the owner
# and every other waiter.
return await asyncio.shield(pending)
try:
val = await fetch()
async with _shared_cache_lock:
_shared_cache[key] = (time.monotonic() + ttl, val)
_shared_cache_pending.pop(key, None)
pending.set_result(val)
return val
except asyncio.CancelledError:
# Cancellation is a BaseException on supported Python versions, so it
# bypasses the Exception handler below. Wake all current waiters while
# allowing a later caller to retry the fetch.
pending.cancel()
raise
except Exception as e:
async with _shared_cache_lock:
_shared_cache_pending.pop(key, None)
pending.set_exception(e)
raise
finally:
# Keep this cleanup synchronous so a second cancellation cannot
# interrupt it and leave a permanently pending Future behind. All
# access runs on the scheduler's event-loop thread.
if _shared_cache_pending.get(key) is pending:
_shared_cache_pending.pop(key, None)
def compute_next_run(schedule: str, scheduled_time: str,