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()