diff --git a/app/worker/__main__.py b/app/worker/__main__.py index 25d3daf..315966f 100644 --- a/app/worker/__main__.py +++ b/app/worker/__main__.py @@ -65,7 +65,15 @@ async def serve(*, once: bool = False) -> None: if not worked: await asyncio.sleep(settings.worker_poll_seconds) finally: - await offsite_worker.close() + # 关闭阶段**不能**让异常盖掉真正的退出原因: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() diff --git a/app/worker/offsite_mail_worker.py b/app/worker/offsite_mail_worker.py index 7f1b5b5..ed9a117 100644 --- a/app/worker/offsite_mail_worker.py +++ b/app/worker/offsite_mail_worker.py @@ -46,6 +46,17 @@ from app.service.offsite_mail_adapter import ( logger = logging.getLogger(__name__) +def _current_task_is_cancelling() -> bool: + """当前任务是否**自己**在被取消(Ctrl+C / 进程关闭),而不是"这一批被放弃"。 + + `Task.cancelling()`(3.11+)返回"已被请求取消的次数",用它区分上面两种 + `CancelledError`:前者必须原样上抛,后者应转成"本批放弃"。 + """ + task = asyncio.current_task() + cancelling = getattr(task, "cancelling", None) + return bool(cancelling()) if callable(cancelling) else False + + class MailReceiver(Protocol): last_scanned_uid: str | None @@ -241,12 +252,32 @@ class OffsiteMailWorker: return None async def _process_batch(self, lease: CursorLease, context: RequestContext) -> bool: - task = asyncio.current_task() - if task is None: - raise RuntimeError("场外 Worker 无法建立当前任务租约") - heartbeat = asyncio.create_task(self._cursor_heartbeat(lease.lease_id, task)) + """处理一批邮件;**租约丢失只放弃这一批,不是杀掉整个 Worker**。 + + ⚠️ 2026-09-14 修:这里此前把 `asyncio.current_task()`(= `__main__.serve()` + 的**主循环任务**)交给心跳,续租失败时心跳 `task.cancel()` 取消的就是主任务 —— + 后果是 `CancelledError` 从 `run_once()` 一路冒到 `serve()`,进程直接退出, + 而且 `finally` 里的 `offsite_worker.close()`(IMAP logout)在网络已断时会抛 + `ConnectionResetError`,把真正的退出原因也盖掉,最终表现为 + **"Worker 自己退了、退出码 1、日志最后一行是 IMAP 连接被重置"** + (2026-09-14 实测踩到,见 `docs/演示用/记忆系统演示文档-2026-09-14.md` 场景五)。 + + 现在把"这一批"跑在**独立任务**里:心跳取消的是它,调用方只是收到 + "本批被放弃",然后照常回到主循环的下一轮。 + """ + batch = asyncio.create_task(self._process_batch_work(lease, context)) + heartbeat = asyncio.create_task(self._cursor_heartbeat(lease.lease_id, batch)) try: - return await self._process_batch_work(lease, context) + return await batch + except asyncio.CancelledError: + # 两种情况必须分开:① 心跳因为租约丢失取消了这一批 → 本批放弃、Worker 继续; + # ② 外层(Ctrl+C / 进程关闭)取消了当前任务 → 必须原样向上抛,不能吞。 + if batch.cancelled() and not _current_task_is_cancelling(): + logger.warning( + "场外邮件批次被放弃(游标租约已失效)lease_id=%s", lease.lease_id + ) + return False + raise finally: heartbeat.cancel() with suppress(asyncio.CancelledError): diff --git a/docs/演示用/记忆系统演示文档-2026-09-14.md b/docs/演示用/记忆系统演示文档-2026-09-14.md index 0a7f2c4..615a812 100644 --- a/docs/演示用/记忆系统演示文档-2026-09-14.md +++ b/docs/演示用/记忆系统演示文档-2026-09-14.md @@ -228,6 +228,14 @@ Milvus 没配时事件**保持 pending**(而不是标成"已同步")。 **这体现什么**:Agent 与记忆都是"**受理 → 排队 → Worker 执行**"的异步三段式。 **停了不会有报错,只会静默不动** —— 所以"看板正常但什么都没发生"时,第一个要查的就是 Worker。 +> ⚠️ **如果发现 Worker 自己退了**(窗口关了 / 退出码 1 / 日志最后一行是 +> `ConnectionResetError [WinError 10054]` 这种 IMAP 连接被重置),那是踩到了 +> 2026-09-14 修掉的一个真缺陷:场外收件的**游标续租失败时,心跳取消的是主循环任务** +> (而不是"这一批"),于是整个 Worker 被取消、退出时 `finally` 里 IMAP `logout()` +> 又把真正的退出原因盖掉了。现在:**续租失败只放弃当前这一批**,Worker 继续跑; +> 关闭阶段的异常只记日志、不改变退出码。 +> 回归守卫:`tests/unit/worker/test_offsite_mail_worker.py::test_cursor_lease_loss_abandons_batch_without_cancelling_the_worker`。 + --- ## 6. 建议的演示顺序与时长 @@ -282,6 +290,11 @@ A:会,但路径是**异步**的:记忆 → `user_facts`(事实层)→ A:会**真的写入**(记忆、事实、画像版本各加一条)。想隔离就用另一个客户号重跑; 演示完不用清理——历史版本本来就是设计的一部分(画像只新增、不覆盖)。 +**Q:Worker 会不会自己中途退出?** +A:正常不会;曾经会 —— 场外收件的游标续租失败时,心跳误把**主循环任务**取消了, +整个 Worker 就退了(退出码 1,日志末尾是 IMAP 连接被重置)。**2026-09-14 已修**: +续租失败只放弃当前这一批。演示前跑一次 §8 自检,看到两个进程都在就没问题。 + --- ## 8. 演示前自检(30 秒) diff --git a/tests/unit/worker/test_offsite_mail_worker.py b/tests/unit/worker/test_offsite_mail_worker.py index 50f409b..c43253e 100644 --- a/tests/unit/worker/test_offsite_mail_worker.py +++ b/tests/unit/worker/test_offsite_mail_worker.py @@ -1,8 +1,10 @@ from __future__ import annotations +import asyncio from datetime import datetime from decimal import Decimal from pathlib import Path +from types import SimpleNamespace from typing import Any import pytest @@ -549,3 +551,66 @@ async def _cursor(maker: async_sessionmaker[AsyncSession]) -> OffsiteMailCursor: cursor = await session.scalar(select(OffsiteMailCursor)) assert cursor is not None return cursor + + +# --------------------------------------------------------------------------- +# 游标租约丢失:只放弃"这一批",不许把整个 Worker 任务取消掉 +# --------------------------------------------------------------------------- + + +class _LeaseLostSession: + """续租时 `rowcount=0`(租约已不在自己名下)的替身会话。""" + + async def __aenter__(self) -> _LeaseLostSession: + return self + + async def __aexit__(self, *args: object) -> None: + return None + + async def execute(self, _statement: object) -> Any: + return SimpleNamespace(rowcount=0) + + def begin(self) -> _LeaseLostSession: + """`async with session.begin()` 要的是**上下文管理器**,不是协程。""" + return self + + +def _lease_lost_session_factory() -> Any: + return _LeaseLostSession + + +@pytest.mark.asyncio +async def test_cursor_lease_loss_abandons_batch_without_cancelling_the_worker() -> None: + """租约丢失 → 本批返回 False、**调用方任务不被取消**。 + + 回归的是 2026-09-14 实测到的一次"Worker 自己退了":心跳拿到的是 + `serve()` 的**主循环任务**,续租失败时 `task.cancel()` 把整个 Worker 干掉了, + 退出码 1,且 `finally` 里 IMAP `logout()` 的 `ConnectionResetError` 把真正原因盖掉。 + """ + worker = OffsiteMailWorker( + # 租约 1 秒 → 心跳每 1/3 秒续一次,几百毫秒内就能触发"续租失败"这条路径。 + _settings(worker_lease_seconds=1), + session_factory=_lease_lost_session_factory(), + ) + started = asyncio.Event() + + async def never_finishes(_lease: object, _context: object) -> bool: + started.set() + await asyncio.sleep(30) # 处理中;只能被心跳取消 + return True + + worker._process_batch_work = never_finishes # type: ignore[method-assign] + lease = SimpleNamespace(lease_id="lease-1", last_uid="0") + caller = asyncio.current_task() + + result = await asyncio.wait_for( + worker._process_batch(lease, RequestContext(user_id="1", trace_id="t")), # type: ignore[arg-type] + timeout=5, + ) + + assert started.is_set(), "批次应该真的跑起来过" + assert result is False, "租约丢失应表达为'本批被放弃'" + assert caller is not None and not caller.cancelling(), ( + "调用方(主循环任务)不能被取消 —— 否则整个 Worker 会退出" + ) + assert not caller.cancelled()