Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9a9b7b2985 | ||
|
|
5b9573d499 | ||
|
|
50444af982 | ||
|
|
2f19fcc972 | ||
|
|
3bd65895ec | ||
|
|
5311d49fa6 |
@@ -34,6 +34,18 @@ Requires Redis with the RedisJSON module (e.g. redis-stack).
|
||||
| -------------------- | ---------------------------------------------------- |
|
||||
| `REDIS_URI` | Redis connection URL (e.g. `redis://localhost:6379`) |
|
||||
|
||||
For key discovery:
|
||||
|
||||
- `list_keys(pattern)` is simple and returns a `list[str]`, but it uses Redis `KEYS` and may block on large datasets.
|
||||
- `scan_keys(pattern, *, count=None)` is preferred for production use and yields keys incrementally via Redis `SCAN`.
|
||||
|
||||
Example:
|
||||
|
||||
```python
|
||||
for key in repo.scan_keys("user:*"):
|
||||
print(key)
|
||||
```
|
||||
|
||||
### MinIO (`ObjectRepositoryInterface`)
|
||||
|
||||
| Environment variable | Description |
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[project]
|
||||
name = "python-repositories"
|
||||
version = "0.5.1"
|
||||
version = "1.1.0"
|
||||
description = "Various python repository interfaces exposed as a python package."
|
||||
authors = [
|
||||
{ name = "Brian Bjarke Jensen", email = "[email protected]" }
|
||||
|
||||
@@ -172,15 +172,12 @@ class MinioAdapter(ObjectRepositoryInterface, ConnectionAwareAdapter):
|
||||
self.logger.warning(
|
||||
f"Object '{object_name}' not found in bucket '{self._bucket_name}'"
|
||||
)
|
||||
else:
|
||||
self.logger.error(repr(exc))
|
||||
except Exception as exc: # pylint: disable=broad-except
|
||||
self.logger.error(repr(exc))
|
||||
return None
|
||||
raise
|
||||
finally:
|
||||
if response is not None:
|
||||
response.close()
|
||||
response.release_conn()
|
||||
return None
|
||||
|
||||
def delete(self, object_name: str) -> None:
|
||||
"""Delete an object from the Minio bucket."""
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
"""Definition of RedisAdapter class."""
|
||||
|
||||
from __future__ import annotations
|
||||
from collections.abc import Iterator
|
||||
from typing import cast
|
||||
|
||||
from python_repositories.adapters.connection_aware_adapter import (
|
||||
@@ -136,11 +137,13 @@ class RedisAdapter(JsonRepositoryInterface, ConnectionAwareAdapter):
|
||||
self._client.json().delete(key)
|
||||
self.logger.debug(f"Deleted {key}")
|
||||
|
||||
def list_keys(self, pattern: str) -> list[str]:
|
||||
"""List keys in Redis matching a pattern."""
|
||||
# Check input
|
||||
def _validate_pattern(self, pattern: str) -> None:
|
||||
if not isinstance(pattern, str) or len(pattern) == 0:
|
||||
raise ValueError("Pattern must be a non-empty string")
|
||||
|
||||
def list_keys(self, pattern: str) -> list[str]:
|
||||
"""List keys in Redis using KEYS; may block on large datasets."""
|
||||
self._validate_pattern(pattern)
|
||||
# Check connection
|
||||
self._require_connected()
|
||||
assert self._client is not None
|
||||
@@ -152,3 +155,29 @@ class RedisAdapter(JsonRepositoryInterface, ConnectionAwareAdapter):
|
||||
keys: list[str] = [key.decode(self.encoding) for key in keys_raw]
|
||||
self.logger.debug(f"Got {keys} matching {pattern}")
|
||||
return keys
|
||||
|
||||
def scan_keys(
|
||||
self,
|
||||
pattern: str,
|
||||
*,
|
||||
count: int | None = None,
|
||||
) -> Iterator[str]:
|
||||
"""Yield keys in Redis using SCAN to avoid blocking large datasets."""
|
||||
self._validate_pattern(pattern)
|
||||
self._require_connected()
|
||||
assert self._client is not None
|
||||
client = self._client
|
||||
|
||||
def _decode(key: bytes | str) -> str:
|
||||
return key if isinstance(key, str) else key.decode(self.encoding)
|
||||
|
||||
def _iter() -> Iterator[str]:
|
||||
scan_iter = (
|
||||
client.scan_iter(match=pattern, count=count)
|
||||
if count is not None
|
||||
else client.scan_iter(match=pattern)
|
||||
)
|
||||
for key_raw in scan_iter:
|
||||
yield _decode(key_raw)
|
||||
|
||||
return _iter()
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
"""Definition of JsonRepositoryInterface abstract base class."""
|
||||
|
||||
from collections.abc import Iterator
|
||||
from abc import ABC, abstractmethod
|
||||
|
||||
|
||||
@@ -25,3 +26,13 @@ class JsonRepositoryInterface(ABC):
|
||||
def list_keys(self, pattern: str) -> list[str]:
|
||||
"""List keys matching a glob pattern."""
|
||||
...
|
||||
|
||||
def scan_keys(
|
||||
self,
|
||||
pattern: str,
|
||||
*,
|
||||
count: int | None = None,
|
||||
) -> Iterator[str]:
|
||||
"""Yield keys matching a glob pattern incrementally."""
|
||||
del count
|
||||
yield from self.list_keys(pattern)
|
||||
|
||||
@@ -9,7 +9,11 @@ class ObjectRepositoryInterface(ABC):
|
||||
|
||||
@abstractmethod
|
||||
def get(self, object_name: str) -> BytesIO | None:
|
||||
"""Get an object by name."""
|
||||
"""Get an object by name.
|
||||
|
||||
Returns None when the object does not exist. Raises ConnectionError when
|
||||
not connected. Other backend errors propagate to the caller.
|
||||
"""
|
||||
...
|
||||
|
||||
@abstractmethod
|
||||
|
||||
@@ -218,10 +218,8 @@ def test_should_log_warning_when_getting_nonexistent_object(
|
||||
)
|
||||
|
||||
|
||||
def test_should_log_error_when_getting_with_s3error_other_than_no_such_key(
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""Test that the MinioAdapter logs an error for unhandled S3 errors."""
|
||||
def test_should_reraise_s3error_other_than_no_such_key() -> None:
|
||||
"""Test that the MinioAdapter re-raises unhandled S3 errors."""
|
||||
mock_client = MagicMock(spec=Minio)
|
||||
mock_client.bucket_exists.return_value = True
|
||||
other_s3error = S3Error(
|
||||
@@ -236,25 +234,20 @@ def test_should_log_error_when_getting_with_s3error_other_than_no_such_key(
|
||||
)
|
||||
mock_client.get_object.side_effect = other_s3error
|
||||
adapter = MinioAdapter(config=TEST_MINIO_CONFIG, client=mock_client)
|
||||
with caplog.at_level("ERROR"):
|
||||
result = adapter.get("missing-object")
|
||||
assert result is None
|
||||
assert repr(other_s3error) in caplog.text
|
||||
with pytest.raises(S3Error) as exc_info:
|
||||
adapter.get("missing-object")
|
||||
assert exc_info.value.code == "UnhandledError"
|
||||
|
||||
|
||||
def test_should_log_error_when_getting_with_general_exception(
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
) -> None:
|
||||
"""Test that the MinioAdapter logs an error on general exceptions during get."""
|
||||
def test_should_reraise_general_exception() -> None:
|
||||
"""Test that the MinioAdapter re-raises general exceptions during get."""
|
||||
mock_client = MagicMock(spec=Minio)
|
||||
mock_client.bucket_exists.return_value = True
|
||||
general_exception = Exception("General failure")
|
||||
mock_client.get_object.side_effect = general_exception
|
||||
adapter = MinioAdapter(config=TEST_MINIO_CONFIG, client=mock_client)
|
||||
with caplog.at_level("ERROR"):
|
||||
result = adapter.get("missing-object")
|
||||
assert result is None
|
||||
assert repr(general_exception) in caplog.text
|
||||
with pytest.raises(Exception, match="General failure"):
|
||||
adapter.get("missing-object")
|
||||
|
||||
|
||||
def test_should_put_data(
|
||||
|
||||
@@ -254,5 +254,34 @@ def test_should_raise_connection_error_on_list_keys_when_not_connected(
|
||||
adapter.list_keys("some_pattern")
|
||||
|
||||
|
||||
def test_should_scan_keys(
|
||||
redis_adapter: RedisAdapter,
|
||||
) -> None:
|
||||
"""Test scanning keys matching a pattern returns correct keys."""
|
||||
redis_adapter.set("key1", {"a": 1})
|
||||
redis_adapter.set("key2", {"b": 2})
|
||||
keys = set(redis_adapter.scan_keys("key*"))
|
||||
assert keys == {"key1", "key2"}
|
||||
|
||||
|
||||
def test_should_raise_value_error_on_invalid_scan_keys_pattern(
|
||||
redis_adapter: RedisAdapter,
|
||||
) -> None:
|
||||
"""Test that the RedisAdapter raises ValueError when scanning keys with an invalid pattern."""
|
||||
invalid_patterns = ["", 123, None]
|
||||
for pattern in invalid_patterns:
|
||||
with pytest.raises(ValueError):
|
||||
list(redis_adapter.scan_keys(pattern)) # type: ignore[arg-type]
|
||||
|
||||
|
||||
def test_should_raise_connection_error_on_scan_keys_when_not_connected(
|
||||
redis_config: RedisConfig,
|
||||
) -> None:
|
||||
"""Test that the RedisAdapter raises ConnectionError when scanning keys while not connected."""
|
||||
adapter = RedisAdapter(config=redis_config)
|
||||
with pytest.raises(ConnectionError):
|
||||
list(adapter.scan_keys("some_pattern"))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
pytest.main(["-s", "-v", __file__])
|
||||
|
||||
@@ -80,3 +80,24 @@ def test_instantiation_fails_when_list_keys_not_implemented() -> None:
|
||||
|
||||
with pytest.raises(TypeError):
|
||||
_ = Incomplete() # type: ignore
|
||||
|
||||
|
||||
def test_scan_keys_defaults_to_list_keys() -> None:
|
||||
"""Test that the default scan_keys implementation delegates to list_keys."""
|
||||
|
||||
class Complete(JsonRepositoryInterface):
|
||||
def get(self, key: str) -> dict | None:
|
||||
return None
|
||||
|
||||
def set(self, key: str, data: dict) -> None:
|
||||
pass
|
||||
|
||||
def delete(self, key: str) -> None:
|
||||
pass
|
||||
|
||||
def list_keys(self, pattern: str) -> list[str]:
|
||||
return [f"{pattern}-1", f"{pattern}-2"]
|
||||
|
||||
repository = Complete()
|
||||
|
||||
assert list(repository.scan_keys("user")) == ["user-1", "user-2"]
|
||||
|
||||
@@ -5,7 +5,8 @@ from __future__ import annotations
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
import pytest
|
||||
from minio import Minio
|
||||
from minio import Minio, S3Error
|
||||
from urllib3.response import BaseHTTPResponse
|
||||
|
||||
from python_repositories.adapters.minio_adapter import MinioAdapter
|
||||
from python_repositories.interfaces import ObjectRepositoryInterface
|
||||
@@ -119,8 +120,63 @@ def test_get_closes_response_when_read_fails() -> None:
|
||||
mock_client.get_object.return_value = mock_response
|
||||
adapter = MinioAdapter(config=TEST_MINIO_CONFIG, client=mock_client)
|
||||
|
||||
result = adapter.get("some-object")
|
||||
with pytest.raises(OSError, match="connection reset"):
|
||||
adapter.get("some-object")
|
||||
|
||||
assert result is None
|
||||
mock_response.close.assert_called_once()
|
||||
mock_response.release_conn.assert_called_once()
|
||||
|
||||
|
||||
def test_get_returns_none_for_no_such_key() -> None:
|
||||
"""get() returns None when the object does not exist."""
|
||||
mock_client = MagicMock(spec=Minio)
|
||||
mock_client.bucket_exists.return_value = True
|
||||
mock_client.get_object.side_effect = S3Error(
|
||||
MagicMock(spec=BaseHTTPResponse),
|
||||
"NoSuchKey",
|
||||
"",
|
||||
"",
|
||||
"",
|
||||
"",
|
||||
bucket_name="test-bucket",
|
||||
object_name="missing-object",
|
||||
)
|
||||
adapter = MinioAdapter(config=TEST_MINIO_CONFIG, client=mock_client)
|
||||
|
||||
result = adapter.get("missing-object")
|
||||
|
||||
assert result is None
|
||||
|
||||
|
||||
def test_get_reraises_other_s3_errors() -> None:
|
||||
"""get() re-raises S3 errors other than NoSuchKey."""
|
||||
mock_client = MagicMock(spec=Minio)
|
||||
mock_client.bucket_exists.return_value = True
|
||||
other_s3error = S3Error(
|
||||
MagicMock(spec=BaseHTTPResponse),
|
||||
"AccessDenied",
|
||||
"",
|
||||
"",
|
||||
"",
|
||||
"",
|
||||
bucket_name="test-bucket",
|
||||
object_name="some-object",
|
||||
)
|
||||
mock_client.get_object.side_effect = other_s3error
|
||||
adapter = MinioAdapter(config=TEST_MINIO_CONFIG, client=mock_client)
|
||||
|
||||
with pytest.raises(S3Error) as exc_info:
|
||||
adapter.get("some-object")
|
||||
|
||||
assert exc_info.value.code == "AccessDenied"
|
||||
|
||||
|
||||
def test_get_reraises_general_exception_from_get_object() -> None:
|
||||
"""get() re-raises unexpected exceptions from get_object."""
|
||||
mock_client = MagicMock(spec=Minio)
|
||||
mock_client.bucket_exists.return_value = True
|
||||
mock_client.get_object.side_effect = Exception("General failure")
|
||||
adapter = MinioAdapter(config=TEST_MINIO_CONFIG, client=mock_client)
|
||||
|
||||
with pytest.raises(Exception, match="General failure"):
|
||||
adapter.get("some-object")
|
||||
|
||||
@@ -114,3 +114,53 @@ def test_connect_closes_existing_non_injected_client(
|
||||
|
||||
stale_client.close.assert_called_once()
|
||||
assert adapter._client is new_client
|
||||
|
||||
|
||||
def test_scan_keys_yields_decoded_keys() -> None:
|
||||
mock_client = MagicMock(spec=redis.Redis)
|
||||
mock_client.scan_iter.return_value = iter([b"key1", b"key2"])
|
||||
adapter = RedisAdapter(config=TEST_REDIS_CONFIG, client=mock_client)
|
||||
|
||||
keys = list(adapter.scan_keys("key*"))
|
||||
|
||||
assert keys == ["key1", "key2"]
|
||||
mock_client.scan_iter.assert_called_once_with(match="key*")
|
||||
mock_client.keys.assert_not_called()
|
||||
|
||||
|
||||
def test_scan_keys_forwards_count() -> None:
|
||||
mock_client = MagicMock(spec=redis.Redis)
|
||||
mock_client.scan_iter.return_value = iter([b"key1"])
|
||||
adapter = RedisAdapter(config=TEST_REDIS_CONFIG, client=mock_client)
|
||||
|
||||
keys = list(adapter.scan_keys("key*", count=50))
|
||||
|
||||
assert keys == ["key1"]
|
||||
mock_client.scan_iter.assert_called_once_with(match="key*", count=50)
|
||||
|
||||
|
||||
def test_list_keys_raises_value_error_on_invalid_pattern() -> None:
|
||||
mock_client = MagicMock(spec=redis.Redis)
|
||||
adapter = RedisAdapter(config=TEST_REDIS_CONFIG, client=mock_client)
|
||||
|
||||
invalid_patterns = ["", 123, None]
|
||||
for pattern in invalid_patterns:
|
||||
with pytest.raises(ValueError, match="Pattern must be a non-empty string"):
|
||||
adapter.list_keys(pattern) # type: ignore[arg-type]
|
||||
|
||||
|
||||
def test_scan_keys_raises_value_error_on_invalid_pattern() -> None:
|
||||
mock_client = MagicMock(spec=redis.Redis)
|
||||
adapter = RedisAdapter(config=TEST_REDIS_CONFIG, client=mock_client)
|
||||
|
||||
invalid_patterns = ["", 123, None]
|
||||
for pattern in invalid_patterns:
|
||||
with pytest.raises(ValueError, match="Pattern must be a non-empty string"):
|
||||
list(adapter.scan_keys(pattern)) # type: ignore[arg-type]
|
||||
|
||||
|
||||
def test_scan_keys_raises_connection_error_when_not_connected() -> None:
|
||||
adapter = RedisAdapter(config=TEST_REDIS_CONFIG)
|
||||
|
||||
with pytest.raises(ConnectionError):
|
||||
list(adapter.scan_keys("key*"))
|
||||
|
||||
Reference in New Issue
Block a user