Files
python-repositories/tests/unit/memory_queue_adapter_test.py
T
Brian Bjarke JensenandCursor a78320b434
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
Move queue helper functions to staticmethods on their adapters.
Co-authored-by: Cursor <[email protected]>
2026-07-16 21:07:35 +02:00

150 lines
4.4 KiB
Python

"""Unit tests for MemoryQueueAdapter."""
from __future__ import annotations
from datetime import UTC, datetime, timedelta
import threading
import pytest
from python_repositories.adapters.memory_queue_adapter import MemoryQueueAdapter
from python_repositories.interfaces.queue_repository_interface import (
QueueRepositoryInterface,
)
def test_implements_interface() -> None:
assert issubclass(MemoryQueueAdapter, QueueRepositoryInterface)
def test_fifo_without_dedup() -> None:
queue = MemoryQueueAdapter()
queue.enqueue([{"id": 1}, {"id": 2}, {"id": 3}])
assert queue.size() == 3
assert queue.dequeue_batch(max_items=2) == [{"id": 1}, {"id": 2}]
assert queue.dequeue_batch(max_items=10) == [{"id": 3}]
assert queue.dequeue_batch() == []
assert queue.size() == 0
def test_empty_enqueue_is_noop() -> None:
queue = MemoryQueueAdapter()
queue.enqueue([])
assert queue.size() == 0
assert queue.enqueue_and_return_added([]) == []
def test_dedup_keeps_first() -> None:
queue = MemoryQueueAdapter(dedup_keys=("id",))
queue.enqueue([{"id": 1, "v": "a"}, {"id": 1, "v": "b"}, {"id": 2, "v": "c"}])
assert queue.size() == 2
assert queue.dequeue_batch(max_items=10) == [
{"id": 1, "v": "a"},
{"id": 2, "v": "c"},
]
def test_missing_dedup_key_raises() -> None:
queue = MemoryQueueAdapter(dedup_keys=("id",))
with pytest.raises(ValueError, match="missing dedup key"):
queue.enqueue([{"name": "x"}])
def test_missing_age_key_raises() -> None:
queue = MemoryQueueAdapter(age_key="created_at")
with pytest.raises(ValueError, match="missing age key"):
queue.enqueue([{"id": 1}])
def test_evict_older_than() -> None:
queue = MemoryQueueAdapter(age_key="created_at", dedup_keys=("id",))
now = datetime.now(UTC)
queue.enqueue(
[
{"id": 1, "created_at": now - timedelta(hours=2)},
{"id": 2, "created_at": now - timedelta(minutes=30)},
{"id": 3, "created_at": (now - timedelta(hours=3)).isoformat()},
]
)
removed = queue.evict_older_than(now - timedelta(hours=1))
assert removed == 2
assert queue.size() == 1
remaining = queue.dequeue_batch(max_items=10)
assert remaining[0]["id"] == 2
def test_evict_without_age_key_is_noop() -> None:
queue = MemoryQueueAdapter()
queue.enqueue([{"id": 1}])
assert queue.evict_older_than(datetime.now(UTC)) == 0
assert queue.size() == 1
def test_parse_age_rejects_unsupported_type() -> None:
with pytest.raises(ValueError, match="age value must be"):
MemoryQueueAdapter._parse_age(123)
def test_parse_age_aware_datetime() -> None:
from datetime import timezone
eastern = timezone(timedelta(hours=-5))
aware = datetime(2024, 1, 1, 12, 0, 0, tzinfo=eastern)
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 MemoryQueueAdapter._parse_age(naive) == datetime(
2024, 1, 1, 12, 0, 0, tzinfo=UTC
)
def test_parse_age_zulu_string() -> None:
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 MemoryQueueAdapter._parse_age("2024-01-01T00:00:00") == datetime(
2024, 1, 1, 0, 0, 0, tzinfo=UTC
)
def test_dequeue_rejects_negative_max_items() -> None:
queue = MemoryQueueAdapter()
with pytest.raises(ValueError, match="max_items"):
queue.dequeue_batch(max_items=-1)
def test_clear_and_snapshot() -> None:
queue = MemoryQueueAdapter()
queue.enqueue([{"id": 1}, {"id": 2}])
assert queue.snapshot() == [{"id": 1}, {"id": 2}]
queue.clear()
assert queue.size() == 0
assert queue.snapshot() == []
def test_thread_safety_smoke() -> None:
queue = MemoryQueueAdapter(dedup_keys=("id",))
errors: list[BaseException] = []
def worker(start: int) -> None:
try:
for i in range(start, start + 50):
queue.enqueue([{"id": i}])
except BaseException as exc: # noqa: BLE001
errors.append(exc)
threads = [threading.Thread(target=worker, args=(i * 50,)) for i in range(4)]
for thread in threads:
thread.start()
for thread in threads:
thread.join()
assert errors == []
assert queue.size() == 200