2026-09-09 21:55:37 +08:00
|
|
|
|
from collections.abc import Awaitable, Callable
|
|
|
|
|
|
from datetime import UTC, datetime, timedelta
|
|
|
|
|
|
from typing import Any
|
|
|
|
|
|
|
|
|
|
|
|
from sqlalchemy import select, true
|
|
|
|
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
|
|
|
|
|
|
|
|
from app.model.platform import DomainEventOutbox, OutboxDelivery
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-11 16:28:19 +08:00
|
|
|
|
class OutboxHandlerError(ValueError):
|
|
|
|
|
|
"""Handler 失败,且失败原因是**可以安全落库**的固定文案。
|
|
|
|
|
|
|
|
|
|
|
|
继承 `ValueError` 而不是 `Exception`:这些失败(run not found、payload 不完整、
|
|
|
|
|
|
mode 非法)本来就是 ValueError 语义,保持继承关系才不会改动既有的 `except
|
|
|
|
|
|
ValueError` 行为与断言。
|
|
|
|
|
|
|
|
|
|
|
|
为什么需要这个类型:`last_error` 默认只记异常类名,因为异常消息可能含凭据、SQL
|
|
|
|
|
|
语句或客户标识(`tests/unit/worker/test_outbox_worker.py` 里那条
|
|
|
|
|
|
`RuntimeError("credential=do-not-log")` 就是守这条的)。
|
|
|
|
|
|
|
|
|
|
|
|
但只记类名又不够 —— `dispatch`(run not found)、`dispatch_run_completed`、
|
|
|
|
|
|
`dispatch_memory_extraction`、`dispatch_profile_rebuild` 抛的全是 ValueError,实测
|
|
|
|
|
|
库里 373 条死信的 `last_error` 都是裸的 `"ValueError"`,分不清是哪一处失败的。
|
|
|
|
|
|
|
|
|
|
|
|
折中办法:handler 想让人看见原因时,抛这个类型,`reason` 由**代码写死**、不含任何
|
|
|
|
|
|
请求数据,于是可以落库;其余异常仍然只记类名。
|
|
|
|
|
|
"""
|
|
|
|
|
|
|
|
|
|
|
|
def __init__(self, reason: str) -> None:
|
|
|
|
|
|
super().__init__(reason)
|
|
|
|
|
|
self.reason = reason
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def safe_error_text(exc: BaseException) -> str:
|
|
|
|
|
|
"""把异常转成可落库的失败原因。
|
|
|
|
|
|
|
|
|
|
|
|
- `OutboxHandlerError` → `类名: 固定文案`(文案由代码写死,安全);
|
|
|
|
|
|
- 其他异常 → **只记类名**(消息可能含凭据/请求数据,不落库)。
|
|
|
|
|
|
"""
|
|
|
|
|
|
if isinstance(exc, OutboxHandlerError):
|
|
|
|
|
|
return f"{type(exc).__name__}: {exc.reason}"[:500]
|
|
|
|
|
|
return type(exc).__name__
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-09 21:55:37 +08:00
|
|
|
|
class OutboxWorker:
|
|
|
|
|
|
"""One dispatcher per event type, not a fan-out transport.
|
|
|
|
|
|
|
|
|
|
|
|
Handlers must use this session for local writes and must not commit it.
|
|
|
|
|
|
External side effects still require their own idempotency protocol.
|
|
|
|
|
|
"""
|
|
|
|
|
|
def __init__(
|
|
|
|
|
|
self,
|
|
|
|
|
|
session: AsyncSession,
|
|
|
|
|
|
handlers: dict[str, Callable[[dict[str, Any]], Awaitable[None]]],
|
|
|
|
|
|
) -> None:
|
|
|
|
|
|
self.session = session
|
|
|
|
|
|
self.handlers = handlers
|
|
|
|
|
|
|
|
|
|
|
|
async def publish_one(self, *, aggregate_id: str | None = None) -> bool:
|
|
|
|
|
|
if not self.handlers:
|
|
|
|
|
|
return False
|
|
|
|
|
|
now = datetime.now(UTC).replace(tzinfo=None)
|
|
|
|
|
|
event = await self.session.scalar(
|
|
|
|
|
|
select(DomainEventOutbox)
|
|
|
|
|
|
.where(
|
|
|
|
|
|
DomainEventOutbox.status.in_({"pending", "failed"}),
|
|
|
|
|
|
DomainEventOutbox.event_type.in_(tuple(self.handlers)),
|
|
|
|
|
|
DomainEventOutbox.aggregate_id == aggregate_id if aggregate_id else true(),
|
|
|
|
|
|
(
|
|
|
|
|
|
DomainEventOutbox.next_retry_at.is_(None)
|
|
|
|
|
|
| (DomainEventOutbox.next_retry_at <= now)
|
|
|
|
|
|
),
|
|
|
|
|
|
)
|
|
|
|
|
|
.order_by(DomainEventOutbox.id)
|
|
|
|
|
|
.limit(1)
|
|
|
|
|
|
.with_for_update(skip_locked=True)
|
|
|
|
|
|
)
|
|
|
|
|
|
if event is None:
|
|
|
|
|
|
await self.session.rollback()
|
|
|
|
|
|
return False
|
|
|
|
|
|
handler = self.handlers.get(event.event_type)
|
|
|
|
|
|
if handler is None:
|
|
|
|
|
|
event.status = "dead"
|
|
|
|
|
|
event.last_error = "no handler registered"
|
|
|
|
|
|
await self._commit()
|
|
|
|
|
|
return False
|
|
|
|
|
|
try:
|
|
|
|
|
|
delivery = await self.session.scalar(
|
|
|
|
|
|
select(OutboxDelivery).where(
|
|
|
|
|
|
OutboxDelivery.event_id == event.event_id,
|
|
|
|
|
|
OutboxDelivery.consumer_name == event.event_type,
|
|
|
|
|
|
)
|
|
|
|
|
|
)
|
|
|
|
|
|
except BaseException:
|
|
|
|
|
|
await self.session.rollback()
|
|
|
|
|
|
raise
|
|
|
|
|
|
try:
|
|
|
|
|
|
if delivery is None:
|
|
|
|
|
|
await handler(event.payload)
|
|
|
|
|
|
self.session.add(OutboxDelivery(
|
|
|
|
|
|
event_id=event.event_id,
|
|
|
|
|
|
consumer_name=event.event_type,
|
|
|
|
|
|
delivered_at=datetime.now(UTC).replace(tzinfo=None),
|
|
|
|
|
|
))
|
|
|
|
|
|
await self.session.flush()
|
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
|
event.retry_count += 1
|
2026-09-11 16:28:19 +08:00
|
|
|
|
# 只让"代码写死的固定文案"落库;异常消息可能含凭据,见 safe_error_text。
|
|
|
|
|
|
event.last_error = safe_error_text(exc)
|
2026-09-09 21:55:37 +08:00
|
|
|
|
event.status = "dead" if event.retry_count >= 5 else "failed"
|
|
|
|
|
|
event.next_retry_at = now + timedelta(seconds=min(300, 2**event.retry_count))
|
|
|
|
|
|
else:
|
|
|
|
|
|
event.status = "published"
|
|
|
|
|
|
event.published_at = datetime.now(UTC).replace(tzinfo=None)
|
|
|
|
|
|
event.last_error = None
|
|
|
|
|
|
event.next_retry_at = None
|
|
|
|
|
|
event.updated_at = datetime.now(UTC).replace(tzinfo=None)
|
|
|
|
|
|
await self._commit()
|
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
|
|
async def _commit(self) -> None:
|
|
|
|
|
|
try:
|
|
|
|
|
|
await self.session.commit()
|
|
|
|
|
|
except BaseException:
|
|
|
|
|
|
await self.session.rollback()
|
|
|
|
|
|
raise
|