Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
8 changes: 8 additions & 0 deletions changelog.d/252.added.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
**`@cache(coalesce=True)` runs the handler once for concurrent misses.**
Requests that miss a key while another request in the same process is
rendering it wait for that request, then are served the stored entry, so a
cold or just-expired key no longer sends every concurrent request to the
handler. When the first response is not stored (a cookie, `private`, an error
status), the waiting requests run the handler themselves, all at once. Off by
default; it needs a positive `ttl` and is rejected with `private` or
`no_cache`. Each worker process still runs the handler once per cold key.
3 changes: 2 additions & 1 deletion docs/APP_CACHE.md
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,8 @@ Complete runnable example: [`examples/app_cache.py`](https://github.com/allen009
When the factory is expensive (a slow database query, a rate-limited upstream
API) and the key is hot, cache expiry turns into simultaneous recomputations.
Distributed stampede protection built on `CacheLock` is on by default; you can
tune or turn it off per call or manager-wide:
tune or turn it off per call or manager-wide (for `@cache` routes, see
[Concurrent misses](HTTP_CACHING.md#concurrent-misses)):

```python
# Per-call protection:
Expand Down
13 changes: 9 additions & 4 deletions docs/CACHE_FLOW.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,9 +9,10 @@ keys), [`_vary.py`](https://github.com/allen0099/FastAPI-CacheX/blob/master/fast
[`_cache_control.py`](https://github.com/allen0099/FastAPI-CacheX/blob/master/fastapi_cachex/_cache_control.py)
(`Cache-Control`),
[`_stored_response.py`](https://github.com/allen0099/FastAPI-CacheX/blob/master/fastapi_cachex/_stored_response.py)
(storing and replaying a response, ETags and 304s) and
(storing and replaying a response, ETags and 304s),
[`_rendering.py`](https://github.com/allen0099/FastAPI-CacheX/blob/master/fastapi_cachex/_rendering.py) (running the
handler and adding its dependencies' headers).
handler and adding its dependencies' headers) and
[`_coalesce.py`](https://github.com/allen0099/FastAPI-CacheX/blob/master/fastapi_cachex/_coalesce.py) (`coalesce=`).

## Overall flow

Expand Down Expand Up @@ -167,8 +168,9 @@ Arguments are validated when the decorator is applied, and a `CacheXError` is
raised if `public` and `private` are both set, if only one of `stale` /
`stale_ttl` is given, if `ttl` is not an `int`, is negative or is larger than
`MAX_TTL`, if `vary` is not a list of header field names, if `sort_query` is
not a `bool` or is passed with a custom `key_builder`, or if `key_builder` is
an `async` callable.
not a `bool` or is passed with a custom `key_builder`, if `key_builder` is
an `async` callable, or if `coalesce` is not a `bool`, or is set without a
positive `ttl` or together with `private` or `no_cache`.

The header value is built once per decorated route:

Expand Down Expand Up @@ -259,6 +261,9 @@ if bypass or (credential and not cache_authorized):
# HEAD: key_builder sees the request with method GET
cache_key = key_builder(request) + vary_components(request) # built only here
entry = await backend.get(cache_key) # expired entries are already skipped here
if coalesce and entry is None and (leader := running_miss(cache_key)):
await leader # GET only leads; HEAD only waits
entry = await backend.get(cache_key) # None: render below, without waiting

if client_etag and no_cache:
fresh = await render() # no-cache: always re-render first
Expand Down
34 changes: 34 additions & 0 deletions docs/HTTP_CACHING.md
Original file line number Diff line number Diff line change
Expand Up @@ -244,6 +244,40 @@ empty body. The server drops the body of every HEAD response.
A route that already accepted HEAD ran the handler on every HEAD request before
0.4.2 and added none of these headers; it now skips the handler on a hit.

### Concurrent misses

By default every request that misses runs the handler: 20 concurrent requests
to a cold key, or arriving just after its entry expired, run it 20 times.
`coalesce=True` runs it once per key in each process:

```python
@app.get("/report")
@cache(ttl=60, coalesce=True)
async def report() -> dict[str, int]: ...
```

The first request to miss renders as usual. Requests that miss the same key
while it runs wait for it to answer, then read the backend again and are served
the stored entry, with an `ETag`, `Age` and 304 handling like any hit (#252).
When the first response was not stored (it set a cookie, was `private` or
`no-store`, had an error status or a status a dependency set, was streamed, or
the handler raised), each waiting request runs the handler itself, all at once
rather than one after another.

- Only requests in one process are coalesced, without a lock in the backend:
each worker still runs the handler once per cold key. For one run across
workers, cache the expensive part with `get_or_set()`, whose
[stampede protection](APP_CACHE.md#stampede-protection) uses a distributed
lock.
- A waiting request waits as long as the first one's handler takes; there is
no separate timeout. A request whose backend read fails does not wait.
- A HEAD request waits for a running GET of the same key but never makes
others wait, since its response is not stored.
- It needs a positive `ttl`. With `private=True` or `no_cache=True`, which
never serve a stored entry, the decorator raises `CacheXError`; with
`no_store=True` it is ignored with the usual warning. Requests with
credentials that bypass the backend are not coalesced.

### Requests with credentials

A single-page app that sends `Authorization` on every request, or a site where
Expand Down
70 changes: 70 additions & 0 deletions fastapi_cachex/_coalesce.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
"""In-process coalescing of concurrent misses for `@cache(coalesce=True)` (#252).

The first request to miss a key leads: it runs the handler as usual. Requests
that miss the same key while it runs follow: they wait until the leader has
answered, then read the backend again. A follower that still finds no entry
(the leader's response was not stored, or the handler failed) renders its own
response without waiting again, so followers never queue behind each other.

Only requests in this process are coalesced. The registry is keyed by event
loop and backend as well as the cache key, so two apps or loops never wait on
each other.
"""

import asyncio

# One future per key being rendered by a leader; done once it has answered.
_IN_FLIGHT: dict[tuple[int, int, str], asyncio.Future[None]] = {}


class _Flight:
"""One request's part in coalescing: leader, follower or neither.

The wrapper creates one per request and calls ``finish()`` once the
request has been answered, however it ended, so a leader always releases
its followers and leaves nothing behind in the registry.
"""

def __init__(self) -> None:
self._key: tuple[int, int, str] | None = None
self._future: asyncio.Future[None] | None = None

def join(
self, backend: object, cache_key: str, *, may_lead: bool
) -> asyncio.Future[None] | None:
"""Return the leader's future to wait on, or ``None`` to render.

With ``may_lead`` and no leader for the key, this request becomes the
leader. Without it (HEAD, whose response is never stored), it renders
without blocking anyone.
"""
loop = asyncio.get_running_loop()
key = (id(loop), id(backend), cache_key)
running = _IN_FLIGHT.get(key)
if running is not None or not may_lead:
return running
self._key = key
self._future = loop.create_future()
_IN_FLIGHT[key] = self._future
return None

def finish(self) -> None:
"""Release the followers, if this request led."""
if self._key is None or self._future is None:
return
# Only the leader registers or removes its key, and followers wait
# through a shield, so the future is still pending. `pop` because a
# test teardown may have cleared the registry already.
_IN_FLIGHT.pop(self._key, None)
self._future.set_result(None)
self._key = None
self._future = None


async def _wait_for(leader: asyncio.Future[None]) -> None:
"""Wait until the leader has answered.

Shielded: a follower that is cancelled stops waiting without cancelling
the leader's future for the other followers.
"""
await asyncio.shield(leader)
92 changes: 81 additions & 11 deletions fastapi_cachex/cache.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,8 @@
from ._cache_control import _cache_control_for
from ._cache_control import _unshareable_reason
from ._cache_control import _with_cache_control
from ._coalesce import _Flight
from ._coalesce import _wait_for
from ._key_builders import _append_key_components
from ._key_builders import _build_key
from ._key_builders import _resolve_key_builder
Expand All @@ -50,13 +52,15 @@
from ._vary import _validate_vary
from ._vary import _vary_components
from .backends.base import MAX_TTL
from .backends.base import BaseCacheBackend
from .directives import DirectiveType
from .exceptions import BackendNotFoundError
from .exceptions import CacheXError
from .exceptions import RequestNotFoundError
from .headers import add_vary
from .proxy import BackendProxy
from .proxy import get_backend_or_fallback
from .types import CacheEntry
from .types import CacheKeyBuilder
from .types import log_ref

Expand Down Expand Up @@ -111,6 +115,22 @@ def _log_backend_failure(
logger.debug("Cache backend %s; key_ref=%s key=%s", what, key_ref, cache_key)


async def _read_entry(
backend: BaseCacheBackend, cache_key: str, request: Request, *, fail_open: bool
) -> tuple[CacheEntry | None, bool]:
"""Read the entry under ``cache_key`` and whether the read worked.

A failed read is a miss if failing open.
"""
try:
return await backend.get(cache_key), True
except Exception as e:
if not fail_open:
raise
_log_backend_failure("read failed; serving uncached", request, cache_key, e)
return None, False


# Where `FastAPICacheXSessionMiddleware` puts the session it loaded;
# `get_session` reads it.
_SESSION_STATE_KEY = "__fastapi_cachex_session"
Expand Down Expand Up @@ -357,6 +377,7 @@ def cache(
cache_authorized: bool = False,
vary: Sequence[str] | None = None,
sort_query: bool | None = None,
coalesce: bool = False,
) -> Callable[[HandlerCallable], AsyncResponseCallable]:
"""Cache decorator for FastAPI route handlers.

Expand Down Expand Up @@ -470,6 +491,19 @@ def cache(
is (``build_cache_key()`` sorts unless told otherwise), and
passing ``sort_query`` with one is rejected. Pass the same value
to ``invalidate()``.
coalesce: Run the handler once for concurrent misses of one key in
this process. The first request to miss renders as usual; requests
that miss the same key meanwhile wait for it to answer, then read
the backend again and are served the stored entry. If the first
response was not stored (it set a cookie, was ``private``, had an
error status, was streamed, or the handler raised), each of them
runs the handler itself, at once. A request whose backend read
fails does not wait. Each worker process still runs
the handler once per cold key, and a waiting request waits as
long as the first one's handler takes. A HEAD request waits for a
running GET but never makes others wait. Requires a positive
``ttl`` and is rejected with ``private`` or ``no_cache``, which
never serve a stored entry.

Returns:
Decorator function that wraps route handlers with caching logic
Expand All @@ -481,7 +515,9 @@ def cache(
negative or is larger than ``MAX_TTL``, or if ``vary`` is not a
list of header field names (a single string is rejected), if
``sort_query`` is not a ``bool`` or is passed with
``key_builder``, or if ``key_builder`` is an ``async`` callable.
``key_builder``, if ``key_builder`` is an ``async`` callable, or
if ``coalesce`` is not a ``bool`` or is set without a positive
``ttl`` or with ``private`` or ``no_cache``.
At request time, if
``key_builder`` returns anything but a ``str``.

Expand Down Expand Up @@ -519,6 +555,18 @@ def decorator(func: HandlerCallable) -> AsyncResponseCallable:
if ttl is not None and ttl > MAX_TTL:
msg = f"ttl must be at most {MAX_TTL} seconds"
raise CacheXError(msg)
if not isinstance(coalesce, bool):
# Unreachable for type checkers; guards untyped callers.
msg = f"coalesce must be a bool, got {type(coalesce).__name__}" # type: ignore[unreachable]
raise CacheXError(msg)
if coalesce and not no_store and (private or no_cache or not ttl):
# Such a route never serves a stored entry, so there would be
# nothing to wait for.
msg = (
"coalesce needs a positive ttl and cannot be combined with "
"private or no_cache, which never serve a stored entry"
)
raise CacheXError(msg)
vary_names = _validate_vary(vary)
builder = _resolve_key_builder(key_builder, sort_query)
if any(name.lower() == "cookie" for name in vary_names):
Expand All @@ -539,6 +587,7 @@ def decorator(func: HandlerCallable) -> AsyncResponseCallable:
("private", private),
("immutable", immutable),
("must_revalidate", must_revalidate),
("coalesce", coalesce),
)
if value is not None and value is not False
]
Expand Down Expand Up @@ -635,6 +684,7 @@ def decorator(func: HandlerCallable) -> AsyncResponseCallable:
)

async def respond(
flight: _Flight,
sub_response: Response | None,
dependency_lines: Sequence[tuple[bytes, bytes]],
dependency_status: int | None,
Expand Down Expand Up @@ -768,13 +818,23 @@ async def respond(
_vary_components(req, vary_names),
)

try:
cached_data = await cache_backend.get(cache_key)
except Exception as e:
if not fail_open:
raise
_log_backend_failure("read failed; serving uncached", req, cache_key, e)
cached_data = None
cached_data, read_ok = await _read_entry(
cache_backend, cache_key, req, fail_open=fail_open
)
# A failed read does not wait: the backend is likely down, so the
# leader could not store its response either. A miss read just
# before a leader stores may still lead a second render; that
# window is one backend round trip.
if coalesce and read_ok and cached_data is None:
leader = flight.join(
cache_backend, cache_key, may_lead=req.method == "GET"
)
if leader is not None:
logger.debug("Waiting for a running miss; key=%s", cache_key)
await _wait_for(leader)
cached_data, _ = await _read_entry(
cache_backend, cache_key, req, fail_open=fail_open
)

current_response: Response | None = None
current_body: bytes | None = None
Expand Down Expand Up @@ -964,9 +1024,19 @@ async def serve(*args: Any, **kwargs: Any) -> Response:
None if sub_response is None else sub_response.status_code
)
req: Request | None = kwargs.get(request_name)
response = await respond(
sub_response, dependency_lines, dependency_status, *args, **kwargs
)
flight = _Flight()
try:
response = await respond(
flight,
sub_response,
dependency_lines,
dependency_status,
*args,
**kwargs,
)
finally:
# Releases the requests waiting on this one (#252).
flight.finish()
return _with_dependency_headers(
response,
sub_response,
Expand Down
2 changes: 1 addition & 1 deletion i18n/zh-TW/docs/APP_CACHE.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,7 @@ await manager.clear_pattern("user:*") # 比對 "myapp:user:*"

## Cache stampede 保護 {#stampede-protection}

當 `factory` 的運算成本很高(例如慢速資料庫查詢、受速率限制的外部 API)且該鍵又是熱門鍵時,快取過期會導致多個請求同時重新計算。以 `CacheLock` 實作的分散式 cache stampede 保護預設開啟;你可以在單次呼叫或 manager 全域調整或關閉它:
當 `factory` 的運算成本很高(例如慢速資料庫查詢、受速率限制的外部 API)且該鍵又是熱門鍵時,快取過期會導致多個請求同時重新計算。以 `CacheLock` 實作的分散式 cache stampede 保護預設開啟;你可以在單次呼叫或 manager 全域調整或關閉它(`@cache` 路由請見[同時發生的未命中](HTTP_CACHING.md#concurrent-misses)):

```python
# 單次呼叫保護:
Expand Down
Loading
Loading