Files
odysseus/tests/test_task_scheduler_cache.py
T
Léo 6e4b3aa5bd 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.
2026-09-30 16:05:09 +02:00

87 lines
2.4 KiB
Python

import asyncio
import pytest
from src import task_scheduler
@pytest.fixture(autouse=True)
def clear_shared_cache():
task_scheduler._shared_cache.clear()
task_scheduler._shared_cache_pending.clear()
yield
task_scheduler._shared_cache.clear()
task_scheduler._shared_cache_pending.clear()
async def test_cached_owner_cancellation_wakes_waiters_and_allows_retry():
key = ("cancelled-owner",)
fetch_started = asyncio.Event()
async def blocked_fetch():
fetch_started.set()
await asyncio.Event().wait()
owner = asyncio.create_task(task_scheduler._cached(key, 60, blocked_fetch))
await fetch_started.wait()
async def unexpected_fetch():
pytest.fail("a waiter must share the owner's fetch")
waiter = asyncio.create_task(task_scheduler._cached(key, 60, unexpected_fetch))
await asyncio.sleep(0)
owner.cancel()
with pytest.raises(asyncio.CancelledError):
await owner
with pytest.raises(asyncio.CancelledError):
await asyncio.wait_for(waiter, timeout=1)
assert key not in task_scheduler._shared_cache_pending
async def retry_fetch():
return "fresh"
result = await asyncio.wait_for(
task_scheduler._cached(key, 60, retry_fetch),
timeout=1,
)
assert result == "fresh"
async def test_cached_waiter_cancellation_does_not_cancel_shared_fetch():
key = ("cancelled-waiter",)
fetch_started = asyncio.Event()
release_fetch = asyncio.Event()
async def blocked_fetch():
fetch_started.set()
await release_fetch.wait()
return "shared"
owner = asyncio.create_task(task_scheduler._cached(key, 60, blocked_fetch))
await fetch_started.wait()
async def unexpected_fetch():
pytest.fail("a waiter must share the owner's fetch")
waiter = asyncio.create_task(task_scheduler._cached(key, 60, unexpected_fetch))
await asyncio.sleep(0)
waiter.cancel()
with pytest.raises(asyncio.CancelledError):
await waiter
pending = task_scheduler._shared_cache_pending[key]
assert not pending.cancelled()
assert not owner.done()
release_fetch.set()
assert await asyncio.wait_for(owner, timeout=1) == "shared"
assert key not in task_scheduler._shared_cache_pending
async def cache_miss():
pytest.fail("the successful owner result should be cached")
assert await task_scheduler._cached(key, 60, cache_miss) == "shared"