diff --git a/python_repositories/adapters/file_backed_queue_adapter.py b/python_repositories/adapters/file_backed_queue_adapter.py index e20abfd..2aa5f09 100644 --- a/python_repositories/adapters/file_backed_queue_adapter.py +++ b/python_repositories/adapters/file_backed_queue_adapter.py @@ -11,22 +11,13 @@ from typing import Any, Self, TextIO import structlog -from python_repositories.adapters.memory_queue_adapter import ( - MemoryQueueAdapter, - _parse_age, -) +from python_repositories.adapters.memory_queue_adapter import MemoryQueueAdapter from python_repositories.config.file_queue_config import FileQueueConfig from python_repositories.interfaces.queue_repository_interface import ( QueueRepositoryInterface, ) -def _json_default(value: object) -> str: - if isinstance(value, datetime): - return value.isoformat() - raise TypeError(f"Object of type {type(value)!r} is not JSON serializable") - - class FileBackedQueueAdapter(QueueRepositoryInterface): """Queue that mirrors an in-memory buffer to an append-friendly JSONL file. @@ -54,6 +45,12 @@ class FileBackedQueueAdapter(QueueRepositoryInterface): self._append_handle: TextIO | None = None self.logger = structlog.get_logger(self.__class__.__name__) + @staticmethod + def _json_default(value: object) -> str: + if isinstance(value, datetime): + return value.isoformat() + raise TypeError(f"Object of type {type(value)!r} is not JSON serializable") + def __enter__(self) -> Self: self.connect() return self @@ -170,7 +167,7 @@ class FileBackedQueueAdapter(QueueRepositoryInterface): ) continue try: - if _parse_age(age_value) < cutoff: + if MemoryQueueAdapter._parse_age(age_value) < cutoff: continue except ValueError: self.logger.warning( @@ -194,7 +191,7 @@ class FileBackedQueueAdapter(QueueRepositoryInterface): raise RuntimeError("append handle is not open") for item in items: self._append_handle.write( - json.dumps(item, default=_json_default, separators=(",", ":")) + json.dumps(item, default=self._json_default, separators=(",", ":")) ) self._append_handle.write("\n") self._append_handle.flush() @@ -210,7 +207,7 @@ class FileBackedQueueAdapter(QueueRepositoryInterface): with tmp_path.open("w", encoding="utf-8") as handle: for item in self._memory.snapshot(): handle.write( - json.dumps(item, default=_json_default, separators=(",", ":")) + json.dumps(item, default=self._json_default, separators=(",", ":")) ) handle.write("\n") handle.flush() diff --git a/python_repositories/adapters/memory_queue_adapter.py b/python_repositories/adapters/memory_queue_adapter.py index b5eab1f..d26f2e5 100644 --- a/python_repositories/adapters/memory_queue_adapter.py +++ b/python_repositories/adapters/memory_queue_adapter.py @@ -13,21 +13,6 @@ from python_repositories.interfaces.queue_repository_interface import ( ) -def _parse_age(value: object) -> datetime: - """Normalize an age field to a timezone-aware UTC datetime.""" - if isinstance(value, datetime): - if value.tzinfo is None: - return value.replace(tzinfo=UTC) - return value.astimezone(UTC) - if isinstance(value, str): - normalized = value.replace("Z", "+00:00") - parsed = datetime.fromisoformat(normalized) - if parsed.tzinfo is None: - return parsed.replace(tzinfo=UTC) - return parsed.astimezone(UTC) - raise ValueError(f"age value must be datetime or ISO-8601 str, got {type(value)!r}") - - class MemoryQueueAdapter(QueueRepositoryInterface): """Thread-safe in-process FIFO queue with optional deduplication.""" @@ -43,6 +28,23 @@ class MemoryQueueAdapter(QueueRepositoryInterface): self._items: OrderedDict[Hashable, dict[str, Any]] = OrderedDict() self._seq = 0 + @staticmethod + def _parse_age(value: object) -> datetime: + """Normalize an age field to a timezone-aware UTC datetime.""" + if isinstance(value, datetime): + if value.tzinfo is None: + return value.replace(tzinfo=UTC) + return value.astimezone(UTC) + if isinstance(value, str): + normalized = value.replace("Z", "+00:00") + parsed = datetime.fromisoformat(normalized) + if parsed.tzinfo is None: + return parsed.replace(tzinfo=UTC) + return parsed.astimezone(UTC) + raise ValueError( + f"age value must be datetime or ISO-8601 str, got {type(value)!r}" + ) + def enqueue(self, items: list[dict[str, Any]]) -> None: """Append items; skip duplicates when ``dedup_keys`` is configured.""" self.enqueue_and_return_added(items) @@ -84,13 +86,13 @@ class MemoryQueueAdapter(QueueRepositoryInterface): """Remove items whose age field is strictly older than ``cutoff``.""" if self._age_key is None: return 0 - cutoff_utc = _parse_age(cutoff) + cutoff_utc = self._parse_age(cutoff) removed = 0 with self._lock: to_remove = [ key for key, item in self._items.items() - if _parse_age(item[self._age_key]) < cutoff_utc + if self._parse_age(item[self._age_key]) < cutoff_utc ] for key in to_remove: del self._items[key] @@ -116,7 +118,7 @@ class MemoryQueueAdapter(QueueRepositoryInterface): if self._age_key is not None: if self._age_key not in item: raise ValueError(f"item missing age key: {self._age_key!r}") - item[self._age_key] = _parse_age(item[self._age_key]) + item[self._age_key] = self._parse_age(item[self._age_key]) return item def _make_key(self, item: dict[str, Any]) -> Hashable: diff --git a/tests/unit/file_backed_queue_adapter_test.py b/tests/unit/file_backed_queue_adapter_test.py index 2972da8..bd9eecb 100644 --- a/tests/unit/file_backed_queue_adapter_test.py +++ b/tests/unit/file_backed_queue_adapter_test.py @@ -207,10 +207,8 @@ def test_reconnect_replays_after_disconnect(tmp_path: Path) -> None: def test_json_default_rejects_unsupported() -> None: - from python_repositories.adapters.file_backed_queue_adapter import _json_default - with pytest.raises(TypeError): - _json_default(object()) + FileBackedQueueAdapter._json_default(object()) def test_disconnect_ignores_fsync_oserror( diff --git a/tests/unit/memory_queue_adapter_test.py b/tests/unit/memory_queue_adapter_test.py index 0e3cb15..84d6ebc 100644 --- a/tests/unit/memory_queue_adapter_test.py +++ b/tests/unit/memory_queue_adapter_test.py @@ -7,10 +7,7 @@ import threading import pytest -from python_repositories.adapters.memory_queue_adapter import ( - MemoryQueueAdapter, - _parse_age, -) +from python_repositories.adapters.memory_queue_adapter import MemoryQueueAdapter from python_repositories.interfaces.queue_repository_interface import ( QueueRepositoryInterface, ) @@ -85,7 +82,7 @@ def test_evict_without_age_key_is_noop() -> None: def test_parse_age_rejects_unsupported_type() -> None: with pytest.raises(ValueError, match="age value must be"): - _parse_age(123) + MemoryQueueAdapter._parse_age(123) def test_parse_age_aware_datetime() -> None: @@ -93,22 +90,26 @@ def test_parse_age_aware_datetime() -> None: eastern = timezone(timedelta(hours=-5)) aware = datetime(2024, 1, 1, 12, 0, 0, tzinfo=eastern) - assert _parse_age(aware) == datetime(2024, 1, 1, 17, 0, 0, tzinfo=UTC) + assert MemoryQueueAdapter._parse_age(aware) == datetime( + 2024, 1, 1, 17, 0, 0, tzinfo=UTC + ) def test_parse_age_naive_datetime() -> None: naive = datetime(2024, 1, 1, 12, 0, 0) - assert _parse_age(naive) == datetime(2024, 1, 1, 12, 0, 0, tzinfo=UTC) + assert MemoryQueueAdapter._parse_age(naive) == datetime( + 2024, 1, 1, 12, 0, 0, tzinfo=UTC + ) def test_parse_age_zulu_string() -> None: - assert _parse_age("2024-01-01T00:00:00Z") == datetime( + assert MemoryQueueAdapter._parse_age("2024-01-01T00:00:00Z") == datetime( 2024, 1, 1, 0, 0, 0, tzinfo=UTC ) def test_parse_age_naive_string() -> None: - assert _parse_age("2024-01-01T00:00:00") == datetime( + assert MemoryQueueAdapter._parse_age("2024-01-01T00:00:00") == datetime( 2024, 1, 1, 0, 0, 0, tzinfo=UTC )