From 5149299c1b6e84eae897798b3761ebef1dff91e7 Mon Sep 17 00:00:00 2001 From: Brian Bjarke Jensen Date: Sun, 5 Jul 2026 15:28:25 +0200 Subject: [PATCH 1/2] Strengthen is_connected with cached health probes. Convert is_connected to a method that verifies backend liveness via TTL-cached ping (Redis) or bucket_exists (MinIO), with cache invalidation on connect/disconnect. Co-authored-by: Cursor --- README.md | 2 + python_repositories/adapters/minio_adapter.py | 57 ++++-- python_repositories/adapters/redis_adapter.py | 54 ++++- .../interfaces/connection_aware_interface.py | 6 +- .../connection_aware_interface_test.py | 2 - tests/integration/minio_adapter_test.py | 6 +- tests/integration/redis_adapter_test.py | 12 +- tests/unit/connection_health_test.py | 184 ++++++++++++++++++ 8 files changed, 289 insertions(+), 34 deletions(-) create mode 100644 tests/unit/connection_health_test.py diff --git a/README.md b/README.md index 5dcfb34..a67d668 100644 --- a/README.md +++ b/README.md @@ -12,6 +12,8 @@ Subclass an adapter in your own repository to add domain-specific methods while | **Adapters** | Technology-specific base classes (`RedisAdapter`, `MinioAdapter`) | | **Your project** | Subclass an adapter and add domain methods | +Connection adapters expose `connect()`, `disconnect()`, and `is_connected()`. The latter verifies backend reachability with a cached health probe (default TTL: 1 second). Subclasses may override `health_check_ttl_seconds`. + ## Optional dependencies Repository **interfaces** import with the base package. **Adapters** require the matching extra; importing an adapter without its extra raises `ImportError` with install instructions. diff --git a/python_repositories/adapters/minio_adapter.py b/python_repositories/adapters/minio_adapter.py index e659c13..579c2be 100644 --- a/python_repositories/adapters/minio_adapter.py +++ b/python_repositories/adapters/minio_adapter.py @@ -2,9 +2,10 @@ from __future__ import annotations import os +import structlog +import time from io import BytesIO from typing import Self -import structlog from python_utils import check_env from python_repositories.interfaces import ( @@ -34,6 +35,7 @@ class MinioAdapter( secret_key_env_var_name: str = "MINIO_SECRET_KEY" bucket_env_var_name: str = "MINIO_BUCKET" chunk_size: int = 5 * 2**20 # 5 MiB + health_check_ttl_seconds: float = 1.0 def __init__(self) -> None: # Setup logger @@ -52,6 +54,8 @@ class MinioAdapter( # Prepare internal variables self._client: minio.Minio | None = None self._bucket_name: str | None = None + self._health_check_at: float | None = None + self._health_check_ok: bool = False def __enter__(self) -> Self: """Enter the context.""" @@ -75,10 +79,11 @@ class MinioAdapter( def connect(self) -> None: """Connect to the Minio server.""" - # Stop if already connected - if self.is_connected: + if self._client is not None and self.is_connected(): self.logger.info("Already connected to Minio") return + if self._client is not None: + self.disconnect() # Prepare arguments endpoint = str(os.getenv(self.endpoint_env_var_name)) access_key = str(os.getenv(self.access_key_env_var_name)) @@ -103,6 +108,7 @@ class MinioAdapter( # Persist information self._client = client self._bucket_name = bucket + self._invalidate_health_cache() def disconnect(self) -> None: """Disconnect from the Minio server.""" @@ -111,13 +117,36 @@ class MinioAdapter( # Reset client self._client = None self._bucket_name = None + self._invalidate_health_cache() + + def _invalidate_health_cache(self) -> None: + self._health_check_at = None + self._health_check_ok = False + + def _probe_connection(self) -> bool: + assert self._client is not None and self._bucket_name is not None + try: + return bool(self._client.bucket_exists(self._bucket_name)) + except Exception: # pylint: disable=broad-except + return False - @property def is_connected(self) -> bool: - """Check if connected to Minio server.""" - res = self._client is not None - self.logger.debug(res) - return res + """Check if connected to the Minio server.""" + if self._client is None or self._bucket_name is None: + return False + now = time.monotonic() + if self._health_check_at is not None: + seconds_since_last_health_check = now - self._health_check_at + cache_is_fresh = ( + seconds_since_last_health_check < self.health_check_ttl_seconds + ) + if cache_is_fresh: + return self._health_check_ok + result = self._probe_connection() + self._health_check_at = now + self._health_check_ok = result + self.logger.debug("Connection status", connected=result) + return result def put( self, @@ -134,8 +163,9 @@ class MinioAdapter( if not isinstance(content_type, str) or len(content_type) == 0: raise ValueError("content_type must be a non-empty string") # Check connection - if self._client is None or self._bucket_name is None or not self.is_connected: + if not self.is_connected(): raise ConnectionError("Not connected to Minio") + assert self._client is not None and self._bucket_name is not None # Prepare buffer for reading num_bytes = data.getbuffer().nbytes data.seek(0) @@ -159,8 +189,9 @@ class MinioAdapter( if not isinstance(object_name, str) or len(object_name) == 0: raise ValueError("object_name must be a non-empty string") # Check connection - if self._client is None or self._bucket_name is None or not self.is_connected: + if not self.is_connected(): raise ConnectionError("Not connected to Minio") + assert self._client is not None and self._bucket_name is not None # Get data from bucket # N.B. bucket name is set when connecting try: @@ -194,8 +225,9 @@ class MinioAdapter( if not isinstance(object_name, str) or len(object_name) == 0: raise ValueError("object_name must be a non-empty string") # Check connection - if self._client is None or self._bucket_name is None or not self.is_connected: + if not self.is_connected(): raise ConnectionError("Not connected to Minio") + assert self._client is not None and self._bucket_name is not None # Delete object from bucket # N.B. bucket name is set when connecting self._client.remove_object( @@ -213,8 +245,9 @@ class MinioAdapter( raise ValueError("prefix must be a string") # Check connection # N.B. bucket name is set when connecting - if self._client is None or self._bucket_name is None or not self.is_connected: + if not self.is_connected(): raise ConnectionError("Not connected to Minio") + assert self._client is not None and self._bucket_name is not None # List objects in bucket objects = self._client.list_objects( bucket_name=self._bucket_name, diff --git a/python_repositories/adapters/redis_adapter.py b/python_repositories/adapters/redis_adapter.py index 3dc5daa..4af6db7 100644 --- a/python_repositories/adapters/redis_adapter.py +++ b/python_repositories/adapters/redis_adapter.py @@ -4,6 +4,7 @@ from __future__ import annotations from typing import Self, cast import os import structlog +import time from python_utils import check_env @@ -33,6 +34,7 @@ class RedisAdapter( uri_env_var_name: str = "REDIS_URI" path: str = "." # JSON root path, updated in __init__ encoding: str = "UTF-8" + health_check_ttl_seconds: float = 1.0 def __init__(self) -> None: # Setup logger @@ -44,6 +46,8 @@ class RedisAdapter( # Prepare internal variables self._client: redis.Redis | None = None self.path: str = RedisPath.root_path() + self._health_check_at: float | None = None + self._health_check_ok: bool = False def __enter__(self) -> Self: """Enter the context.""" @@ -67,6 +71,10 @@ class RedisAdapter( def connect(self) -> None: """Connect to the Redis server.""" + if self._client is not None: + self._client.close() + self._client = None + self._invalidate_health_cache() # Prepare arguments uri = str(os.getenv(self.uri_env_var_name)) # Connect client @@ -81,6 +89,7 @@ class RedisAdapter( raise ConnectionError(f"Could not connect to Redis at {uri}") from exc # Persist client self._client = client + self._invalidate_health_cache() def disconnect(self) -> None: """Disconnect from the Redis server.""" @@ -89,13 +98,36 @@ class RedisAdapter( self._client.close() # Reset client self._client = None + self._invalidate_health_cache() + + def _invalidate_health_cache(self) -> None: + self._health_check_at = None + self._health_check_ok = False + + def _probe_connection(self) -> bool: + assert self._client is not None + try: + return bool(self._client.ping()) + except (redis.ConnectionError, redis.TimeoutError): + return False - @property def is_connected(self) -> bool: - """Check if connected to Redis server.""" - res = self._client is not None - self.logger.debug(res) - return res + """Check if connected to the Redis server.""" + if self._client is None: + return False + now = time.monotonic() + if self._health_check_at is not None: + seconds_since_last_health_check = now - self._health_check_at + cache_is_fresh = ( + seconds_since_last_health_check < self.health_check_ttl_seconds + ) + if cache_is_fresh: + return self._health_check_ok + result = self._probe_connection() + self._health_check_at = now + self._health_check_ok = result + self.logger.debug("Connection status", connected=result) + return result def set(self, key: str, data: dict) -> None: """Set a JSON object in Redis.""" @@ -105,8 +137,9 @@ class RedisAdapter( if not isinstance(data, dict) or len(data) == 0: raise ValueError("Data must be a non-empty dictionary") # Check connection - if self._client is None or not self.is_connected: + if not self.is_connected(): raise ConnectionError("Not connected to Redis") + assert self._client is not None # Set data self._client.json().set(key, self.path, data) self.logger.debug(f"Set {key} to {data}") @@ -117,8 +150,9 @@ class RedisAdapter( if not isinstance(key, str) or len(key) == 0: raise ValueError("Key must be a non-empty string") # Check connection - if self._client is None or not self.is_connected: + if not self.is_connected(): raise ConnectionError("Not connected to Redis") + assert self._client is not None # Get data data = cast( dict | None, @@ -133,8 +167,9 @@ class RedisAdapter( if not isinstance(key, str) or len(key) == 0: raise ValueError("Key must be a non-empty string") # Check connection - if self._client is None or not self.is_connected: + if not self.is_connected(): raise ConnectionError("Not connected to Redis") + assert self._client is not None # Delete data self._client.json().delete(key) self.logger.debug(f"Deleted {key}") @@ -145,8 +180,9 @@ class RedisAdapter( if not isinstance(pattern, str) or len(pattern) == 0: raise ValueError("Pattern must be a non-empty string") # Check connection - if self._client is None or not self.is_connected: + if not self.is_connected(): raise ConnectionError("Not connected to Redis") + assert self._client is not None # List keys keys_raw = cast( list[bytes], diff --git a/python_repositories/interfaces/connection_aware_interface.py b/python_repositories/interfaces/connection_aware_interface.py index a1d1cbb..fd8d1c3 100644 --- a/python_repositories/interfaces/connection_aware_interface.py +++ b/python_repositories/interfaces/connection_aware_interface.py @@ -16,8 +16,10 @@ class ConnectionAwareInterface(ABC): """Disconnect from resource.""" ... - @property @abstractmethod def is_connected(self) -> bool: - """Check if connected to resource.""" + """Return whether the adapter has an active, reachable connection. + + Implementations may perform a cached network probe to verify liveness. + """ ... diff --git a/tests/integration/connection_aware_interface_test.py b/tests/integration/connection_aware_interface_test.py index d7d2e72..f4d3414 100644 --- a/tests/integration/connection_aware_interface_test.py +++ b/tests/integration/connection_aware_interface_test.py @@ -15,7 +15,6 @@ def test_instantiation_fails_when_connect_not_implemented() -> None: def disconnect(self) -> None: pass - @property def is_connected(self) -> bool: return False @@ -32,7 +31,6 @@ def test_instantiation_fails_when_disconnect_not_implemented() -> None: def connect(self) -> None: pass - @property def is_connected(self) -> bool: return False diff --git a/tests/integration/minio_adapter_test.py b/tests/integration/minio_adapter_test.py index 3d8586f..e79e157 100644 --- a/tests/integration/minio_adapter_test.py +++ b/tests/integration/minio_adapter_test.py @@ -115,7 +115,7 @@ def test_should_have_logger_when_instantiated() -> None: def test_should_not_be_connected_when_instantiated() -> None: """Test that the MinioAdapter is not connected when instantiated.""" adapter = MinioAdapter() - assert not adapter.is_connected + assert not adapter.is_connected() def test_should_log_info_when_already_connected( @@ -137,7 +137,7 @@ def test_should_raise_connection_error_when_unable_to_connect( adapter = MinioAdapter() with pytest.raises(ConnectionError): adapter.connect() - assert not adapter.is_connected + assert not adapter.is_connected() def test_should_log_info_when_creating_expected_bucket( @@ -163,7 +163,7 @@ def test_should_log_error_on_exception_during_exit( """Test that the MinioAdapter logs an error if an exception occurs during __exit__.""" try: with MinioAdapter() as adapter: - assert adapter.is_connected + assert adapter.is_connected() raise ValueError("Simulated error") except ValueError: pass # Expected diff --git a/tests/integration/redis_adapter_test.py b/tests/integration/redis_adapter_test.py index 95ab9bb..7cfe8fd 100644 --- a/tests/integration/redis_adapter_test.py +++ b/tests/integration/redis_adapter_test.py @@ -9,7 +9,7 @@ from python_repositories.interfaces import JsonRepositoryInterface @pytest.fixture(scope="module") -def data() -> Generator[dict[str, str]]: +def data() -> Generator[dict[str, str], None, None]: """Provide a sample data dictionary for tests.""" yield {"foo": "bar"} @@ -18,7 +18,7 @@ def data() -> Generator[dict[str, str]]: def data_in_redis( raw_redis_client: redis.Redis, data: dict[str, str], -) -> Generator[tuple[str, dict[str, str]]]: +) -> Generator[tuple[str, dict[str, str]], None, None]: """Fixture to set up a known value in Redis before each test.""" key = "test_key" path = RedisPath.root_path() @@ -31,7 +31,7 @@ def data_in_redis( @pytest.fixture(scope="module") -def redis_adapter(redis_container: str) -> Generator[RedisAdapter]: +def redis_adapter(redis_container: str) -> Generator[RedisAdapter, None, None]: """Fixture to provide a connected RedisAdapter instance.""" adapter = RedisAdapter() adapter.connect() @@ -63,7 +63,7 @@ def test_should_not_be_connected_when_instantiated(redis_container: str) -> None """Test that the RedisAdapter is not connected when instantiated.""" adapter = RedisAdapter() assert adapter._client is None - assert not adapter.is_connected + assert not adapter.is_connected() def test_should_raise_connection_error_when_unable_to_connect( @@ -77,7 +77,7 @@ def test_should_raise_connection_error_when_unable_to_connect( with pytest.raises(ConnectionError): adapter.connect() assert adapter._client is None - assert not adapter.is_connected + assert not adapter.is_connected() def test_connect_raises_connection_error_when_unable_to_ping( @@ -109,7 +109,7 @@ def test_should_log_error_on_exception_during_exit( """Test that the RedisAdapter logs an error when an exception occurs during context exit.""" try: with RedisAdapter() as adapter: - assert adapter.is_connected + assert adapter.is_connected() raise ValueError("Simulated error") except ValueError: pass # Expected diff --git a/tests/unit/connection_health_test.py b/tests/unit/connection_health_test.py new file mode 100644 index 0000000..1dbccbb --- /dev/null +++ b/tests/unit/connection_health_test.py @@ -0,0 +1,184 @@ +"""Tests for TTL-cached connection health checks on adapters.""" + +from __future__ import annotations + +from unittest.mock import MagicMock, patch + +import pytest +import redis + +from python_repositories.adapters.minio_adapter import MinioAdapter +from python_repositories.adapters.redis_adapter import RedisAdapter + + +@pytest.fixture +def redis_adapter(monkeypatch: pytest.MonkeyPatch) -> RedisAdapter: + monkeypatch.setenv("REDIS_URI", "redis://localhost:6379") + return RedisAdapter() + + +@pytest.fixture +def minio_adapter(monkeypatch: pytest.MonkeyPatch) -> MinioAdapter: + monkeypatch.setenv("MINIO_ENDPOINT", "localhost:9000") + monkeypatch.setenv("MINIO_ACCESS_KEY", "minioadmin") + monkeypatch.setenv("MINIO_SECRET_KEY", "minioadmin") + monkeypatch.setenv("MINIO_BUCKET", "test-bucket") + return MinioAdapter() + + +class TestRedisConnectionHealth: + def test_not_connected_when_no_client(self, redis_adapter: RedisAdapter) -> None: + assert not redis_adapter.is_connected() + + def test_connected_when_probe_succeeds(self, redis_adapter: RedisAdapter) -> None: + mock_client = MagicMock(spec=redis.Redis) + mock_client.ping.return_value = True + redis_adapter._client = mock_client + + assert redis_adapter.is_connected() + mock_client.ping.assert_called_once() + + def test_stale_connection_when_probe_fails( + self, redis_adapter: RedisAdapter + ) -> None: + mock_client = MagicMock(spec=redis.Redis) + mock_client.ping.side_effect = redis.ConnectionError("connection lost") + redis_adapter._client = mock_client + + assert not redis_adapter.is_connected() + + def test_cache_hit_avoids_second_probe(self, redis_adapter: RedisAdapter) -> None: + mock_client = MagicMock(spec=redis.Redis) + mock_client.ping.return_value = True + redis_adapter._client = mock_client + + with patch( + "python_repositories.adapters.redis_adapter.time.monotonic", + return_value=100.0, + ): + assert redis_adapter.is_connected() + assert redis_adapter.is_connected() + + mock_client.ping.assert_called_once() + + def test_cache_miss_runs_probe_again(self, redis_adapter: RedisAdapter) -> None: + mock_client = MagicMock(spec=redis.Redis) + mock_client.ping.return_value = True + redis_adapter._client = mock_client + + with patch( + "python_repositories.adapters.redis_adapter.time.monotonic", + side_effect=[100.0, 102.0], + ): + assert redis_adapter.is_connected() + assert redis_adapter.is_connected() + + assert mock_client.ping.call_count == 2 + + def test_disconnect_clears_cache(self, redis_adapter: RedisAdapter) -> None: + mock_client = MagicMock(spec=redis.Redis) + mock_client.ping.return_value = True + redis_adapter._client = mock_client + + with patch( + "python_repositories.adapters.redis_adapter.time.monotonic", + return_value=100.0, + ): + assert redis_adapter.is_connected() + + redis_adapter.disconnect() + redis_adapter._client = mock_client + + with patch( + "python_repositories.adapters.redis_adapter.time.monotonic", + return_value=100.0, + ): + assert redis_adapter.is_connected() + + assert mock_client.ping.call_count == 2 + + +class TestMinioConnectionHealth: + def test_not_connected_when_no_client(self, minio_adapter: MinioAdapter) -> None: + assert not minio_adapter.is_connected() + + def test_not_connected_when_bucket_name_missing( + self, minio_adapter: MinioAdapter + ) -> None: + minio_adapter._client = MagicMock() + minio_adapter._bucket_name = None + + assert not minio_adapter.is_connected() + + def test_connected_when_probe_succeeds(self, minio_adapter: MinioAdapter) -> None: + mock_client = MagicMock() + mock_client.bucket_exists.return_value = True + minio_adapter._client = mock_client + minio_adapter._bucket_name = "test-bucket" + + assert minio_adapter.is_connected() + mock_client.bucket_exists.assert_called_once_with("test-bucket") + + def test_stale_connection_when_probe_fails( + self, minio_adapter: MinioAdapter + ) -> None: + mock_client = MagicMock() + mock_client.bucket_exists.side_effect = Exception("connection lost") + minio_adapter._client = mock_client + minio_adapter._bucket_name = "test-bucket" + + assert not minio_adapter.is_connected() + + def test_cache_hit_avoids_second_probe(self, minio_adapter: MinioAdapter) -> None: + mock_client = MagicMock() + mock_client.bucket_exists.return_value = True + minio_adapter._client = mock_client + minio_adapter._bucket_name = "test-bucket" + + with patch( + "python_repositories.adapters.minio_adapter.time.monotonic", + return_value=100.0, + ): + assert minio_adapter.is_connected() + assert minio_adapter.is_connected() + + mock_client.bucket_exists.assert_called_once() + + def test_cache_miss_runs_probe_again(self, minio_adapter: MinioAdapter) -> None: + mock_client = MagicMock() + mock_client.bucket_exists.return_value = True + minio_adapter._client = mock_client + minio_adapter._bucket_name = "test-bucket" + + with patch( + "python_repositories.adapters.minio_adapter.time.monotonic", + side_effect=[100.0, 102.0], + ): + assert minio_adapter.is_connected() + assert minio_adapter.is_connected() + + assert mock_client.bucket_exists.call_count == 2 + + def test_disconnect_clears_cache(self, minio_adapter: MinioAdapter) -> None: + mock_client = MagicMock() + mock_client.bucket_exists.return_value = True + minio_adapter._client = mock_client + minio_adapter._bucket_name = "test-bucket" + + with patch( + "python_repositories.adapters.minio_adapter.time.monotonic", + return_value=100.0, + ): + assert minio_adapter.is_connected() + + minio_adapter.disconnect() + minio_adapter._client = mock_client + minio_adapter._bucket_name = "test-bucket" + + with patch( + "python_repositories.adapters.minio_adapter.time.monotonic", + return_value=100.0, + ): + assert minio_adapter.is_connected() + + assert mock_client.bucket_exists.call_count == 2 From 34fc0468686212145f85855498bb3ac2fa99d2c2 Mon Sep 17 00:00:00 2001 From: Brian Bjarke Jensen Date: Sun, 5 Jul 2026 15:38:49 +0200 Subject: [PATCH 2/2] Sync uv.lock with pyproject.toml version 0.3.2. Co-authored-by: Cursor --- uv.lock | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/uv.lock b/uv.lock index d7f05c6..bcb1909 100644 --- a/uv.lock +++ b/uv.lock @@ -1053,7 +1053,7 @@ wheels = [ [[package]] name = "python-repositories" -version = "0.3.1" +version = "0.3.2" source = { editable = "." } dependencies = [ { name = "python-utils" },