Add QueueRepositoryInterface with memory and file-backed adapters.
PR Title Check / check-title (pull_request) Successful in 9s
Code Quality Pipeline / code-quality (pull_request) Failing after 53s
Test Python Package / unit-tests (pull_request) Successful in 1m1s
Test Python Package / integration-tests (pull_request) Successful in 1m44s
Test Python Package / coverage-report (pull_request) Successful in 13s
PR Title Check / check-title (pull_request) Successful in 9s
Code Quality Pipeline / code-quality (pull_request) Failing after 53s
Test Python Package / unit-tests (pull_request) Successful in 1m1s
Test Python Package / integration-tests (pull_request) Successful in 1m44s
Test Python Package / coverage-report (pull_request) Successful in 13s
Provide a generic disk-backed FIFO queue with configurable path, retention, and dedup keys so consumers can buffer items across restarts without optional extras. Co-authored-by: Cursor <[email protected]>
This commit is contained in:
co-authored by
Cursor
parent
88ea7f06ab
commit
4f58c32dd6
@@ -6,17 +6,28 @@ from typing import TYPE_CHECKING
|
||||
|
||||
# Interfaces are always available; they have no optional backend dependencies.
|
||||
from . import adapters
|
||||
from .config import MinioConfig, PostgresConfig, RedisConfig, load_dotenv
|
||||
from .config import (
|
||||
FileQueueConfig,
|
||||
MinioConfig,
|
||||
PostgresConfig,
|
||||
RedisConfig,
|
||||
load_dotenv,
|
||||
)
|
||||
from .interfaces import (
|
||||
ConnectionAwareInterface,
|
||||
ContextAwareInterface,
|
||||
JsonRepositoryInterface,
|
||||
ObjectRepositoryInterface,
|
||||
QueueRepositoryInterface,
|
||||
TableRepositoryInterface,
|
||||
)
|
||||
|
||||
# Adapters are imported only for static type checkers; runtime loading is delegated below.
|
||||
if TYPE_CHECKING:
|
||||
from .adapters.file_backed_queue_adapter import (
|
||||
FileBackedQueueAdapter as FileBackedQueueAdapter,
|
||||
)
|
||||
from .adapters.memory_queue_adapter import MemoryQueueAdapter as MemoryQueueAdapter
|
||||
from .adapters.minio_adapter import MinioAdapter as MinioAdapter
|
||||
from .adapters.postgres_adapter import PostgresAdapter as PostgresAdapter
|
||||
from .adapters.redis_adapter import RedisAdapter as RedisAdapter
|
||||
@@ -24,10 +35,12 @@ if TYPE_CHECKING:
|
||||
__all__ = [
|
||||
"ConnectionAwareInterface",
|
||||
"ContextAwareInterface",
|
||||
"FileQueueConfig",
|
||||
"JsonRepositoryInterface",
|
||||
"MinioConfig",
|
||||
"ObjectRepositoryInterface",
|
||||
"PostgresConfig",
|
||||
"QueueRepositoryInterface",
|
||||
"RedisConfig",
|
||||
"TableRepositoryInterface",
|
||||
"load_dotenv",
|
||||
|
||||
@@ -11,6 +11,8 @@ from typing import TYPE_CHECKING
|
||||
|
||||
# Adapters are imported only for static type checkers; runtime loading is deferred below.
|
||||
if TYPE_CHECKING:
|
||||
from .file_backed_queue_adapter import FileBackedQueueAdapter
|
||||
from .memory_queue_adapter import MemoryQueueAdapter
|
||||
from .minio_adapter import MinioAdapter
|
||||
from .postgres_adapter import PostgresAdapter
|
||||
from .redis_adapter import RedisAdapter
|
||||
@@ -22,12 +24,19 @@ _LAZY_EXPORTS = {
|
||||
"RedisAdapter": (".redis_adapter", "RedisAdapter"),
|
||||
"MinioAdapter": (".minio_adapter", "MinioAdapter"),
|
||||
"PostgresAdapter": (".postgres_adapter", "PostgresAdapter"),
|
||||
"MemoryQueueAdapter": (".memory_queue_adapter", "MemoryQueueAdapter"),
|
||||
"FileBackedQueueAdapter": (
|
||||
".file_backed_queue_adapter",
|
||||
"FileBackedQueueAdapter",
|
||||
),
|
||||
}
|
||||
|
||||
__all__ = [
|
||||
"RedisAdapter",
|
||||
"MinioAdapter",
|
||||
"PostgresAdapter",
|
||||
"MemoryQueueAdapter",
|
||||
"FileBackedQueueAdapter",
|
||||
]
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,220 @@
|
||||
"""File-backed queue adapter: in-memory hot buffer mirrored to JSONL on disk."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import UTC, datetime, timedelta
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import threading
|
||||
from typing import Any, Self, TextIO
|
||||
|
||||
import structlog
|
||||
|
||||
from python_repositories.adapters.memory_queue_adapter import (
|
||||
MemoryQueueAdapter,
|
||||
_parse_age,
|
||||
)
|
||||
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.
|
||||
|
||||
Does not extend ``ConnectionAwareAdapter``. ``connect()`` replays the JSONL
|
||||
file into memory; ``disconnect()`` flushes any open handle.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
config: FileQueueConfig | None = None,
|
||||
memory: MemoryQueueAdapter | None = None,
|
||||
) -> None:
|
||||
if config is None:
|
||||
config = FileQueueConfig.from_env()
|
||||
if memory is None:
|
||||
memory = MemoryQueueAdapter(
|
||||
dedup_keys=config.dedup_keys,
|
||||
age_key=config.age_key,
|
||||
)
|
||||
self._config = config
|
||||
self._memory = memory
|
||||
self._lock = threading.RLock()
|
||||
self._connected = False
|
||||
self._append_handle: TextIO | None = None
|
||||
self.logger = structlog.get_logger(self.__class__.__name__)
|
||||
|
||||
def __enter__(self) -> Self:
|
||||
self.connect()
|
||||
return self
|
||||
|
||||
def __exit__(
|
||||
self, exc_type: type | None, exc_val: object | None, exc_tb: object | None
|
||||
) -> None:
|
||||
del exc_type, exc_val, exc_tb
|
||||
self.disconnect()
|
||||
|
||||
def connect(self) -> None:
|
||||
"""Create parent directories and replay the JSONL file into memory."""
|
||||
with self._lock:
|
||||
if self._connected:
|
||||
return
|
||||
path = self._config.path
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
self._memory.clear()
|
||||
if path.is_file():
|
||||
self._replay_file(path)
|
||||
self._append_handle = path.open("a", encoding="utf-8")
|
||||
self._connected = True
|
||||
|
||||
def disconnect(self) -> None:
|
||||
"""Flush and close the append handle if open."""
|
||||
with self._lock:
|
||||
if self._append_handle is not None:
|
||||
self._append_handle.flush()
|
||||
try:
|
||||
os.fsync(self._append_handle.fileno())
|
||||
except OSError:
|
||||
pass
|
||||
self._append_handle.close()
|
||||
self._append_handle = None
|
||||
self._connected = False
|
||||
|
||||
def is_connected(self) -> bool:
|
||||
"""Return whether ``connect()`` has completed successfully."""
|
||||
return self._connected
|
||||
|
||||
def enqueue(self, items: list[dict[str, Any]]) -> None:
|
||||
"""Append items to memory and the JSONL file; optionally age-evict."""
|
||||
self._require_connected()
|
||||
with self._lock:
|
||||
added = self._memory.enqueue_and_return_added(items)
|
||||
if added:
|
||||
self._append_items(added)
|
||||
if self._config.age_key is not None:
|
||||
cutoff = datetime.now(UTC) - timedelta(
|
||||
hours=self._config.max_age_hours
|
||||
)
|
||||
if self._memory.evict_older_than(cutoff):
|
||||
self._compact()
|
||||
|
||||
def dequeue_batch(self, *, max_items: int = 1000) -> list[dict[str, Any]]:
|
||||
"""Dequeue from memory and compact the JSONL file to match."""
|
||||
self._require_connected()
|
||||
with self._lock:
|
||||
batch = self._memory.dequeue_batch(max_items=max_items)
|
||||
if batch:
|
||||
self._compact()
|
||||
return batch
|
||||
|
||||
def size(self) -> int:
|
||||
"""Return the number of items currently in the in-memory buffer."""
|
||||
self._require_connected()
|
||||
return self._memory.size()
|
||||
|
||||
def evict_older_than(self, cutoff: datetime) -> int:
|
||||
"""Evict aged items from memory and compact the JSONL file."""
|
||||
self._require_connected()
|
||||
with self._lock:
|
||||
removed = self._memory.evict_older_than(cutoff)
|
||||
if removed:
|
||||
self._compact()
|
||||
return removed
|
||||
|
||||
def _require_connected(self) -> None:
|
||||
if not self._connected:
|
||||
raise RuntimeError("FileBackedQueueAdapter is not connected")
|
||||
|
||||
def _replay_file(self, path: Path) -> None:
|
||||
cutoff: datetime | None = None
|
||||
if self._config.age_key is not None:
|
||||
cutoff = datetime.now(UTC) - timedelta(hours=self._config.max_age_hours)
|
||||
with path.open(encoding="utf-8") as handle:
|
||||
for line_number, line in enumerate(handle, start=1):
|
||||
stripped = line.strip()
|
||||
if not stripped:
|
||||
continue
|
||||
try:
|
||||
payload = json.loads(stripped)
|
||||
except json.JSONDecodeError:
|
||||
self.logger.warning(
|
||||
"Skipping corrupt JSONL line",
|
||||
path=str(path),
|
||||
line_number=line_number,
|
||||
)
|
||||
continue
|
||||
if not isinstance(payload, dict):
|
||||
self.logger.warning(
|
||||
"Skipping non-object JSONL line",
|
||||
path=str(path),
|
||||
line_number=line_number,
|
||||
)
|
||||
continue
|
||||
item: dict[str, Any] = payload
|
||||
if cutoff is not None and self._config.age_key is not None:
|
||||
age_value = item.get(self._config.age_key)
|
||||
if age_value is None:
|
||||
self.logger.warning(
|
||||
"Skipping item missing age key on replay",
|
||||
path=str(path),
|
||||
line_number=line_number,
|
||||
age_key=self._config.age_key,
|
||||
)
|
||||
continue
|
||||
try:
|
||||
if _parse_age(age_value) < cutoff:
|
||||
continue
|
||||
except ValueError:
|
||||
self.logger.warning(
|
||||
"Skipping item with invalid age on replay",
|
||||
path=str(path),
|
||||
line_number=line_number,
|
||||
)
|
||||
continue
|
||||
try:
|
||||
self._memory.enqueue_and_return_added([item])
|
||||
except ValueError as exc:
|
||||
self.logger.warning(
|
||||
"Skipping invalid item on replay",
|
||||
path=str(path),
|
||||
line_number=line_number,
|
||||
error=str(exc),
|
||||
)
|
||||
|
||||
def _append_items(self, items: list[dict[str, Any]]) -> None:
|
||||
if self._append_handle is None:
|
||||
raise RuntimeError("append handle is not open")
|
||||
for item in items:
|
||||
self._append_handle.write(
|
||||
json.dumps(item, default=_json_default, separators=(",", ":"))
|
||||
)
|
||||
self._append_handle.write("\n")
|
||||
self._append_handle.flush()
|
||||
|
||||
def _compact(self) -> None:
|
||||
"""Rewrite the JSONL file from the current in-memory snapshot."""
|
||||
path = self._config.path
|
||||
tmp_path = path.with_suffix(path.suffix + ".tmp")
|
||||
if self._append_handle is not None:
|
||||
self._append_handle.flush()
|
||||
self._append_handle.close()
|
||||
self._append_handle = None
|
||||
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=(",", ":"))
|
||||
)
|
||||
handle.write("\n")
|
||||
handle.flush()
|
||||
tmp_path.replace(path)
|
||||
self._append_handle = path.open("a", encoding="utf-8")
|
||||
@@ -0,0 +1,126 @@
|
||||
"""In-memory queue adapter with optional dedup and age-based eviction."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections import OrderedDict
|
||||
from collections.abc import Hashable
|
||||
from datetime import UTC, datetime
|
||||
import threading
|
||||
from typing import Any
|
||||
|
||||
from python_repositories.interfaces.queue_repository_interface import (
|
||||
QueueRepositoryInterface,
|
||||
)
|
||||
|
||||
|
||||
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."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
dedup_keys: tuple[str, ...] = (),
|
||||
age_key: str | None = None,
|
||||
) -> None:
|
||||
self._dedup_keys = dedup_keys
|
||||
self._age_key = age_key
|
||||
self._lock = threading.RLock()
|
||||
self._items: OrderedDict[Hashable, dict[str, Any]] = OrderedDict()
|
||||
self._seq = 0
|
||||
|
||||
def enqueue(self, items: list[dict[str, Any]]) -> None:
|
||||
"""Append items; skip duplicates when ``dedup_keys`` is configured."""
|
||||
self.enqueue_and_return_added(items)
|
||||
|
||||
def enqueue_and_return_added(
|
||||
self, items: list[dict[str, Any]]
|
||||
) -> list[dict[str, Any]]:
|
||||
"""Enqueue items and return the subset that was newly stored."""
|
||||
if not items:
|
||||
return []
|
||||
added: list[dict[str, Any]] = []
|
||||
with self._lock:
|
||||
for raw in items:
|
||||
item = self._normalize_item(raw)
|
||||
key = self._make_key(item)
|
||||
if self._dedup_keys and key in self._items:
|
||||
continue
|
||||
self._items[key] = item
|
||||
added.append(item)
|
||||
return added
|
||||
|
||||
def dequeue_batch(self, *, max_items: int = 1000) -> list[dict[str, Any]]:
|
||||
"""Remove and return up to ``max_items`` items in FIFO order."""
|
||||
if max_items < 0:
|
||||
raise ValueError("max_items must be >= 0")
|
||||
with self._lock:
|
||||
batch: list[dict[str, Any]] = []
|
||||
for _ in range(min(max_items, len(self._items))):
|
||||
_key, item = self._items.popitem(last=False)
|
||||
batch.append(item)
|
||||
return batch
|
||||
|
||||
def size(self) -> int:
|
||||
"""Return the number of items currently in the queue."""
|
||||
with self._lock:
|
||||
return len(self._items)
|
||||
|
||||
def evict_older_than(self, cutoff: datetime) -> int:
|
||||
"""Remove items whose age field is strictly older than ``cutoff``."""
|
||||
if self._age_key is None:
|
||||
return 0
|
||||
cutoff_utc = _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
|
||||
]
|
||||
for key in to_remove:
|
||||
del self._items[key]
|
||||
removed += 1
|
||||
return removed
|
||||
|
||||
def clear(self) -> None:
|
||||
"""Remove all items from the queue."""
|
||||
with self._lock:
|
||||
self._items.clear()
|
||||
|
||||
def snapshot(self) -> list[dict[str, Any]]:
|
||||
"""Return a shallow copy of queued items in FIFO order."""
|
||||
with self._lock:
|
||||
return [dict(item) for item in self._items.values()]
|
||||
|
||||
def _normalize_item(self, raw: dict[str, Any]) -> dict[str, Any]:
|
||||
item = dict(raw)
|
||||
if self._dedup_keys:
|
||||
missing = [key for key in self._dedup_keys if key not in item]
|
||||
if missing:
|
||||
raise ValueError(f"item missing dedup key(s): {missing}")
|
||||
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])
|
||||
return item
|
||||
|
||||
def _make_key(self, item: dict[str, Any]) -> Hashable:
|
||||
if not self._dedup_keys:
|
||||
self._seq += 1
|
||||
return self._seq
|
||||
return tuple(item[key] for key in self._dedup_keys)
|
||||
@@ -1,9 +1,11 @@
|
||||
from .dotenv_loader import load_dotenv as load_dotenv
|
||||
from .file_queue_config import FileQueueConfig as FileQueueConfig
|
||||
from .minio_config import MinioConfig as MinioConfig
|
||||
from .postgres_config import PostgresConfig as PostgresConfig
|
||||
from .redis_config import RedisConfig as RedisConfig
|
||||
|
||||
__all__ = [
|
||||
"FileQueueConfig",
|
||||
"MinioConfig",
|
||||
"PostgresConfig",
|
||||
"RedisConfig",
|
||||
|
||||
@@ -0,0 +1,56 @@
|
||||
"""File-backed queue configuration."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass
|
||||
import os
|
||||
from pathlib import Path
|
||||
|
||||
from python_utils import check_env
|
||||
|
||||
from python_repositories.config.dotenv_loader import load_dotenv
|
||||
|
||||
_DEFAULT_MAX_AGE_HOURS = 24
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class FileQueueConfig:
|
||||
"""Configuration for a JSONL file-backed queue.
|
||||
|
||||
Path and key schema are owned by the caller. This package does not assume
|
||||
any particular directory layout or item field names.
|
||||
"""
|
||||
|
||||
path: Path
|
||||
max_age_hours: int = _DEFAULT_MAX_AGE_HOURS
|
||||
dedup_keys: tuple[str, ...] = ()
|
||||
age_key: str | None = None
|
||||
|
||||
@classmethod
|
||||
def from_env(
|
||||
cls,
|
||||
path_env_var_name: str = "FILE_QUEUE_PATH",
|
||||
*,
|
||||
max_age_hours_env_var_name: str = "FILE_QUEUE_MAX_AGE_HOURS",
|
||||
dedup_keys: tuple[str, ...] = (),
|
||||
age_key: str | None = None,
|
||||
use_dotenv: bool = True,
|
||||
) -> FileQueueConfig:
|
||||
"""Load path and optional max age from environment variables.
|
||||
|
||||
``dedup_keys`` and ``age_key`` are domain-specific and must be passed
|
||||
explicitly; they are not read from the environment.
|
||||
"""
|
||||
if use_dotenv:
|
||||
load_dotenv()
|
||||
check_env(path_env_var_name)
|
||||
raw_max_age = os.getenv(max_age_hours_env_var_name)
|
||||
max_age_hours = (
|
||||
_DEFAULT_MAX_AGE_HOURS if raw_max_age is None else int(raw_max_age)
|
||||
)
|
||||
return cls(
|
||||
path=Path(str(os.getenv(path_env_var_name))),
|
||||
max_age_hours=max_age_hours,
|
||||
dedup_keys=dedup_keys,
|
||||
age_key=age_key,
|
||||
)
|
||||
@@ -8,6 +8,9 @@ from .json_repository_interface import (
|
||||
from .object_repository_interface import (
|
||||
ObjectRepositoryInterface as ObjectRepositoryInterface,
|
||||
)
|
||||
from .queue_repository_interface import (
|
||||
QueueRepositoryInterface as QueueRepositoryInterface,
|
||||
)
|
||||
from .table_repository_interface import (
|
||||
TableRepositoryInterface as TableRepositoryInterface,
|
||||
)
|
||||
@@ -17,5 +20,6 @@ __all__ = [
|
||||
"ContextAwareInterface",
|
||||
"JsonRepositoryInterface",
|
||||
"ObjectRepositoryInterface",
|
||||
"QueueRepositoryInterface",
|
||||
"TableRepositoryInterface",
|
||||
]
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
"""Definition of QueueRepositoryInterface protocol."""
|
||||
|
||||
from abc import abstractmethod
|
||||
from datetime import datetime
|
||||
from typing import Any, Protocol, runtime_checkable
|
||||
|
||||
|
||||
@runtime_checkable
|
||||
class QueueRepositoryInterface(Protocol):
|
||||
"""Interface that defines buffered queue enqueue/dequeue methods."""
|
||||
|
||||
@abstractmethod
|
||||
def enqueue(self, items: list[dict[str, Any]]) -> None:
|
||||
"""Append items to the queue.
|
||||
|
||||
Duplicates may be skipped when the implementation is configured with
|
||||
dedup keys. An empty list is a no-op.
|
||||
"""
|
||||
...
|
||||
|
||||
@abstractmethod
|
||||
def dequeue_batch(self, *, max_items: int = 1000) -> list[dict[str, Any]]:
|
||||
"""Remove and return up to ``max_items`` items in FIFO order.
|
||||
|
||||
Returns an empty list when the queue is empty.
|
||||
"""
|
||||
...
|
||||
|
||||
@abstractmethod
|
||||
def size(self) -> int:
|
||||
"""Return the number of items currently in the queue."""
|
||||
...
|
||||
|
||||
@abstractmethod
|
||||
def evict_older_than(self, cutoff: datetime) -> int:
|
||||
"""Remove items older than ``cutoff`` and return how many were removed.
|
||||
|
||||
Age is determined by an implementation-specific item field. When no age
|
||||
field is configured, this is a no-op that returns ``0``.
|
||||
"""
|
||||
...
|
||||
Reference in New Issue
Block a user