共同祖先 bbf623a;主干 54 个提交、118 个文件;本线 25 个文件;9 个冲突文件。 主干这次把 **ZSY 的整条投影实现合进来了(PR #7)**,而本线此前的提交正是 移植并修正同一套代码 —— 因此冲突的本质是"同一功能两份实现并存",取舍错了会把 已修好的缺陷又带回来。逐项取舍与理由见 `docs/39-主干合并对策记录.md`。 ## 取舍(9 个冲突) 取本线: - `app/infrastructure/milvus_profile_projection.py` —— 主干是 ZSY 原版,含两处必炸点: ① `customer_id` 要求 int 而本仓所有生产者都写 `str` ⇒ 每个事件必然失败; ② 不可投影的 `memory_key` 直接 raise ⇒ 一条 `constraint:` 记忆毒死整客户整批。 本线版已放宽为「接受纯数字字符串」与「跳过并留痕」。 - `memory_sync_outbox_worker.py` / `conversation_privacy.py` / `risk_questionnaire.py` —— 代码逐行一致,仅注释与说明文字详略不同(`risk_questionnaire.py` 两边**独立做了 完全相同的修复**,都改成 re-export `app.model.profile`)。 - 两个投影测试文件 —— 本线是他那份的**超集**(4→10、4→5 例,包含他全部用例)。 两边合并: - `app/worker/runtime.py`:`__init__` 两边各加一个参数,都要。 - `app/service/agent/implementations/customer_service.py`:import 取并集; `COMPANY` 取主干的「奶龙基金责任有限公司」("奶龙"是本项目实际品牌名,主干多处出现), `HOTLINE`/`SERVICE_HOURS` **取本线的修复**(主干仍是占位符 `400-XXX-XXXX`, 本线已改为引用 `customer_service_rules` 的唯一来源 —— 这是 A1 缺陷修复, 否则同一客服给客户两个不同号码)。 - `AGENTS.md`:表数/Agent 清单取主干(90/89、7 个 Agent),本线的 `-X utf8` 与两条 outbox 易错点保留,测试基线按合并后实测重算。 ## 消费端只保留一套(本次最重要的一处) 合并后曾出现**两套消费者读同一个 `memory_sync_outbox`**:`__main__.py`(PR #7) 与 `runtime.consume_profile_projections()`(本线),而**两者的 neo4j handler 不同** —— 前者用 ZSY 的 `Neo4jProfileProjection`(按客户各建私有节点), 后者用主干 `ProfileGraphProjectionService`(共享 tag 节点、只投影已确认事实)。 同一事件被谁领到结果不定,等于"同一事实在图里有两种说法",正是**方案 A 要避免的状态**。 现只保留 runtime 那一套(带 `memory_sources` 兜底、neo4j 复用主干服务), 删除 `__main__.py` 的重复接线;装配入口职责仍在该文件(注入 `relationships` / `projection_cleaner`),Milvus 客户端由 `bootstrap` 工厂惰性构造、缺配置时显式降级。 副作用:`app/infrastructure/neo4j_profile_projection.py` 不再被生产代码引用,成为 **死代码**(本线未删,属架构师线,其单测仍在)—— 待架构师决定删或明确分工。 ## 顺带修掉的 3 个继承缺陷(主干同样存在,PR #7 后未整套复跑故未发现) 1. `tools/seed_test_rbac.py` **少建 `review_t`(9004)账号** —— 两个集成测试都依赖它 ("账号存在但无权限应返回 200 空集而非 404"、`PLACEHOLDER_ACCOUNTS`)。 同时把用户↔角色绑定从 `zip(..., strict=True)` 改为**显式配对表**:原写法隐含 "USERS 与 ROLES 一一对应",一加不绑角色的账号就 ValueError、整个种子跑不完 (commit 在最后,外部表现是"什么都没发生")。 2. `CustomerProfileCandidateService._write_profile_snapshot` **漏写 `current_customer_id`** —— 该列不是生成列而是普通可空列 + 唯一键 `uk_profile_snapshot_current`, 不写则唯一键形同虚设(多个 NULL 不冲突),且旧当前版本也没清该列、补写就会撞键。 现旧值置 None、新值显式写入(与 `ProfileGenerationService._clear_current` 一致)。 3. 集成测试前置未记录 —— 13 个登录/RBAC 用例因 401 而红,实为"测试账号不存在", 跑 `seed_test_rbac.py` + `set_user_password.py` 后转绿;已在 `AGENTS.md` 记明, 避免被误判成代码缺陷。 ## 文档 - 新增 `docs/39-主干合并对策记录.md`(逐文件取舍 + 理由 + 遗留) - `docs/37` 订正一处过时说法:曾写 `current_customer_id` 无人使用且故意不映射, 实际 `app/model/profile.py` 已映射且有人使用(详见该文档 §6.2 的订正块) - 文档编号:主干已占 29–36,本线两份文档让号至 `docs/37`、`docs/38` ## 验证(合并后实测) - `pytest tests`(全量)→ `2 failed, 1396 passed, 2 skipped` - `pytest tests/integration` → `102 passed, 1 skipped`(修上述 1、2 后从 15 failed 归零) - `mypy app` → `Success: no issues found in 245 source files` - `tools/audit_schema.py` → 89 张业务表无缺失/意外(未改动任何表结构) - `tools/check_authoritative_docs.py` → 52 份文档无编号冲突 - `tools/check_rbac_seed_consistency.py` → 通过 那 2 个失败是既有环境项(`test_offsite_document_recognition_adapter.py` 断言请求体 中文原文而 httpx 序列化成 `\uXXXX`),与本次合并无关。
85 lines
4.2 KiB
Python
85 lines
4.2 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:
|
||
await offsite_worker.close()
|
||
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()
|