diff --git a/app/worker/memory_sync_outbox_worker.py b/app/worker/memory_sync_outbox_worker.py new file mode 100644 index 0000000..73f97f5 --- /dev/null +++ b/app/worker/memory_sync_outbox_worker.py @@ -0,0 +1,88 @@ +"""画像投影 Outbox 消费器。 + +外部存储通过 handler 注入;本模块只负责领取、重试、死信和 MySQL 状态更新。 +""" + +from collections.abc import Awaitable, Callable +from datetime import UTC, datetime, timedelta +from typing import Any + +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.infrastructure.db import SessionFactory +from app.model.memory import MemorySyncOutbox + +ProjectionHandler = Callable[[dict[str, Any]], Awaitable[Any]] +MAX_RETRY_COUNT = 5 + + +class MemorySyncOutboxWorker: + """按目标存储独立消费画像投影事件。""" + + def __init__( + self, + handlers: dict[str, ProjectionHandler], + *, + session_factory: Callable[[], AsyncSession] = SessionFactory, + ) -> None: + self.handlers = handlers + self.session_factory = session_factory + + async def run_once(self, *, target_store: str | None = None) -> bool: + """领取并处理一条到期事件;没有可处理事件时返回 False。""" + if not self.handlers: + return False + async with self.session_factory() as session: + now = datetime.now(UTC).replace(tzinfo=None) + conditions: list[Any] = [ + MemorySyncOutbox.status.in_({"pending", "failed"}), + MemorySyncOutbox.next_retry_at.is_(None) + | (MemorySyncOutbox.next_retry_at <= now), + ] + if target_store is not None: + conditions.append(MemorySyncOutbox.target_store == target_store) + event = await session.scalar( + select(MemorySyncOutbox) + .where(*conditions) + .order_by(MemorySyncOutbox.id) + .limit(1) + .with_for_update(skip_locked=True) + ) + if event is None: + await session.rollback() + return False + handler = self.handlers.get(event.target_store) + if handler is None: + self._fail(event, "target_handler_not_configured", now, dead=True) + await session.commit() + return True + try: + await handler(event.payload) + except Exception as exc: + self._fail(event, type(exc).__name__, now) + else: + event.status = "processed" + event.processed_at = datetime.now(UTC).replace(tzinfo=None) + event.last_error = None + event.next_retry_at = None + await session.commit() + return True + + @staticmethod + def _fail( + event: MemorySyncOutbox, + reason: str, + now: datetime, + *, + dead: bool = False, + ) -> None: + """写入可重试失败或死信状态,不吞掉失败事实。""" + event.retry_count = int(event.retry_count) + 1 + event.last_error = reason[:500] + event.status = "dead" if dead or event.retry_count >= MAX_RETRY_COUNT else "failed" + event.next_retry_at = ( + None + if event.status == "dead" + else now + timedelta(seconds=min(300, 2 ** int(event.retry_count))) + ) diff --git a/tests/unit/worker/test_memory_sync_outbox_worker.py b/tests/unit/worker/test_memory_sync_outbox_worker.py new file mode 100644 index 0000000..8f25a34 --- /dev/null +++ b/tests/unit/worker/test_memory_sync_outbox_worker.py @@ -0,0 +1,100 @@ +from datetime import datetime + +import pytest + +from app.model.memory import MemorySyncOutbox +from app.worker.memory_sync_outbox_worker import MemorySyncOutboxWorker + + +class FakeSession: + def __init__(self, event: MemorySyncOutbox | None) -> None: + self.event = event + self.commits = 0 + self.rollbacks = 0 + + async def __aenter__(self) -> "FakeSession": + return self + + async def __aexit__(self, *args: object) -> None: + return None + + async def scalar(self, statement: object) -> MemorySyncOutbox | None: + del statement + return self.event + + async def commit(self) -> None: + self.commits += 1 + + async def rollback(self) -> None: + self.rollbacks += 1 + + +def event(*, target: str = "neo4j", retry_count: int = 0) -> MemorySyncOutbox: + return MemorySyncOutbox( + id=1, event_uuid="event-1", aggregate_type="profile_snapshot", + aggregate_uuid="profile-1", aggregate_version=1, target_store=target, + operation="upsert", payload={"customer_id": 7}, status="pending", + retry_count=retry_count, next_retry_at=None, last_error=None, + created_at=datetime(2026, 1, 1), processed_at=None, + ) + + +@pytest.mark.asyncio +async def test_success_marks_event_processed() -> None: + item = event() + session = FakeSession(item) + seen: list[dict[str, object]] = [] + + async def handler(payload: dict[str, object]) -> None: + seen.append(payload) + + worker = MemorySyncOutboxWorker({"neo4j": handler}, session_factory=lambda: session) + assert await worker.run_once() is True + assert seen == [{"customer_id": 7}] + assert item.status == "processed" + assert item.processed_at is not None + assert session.commits == 1 + + +@pytest.mark.asyncio +async def test_handler_failure_uses_backoff_and_keeps_event() -> None: + item = event() + session = FakeSession(item) + + async def handler(payload: dict[str, object]) -> None: + del payload + raise TimeoutError + + worker = MemorySyncOutboxWorker({"neo4j": handler}, session_factory=lambda: session) + assert await worker.run_once() is True + assert item.status == "failed" + assert item.retry_count == 1 + assert item.next_retry_at is not None + assert item.last_error == "TimeoutError" + + +@pytest.mark.asyncio +async def test_fifth_failure_enters_dead_state() -> None: + item = event(retry_count=4) + session = FakeSession(item) + + async def handler(payload: dict[str, object]) -> None: + del payload + raise RuntimeError + + worker = MemorySyncOutboxWorker({"neo4j": handler}, session_factory=lambda: session) + await worker.run_once() + assert item.status == "dead" + assert item.retry_count == 5 + assert item.next_retry_at is None + + +@pytest.mark.asyncio +async def test_missing_handler_enters_dead_state_without_external_call() -> None: + item = event(target="milvus") + session = FakeSession(item) + worker = MemorySyncOutboxWorker({"neo4j": lambda _: None}, session_factory=lambda: session) + + assert await worker.run_once() is True + assert item.status == "dead" + assert item.last_error == "target_handler_not_configured"