feat: add memory sync outbox worker

This commit is contained in:
张胜宇
2026-09-11 20:07:55 +08:00
parent f167390a8a
commit 6c605ec3c1
2 changed files with 188 additions and 0 deletions
+88
View File
@@ -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)))
)
@@ -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"