## 现象(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 会不会自己中途退出"。
93 lines
4.8 KiB
Python
93 lines
4.8 KiB
Python
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()
|