Files
lzf_0626 01e4e6a687 修掉"Worker 会自己退出"的真缺陷(场外游标续租失败误取消主任务)
## 现象(2026-09-14 实测,非人为停止)
Worker 进程自己退出,退出码 1,日志末尾是
`ConnectionResetError: [WinError 10054] 远程主机强迫关闭了一个现有的连接`
(发生在 `offsite_worker.close()` 的 IMAP `logout()`),
而真正的起点是更上面那句 `asyncio.exceptions.CancelledError`。

## 根因(两处,都是真缺陷)
1. **取消错了对象**:`OffsiteMailWorker._process_batch` 把
   `asyncio.current_task()`(= `__main__.serve()` 的**主循环任务**)交给游标心跳,
   心跳在续租失败(`rowcount != 1`)或续租抛异常时执行 `task.cancel()` ——
   于是"放弃这一批"变成了"**杀掉整个 Worker**"。`CancelledError` 从 `run_once()`
   一路冒到 `serve()`,主循环直接结束。
2. **收尾异常盖掉退出原因**:`serve()` 的 `finally` 里 `await offsite_worker.close()`
   在网络已断时抛 `ConnectionResetError`,把 `CancelledError` 顶掉,
   表现为"关闭流程崩了",看不出真实原因。

## 改法
- `_process_batch` 把"这一批"跑在**独立任务** `batch` 里,心跳只取消 `batch`;
  调用方捕获 `CancelledError` 后区分两种情况:
  **本批被放弃**(`batch.cancelled()` 且当前任务自己没有在取消)→ 记 warning、返回 `False`、
  Worker 继续下一轮;**外层在取消当前任务**(Ctrl+C / 进程关闭)→ 原样上抛,绝不吞掉。
- `serve()` 的 `finally` 里关闭场外 Worker 包 try/except:收尾失败只记日志,
  **不改变退出码与退出原因**。
- 新增 `_current_task_is_cancelling()` 用 `Task.cancelling()` 做这个区分(3.11+)。

## 守卫
`tests/unit/worker/test_offsite_mail_worker.py::test_cursor_lease_loss_abandons_batch_without_cancelling_the_worker`
—— 续租失败(rowcount=0)时断言:批次确实跑起来过、返回 `False`、
**调用方任务没有被取消**。修复前这条用例会挂在"调用方被取消"上。

## 验证
- `pytest tests/unit/worker` → 112 passed
- `pytest tests/unit tests/contract` → 见下方(0 failed)
- `mypy app/worker/offsite_mail_worker.py app/worker/__main__.py` → 0 错(除组员文件里既有的 1 个)
- 重启 Worker 后跑记忆演示链路:候选 → verified → active → **2 秒**收敛,进程稳定

## 文档
`docs/演示用/记忆系统演示文档-2026-09-14.md`:
场景五补"如果发现 Worker 自己退了是怎么回事",问答补"Worker 会不会自己中途退出"。
2026-09-15 08:01:00 +08:00

93 lines
4.8 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.
import argparse
import asyncio
import logging
from app.core.config import get_settings
from app.infrastructure.db import engine
from app.service.agent.bootstrap import get_relationship_service
from app.service.projection_cleanup_service import ProjectionCleanupService
from app.worker.offsite_mail_worker import OffsiteMailWorker
from app.worker.runtime import WorkerRuntime
logger = logging.getLogger(__name__)
async def serve(*, once: bool = False) -> None:
settings = get_settings()
relationships = get_relationship_service()
# 在**组装层**注入投影删除客户端:记忆失效/销户时清理图库与向量库里的派生数据。
# 放在这里而不是 runtime 内部兜底,是为了保留"未注入即显式降级并留痕"的语义
# (有单测守着这一点),也让"生产装配了什么"在入口处一眼可见。
# 该服务返回自己模块里的 ProjectionCleanupOutcome(字段与 runtime 的同名结构一致),
# 结构化契约成立但名义类型不同,故显式忽略:为此把结构体抽到共享模块会造成
# service 与 worker 两个层次互相导入,不值得。
runtime = WorkerRuntime(
settings=settings,
relationships=relationships,
projection_cleaner=ProjectionCleanupService( # type: ignore[arg-type]
relationships=relationships
),
)
# 画像投影(`memory_sync_outbox`)的消费者**只有一套**,在 `WorkerRuntime.run_once()`
# 内部(`consume_profile_projections`),本入口**不再**另起一个 worker。
#
# ⚠️ 为什么这里不能另装一个:合并主干 PR #7 后,本入口曾有一个
# `MemorySyncOutboxWorker(memory_sync_handlers)` 与 runtime 内那套**同时读同一个队列**,
# 而两套的 handler 并不相同 —— 入口那套的 `neo4j` 指向 ZSY 的
# `Neo4jProfileProjection`(按客户各建私有节点),runtime 那套指向主干的
# `ProfileGraphProjectionService`(共享 tag 节点、只投影已确认事实)。
# 同一事件被哪套领到结果不定,等于**同一事实在图里有两种说法**。
# 现统一走 runtime 那套,理由:它带 `memory_sources` 缺失兜底,且 neo4j 复用主干服务
# (方案 A:不引入第二套图投影)。装配入口的职责仍在本文件 —— 注入 `relationships`
# 与 `projection_cleaner`;Milvus 客户端由 `bootstrap.get_milvus_profile_vector_client()`
# 惰性构造(缺配置时返回 None,runtime 显式降级、事件保持 pending)。
#
# 场外收件 Worker 必须与底座 Worker 同进程同入口:2026-09-11 01:45 的一次批量
# 文件覆盖把这处接线删掉了,导致邮件 Worker 完全不再运行、邮箱无人收取。
offsite_worker = OffsiteMailWorker(settings)
try:
while True:
try:
worked = await runtime.run_once()
worked = await offsite_worker.run_once() or worked
except Exception:
# 常驻 Worker 不能因为"某一轮"的异常就整体退出:数据库抖动、
# 迁移期间锁表、外部依赖瞬断都会命中这里,而 `run_once` 里的
# 裸查询没有兜底。单轮失败记录堆栈后退避重试;`--once` 模式
# 保持抛出,便于诊断一次性运行的真实问题。
logger.warning("worker round failed; retrying after backoff", exc_info=True)
if once:
raise
await asyncio.sleep(settings.worker_poll_seconds)
continue
if once:
return
if not worked:
await asyncio.sleep(settings.worker_poll_seconds)
finally:
# 关闭阶段**不能**让异常盖掉真正的退出原因:2026-09-14 实测过一次
# "Worker 自己退了",日志最后一屏是 `offsite_worker.close()` 里 IMAP
# `logout()` 抛的 `ConnectionResetError [WinError 10054]`(网络已断),
# 把前面的 `CancelledError` 顶掉,看起来像是关闭流程崩了。
# 收尾失败只记日志,退出码与原因保持原样。
try:
await offsite_worker.close()
except Exception: # noqa: BLE001 - 收尾失败不该改变退出语义
logger.warning("场外收件 Worker 关闭失败(不影响退出原因)", exc_info=True)
await engine.dispose()
def main() -> None:
parser = argparse.ArgumentParser(description="Agent 底座独立 Worker")
parser.add_argument("--once", action="store_true", help="消费一次后退出")
args = parser.parse_args()
logging.basicConfig(level=get_settings().log_level)
try:
asyncio.run(serve(once=args.once))
except KeyboardInterrupt:
pass
if __name__ == "__main__":
main()