Files
group_fqcd_jr/app/worker/outbox_worker.py
lzf_0626 8b8883ccca Worker 失败原因不再只留类名:区分"可落库的固定文案"与"异常消息"
起因是查 Worker 运行状态时发现库里 373 条死信的 last_error 全是裸的 "ValueError"
(工具 tools/probe_worker_state.py,证据 docs/evidence/worker-state.json)。
dispatch(run not found)、dispatch_run_completed、dispatch_memory_extraction、
dispatch_profile_rebuild 抛的都是 ValueError,只记类名等于把"哪一处失败"也一起丢了。

但"直接存 str(exc)"是错的:tests/unit/worker/test_outbox_worker.py 那条
RuntimeError("credential=do-not-log") 断言异常消息不得落库 —— 它可能含凭据、SQL 或
客户标识。第一版改动就是这么写的,被这个测试当场拦下(这测试写得值)。

折中:

- 新增 OutboxHandlerError(继承 ValueError,这些失败本就是 ValueError 语义,保持
  继承关系才不会改动既有的 except ValueError 行为与断言)。它的 reason 由代码写死、
  不含任何请求数据,因此可以落库;
- safe_error_text:OutboxHandlerError → "类名: 固定文案"(截断 500 字符),
  其余异常 → 仍只记类名;
- runtime.py 的 5 处 handler 失败改抛 OutboxHandlerError。

测试:新增"固定文案落库"用例;并把既有用例的断言收紧为 last_error == "RuntimeError"
(原先只断言"不含 do-not-log",太松,漏掉的情况测不出来)。

顺带产出 tools/probe_worker_state.py(只读):outbox / agent_run 各状态计数、按事件
类型分组、死信原因聚合。当前环境实测 pending 347、dead 373、published 410、
agent_run 无 queued/running。

门禁:ruff 干净 / mypy 138 文件 / 697 unit+contract / 33 integration。
2026-09-11 16:28:19 +08:00

128 lines
5.0 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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
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__
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
# 只让"代码写死的固定文案"落库;异常消息可能含凭据,见 safe_error_text。
event.last_error = safe_error_text(exc)
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