Files
GaoYiYuan_0626 fffb78a11c docs: 记忆文件同步架构改进交付态(纯文档,零代码)
- 架构设计说明书.md:同步桌面终版——§2.1 补 trace 顺序测试守卫、§4.2.1 无 key 告警闭环、
  §4.3.4 重写为 Redis 双层锁(TTL 30s/Lua 释放/进程内备)、§5.3/§6.10/§7.1/§7.2/§7.3/§7.4/§9.4
  锁与告警状态更新、基线 510
- FRAMEWORK/FLOW/REQUIREMENTS:风控 Tool 四→五(补 query_overdue_alerts)、词表 42→45 口径
- README:risk-m3 tag 注记 + API 契约表改五只读 Tool
- 注:交接文档.md v2.2 同步更新(该文件被 .gitignore 忽略,仅本地)
2026-09-09 20:06:53 +08:00

41 KiB
Raw Permalink Blame History

XingHuo 智能财富管家 · 架构设计说明书

代码基线:分支 risk-control-agent(2026-09-09 架构改进 T-101~T-202 落地后),全量 pytest 510 绿(503 + 锁测试 6 例 + 中间件顺序守卫 1 例)。 编写原则:本文所有论断均落到具体文件、类、函数、行号;凡文档口径与代码实测不一致处,以代码为准并在 §9.4 列出。 与既有文档的关系:FRAMEWORK.md 是选型与状态的速查表,FLOW.md 是端到端链路与 bootstrap,本文是设计决策与理由的系统阐述——回答"为什么是这样"以及"换成别的会怎样"。


0. 阅读指南

你想知道 直接看
系统整体长什么样 §1
一次请求从头到尾怎么走 §2
某个目录/层是干什么的 §3
某个核心机制怎么实现、为什么 §4(按子系统分)
模块之间怎么协作、往哪扩展 §5
设计取舍与替代方案代价 §6
这套设计在什么前提下成立、哪里会崩 §7
有多种做法时代码选了哪种 §8

1. 系统全景

1.1 一句话骨架

Gateway 鉴权 → FastAPI api → service(LangGraph 编排)→ tool / repository → MySQL 双库 / Redis / Milvus / Neo4j → 审计落库

展开成进程内分层(app/):

┌──────────────────────────────────────────────────────────────┐
│  main.py          装配:lifespan + 中间件栈 + 路由 + 错误 handler │
├──────────────────────────────────────────────────────────────┤
│  api/             路由 · 鉴权工厂 · 中间件(薄,无业务逻辑)       │
├──────────────────────────────────────────────────────────────┤
│  service/         编排与业务内核(LangGraph · Tool · 风控引擎)   │
├──────────────┬───────────────────────────────────────────────┤
│  tool/        │ 纯查询适配器(无业务流程,注册表白名单)           │
│  repository/  │ 数据访问(core_ro 只读 / session / risk)         │
│  gateway/     │ 模拟外部交易系统(唯一允许写 Core 的例外)         │
├──────────────┴───────────────────────────────────────────────┤
│  model/(占位)  config/  utils/(trace·db·response·authz)      │
└──────────────────────────────────────────────────────────────┘

1.2 技术选型与硬约束

维度 选型 约束来源
运行时 Python 3.13.14 + FastAPI FRAMEWORK.md §1 已定
Agent 编排 LangGraph 1.2.x + langchain-openai 同上
关系库 MySQL 8.0,双库 jinrong_agent + jinrong_core 同上
缓存 Redis 8.10.1 同上
向量库 Milvus Lite + pymilvus 3.0.1,1024 维 内存 15.4GB 的妥协
Embedding Ollama bge-m3(本地,不出内网) 数据合规
生成 DeepSeek API —

五条不可触碰的红线(MEMORY.md §3,代码层面均有落点):

  1. Core 正式 C1~C5 不可被画像覆盖 → 靠双库物理隔离 + core_ro 纯 SELECT 约定(§4.6.1)
  2. 审计表只 INSERT → 仓储层只暴露 insert_audit_log / insert_input_guard_log(§4.6.4)
  3. 代理人草稿不外发客户 → 无外发通道
  4. 不自动冻户、不自动改风险等级 → 引擎无写 core_customer_risk 的路径
  5. 仅 R-02 适当性可阻断交易 → 代码级物理隔离,见 §4.3.3

1.3 双库与数据分层

层 存储 内容 谁写
L0 jinrong_core 客户主档、风评、持仓、流水、交易 只读(唯一例外见下)
L3 jinrong_agent 风控自有监测层 customer_profile_l3 风控引擎
会话 Redis + jinrong_agent 窗口缓存 + 永久消息 对话链路
RAG Milvus + data/kb/ kb_product_rules 产品规则 scripts/kb/build_kb.py
关系 Neo4j 客户-产品-代理人图谱 scripts/sync/sync_neo4j.py

表规模:jinrong_core 12 张 + jinrong_agent 17 张(共用 11 + Agent 专用 6)。


2. 请求生命周期

2.1 中间件栈:实际执行顺序及其推导

app/main.py 的注册序是 audit(第 81 行)在前、trace(第 87 行)在后。

Starlette 的 add_middleware() 内部执行 self.user_middleware.insert(0, ...),而 build_middleware_stack() 用 reversed(middleware) 由内向外包裹——数组中越靠前越外层。

因此实际栈序(外 → 内):

ServerErrorMiddleware → trace_middleware → audit_middleware → ExceptionMiddleware → 路由

trace 在外层先执行,audit 在内层后执行。

代码自证:audit_middleware.py:9 注释写"注册顺序:trace 中间件之后注册(即执行序在 trace 之内),保证审计时 trace_id/request_id 已绑定"。

⚠️ 这是一处隐式依赖:audit 的正确性完全依赖注册顺序,而顺序由"谁写在文件后面"决定。原靠注释固化(两处都有),现已由测试守卫锁定:tests/test_audit_middleware.py::test_audit_middleware_runs_inside_trace_middleware 断言 audit_log.trace_id 非空——调换两个装饰器的注册顺序该测试即变红(已实测验证)。测试即文档,比注释更可靠。

2.2 四个 ID,各管一层

ID 层级 生成 稳定性
actor_id 人 JWT claims.sub(deps.py:189) 跨请求永不变
session_id 一次聊天 sess-{uuid4().hex[:16]}(chat.py:223) 跨请求复用
trace_id 单次请求 trc-{uuid4().hex[:16]}(trace.py:26) 每请求新建
request_id 单次请求 req-{uuid4().hex[:16]}(trace.py:54) 每请求新建

传播机制是 contextvars(trace.py:18)——中间件 set 一次,之后任何层直接 current_trace() 读取,不需要把 trace_id 当参数层层传递。这是"全链路"的实现基础。

trace_id 与 request_id 分离是 B7 复审 P3-4 的整改结果:前者贯通链路,后者标识单次 HTTP 请求用于幂等对账,注释明确"两者不再互用"(trace.py:19-20)。

2.3 三条主链路

(a) 同步对话 POST /api/chat

trace 中间件(set trace/request)
  → audit 中间件
  → deps.get_auth_context(JWT 验签 + X-Agent-Type 准入)
  → _guard_request  准入 → 空白 → 限流429 → 注入/超长400
  → _prepare_turn   会话解析/创建 + 归属断言
  → memory_service.get_recent(Redis 窗口,miss 回源 MySQL)
  → agent_service.chat(LangGraph: tool → llm → guard)
  → session_repository.insert_turn(user+assistant 同事务)
  → memory_service.append_window
  → audit 中间件落 http_access

(b) SSE 流式 POST /api/chat/stream

同步部分完全一致,差别在最后:_events() 生成器逐块 yield,整轮收完 done 才一次性落库(chat.py:440 insert_turn)。

关键细节:StreamingResponse 的生成器是在中间件 finally 已经 reset_trace() 之后才被消费的,此时 current_trace() 已返回空串。代码在 chat.py:397 提前快照:

trace_id = current_trace() or ""   # 在生成器外面,上下文还在时取值

生成器闭包捕获该变量,后续 insert_turn(..., trace_id=trace_id) 用快照而非现取。若按常规写法在生成器里现取,SSE 落库的 trace_id 会全是空串,消息历史直接断链。

(c) 交易事件 POST /api/simulate/trade

gateway.submit_trade
  → suitability_check(R-02,唯一阻断点;阻断则 return,交易不落库)
  → gateway_repository.insert_trade(写 core_trade)
  → engine.process_trade_event(规则评估 → 预警落库 → L3 → AML)

2.4 统一错误体契约

utils/response.py:28 error_body() 固定输出:

{ "error_code": "...", "message": "...", "trace_id": "...", "request_id": "..." }
  • 4xx:路由 raise ApiError / PermissionDenied(utils/exceptions.py),由 register_error_handlers 转 JSON。
  • 500:由 trace_middleware 兜底(main.py:102-106)——异常发生在中间件以内时,Starlette 的 ServerErrorMiddleware 生成的 500 响应不经过用户中间件,会导致 trace 头丢失;所以这里 catch 后直接产出统一错误体。

3. 分层职责逐一讲解

3.1 api/ — 路由与接缝

文件 职责 状态
deps.py 鉴权工厂 get_auth_context、准入矩阵、归属断言 已实现
chat.py 对话四端点(同步 + 拉侧三端点 + SSE) 已实现
risk.py / simulate.py 风控 API / 模拟网关路由 已实现
audit_middleware.py 每请求 http_access 审计 已实现
auth_adapter.py 宿主→模块 AuthContext 适配 预制件,未接线
admin.py / knowledge.py 仅一行 docstring 空壳

admin.py / knowledge.py 为什么是空的:T-21 拍板"一期只做脚本入库,上传/重建端点不做",知识库入库走 scripts/kb/build_kb.py 离线脚本。二者未在 main.py 中 include_router,不暴露任何端点,也不影响启动。注意 knowledge.py 空 ≠ RAG 能力缺失——检索能力在 service/rag_service.py + tool/kb_tools.py,已可用。

分层纪律:api 允许调 service;禁止直连 Milvus / 在路由里写复杂 SQL。

3.2 service/ — 编排与业务内核

分三类:

  1. 对话编排:agent_service.py(LangGraph 图)、tool_service.py(Tool 编排)、memory_service.py(记忆)
  2. 公共底座:auth_service.py(JWT)、input_guard.py(输入防护)、rag_service.py / milvus_service.py / embedding.py(RAG)
  3. 风控业务:service/risk/*(引擎、规则、预警、AML、L3、评分、锁、Redis 网关、对话 Tool)、service/suitability.py

3.3 tool/ — 纯查询适配器

边界定义:Tool 只做"参数校验 → 调 service/repository → 归一化输出",不承载业务流程。

  • core_tools.py:三个 Core 只读 Tool(query_customer_profile / query_holdings / query_recent_trades)
  • kb_tools.py:search_knowledge
  • document_parser.py / embedding_tool.py / milvus_tool.py:早期占位 stub(各不足 90 字节),T-21 的实际实现落在 service/ 层

3.4 repository/ — 数据访问

  • core_ro.py:CoreReadOnlyRepository,只 SELECT
  • session_repository.py:agent_session / agent_message / agent_tool_call
  • risk_repository.py:风控四表 + 两张审计表(仅 INSERT)

3.5 gateway/ — 模拟外部交易系统

为什么单独成包:它是全项目唯一允许写 Core 的模块(FRAMEWORK.md §3 明确列为分层例外)。生产环境由真实交易系统回调替代,整个包可整体退役。替换边界就是 submit_trade(req) 的契约。

3.6 model/ — 占位层

schemas.py 仅 57 字节 docstring;entities.py 定义了 31 个 ORM 类但运行时不使用——docstring 明确"禁止 create_all、仅 IDE 参照、与 SQL 人工同步"。

实际后果:风控全程用 dict 传参(insert_alert(alert: dict) 等)。见 §6.6 的代价分析。

3.7 config/ 与 utils/

  • config/settings.py:BaseSettings + .env,进程单例 settings = Settings()
  • utils/db.py:get_engine(database) 按库名缓存单例(§4.6.2)
  • utils/authz.py:record_authz_denial,下沉到 utils 的原因是 service 层不得反向依赖 api
  • utils/trace.py、utils/response.py、utils/exceptions.py、utils/desensitize.py
  • utils/logger.py:仅一行 docstring,完全占位——全站日志无配置、无落盘、无结构化(见 §7.3)

4. 核心子系统深讲

4.1 身份与准入体系

职责与位置:api/deps.py + service/auth_service.py + utils/authz.py,处于 api 层最前,是所有业务的前置闸门。

核心实现(deps.py:225-279):

1. Authorization: Bearer 存在 → auth_service.verify_token(验签 + 必填 claims + jti 吊销)
   · Bearer 为全环境首选,大小写不敏感(auth_header[:7].lower() == "bearer ")
2. X-Agent-Type 交叉校验 → 缺失 401 / 域外 400 / 不符矩阵 403(均落审计)
3. 无 Bearer 时:仅当 app_env == "development" 且 jwt_public_key_path 为空
   → 才走 X-Debug-Role / X-Debug-Actor 兜底(双闸门)

准入矩阵 AGENT_ACCESS_MATRIX(deps.py:46-57):

Agent token_type roles
customer customer customer
advisor staff advisor / compliance / ops
analyst staff analyst / compliance
risk staff, service risk_officer / risk_manager / service_risk

deny() 的关键性质(deps.py:119-129):先落审计再抛异常。审计下沉到 utils/authz.record_authz_denial,双写 audit_log + input_guard_log,然后 raise PermissionDenied。即 403 一律留痕,语义 fail-closed。

C5 红线:risk_manager 禁止走对话线(三重保险)

层 位置 行为
① 矩阵层 deps.py:53-56 risk 行 roles 含 risk_manager → HTTP 台账放行
② 对话入口 chat.py:122-124 if agent_type == "risk" and "risk_manager" in auth.roles: deny(...)
③ Tool 层 tool_service.py:138-156 assert_tool_access 只放行 risk_officer,其余 AUTH_403_SCOPE

为什么分三处而不是一处:① 是能力声明(这个角色"能进哪个 Agent"),② 是业务线约束(同一 Agent 的两条线权限不同),③ 是纵深兜底(前两层被绕过时仍 fail-closed)。代价是权限逻辑分散,改动需三处同步(§7.4)。

权衡

决策 选了 替代方案 代价/收益
身份来源 JWT RS256(生产)+ HS256(dev) 每次查库 session JWT 免去 DB 往返,但吊销依赖 Redis jti 黑名单,Redis 挂则吊销失效(fail-open)
dev 通道 debug 头兜底 一律强制 JWT 演示/CI 免签方便;但双闸门一旦误配(生产设 app_env=development 且公钥为空)即放开无签名身份
审计失败 fail-open(仅 warning) fail-closed 可用性优先;合规强场景应切 fail-closed,已登记待决

4.2 对话编排子系统

4.2.1 LangGraph 图结构

agent_service.py:173-183:

def build_graph():
    graph = StateGraph(ChatState)
    graph.add_node("tool", tool_node)
    graph.add_node("llm", llm_node)
    graph.add_node("guard", guard_node)
    graph.add_edge(START, "tool")
    graph.add_edge("tool", "llm")
    graph.add_edge("llm", "guard")
    graph.add_edge("guard", END)
    return graph.compile()

这是一个线性图,没有条件边。四 Agent 的差异不是靠图分支,而是靠数据驱动:

  • _SYSTEM_PROMPTS(:40-56)按 agent_type 固化四种角色边界
  • _INTENT_KEYWORDS(tool_service.py:73-96)按 agent 分组决定命中哪些 Tool
  • 结果:analyst 无关键词 → Tool 节点空转 → 纯 LLM;risk 命中风控 Tool;customer/advisor 命中 Core RO + KB

ChatState 字段(:59-77):agent_type / history / user_message / messages(Annotated + operator.add) / reply / has_disclaimer / session_id / trace_id / actor / customer_id / tool_results。messages 用 operator.add 归并以避免节点覆盖历史。

多解说明:LangGraph 支持 add_conditional_edges 做真正的图分支。本项目选择了数据驱动而非图分支——理由是四 Agent 共享相同的"取数 → 生成 → 加免责声明"骨架,差异只在提示词与可用工具集,抽成一张图比维护四张子图成本低。代价是无法表达"风控对话需要额外的复核节点"这类结构性差异,若将来某 Agent 需要独立节点拓扑,需拆图。

无 key 降级(:148-150):if not settings.deepseek_api_key: reply = _degraded_reply(state),降级回复携带 Tool 查询摘要并加 _DEGRADED_PREFIX,不抛异常。好处是演示链路不断;原风险是生产漏配 key 会"看起来正常"——该风险已由 T-107 闭环:main.py 的 lifespan 在 DEEPSEEK_API_KEY 缺失时打印显式告警(函数内延迟导入 _DEGRADED_PREFIX 防循环导入),启动日志即提示将走降级回复,不再是静默失败。改进已完成,见 §7.3 P0-3。

4.2.2 流式执行时序

stream_chat(:275-319):Tool 同步跑完,再流式推 LLM 文本(非交错)。

tool_node(state) → tool_results
_compose_messages(把 Tool 结果注入上下文)
llm.stream(messages) → 逐块 yield ("delta", text)
yield ("done", 完整正文)

为什么不交错:Tool 结果要作为 LLM 输入且需落 agent_tool_call 留痕,必须同步完成;交错会让 LLM 在 Tool 未返回时"瞎编"——金融场景不可接受。代价是首字延迟 = Tool 耗时 + LLM 首字(§7.3)。

4.2.3 Tool 编排

注册表三层懒加载(tool_service.py:115-132):

core_tools.TOOL_REGISTRY  →  chat_tools.RISK_TOOL_REGISTRY  →  kb_tools.KB_TOOL_REGISTRY

懒导入避免模块加载耦合。ToolSpec 字段:func / description / requires_customer / param_whitelist / int_bounds [/ skip_access_check]。

意图识别是关键词匹配,不是 LLM(match_intent:100-112):按 agent 分组遍历 (tool_name, (kws...)),any(k in message for k in kws),命中即停(单意图)。

权衡:LLM 路由更灵活但增加一次调用延迟与成本,且意图判定不可复现(不利于审计);关键词零延迟、可复现、单测可锁。代价是用户换说法就漏触("我的基金" vs "持仓"),且不支持多意图。

归属校验(assert_tool_access:138-156):

角色 规则
risk_officer 全量放行
customer actor_id == customer_id
advisor core_ro.is_advisor_assigned(actor_id, customer_id)
其余 AUTH_403_SCOPE fail-closed

skip_access_check=True 仅 KB Tool(公开知识无客户对象),开放范围由意图层约束(仅 customer/advisor 配词)。

入参归一化(_normalize_params:162-196):白名单 + 整数边界钳制,未知键直接报 TOOL_BAD_PARAM 而不静默丢弃——静默丢弃会让 LLM 以为参数生效了。

留痕:run_tool 执行后无条件调 _audit_tool_call 写 agent_tool_call(含 tool_input / tool_output / status / latency_ms),落库失败降级 warning 不阻塞。

4.2.4 记忆分层

  • Redis 窗口:key sess:{agent}:{session_id}:msgs,WINDOW_SIZE=20,TTL=2h,RPUSH + LTRIM 保最近 20 条
  • MySQL 权威:insert_turn 落 agent_message
  • 回源:get_recent 先 lrange,miss 或异常 fallback list_messages

设计要点:Redis 是缓存不是权威,丢了可重建。这与 Redis 在限流、发布订阅处的 fail-open 口径一致。

4.2.5 输入防护(T-03)

顺序(chat.py:144-206 _guard_request):

Agent 准入 → 空白 → 限流 429 → 注入/超长 400 → 归属 → 会话
  • 注入词表 45 条(input_guard.py:43-94;注:旧文档写 42 条,实测 45 条 = 指令覆盖 18 + 角色重置 11 + 系统提示泄露 9 + 越权诱导 7)。短语精确子串匹配,不做模糊正则(防误杀"忽略这只股票"这类业务句)。
  • oversize:MESSAGE_MAX_LENGTH = 4000
  • 限流:Redis 固定窗口,actor 级 ratelimit:{agent}:{actor},默认 30/min,fail-open

为什么限流 fail-open 而鉴权 fail-closed:限流是可用性保护,不是安全边界;安全边界(鉴权/归属/注入)必须 fail-closed。这个区分是刻意的。代价是 Redis 故障时恶意用户可绕过限流(§7.3)。

为什么限流排在内容防护之前:计数需覆盖全部请求(含将被注入拦截的),让重复攻击者快速收敛到 429。

4.3 风控事件线

4.3.1 完整调用链

app/gateway/trade_gateway.py:86  submit_trade
  ├─ :122  suitability_check          ← R-02 唯一阻断点
  │        └─ blocked → record_suitability_alert → return(交易不落库)
  ├─ :167  gateway_repository.insert_trade(写 core_trade)
  └─ :171  engine.process_trade_event(同步)
              engine.py:114
              ├─ core.list_trades_range(取当日流水)
              ├─ :139  run_rules → RISK-001~005 评估
              ├─ :145  rule_concentration(RISK-006,输入域为持仓,故不并入 run_rules)
              ├─ :150  record_trade_alerts(聚合落 risk_alert)
              ├─ :160  upsert_profile_l3(L3 监测层)
              ├─ :174  match_customer(AML)
              └─ :177  record_aml_alert

4.3.2 规则与评分

service/risk/rules.py 全部为纯函数,run_rules:201 编排:

规则 函数 阈值来源
RISK-001 大额 rule_large_amount:91 risk_large_amount = 500000
RISK-002 当日累计 rule_daily_total:102 risk_daily_total = 500000
RISK-003 频繁交易 rule_freq_trade:113 risk_freq_count = 3
RISK-004 试探性 rule_probe_pattern:127 5 分钟内 3 笔 × 40 万
RISK-005 小额铺垫 rule_small_then_large:145 1 万 × 3 笔
RISK-006 集中度(C4) rule_concentration:164 concentration_threshold = 0.80

评分是静态映射(RULE_SCORES:18):001/002 = 70、003 = 50、004/005 = 80、006 = 60;alert_service.py:22 TIER_SCORE sustainability = 90、aml = 95。

⚠️ service/risk/scoring.py 不是实际评分器——recompute_customer_score 直接 raise NotImplementedError("R-05 dynamic scoring not implemented"),是签名冻结的预留桩。L3 的 risk_score 一期恒为 NULL。

4.3.3 为什么只有 R-02 能阻断

代码级物理隔离,两套路径不共用出口:

  • R-02 在网关层(submit_trade:126):if result.blocked: ...; return {blocked: True},此时 insert_trade(:167)尚未执行
  • RISK-001~006 全在 process_trade_event 内,命中只调 record_trade_alerts 出 risk_alert(status=pending_review);engine.py:122 注释明确事件线全程 return {"blocked": False, ...},引擎无任何阻断出口

判定权唯一归属 core_ro.check_suitability 返回的 blocked 字段。

4.3.4 L3 监测层

customer_profile_l3 表,字段 monitor_tier(normal < watch < high)、monitor_tags(追加合并)、risk_score(一期 NULL)、last_alert_id、computed_at。

并发控制:进程内锁 + 乐观锁 update_l3(expected_computed_at) + 3 次重试。

⚠️ service/risk/locks.py 已完成「双层锁」改造(T-201):run_locked(key, fn) 优先用 Redis SET NX EX 抢分布式锁(redis_gateway.acquire_lock / release_lock 由 Lua 原子释放、只删自己的锁),抢到则 fn(locked=True);Redis 等待超时(2.0s)→ fn(locked=False) 沿用原降级语义(绝不抛异常);Redis 不可用(含测试 Fake 缺方法的 AttributeError)→ 安全退回进程内 threading.Lock,最坏退化为单进程行为,不会更差。常量:LOCK_TTL_SECONDS=30 / LOCK_TIMEOUT_SECONDS=2.0 / _LOCK_KEY_PREFIX="lock:" / _RETRY_INTERVAL=0.05。三处调用点(alert_service 聚合 ×2 含 R-02 阻断点、profile_l3 ×1)签名与调用方式不变,多实例部署不再失效(§7.2 已更新)。

4.3.5 C4 / C5 / C6 追加需求

需求 实现 核心逻辑
C4 集中度 rules.py:164 + engine.py:144 R4+R5 市值占比 ≥ 0.80;holdings_truncated 视同达标(保守口径)
C5 时效升级 escalation_service.py:59/84 普通 4h/24h、AML 1h/4h;只写 payload.escalation_level,不改 status;幂等仅升不降
C6 行为链 agent_behavior_service.py:108/289 A 诱导调仓(赎回 2h 内申购不同产品,3 次)/ B 越权试探(5 次/72h)/ C 越权查询(10 次/24h);出代理人维度独立单

C6 的单据被 find_pending_event_alert 显式排除,避免误并入客户单。

4.3.6 引擎异常补偿

trade_gateway.py:185-203:交易已提交但引擎异常 → 记 risk_engine_error 审计 + engine_error=true,由 scripts/demo/rebuild_alerts.py 补偿重放。

这是最终一致而非原子:写 core_trade(core 库)与写 risk_alert(agent 库)是两个独立事务,双库无法纳入同一事务(§6.2)。

4.4 风控对话线

service/risk/chat_tools.py:328 RISK_TOOL_REGISTRY 实际注册 5 个:

Tool 说明 requires_customer
query_overdue_alerts 超期未处置预警(C5) False
alert_query 预警查询(risk_officer 可查全量) False
customer_context 客户风控上下文 —
suitability_check 适当性校验(只落 risk_suitability_log) —
aml_lookup AML 名单查询 —

(query_agent_behavior:276 已定义但未注册,仅由意图层使用。旧文档"四个 Tool"为过时口径。)

只读边界:全部无 insert alert / 处置 / 改表动作。事件线独占出预警单、人工处置、AML 扫描。suitability_check 复用 service/suitability.py,只落审计日志。

4.5 RAG 链路

环节 实现
Embedding service/embedding.py,Ollama bge-m3,1024 维
向量库 service/milvus_service.py,collection kb_product_rules,schema 含溯源字段(source_doc_id / source_version / effective_date)
检索 search_kb:119-166,filter 内建 effective_date <= today(只返回已生效文档)
输出 rag_service.search_knowledge:53-77 返回 chunks + source_refs(去重溯源清单)

Ollama 失败 = 报错,不降级(embedding.py:7-10):禁止零向量或截断,一律抛 EmbeddingError,由 run_tool 转 TOOL_ERROR。

为什么这里反而 fail-closed:假装"没有结果"会让 LLM 编造产品规则——在金融合规场景下,暴露"查询未完成"比给出错误答案安全得多。与限流的 fail-open 形成对比,说明降级策略是按"失败后果"逐项决定的,不是全局统一口径。

4.6 数据访问与审计

4.6.1 core_ro 只读契约

CoreReadOnlyRepository(core_ro.py:51-498),绑定 get_engine(settings.mysql_core_database),全部方法用 self._engine.connect() 执行 text() SELECT。

⚠️ 只读靠约定,不靠强制:类注释声明"仅 SELECT",但没有 DB 级只读账号,没有 SQL 层拦截。任何持有 core engine 的代码都能 INSERT。改进见 §7.1。

check_suitability(:80-199)—— 全系统唯一的阻断判定核:

SQL 取回 core_suitability_rule 的 matrix_match_result(LEFT JOIN ... ON sr.customer_risk_code = r.risk_code AND sr.product_risk_code = p.min_risk_code),然后按序判定:

1. 客户/产品缺失           → forbidden / SUIT_NOT_FOUND
2. risk_is_expired (FM-03) → 阻断 / SUIT_RISK_EXPIRED
3. professional 投资者      → 豁免 / professional_exempt (JR-AST-PRO)
4. matrix == forbidden/null → 阻断 / SUIT_RISK_MISMATCH (JR-AST-012)
5. 否则                    → allowed / allowed_with_disclosure
6. age >= 70 且产品 >= R3  → 阻断 / SUIT_AGE_CONFIRM (FM-01)

结果经 _suitability_result(:201-246,21 字段)与 risk_suitability_log 一一对应。

矩阵判定是数据驱动的(查表而非硬编码 if-else),新增风险等级组合只需改 core_suitability_rule 表数据。这是阶段一 AL-01~08 的整改成果(旧的 SUIT-001~008 纯函数矩阵已退役)。

4.6.2 双库实现与代价

utils/db.py:23-36 get_engine(database) 按库名缓存 Engine 单例(字典 + threading.Lock)。两个库各自持有独立 engine。

代价:跨库无法 JOIN。代码从不跨库联表,而是在 Python 层分别取数合并(如 submit_trade 同时调 core_ro 与 risk_repo)。更进一步,两库的写操作无法纳入同一事务——这是 §4.3.6 最终一致性的根因。

4.6.3 会话三表与 insert_turn

  • agent_session:session_id(UK) / trace_id / agent_type / actor_id / customer_id / status(active|closed|blocked) / metadata
  • agent_message:session_id / trace_id / seq_no / role / content / has_disclaimer
  • agent_tool_call:session_id / trace_id / tool_name / tool_input / tool_output / status / latency_ms

insert_turn(:218-275)同事务落库:with self._engine.begin() 内先取 MAX(seq_no)+1,再连续 INSERT user / assistant。杜绝"user 落了 assistant 没落"的半截历史污染 LLM 上下文。

close_session(:78)用条件 UPDATE ... WHERE status='active',幂等返回 rowcount > 0。

⚠️ agent_message 的 (session_id, seq_no) 是普通索引(KEY idx_session_seq)而非唯一索引,且 insert_turn 取 seq 的 SELECT 无 FOR UPDATE,并发同会话会静默产生重号消息。见 §7.3。

4.6.4 审计表

  • audit_log:trace_id / event_type / agent_type / actor_id / customer_id / rule_id / input_summary(JSON) / decision / risk_score / handler_* / created_at —— 只 INSERT
  • input_guard_log:trace_id / session_id / agent_type / actor_id / guard_type(ENUM 四值) / raw_excerpt / action(blocked|sanitized|passed)

"双写"在项目里有两处含义,注意区分:

  1. 越权双写:deny 同时写 audit_log + input_guard_log(utils/authz.py:44-76,无事务)
  2. 事件双写:每笔交易同时落 risk_alert + audit_log;处置时 handle_alert_with_audit:309 状态变更与审计同事务(防无痕状态变更)

5. 模块协作:依赖方向、接口边界与扩展点

5.1 依赖方向规则

允许:  api → service → tool / repository / model / config
禁止:  api 直连 Milvus;api 写复杂 SQL;tool 写业务流程;repository 写 Core
例外:  app/gateway/ 为模拟外部系统模块,其 gateway_repository 可 INSERT core_trade
反向依赖禁止:service 不得依赖 api(这就是 utils/authz.py 存在的原因)

5.2 接口边界速查

边界 契约 扩展方式
宿主 → 模块鉴权 auth_adapter.from_host_auth(未接线) AL-09 合并后接入
api → service Pydantic 请求模型 + AuthContext 新增端点只需加路由
对话 → Tool get_registered_tool 注册表白名单 新增 Tool = 注册表加一项 + 意图词表加词
事件 → 规则 rules.py 纯函数 + RiskThresholds.from_settings 阈值改 .env,新规则改代码
模块 → Core core_ro 纯 SELECT 只读,不动
阻断判定 core_suitability_rule 表数据 改数据即可,不改代码

5.3 主要扩展点

  1. 新增对话 Tool:注册表加一项(core / risk / kb 三选一)+ _INTENT_KEYWORDS 配词
  2. 新增风控规则:rules.py 加纯函数 + run_rules 编排 + RULE_SCORES 映射
  3. 接入动态评分:实现 scoring.py:recompute_customer_score(签名已冻结),同步放开 L3 risk_score
  4. 替换真实交易系统:整体退役 app/gateway/,实现 submit_trade 契约
  5. 切分布式锁:locks.py 换 Redis SET NX(接口不变,注释已预留)—— T-201 已完成:现为 Redis 双层锁,三处调用点未改动(见 §4.3.4)
  6. 前端接入:后端接入面已就绪(方案 B 三端点 + 方案 C SSE),web/ 待 init

6. 设计权衡总表

# 设计决策 采用方案 替代方案 采用方案的收益 替代方案的代价
6.1 四 Agent 差异表达 一张图 + 数据驱动(提示词+意图表) 四张子图 + 条件边 维护成本低、骨架统一 无法表达结构性差异(如风控复核节点)
6.2 双库 jinrong_core + jinrong_agent 物理隔离 单库 + schema 前缀 Core 红线靠物理隔离兜底 跨库不能 JOIN、无分布式事务 → 最终一致
6.3 意图识别 关键词匹配 LLM function calling 零延迟、可复现、可单测 换说法漏触、不支持多意图
6.4 流式与 Tool Tool 同步跑完再流式推 Tool 与 LLM 交错 结果确定性、可留痕 首字延迟 = Tool 耗时 + LLM 首字
6.5 只读保障 约定 + 分库 DB 只读账号 + ORM 强制 实现简单 无强制,任何持有 engine 的代码都能写
6.6 数据传参 dict Pydantic / TypedDict / ORM 敏捷,快速迭代 无编译期校验,拼写错误运行时才暴露;entities.py 31 个类闲置
6.7 限流降级 fail-open fail-closed Redis 故障不阻断业务 故障时恶意用户可绕过
6.8 Embedding 降级 fail-closed(抛错) 返回空结果 避免 LLM 编造产品规则 依赖可用性,Ollama/Milvus 故障直接报错
6.9 会话落库 整轮一次性落(同事务) 边生成边落 无半截历史 断连/异常整轮丢失(Tool 留痕仍在)
6.10 锁 双层:Redis SET NX EX 为主(TTL 30s、超时 2s 降级 fn(locked=False)),进程内 threading.Lock 为备(Redis 不可用时退回) 纯进程内锁 多实例部署不再失效,Redis 故障不阻塞业务、不抛异常 引入 Redis 依赖;需防临界区超 TTL 自动解锁(已用 30s TTL 缓解)

7. 前提假设、适用边界与改进空间

7.1 前提假设(不成立则设计失效)

  1. 单进程部署——locks.py 已升级为 Redis 双层锁(多实例不再失效),但 insert_turn 的 MAX(seq_no)+1 竞态(注释承认"并发写锁归后续")仍以此为前提
  2. Core 是模拟库且无人写入——只读靠约定,无 DB 级账号保护
  3. 用户用语高度收敛——关键词意图识别的前提
  4. 演示/内网环境——dev debug 通道、_degraded_reply 静默降级均在公网生产下变危险
  5. SQLite 测试与 MySQL 生产行为一致——项目大量 Python 端垫片(_as_date / _is_expired / str(amount) / _jsonable)源于此

7.2 适用边界

能力 边界
分布式锁 ✅ 双层(Redis 主 + 进程内备,T-201 已落地);⚠️ insert_turn 的 seq 竞态仍在(§7.3 P0-2)
动态风险评分 ❌ scoring.py 是桩,L3 risk_score 恒 NULL
跨库事务 ❌ 双库最终一致,靠 rebuild_alerts.py 补偿
知识库管理 API ❌ 一期只做脚本入库
图数据库 ⚠️ Neo4j 仅由 sync_neo4j.py 灌图,运行时未参与主链路
Milvus Lite ⚠️ 演示规模,不适合生产 HA
convert 交易 ❌ 网关主动拒绝(DB ENUM 支持,代码收窄)

7.3 主要改进空间(按优先级)

P0 — 安全与正确

  1. 为 core_ro 配 DB 级只读账号:当前红线仅靠约定,配只读账号是从根基卡死(§4.6.1)
  2. agent_message 并发重号:索引非唯一 + SELECT 无锁 → 并发同会话静默重号。修法:加 UNIQUE(session_id, seq_no) 并同步给 insert_turn 加 IntegrityError 重试(只加索引会让并发变 500)
  3. ✅ 无 key 降级改显式告警:_degraded_reply 静默返回,生产漏配会"看似正常" —— T-107 已完成:main.py 启动即打印告警(见 §4.2.1)
  4. 审计/限流 fail-open 策略可配置:合规强场景应能切 fail-closed

P1 — 可维护性

  1. ✅ 中间件顺序加测试守卫:原靠注释固化,调换装饰器会静默丢 trace_id —— T-202 已完成:test_audit_middleware_runs_inside_trace_middleware 断言 audit_log.trace_id 非空,调换注册顺序即变红(§2.1)
  2. 权限逻辑收拢:C5 红线分散在矩阵 / _assert_chat_entry / assert_tool_access 三处
  3. 注册表合并:TOOL_REGISTRY 与 KB_TOOL_REGISTRY 分裂
  4. model/ 层启用:逐步用 Pydantic / TypedDict 约束仓储入参

P2 — 体验与能力

  1. 意图识别升级:关键词 → 同义词扩展 + 多意图(谨慎用 LLM,会损失可复现性)
  2. 流式首字优化:Tool 执行期间先推"查询中"占位帧
  3. 注入词表运营化:45 条硬编码,宜配置化 + 持续红队补充
  4. 限流改滑动窗口:当前 INCR + 首命中 EXPIRE 非原子
  5. 日志基建:utils/logger.py 仅一行 docstring,需 basicConfig + 落盘 + trace_id 注入

7.4 已知待决事项

  • 审计降级 fail-open 是否切 fail-closed(已登记)
  • chat 链路 risk_suitability_log.actor_id 落 SYSTEM 待评估
  • 前端 React 多 Agent 入口 web/ 归属待拍板

8. 多解处:代码实际采用哪一种

议题 候选 实际采用 原因
四 Agent 是否各用一张 LangGraph 图 A 四张子图 / B 一张图数据驱动 B 骨架相同,差异仅在提示词与工具集
Tool 与 LLM 是否交错流式 A 交错 / B Tool 先同步跑完 B 金融场景不接受 LLM 在 Tool 未返回时编造
意图识别用 LLM 还是规则 A LLM function calling / B 关键词 B 零延迟、可复现、可单测锁定
适当性矩阵硬编码还是查表 A 硬编码 if-else / B 查 core_suitability_rule B 阶段一 AL-01~08 整改成果,新增组合只改数据
失败降级统一 fail-open 还是 fail-closed A 统一 / B 按失败后果逐项定 B 限流 fail-open(可用性保护),embedding fail-closed(防编造)
阻断逻辑放引擎还是网关 A 引擎统一判定 / B 网关专用分支 B 只有 R-02 能阻断,物理隔离防误扩权
数据传参用 ORM 还是 dict A ORM / TypedDict / B dict B 快速迭代优先,代价是失去编译期校验
新增/变更双写是否包事务 A 统一包 / B 按后果定 B 新增可丢留痕;变更(处置)必须同事务,否则无痕状态变更

9. 附录

9.1 关键文件索引

关注点 文件
装配与中间件栈 app/main.py
鉴权与准入 app/api/deps.py、app/service/auth_service.py、app/utils/authz.py
对话四端点 app/api/chat.py
LangGraph 编排 app/service/agent_service.py
Tool 编排 app/service/tool_service.py、app/tool/core_tools.py、app/service/risk/chat_tools.py、app/tool/kb_tools.py
输入防护 app/service/input_guard.py
风控事件线 app/gateway/trade_gateway.py、app/service/risk/engine.py、app/service/risk/rules.py
适当性判定 app/repository/core_ro.py::check_suitability、app/service/suitability.py
数据访问 app/repository/core_ro.py、session_repository.py、risk_repository.py
双 ID 贯通 app/utils/trace.py
并发原语 app/service/risk/locks.py、app/service/risk/redis_gateway.py

9.2 风控阈值清单(config/settings.py:48-60)

risk_large_amount=500000、risk_daily_total=500000、risk_freq_count=3、risk_probe_window_minutes=5 / risk_probe_count=3 / risk_probe_amount=400000、risk_small_amount=10000 / risk_small_count=3、risk_aml_default_threshold=0.85、concentration_threshold=0.80、guard_rate_limit_max=30/60s

9.3 表清单

  • 共用 11 张:agent_session / agent_message / agent_tool_call / audit_log / input_guard_log / customer_advisor_rel / customer_profile_l1 / customer_profile_l2 / customer_profile_l3 / risk_alert / risk_suitability_log
  • Agent 专用 6 张:customer_threshold_config / customer_notify_log / advisor_draft / compliance_hit_log / analytics_query_log / risk_aml_list
  • Core 12 张:scripts/core/01-ddl.sql

9.4 文档口径与代码实测不一致处(以代码为准)

项 旧文档口径 代码实测
注入词表条数 42 条 45 条(input_guard.py:43-94)
风控对话 Tool 4 个 5 个(chat_tools.py:328,另有 1 个定义未注册)
Agent 专用表 5 张 6 张(SQL 注释写 5,实建 6)
scoring.py 未明确 NotImplementedError 预留桩,非实际评分器
locks.py 未明确 双层锁(T-201 已落地):Redis SET NX EX 为主、进程内 threading.Lock 为备,三处调用点不变
convert 交易 未提及 DB ENUM 支持,网关主动拒绝

本文基于 risk-control-agent 分支的代码实测撰写,已同步 2026-09-09 架构改进(T-101~T-202:双层锁 / 无 key 告警 / 中间件顺序守卫 / 口径勘误,pytest 510 绿)。后续改动请先更新 docs/memory/ 再同步本文。