diff --git a/CHANGELOG.md b/CHANGELOG.md index 4109075..74fefd5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -25,6 +25,15 @@ Note that 0.3.3 was never released; 0.3.4 follows 0.3.2. compare-and-delete, Memcached `ADD` and `GETS` + `CAS`, memory its lock. Third-party backends inherit non-atomic fallbacks. ([#62](https://github.com/allen0099/FastAPI-CacheX/issues/62)) +### Fixed + +- Concurrent first requests to the `AppCache` dependency, in an app that never + configured a backend, no longer each build their own `MemoryBackend` and + `CacheManager`. `get_app_cache` runs in worker threads, so the last one + registered replaced the others and, for a while, requests used caches that + could not see each other's entries. The lazy set-up and the `@cache` + fallback now happen under one lock. ([#76](https://github.com/allen0099/FastAPI-CacheX/issues/76)) + ### Documentation - The documentation is published at , diff --git a/fastapi_cachex/cache.py b/fastapi_cachex/cache.py index cd075da..52ed75a 100644 --- a/fastapi_cachex/cache.py +++ b/fastapi_cachex/cache.py @@ -26,12 +26,12 @@ from starlette.status import HTTP_300_MULTIPLE_CHOICES from starlette.status import HTTP_304_NOT_MODIFIED -from .backends import MemoryBackend from .directives import DirectiveType from .exceptions import BackendNotFoundError from .exceptions import CacheXError from .exceptions import RequestNotFoundError from .proxy import BackendProxy +from .proxy import get_backend_or_fallback from .types import CACHE_KEY_SEPARATOR from .types import CacheEntry from .types import CacheKeyBuilder @@ -470,12 +470,7 @@ def build_cache_control() -> str: @wraps(func) async def wrapper(*args: Any, **kwargs: Any) -> Response: # Resolve backend on every request to support lifespan-configured backends - try: - cache_backend = BackendProxy.get() - except BackendNotFoundError: - cache_backend = MemoryBackend() - BackendProxy.set(cache_backend) - logger.debug("No backend configured; using MemoryBackend fallback") + cache_backend = get_backend_or_fallback() if found_request: req: Request | None = kwargs.get(request_name) diff --git a/fastapi_cachex/dependencies.py b/fastapi_cachex/dependencies.py index 0864061..d720ea7 100644 --- a/fastapi_cachex/dependencies.py +++ b/fastapi_cachex/dependencies.py @@ -1,15 +1,18 @@ """FastAPI dependency injection utilities for cache control.""" +import threading from typing import Annotated from fastapi import Depends -from .backends import MemoryBackend from .backends.base import BaseCacheBackend from .exceptions import BackendNotFoundError from .manager import CacheManager from .manager_proxy import CacheManagerProxy from .proxy import BackendProxy +from .proxy import get_backend_or_fallback + +_manager_lock = threading.Lock() def get_cache_backend() -> BaseCacheBackend: @@ -36,14 +39,17 @@ def get_app_cache() -> CacheManager: try: return CacheManagerProxy.get() except BackendNotFoundError: + pass + # Checked again under the lock: FastAPI runs this sync dependency in a + # worker thread, so concurrent first requests would otherwise each build + # and register their own manager (and fallback backend). + with _manager_lock: try: - backend = BackendProxy.get() + return CacheManagerProxy.get() except BackendNotFoundError: - backend = MemoryBackend() - BackendProxy.set(backend) - manager = CacheManager(backend=backend) - CacheManagerProxy.set(manager) - return manager + manager = CacheManager(backend=get_backend_or_fallback()) + CacheManagerProxy.set(manager) + return manager AppCache = Annotated[CacheManager, Depends(get_app_cache)] diff --git a/fastapi_cachex/proxy.py b/fastapi_cachex/proxy.py index ee15e60..613cd46 100644 --- a/fastapi_cachex/proxy.py +++ b/fastapi_cachex/proxy.py @@ -2,6 +2,7 @@ from __future__ import annotations +import threading import warnings from logging import getLogger from typing import Generic @@ -9,12 +10,17 @@ from typing import TypeVar from .backends import BaseCacheBackend +from .backends import MemoryBackend from .exceptions import BackendNotFoundError ProxyInstance = TypeVar("ProxyInstance") logger = getLogger(__name__) +# Serialises the lazy fallback below. `get_app_cache` is a sync dependency that +# FastAPI runs in a worker thread, so two first requests can reach it at once. +_fallback_lock = threading.Lock() + class ProxyMeta(type): """Metaclass for BackendProxy to prevent instantiation.""" @@ -94,3 +100,25 @@ def set_backend(backend: BaseCacheBackend | None) -> None: stacklevel=2, ) BackendProxy.set(backend) + + +def get_backend_or_fallback() -> BaseCacheBackend: + """Return the configured backend, registering a `MemoryBackend` if none is. + + Used by `@cache` and the `AppCache` dependency. The check and the + registration happen under one lock, so concurrent first callers — including + ones on worker threads — all end up with the same fallback instead of each + installing its own and overwriting the others. + """ + try: + return BackendProxy.get() + except BackendNotFoundError: + pass + with _fallback_lock: + try: + return BackendProxy.get() + except BackendNotFoundError: + backend = MemoryBackend() + BackendProxy.set(backend) + logger.debug("No backend configured; using MemoryBackend fallback") + return backend diff --git a/tests/test_dependencies.py b/tests/test_dependencies.py index 637d7b0..d017e58 100644 --- a/tests/test_dependencies.py +++ b/tests/test_dependencies.py @@ -7,6 +7,7 @@ from fastapi_cachex.backends import MemoryBackend from fastapi_cachex.dependencies import get_app_cache from fastapi_cachex.exceptions import BackendNotFoundError +from fastapi_cachex.manager import CacheManager from fastapi_cachex.manager_proxy import CacheManagerProxy # Setup FastAPI application @@ -67,3 +68,42 @@ async def test_get_app_cache_uses_the_configured_backend(): manager = get_app_cache() assert manager.backend is backend + + +def test_get_app_cache_concurrent_first_calls_share_one_backend(monkeypatch): + """Concurrent first calls must agree on one backend and one manager. + + FastAPI runs the sync `get_app_cache` in worker threads. Without a lock, + two first requests each built a `MemoryBackend` and a `CacheManager`, and + the later `set()` replaced the earlier, leaving two caches that could not + see each other's entries. + """ + import threading + import time + from concurrent.futures import ThreadPoolExecutor + + import fastapi_cachex.proxy + + class SlowMemoryBackend(MemoryBackend): + def __init__(self) -> None: + time.sleep(0.05) # widen the window between the check and the set + super().__init__() + + monkeypatch.setattr(fastapi_cachex.proxy, "MemoryBackend", SlowMemoryBackend) + BackendProxy.set(None) + CacheManagerProxy.set(None) + + workers = 8 + barrier = threading.Barrier(workers) + + def first_call() -> CacheManager: + barrier.wait() + return get_app_cache() + + with ThreadPoolExecutor(max_workers=workers) as pool: + managers = list(pool.map(lambda _: first_call(), range(workers))) + + assert len({id(m) for m in managers}) == 1 + assert len({id(m.backend) for m in managers}) == 1 + assert BackendProxy.get() is managers[0].backend + assert CacheManagerProxy.get() is managers[0]