一、装配 projection_cleaner 新增 app/service/projection_cleanup_service.py,并在**组装层**(app/worker/__main__.py) 注入。此前该客户端一直未装配,清理链路只记录降级(skipped_no_client),实际后果是 **销户后图里仍留着偏好关系**——投顾仍能通过关系网络看到这个人。 清理策略是"以画像为准"而不是按标识直接删边:先删掉该记忆对应的长期事实,再重建画像, 最后用对账修复让图跟着画像收敛。这样即使一条记忆影响多条派生边也能删干净,不依赖 "记得它当初投影成了什么"。Milvus 侧按 memory_uuid 删除;集合不存在(语义召回未启用) 时视为无需清理,客户端不可用则如实报告未清理,绝不伪造成功。 放在组装层而不是 runtime 内部兜底:组件内部给默认实现会把"尚未装配"这一事实悄悄盖住, 而"未注入即显式降级并留痕"是 runtime 刻意保留的语义。最初的改法写成内部兜底,被 5 个 既有单测拦下——那些测试是对的,因此改为在入口处注入。 二、顺带修复:事实消失后画像字段残留 实测:让 preference:horizon 记忆失效并清理后,user_facts 与图边都正确清除,但 fin_customer_profile.investment_horizon 仍是"长期(5年以上)"。根因是 rebuild_profile 只写"本轮有新值"的字段;事实被删后该字段没有新值,旧值就留在画像里——于是记忆已经作废, 投顾还能看到一个客户从未授权继续生效的投资期限。 改为:本轮没有对应事实时**清空**该字段;investor_type 例外——问卷是唯一权威,本轮没有 问卷记录时保持原值(该列 NOT NULL,且清空会在重测空档抹掉开户时的等级)。 三、实测 · 让 preference:horizon 失效后清理:cleaned=True, detail=profile_rebuilt; graph_cleaned; vector_collection_absent; · user_facts 2→1;图边由 ['HAS_GOAL','PREFERS'] 变为 ['PREFERS']; · 画像 investment_horizon 由"长期(5年以上)"变为 None;investor_type 保持 C2、 risk_tags 不受影响(与 horizon 无关); · 恢复记忆状态后重建,画像与图边均正确复原; · ruff 通过、mypy 113 文件无错、unit+contract 447 passed、integration 29 passed。
65 lines
2.6 KiB
Python
65 lines
2.6 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.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
|
|
),
|
|
)
|
|
try:
|
|
while True:
|
|
try:
|
|
worked = await runtime.run_once()
|
|
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 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()
|