Align Redis and MinIO connect() idempotency in ConnectionAwareAdapter.
PR Title Check / check-title (pull_request) Successful in 5s
Test Python Package / unit-tests (pull_request) Successful in 13s
Test Python Package / integration-tests (pull_request) Successful in 30s
Test Python Package / coverage-report (pull_request) Failing after 15s
Code Quality Pipeline / code-quality (pull_request) Successful in 53s

Move shared connect orchestration into the base adapter so both backends short-circuit when already healthy and only reconnect after a failed probe.

Co-authored-by: Cursor <[email protected]>
This commit is contained in:
Brian Bjarke Jensen
2026-07-09 15:46:39 +02:00
co-authored by Cursor
parent c2e6a96c5d
commit 67bb88fbb4
8 changed files with 103 additions and 39 deletions
+1 -1
View File
@@ -12,7 +12,7 @@ Subclass an adapter in your own repository to add domain-specific methods while
| **Adapters** | Technology-specific base classes (`RedisAdapter`, `MinioAdapter`) | | **Adapters** | Technology-specific base classes (`RedisAdapter`, `MinioAdapter`) |
| **Your project** | Subclass an adapter and add domain methods | | **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`. 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`. `connect()` is idempotent: calling it while already connected and healthy is a no-op.
## Optional dependencies ## Optional dependencies
@@ -22,6 +22,7 @@ class ConnectionAwareAdapter(ConnectionAwareInterface, ContextAwareInterface):
def __init__(self) -> None: def __init__(self) -> None:
self.logger = structlog.get_logger(self.__class__.__name__) self.logger = structlog.get_logger(self.__class__.__name__)
self._client_injected = False
self._health_check_at: float | None = None self._health_check_at: float | None = None
self._health_check_ok: bool = False self._health_check_ok: bool = False
@@ -57,6 +58,27 @@ class ConnectionAwareAdapter(ConnectionAwareInterface, ContextAwareInterface):
def _probe_connection(self) -> bool: def _probe_connection(self) -> bool:
"""Backend-specific liveness check; called only when client is ready.""" """Backend-specific liveness check; called only when client is ready."""
@abstractmethod
def _validate_injected_client(self) -> None:
"""Verify an injected client is reachable; raise ConnectionError on failure."""
@abstractmethod
def _establish_connection(self) -> None:
"""Create a backend client and set internal connection state."""
def connect(self) -> None:
"""Connect to the backend; idempotent when already connected and healthy."""
if self._client_injected:
self._validate_injected_client()
self._invalidate_health_cache()
return
if self._is_client_ready() and self.is_connected():
self.logger.info(f"Already connected to {self.connection_name}")
return
self.disconnect()
self._establish_connection()
self._invalidate_health_cache()
def is_connected(self) -> bool: def is_connected(self) -> bool:
"""Check if connected to the backend.""" """Check if connected to the backend."""
if not self._is_client_ready(): if not self._is_client_ready():
+5 -12
View File
@@ -58,23 +58,17 @@ class MinioAdapter(ObjectRepositoryInterface, ConnectionAwareAdapter):
def _is_client_ready(self) -> bool: def _is_client_ready(self) -> bool:
return self._client is not None and self._bucket_name is not None return self._client is not None and self._bucket_name is not None
def connect(self) -> None: def _validate_injected_client(self) -> None:
"""Connect to the Minio server.""" if self._client is None:
if self._client_injected: return
if self._client is not None:
try: try:
_ = self._client.list_buckets() _ = self._client.list_buckets()
except Exception as exc: # pylint: disable=broad-except except Exception as exc: # pylint: disable=broad-except
raise ConnectionError( raise ConnectionError(
f"Could not connect to Minio at {self._config.endpoint}" f"Could not connect to Minio at {self._config.endpoint}"
) from exc ) from exc
self._invalidate_health_cache()
return def _establish_connection(self) -> None:
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()
endpoint = self._config.endpoint endpoint = self._config.endpoint
access_key = self._config.access_key access_key = self._config.access_key
secret_key = self._config.secret_key secret_key = self._config.secret_key
@@ -98,7 +92,6 @@ class MinioAdapter(ObjectRepositoryInterface, ConnectionAwareAdapter):
client.make_bucket(bucket) client.make_bucket(bucket)
self._client = client self._client = client
self._bucket_name = bucket self._bucket_name = bucket
self._invalidate_health_cache()
def disconnect(self) -> None: def disconnect(self) -> None:
"""Disconnect from the Minio server.""" """Disconnect from the Minio server."""
+5 -11
View File
@@ -49,10 +49,9 @@ class RedisAdapter(JsonRepositoryInterface, ConnectionAwareAdapter):
def _is_client_ready(self) -> bool: def _is_client_ready(self) -> bool:
return self._client is not None return self._client is not None
def connect(self) -> None: def _validate_injected_client(self) -> None:
"""Connect to the Redis server.""" if self._client is None:
if self._client_injected: return
if self._client is not None:
try: try:
if not self._client.ping(): if not self._client.ping():
raise ConnectionError( raise ConnectionError(
@@ -62,12 +61,8 @@ class RedisAdapter(JsonRepositoryInterface, ConnectionAwareAdapter):
raise ConnectionError( raise ConnectionError(
f"Could not connect to Redis at {self._config.uri}" f"Could not connect to Redis at {self._config.uri}"
) from exc ) from exc
self._invalidate_health_cache()
return def _establish_connection(self) -> None:
if self._client is not None:
self._client.close()
self._client = None
self._invalidate_health_cache()
uri = self._config.uri uri = self._config.uri
try: try:
client = redis.Redis.from_url( client = redis.Redis.from_url(
@@ -79,7 +74,6 @@ class RedisAdapter(JsonRepositoryInterface, ConnectionAwareAdapter):
except (redis.ConnectionError, redis.TimeoutError) as exc: except (redis.ConnectionError, redis.TimeoutError) as exc:
raise ConnectionError(f"Could not connect to Redis at {uri}") from exc raise ConnectionError(f"Could not connect to Redis at {uri}") from exc
self._client = client self._client = client
self._invalidate_health_cache()
def disconnect(self) -> None: def disconnect(self) -> None:
"""Disconnect from the Redis server.""" """Disconnect from the Redis server."""
@@ -8,7 +8,11 @@ class ConnectionAwareInterface(ABC):
@abstractmethod @abstractmethod
def connect(self) -> None: def connect(self) -> None:
"""Connect to resource.""" """Connect to resource.
Implementations should be idempotent: calling connect while already
connected and healthy is a no-op.
"""
... ...
@abstractmethod @abstractmethod
+11
View File
@@ -1,6 +1,7 @@
"""Integration tests for the RedisAdapter.""" """Integration tests for the RedisAdapter."""
from collections.abc import Generator from collections.abc import Generator
import logging
import pytest import pytest
import redis import redis
@@ -48,6 +49,16 @@ def clear_redis(raw_redis_client: redis.Redis) -> None:
raw_redis_client.flushall() raw_redis_client.flushall()
def test_should_log_info_when_already_connected(
redis_adapter: RedisAdapter,
caplog: pytest.LogCaptureFixture,
) -> None:
"""Test that the RedisAdapter logs info when connect is called while already connected."""
with caplog.at_level(logging.INFO):
redis_adapter.connect()
assert "Already connected to Redis" in caplog.text
def test_should_raise_connection_error_when_unable_to_connect() -> None: def test_should_raise_connection_error_when_unable_to_connect() -> None:
"""Test that the RedisAdapter raises ConnectionError when unable to connect.""" """Test that the RedisAdapter raises ConnectionError when unable to connect."""
adapter = RedisAdapter(config=RedisConfig(uri="redis://invalid:6379")) adapter = RedisAdapter(config=RedisConfig(uri="redis://invalid:6379"))
+21
View File
@@ -95,6 +95,27 @@ def test_connect_disconnects_before_reconnect(
new_client.list_buckets.assert_called_once() new_client.list_buckets.assert_called_once()
def test_connect_skips_reconnect_when_already_connected(
monkeypatch: pytest.MonkeyPatch,
) -> None:
stale_client = MagicMock(spec=Minio)
stale_client.bucket_exists.return_value = True
adapter = MinioAdapter(config=TEST_MINIO_CONFIG)
adapter._client = stale_client
adapter._bucket_name = TEST_MINIO_CONFIG.bucket
minio_ctor = MagicMock()
monkeypatch.setattr(
"python_repositories.adapters.minio_adapter.minio.Minio",
minio_ctor,
)
adapter.connect()
minio_ctor.assert_not_called()
assert adapter._client is stale_client
def test_connect_raises_when_bucket_missing( def test_connect_raises_when_bucket_missing(
monkeypatch: pytest.MonkeyPatch, monkeypatch: pytest.MonkeyPatch,
) -> None: ) -> None:
+20 -1
View File
@@ -99,10 +99,11 @@ def test_connect_with_injected_client_raises_on_redis_error() -> None:
adapter.connect() adapter.connect()
def test_connect_closes_existing_non_injected_client( def test_connect_reconnects_when_existing_client_unhealthy(
monkeypatch: pytest.MonkeyPatch, monkeypatch: pytest.MonkeyPatch,
) -> None: ) -> None:
stale_client = MagicMock(spec=redis.Redis) stale_client = MagicMock(spec=redis.Redis)
stale_client.ping.side_effect = redis.ConnectionError("connection lost")
new_client = MagicMock(spec=redis.Redis) new_client = MagicMock(spec=redis.Redis)
new_client.ping.return_value = True new_client.ping.return_value = True
adapter = RedisAdapter(config=TEST_REDIS_CONFIG) adapter = RedisAdapter(config=TEST_REDIS_CONFIG)
@@ -116,6 +117,24 @@ def test_connect_closes_existing_non_injected_client(
assert adapter._client is new_client assert adapter._client is new_client
def test_connect_skips_reconnect_when_already_connected(
monkeypatch: pytest.MonkeyPatch,
) -> None:
stale_client = MagicMock(spec=redis.Redis)
stale_client.ping.return_value = True
adapter = RedisAdapter(config=TEST_REDIS_CONFIG)
adapter._client = stale_client
from_url = MagicMock()
monkeypatch.setattr("redis.Redis.from_url", from_url)
adapter.connect()
stale_client.close.assert_not_called()
from_url.assert_not_called()
assert adapter._client is stale_client
def test_scan_keys_yields_decoded_keys() -> None: def test_scan_keys_yields_decoded_keys() -> None:
mock_client = MagicMock(spec=redis.Redis) mock_client = MagicMock(spec=redis.Redis)
mock_client.scan_iter.return_value = iter([b"key1", b"key2"]) mock_client.scan_iter.return_value = iter([b"key1", b"key2"])