Move queue helper functions to staticmethods on their adapters.
PR Title Check / check-title (pull_request) Successful in 8s
Test Python Package / unit-tests (pull_request) Successful in 20s
Code Quality Pipeline / code-quality (pull_request) Failing after 32s
Test Python Package / integration-tests (pull_request) Successful in 1m33s
Test Python Package / coverage-report (pull_request) Successful in 16s

Co-authored-by: Cursor <[email protected]>
This commit is contained in:
Brian Bjarke Jensen
2026-07-16 21:07:35 +02:00
co-authored by Cursor
parent 1822d886b1
commit a78320b434
4 changed files with 41 additions and 43 deletions
@@ -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()
@@ -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: