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
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:
co-authored by
Cursor
parent
c2e6a96c5d
commit
67bb88fbb4
@@ -22,6 +22,7 @@ class ConnectionAwareAdapter(ConnectionAwareInterface, ContextAwareInterface):
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.logger = structlog.get_logger(self.__class__.__name__)
|
||||
self._client_injected = False
|
||||
self._health_check_at: float | None = None
|
||||
self._health_check_ok: bool = False
|
||||
|
||||
@@ -57,6 +58,27 @@ class ConnectionAwareAdapter(ConnectionAwareInterface, ContextAwareInterface):
|
||||
def _probe_connection(self) -> bool:
|
||||
"""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:
|
||||
"""Check if connected to the backend."""
|
||||
if not self._is_client_ready():
|
||||
|
||||
@@ -58,23 +58,17 @@ class MinioAdapter(ObjectRepositoryInterface, ConnectionAwareAdapter):
|
||||
def _is_client_ready(self) -> bool:
|
||||
return self._client is not None and self._bucket_name is not None
|
||||
|
||||
def connect(self) -> None:
|
||||
"""Connect to the Minio server."""
|
||||
if self._client_injected:
|
||||
if self._client is not None:
|
||||
try:
|
||||
_ = self._client.list_buckets()
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
raise ConnectionError(
|
||||
f"Could not connect to Minio at {self._config.endpoint}"
|
||||
) from exc
|
||||
self._invalidate_health_cache()
|
||||
def _validate_injected_client(self) -> None:
|
||||
if self._client is None:
|
||||
return
|
||||
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()
|
||||
try:
|
||||
_ = self._client.list_buckets()
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
raise ConnectionError(
|
||||
f"Could not connect to Minio at {self._config.endpoint}"
|
||||
) from exc
|
||||
|
||||
def _establish_connection(self) -> None:
|
||||
endpoint = self._config.endpoint
|
||||
access_key = self._config.access_key
|
||||
secret_key = self._config.secret_key
|
||||
@@ -98,7 +92,6 @@ class MinioAdapter(ObjectRepositoryInterface, ConnectionAwareAdapter):
|
||||
client.make_bucket(bucket)
|
||||
self._client = client
|
||||
self._bucket_name = bucket
|
||||
self._invalidate_health_cache()
|
||||
|
||||
def disconnect(self) -> None:
|
||||
"""Disconnect from the Minio server."""
|
||||
|
||||
@@ -49,25 +49,20 @@ class RedisAdapter(JsonRepositoryInterface, ConnectionAwareAdapter):
|
||||
def _is_client_ready(self) -> bool:
|
||||
return self._client is not None
|
||||
|
||||
def connect(self) -> None:
|
||||
"""Connect to the Redis server."""
|
||||
if self._client_injected:
|
||||
if self._client is not None:
|
||||
try:
|
||||
if not self._client.ping():
|
||||
raise ConnectionError(
|
||||
f"Could not connect to Redis at {self._config.uri}"
|
||||
)
|
||||
except (redis.ConnectionError, redis.TimeoutError) as exc:
|
||||
raise ConnectionError(
|
||||
f"Could not connect to Redis at {self._config.uri}"
|
||||
) from exc
|
||||
self._invalidate_health_cache()
|
||||
def _validate_injected_client(self) -> None:
|
||||
if self._client is None:
|
||||
return
|
||||
if self._client is not None:
|
||||
self._client.close()
|
||||
self._client = None
|
||||
self._invalidate_health_cache()
|
||||
try:
|
||||
if not self._client.ping():
|
||||
raise ConnectionError(
|
||||
f"Could not connect to Redis at {self._config.uri}"
|
||||
)
|
||||
except (redis.ConnectionError, redis.TimeoutError) as exc:
|
||||
raise ConnectionError(
|
||||
f"Could not connect to Redis at {self._config.uri}"
|
||||
) from exc
|
||||
|
||||
def _establish_connection(self) -> None:
|
||||
uri = self._config.uri
|
||||
try:
|
||||
client = redis.Redis.from_url(
|
||||
@@ -79,7 +74,6 @@ class RedisAdapter(JsonRepositoryInterface, ConnectionAwareAdapter):
|
||||
except (redis.ConnectionError, redis.TimeoutError) as exc:
|
||||
raise ConnectionError(f"Could not connect to Redis at {uri}") from exc
|
||||
self._client = client
|
||||
self._invalidate_health_cache()
|
||||
|
||||
def disconnect(self) -> None:
|
||||
"""Disconnect from the Redis server."""
|
||||
|
||||
Reference in New Issue
Block a user