diff --git a/changelog.d/243.added.md b/changelog.d/243.added.md new file mode 100644 index 0000000..fa1d52e --- /dev/null +++ b/changelog.d/243.added.md @@ -0,0 +1,8 @@ +**Every backend can be closed with `aclose()` and used with `async with`.** +`AsyncRedisCacheBackend.aclose()` closes the redis-py client and the connection +pool it created, and `MemcachedBackend.aclose()` closes every pooled socket, so +the lifespan pattern in the backends guide works for all three backends instead +of raising `AttributeError` on Redis and Memcached. `BaseCacheBackend` gains a +no-op `aclose()` that custom backends may override, and `__aenter__`/`__aexit__` +that close the backend when the block ends. Calling `aclose()` twice is safe; a +Redis `connection_pool=` you pass in is left for you to close. diff --git a/docs/BACKENDS.md b/docs/BACKENDS.md index 0663dd6..a5ed325 100644 --- a/docs/BACKENDS.md +++ b/docs/BACKENDS.md @@ -39,22 +39,8 @@ BackendProxy.set(backend) The cleanup task starts on the event loop of the first cache call. If a later call runs on a different loop, for example after the first loop was closed, the task is started again there. To stop it on shutdown, `await backend.aclose()` cancels the -task and waits until it has finished. `stop_cleanup()` only requests cancellation. - -```python -from contextlib import asynccontextmanager - -from fastapi import FastAPI - - -@asynccontextmanager -async def lifespan(app: FastAPI): - yield - await backend.aclose() - - -app = FastAPI(lifespan=lifespan) -``` +task and waits until it has finished (see [Closing a backend](#closing-a-backend)). +`stop_cleanup()` only requests cancellation. `clear_pattern()` matches whole keys case-sensitively on every platform, like Redis. The glob syntax is Python's `fnmatch`, which differs from Redis in two places: negate @@ -171,6 +157,60 @@ was lost), `increment()` a fresh counter of 0. With several servers, a failed on taken out of rotation at once and its keys go to the remaining servers until it answers again. +## Closing a backend + +Every backend has `aclose()`, which releases what it holds open. Call it on +shutdown, at the end of the FastAPI lifespan: + +```python +from contextlib import asynccontextmanager + +from fastapi import FastAPI + +from fastapi_cachex import BackendProxy +from fastapi_cachex.backends import AsyncRedisCacheBackend + + +@asynccontextmanager +async def lifespan(app: FastAPI): + backend = AsyncRedisCacheBackend(host="127.0.0.1", port=6379) + BackendProxy.set(backend) + try: + yield + finally: + BackendProxy.set(None) + await backend.aclose() + + +app = FastAPI(lifespan=lifespan) +``` + +The same lifespan works for every backend. A backend is also an async context +manager, so `async with MemcachedBackend(servers=[...]) as backend:` closes it +when the block ends, including when the block raises. + +| Backend | `aclose()` | +|---------|------------| +| `MemoryBackend` | Cancels the cleanup task and waits until it has finished | +| `AsyncRedisCacheBackend` | Closes the redis-py client and the connection pool it created | +| `MemcachedBackend` | Closes every pooled socket to every server, in a worker thread | + +- Calling `aclose()` more than once is safe. +- The backend owns the client it creates, so `aclose()` closes `backend.client` + too; close the backend rather than reaching into the client. A Redis + connection pool you create yourself and pass as `connection_pool=` is left + open, as redis-py does: whoever creates a pool closes it. +- A closed backend is not locked: the Redis and Memcached clients reconnect on + the next call, and `MemoryBackend` restarts its cleanup task, so a backend used + after `aclose()` needs another `aclose()`. +- `BackendProxy.set(None)` only unregisters the backend; it does not close it. +- A custom backend inherits a no-op `aclose()` from `BaseCacheBackend`; override + it when the backend holds connections or tasks. +- Without `aclose()`, open sockets are closed only by the garbage collector, + which may emit `ResourceWarning`s at shutdown. Before 0.3.9 only + `MemoryBackend` had `aclose()`, and the Redis and Memcached backends were + closed through `backend.client`. + ## Atomic backend primitives Every backend exposes atomic operations on top of `get`/`set`/`delete`, for diff --git a/examples/redis_backend.py b/examples/redis_backend.py index 3fcc8ab..cbf0c9b 100644 --- a/examples/redis_backend.py +++ b/examples/redis_backend.py @@ -40,8 +40,7 @@ async def lifespan(_app: FastAPI) -> AsyncIterator[None]: yield finally: BackendProxy.set(None) - # The types-redis stubs predate aclose() (redis-py 5.0.1+). - await backend.client.aclose() # type: ignore[attr-defined] + await backend.aclose() app = FastAPI(lifespan=lifespan) diff --git a/fastapi_cachex/backends/base.py b/fastapi_cachex/backends/base.py index 83c187c..da2d30a 100644 --- a/fastapi_cachex/backends/base.py +++ b/fastapi_cachex/backends/base.py @@ -4,6 +4,8 @@ from abc import ABC from abc import abstractmethod from collections.abc import Iterable +from types import TracebackType +from typing import TYPE_CHECKING from typing import Any from fastapi_cachex.types import CACHE_KEY_SEPARATOR @@ -11,6 +13,11 @@ from fastapi_cachex.types import counter_entry from fastapi_cachex.types import counter_value +if TYPE_CHECKING: + # Type-only: typing.Self is 3.11+, and typing_extensions is not a runtime + # dependency. + from typing_extensions import Self + def warn_if_path_shaped(pattern: str, cleared: int) -> None: """Warn when a ``clear_pattern`` that cleared nothing was written as a path. @@ -89,7 +96,33 @@ def validate_delta(delta: int) -> int: class BaseCacheBackend(ABC): - """Base class for all cache backends.""" + """Base class for all cache backends. + + Every backend is an async context manager: ``async with`` returns the + backend itself and calls ``aclose()`` on the way out. + """ + + async def aclose(self) -> None: # noqa: B027 - a no-op default, not abstract + """Release what the backend holds open: connections, background tasks. + + Call it once on shutdown, typically at the end of a FastAPI lifespan. + It is safe to call more than once. The base implementation does + nothing, for backends with nothing to release; the built-in backends + override it. + """ + + async def __aenter__(self) -> "Self": + """Return the backend itself.""" + return self + + async def __aexit__( + self, + exc_type: type[BaseException] | None, + exc_value: BaseException | None, + traceback: TracebackType | None, + ) -> None: + """Close the backend with ``aclose()``.""" + await self.aclose() @abstractmethod async def get(self, key: str) -> CacheEntry | None: diff --git a/fastapi_cachex/backends/memcached.py b/fastapi_cachex/backends/memcached.py index cfde919..0edc2c3 100644 --- a/fastapi_cachex/backends/memcached.py +++ b/fastapi_cachex/backends/memcached.py @@ -126,6 +126,15 @@ def __init__( ) self.key_prefix = key_prefix + async def aclose(self) -> None: + """Close every pooled connection to every server. + + Safe to call more than once. The pymemcache client reconnects on the + next call, so a backend used after ``aclose()`` opens new sockets that + need another ``aclose()``. + """ + await asyncio.to_thread(self.client.close) + def _make_key(self, key: str) -> str: """Namespace a cache key, hashing it when Memcached would refuse it. diff --git a/fastapi_cachex/backends/redis.py b/fastapi_cachex/backends/redis.py index c15da04..d7f2422 100644 --- a/fastapi_cachex/backends/redis.py +++ b/fastapi_cachex/backends/redis.py @@ -203,6 +203,19 @@ def load_from_config(config: RedisConfig) -> "AsyncRedisCacheBackend": protocol=config.protocol, ) + async def aclose(self) -> None: + """Close the Redis client and the connection pool it created. + + Safe to call more than once. redis-py reconnects on the next command, + so a backend used after ``aclose()`` needs another ``aclose()``. A + pool passed in through ``connection_pool=`` is left open: redis-py + leaves a pool it did not create to whoever created it, and so does + this backend. + """ + # The types-redis stubs predate aclose() (redis-py 5.0.1+); close() + # is its deprecated alias and warns. + await self.client.aclose() # type: ignore[attr-defined] + def _make_key(self, key: str) -> str: """Add prefix to cache key.""" return f"{self.key_prefix}{key}" diff --git a/i18n/zh-TW/docs/BACKENDS.md b/i18n/zh-TW/docs/BACKENDS.md index b4f0c06..21bad85 100644 --- a/i18n/zh-TW/docs/BACKENDS.md +++ b/i18n/zh-TW/docs/BACKENDS.md @@ -28,22 +28,7 @@ BackendProxy.set(backend) > [!NOTE] > 記憶體快取不適合用於多行程的正式環境。每個行程都各自維護獨立的快取。 -清理 task 會在第一次快取呼叫所在的事件迴圈(event loop)上啟動。若之後的呼叫在另一個迴圈上執行(例如第一個迴圈已關閉),task 會在新的迴圈上重新啟動。關閉應用程式時,`await backend.aclose()` 會取消 task 並等待它結束;`stop_cleanup()` 只會要求取消。 - -```python -from contextlib import asynccontextmanager - -from fastapi import FastAPI - - -@asynccontextmanager -async def lifespan(app: FastAPI): - yield - await backend.aclose() - - -app = FastAPI(lifespan=lifespan) -``` +清理 task 會在第一次快取呼叫所在的事件迴圈(event loop)上啟動。若之後的呼叫在另一個迴圈上執行(例如第一個迴圈已關閉),task 會在新的迴圈上重新啟動。關閉應用程式時,`await backend.aclose()` 會取消 task 並等待它結束(參見[關閉後端](#closing-a-backend));`stop_cleanup()` 只會要求取消。 `clear_pattern()` 在所有平台上都以區分大小寫的方式比對完整的鍵,與 Redis 相同。萬用字元語法採用 Python 的 `fnmatch`,與 Redis 有兩處不同:否定字元類別要寫 `[!...]`(Redis 為 `[^...]`);跳脫特殊字元要放進中括號,例如 `[*]`(Redis 另外也接受 `\*`)。`*`、`?` 與 `[abc]` 在兩者上的行為相同。 @@ -120,6 +105,48 @@ BackendProxy.set(backend) 伺服器無法連線時,所有要送往它的呼叫都會拋出錯誤,一秒後會再嘗試連線。0.3.8 以前,失敗後一秒內的呼叫會回傳虛構的結果:`get()` 當成未命中、`set()` 沒有任何反應(寫入遺失)、`increment()` 當成新的計數器並回傳 0。設定多台伺服器時,失敗的那台會立即移出輪替,它的鍵會改由其餘伺服器處理,直到它恢復回應。 +## 關閉後端 {#closing-a-backend} + +每個後端都有 `aclose()`,用來釋放它持有的連線與背景工作。請在關閉應用程式時,於 FastAPI lifespan 的結尾呼叫它: + +```python +from contextlib import asynccontextmanager + +from fastapi import FastAPI + +from fastapi_cachex import BackendProxy +from fastapi_cachex.backends import AsyncRedisCacheBackend + + +@asynccontextmanager +async def lifespan(app: FastAPI): + backend = AsyncRedisCacheBackend(host="127.0.0.1", port=6379) + BackendProxy.set(backend) + try: + yield + finally: + BackendProxy.set(None) + await backend.aclose() + + +app = FastAPI(lifespan=lifespan) +``` + +每種後端都適用同一個 lifespan。後端也是非同步 context manager,因此 `async with MemcachedBackend(servers=[...]) as backend:` 會在區塊結束時關閉它,區塊拋出例外時也一樣。 + +| 後端 | `aclose()` | +|------|------------| +| `MemoryBackend` | 取消清理 task 並等待它結束 | +| `AsyncRedisCacheBackend` | 關閉 redis-py 用戶端,以及它自己建立的連線池 | +| `MemcachedBackend` | 在工作執行緒中關閉連往每台伺服器的所有連線池 socket | + +- 多次呼叫 `aclose()` 是安全的。 +- 後端擁有它自己建立的用戶端,因此 `aclose()` 也會關閉 `backend.client`;請關閉後端,而不要直接操作用戶端。你自行建立並以 `connection_pool=` 傳入的 Redis 連線池不會被關閉,與 redis-py 的做法相同:誰建立連線池,就由誰關閉。 +- 關閉後的後端並不會被鎖住:Redis 與 Memcached 用戶端會在下一次呼叫時重新連線,`MemoryBackend` 也會重新啟動清理 task,因此在 `aclose()` 之後又使用的後端需要再呼叫一次 `aclose()`。 +- `BackendProxy.set(None)` 只會取消註冊後端,不會關閉它。 +- 自訂後端會從 `BaseCacheBackend` 繼承一個什麼都不做的 `aclose()`;若後端持有連線或背景工作,請覆寫它。 +- 若不呼叫 `aclose()`,開啟中的 socket 只會由垃圾回收器關閉,關閉應用程式時可能會產生 `ResourceWarning`。0.3.9 之前只有 `MemoryBackend` 有 `aclose()`,Redis 與 Memcached 後端要透過 `backend.client` 關閉。 + ## 後端的原子操作 {#atomic-backend-primitives} 每個後端都在 `get`/`set`/`delete` 之上提供原子操作,供會被許多並行請求讀寫的值使用: diff --git a/tests/backends/test_base.py b/tests/backends/test_base.py index f5f7422..4ccacd6 100644 --- a/tests/backends/test_base.py +++ b/tests/backends/test_base.py @@ -121,3 +121,47 @@ async def test_expire_if_equals_fallback_updates_ttl_only_when_matching( assert await backend.expire_if_equals("slot", theirs, ttl=60) is True assert backend.store["slot"] == (theirs, 60) assert await backend.expire_if_equals("missing", theirs, ttl=60) is False + + +async def test_aclose_is_a_no_op_by_default(backend: DictBackend) -> None: + await backend.set("key", CacheEntry(fingerprint="f", content=b"v")) + + await backend.aclose() + await backend.aclose() + + assert "key" in backend.store + + +class ClosingBackend(DictBackend): + """Counts `aclose()` calls.""" + + def __init__(self) -> None: + super().__init__() + self.closed = 0 + + async def aclose(self) -> None: + self.closed += 1 + + +async def test_async_with_returns_the_backend_and_closes_it() -> None: + backend = ClosingBackend() + + async with backend as entered: + assert entered is backend + assert backend.closed == 0 + + assert backend.closed == 1 + + +async def test_async_with_closes_the_backend_when_the_body_raises() -> None: + backend = ClosingBackend() + + async def fail_inside() -> None: + async with backend: + msg = "boom" + raise RuntimeError(msg) + + with pytest.raises(RuntimeError, match="boom"): + await fail_inside() + + assert backend.closed == 1 diff --git a/tests/backends/test_memcached.py b/tests/backends/test_memcached.py index 42c426e..d062432 100644 --- a/tests/backends/test_memcached.py +++ b/tests/backends/test_memcached.py @@ -997,3 +997,43 @@ async def counting_to_thread(func, /, *a, **kw): assert hops == 1 assert len(backend.client.method_calls) > 1 + + +async def test_memcached_aclose_closes_the_client() -> None: + backend = stubbed_backend() + + async with backend as entered: + assert entered is backend + await backend.aclose() + + assert backend.client.close.call_count == 2 + + +async def test_memcached_aclose_is_safe_to_repeat_without_a_connection() -> None: + backend = MemcachedBackend([MEMCACHED_SERVER]) + + await backend.aclose() + await backend.aclose() + + +@requires_memcached +async def test_memcached_aclose_closes_every_pooled_socket() -> None: + backend = MemcachedBackend([MEMCACHED_SERVER], key_prefix="cachex_aclose_test:") + await backend.set("key", CacheEntry(fingerprint="f", content=b"v")) + await backend.delete("key") + pooled = [ + conn + for server in backend.client.clients.values() + for conn in server.client_pool.free + ] + assert pooled + assert all(conn.sock is not None for conn in pooled) + + await backend.aclose() + await backend.aclose() + + assert all(conn.sock is None for conn in pooled) + assert all( + not server.client_pool.free and not server.client_pool.used + for server in backend.client.clients.values() + ) diff --git a/tests/backends/test_memory.py b/tests/backends/test_memory.py index 197cec3..bfa9161 100644 --- a/tests/backends/test_memory.py +++ b/tests/backends/test_memory.py @@ -240,6 +240,16 @@ async def test_aclose_waits_for_the_cleanup_task(memory_backend: MemoryBackend): await memory_backend.aclose() # a second call is a no-op +async def test_async_with_stops_the_cleanup_task(): + async with MemoryBackend() as backend: + await backend.set("key", CacheEntry(fingerprint="f", content=b"v")) + task = backend._cleanup_task + assert task is not None + + assert task.done() + assert backend._cleanup_task is None + + def test_aclose_only_cancels_a_task_on_another_loop(): backend = MemoryBackend() other = asyncio.new_event_loop() diff --git a/tests/backends/test_redis.py b/tests/backends/test_redis.py index db69501..e42818b 100644 --- a/tests/backends/test_redis.py +++ b/tests/backends/test_redis.py @@ -1204,8 +1204,7 @@ async def ttl_zero() -> dict[str, str]: ) # Clean up inside the client's event loop, which owns the pool. client.portal.call(backend.clear) # type: ignore[union-attr] - # types-redis predates aclose() (redis-py 5.0.1). - client.portal.call(backend.client.aclose) # type: ignore[union-attr,call-arg,attr-defined] + client.portal.call(backend.aclose) # type: ignore[union-attr] assert first.status_code == 200 assert first.headers["Cache-Control"] == "max-age=0" @@ -1279,3 +1278,66 @@ async def test_lock_lifecycle_with_redis( assert await lock2.acquire(blocking=False) is False assert await lock1.extend(60) is True assert await lock1.release() is True + + +@requires_redis_package +async def test_redis_aclose_closes_the_client() -> None: + from unittest.mock import AsyncMock + from unittest.mock import MagicMock + + backend = AsyncRedisCacheBackend(host=REDIS_HOST, port=UNCONNECTED_PORT) + real_client = backend.client + backend.client = MagicMock() + backend.client.aclose = AsyncMock() + + async with backend as entered: + assert entered is backend + await backend.aclose() + + assert backend.client.aclose.await_count == 2 + backend.client = real_client + + +@requires_redis_package +async def test_redis_aclose_is_safe_to_repeat_without_a_connection() -> None: + backend = AsyncRedisCacheBackend(host=REDIS_HOST, port=UNCONNECTED_PORT) + + await backend.aclose() + await backend.aclose() + + +@requires_redis_package +async def test_redis_aclose_leaves_a_caller_supplied_pool_open() -> None: + from unittest.mock import AsyncMock + + from redis.asyncio import ConnectionPool + + pool = ConnectionPool(host=REDIS_HOST, port=UNCONNECTED_PORT, decode_responses=True) + backend = AsyncRedisCacheBackend(connection_pool=pool) + disconnect = AsyncMock() + try: + with pytest.MonkeyPatch.context() as mp: + mp.setattr(pool, "disconnect", disconnect) + await backend.aclose() + disconnect.assert_not_awaited() + finally: + await pool.disconnect() + + +@requires_redis +async def test_redis_aclose_disconnects_every_pooled_connection() -> None: + backend = AsyncRedisCacheBackend( + host=REDIS_HOST, port=REDIS_PORT, key_prefix="cachex_aclose_test:" + ) + await backend.set("key", CacheEntry(fingerprint="f", content=b"v")) + assert await backend.get("key") is not None + await backend.delete("key") + pool = backend.client.connection_pool + connections = list(pool._available_connections) + assert connections + assert all(conn.is_connected for conn in connections) + + await backend.aclose() + await backend.aclose() + + assert not any(conn.is_connected for conn in connections) diff --git a/tests/conftest.py b/tests/conftest.py index e987bfd..76dc170 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -52,9 +52,9 @@ async def close_network_clients( ) -> AsyncGenerator[None, None]: """Close every Redis and Memcached client a test opens, when it ends. - The backends have no close method of their own, so a client a test drops - keeps its sockets until the garbage collector finds it, often during a - later test. The `ResourceWarning` then fails whichever test that is. + A client a test drops without `aclose()` keeps its sockets until the + garbage collector finds it, often during a later test. The + `ResourceWarning` then fails whichever test that is. Recording each backend as it is built and closing it here keeps every socket inside the test that opened it. @@ -77,12 +77,9 @@ def record(cls: type[Any], *args: Any, **kwargs: Any) -> Any: yield for backend in opened: # Missing if the constructor raised; a test may swap in a mock. - client: Any = getattr(backend, "client", None) - module = type(client).__module__ - if module.startswith("redis."): - await client.aclose() - elif module.startswith("pymemcache."): - client.close() + module = type(getattr(backend, "client", None)).__module__ + if module.startswith(("redis.", "pymemcache.")): + await backend.aclose() @pytest.fixture(autouse=True)