# 实现方案 · 风控追加需求 v1.1(FR-8 / FR-9 / FR-10 → C4~C6) > 依据:PRD v1.1 §4A(已冻结并入)· 附-风控规则表 v1.1 · 开发计划 v1.2 C4~C6 任务行 + 挂账 #1~#9 > 分支 `risk-control-agent`(旧称 feature/risk 已过时)· 基线:T-21 完成后(425 用例 collect 已验证,全量执行绿待跑) > 本文为编码依据,函数签名/文件路径可直接照抄;与 PRD 冲突时以 PRD 为准。 > **评审修订记录(v1.1 · 2026-09-07)**:经独立 AI 评审 FAIL → 修订本文。P0-1 C4 聚合锚点与 C6 行为链单碰撞(§2.1/§2.5 锚点排除 agent_behavior);P1-1 升级 payload 写入统一读改写(§3.2);P1-2 回归机制改 autouse monkeypatch settings + `from_settings` 补读字段(§2.3/§6.2);P1-3 追加时 alert_type 按合并后规则重算(§2.5);P1-4 L2 推送通知链累积(§3.3);P2-1~P2-6 顺手收口(§2.2/§3.5/§4.2/§4.3/§6.2)。 --- ## 0. 三个需求一句话 + 落库前提核实结论 | 需求 | 规则 | 一句话 | 触发方式 | 优先级 | | --- | --- | --- | --- | --- | | FR-8 | RISK-006 | 客户持仓 R4+R5 市值占比 ≥ 阈值(默认 80%)→ 并入当日事件类预警单 | 交易落库后**同步**(引擎内追加) | P1 | | FR-9 | RISK-007 | pending 单超时(4h/24h,AML 1h/4h)→ payload 写升级标记 + 推送升级,**不改 status** | **定时扫描** 15 分钟 | **P0** | | FR-10 | RISK-008 | 代理人行为链(A 诱导调仓 / B 越权试探 / C 越权查询)→ 代理人维度独立预警单 | **定时扫描** 30 分钟 | P1 | **DDL 核实结论(决定方案形态,均已确认):** | 项 | 结论 | 对方案的影响 | | --- | --- | --- | | `risk_alert.alert_type` | ENUM 五值,无 concentration/agent_behavior | 一律落 `pattern` + `payload.alert_subtype` 区分 | | `risk_alert.status` | ENUM 四值,无 `escalated` | 升级信息全由 `payload.escalation_level/escalated_at` 承载 | | `risk_alert.payload` | JSON NOT NULL | 直接写 dict(repo `_dump_alert` 已序列化) | | `audit_log.event_type` | **VARCHAR(64)**(非 ENUM) | `risk_concentration`/`alert_escalation`/`agent_behavior_detected` 三个新事件类型**可直接写入,无需改表** | | `core_product.min_risk_code` | CHAR(2),取值 R1~R5 | RISK-006 判定 `min_risk_code in ("R4","R5")`(sqlite 测试 DDL 为 VARCHAR(8),同口径兼容) | | `agent_session`/Agent 准入 | `AGENT_ACCESS_MATRIX`["risk"].roles 现为 ("risk_officer","service_risk") | risk_manager 需增补(见 §3.3,注意对话线例外) | --- ## 1. 模块与文件划分总表 ### 1.1 新增文件(6 个) | 文件 | 职责 | 所属任务 | | --- | --- | --- | | `app/service/risk/escalation_service.py` | RISK-007 扫描 + 升级决策 + 幂等 payload 写入 + 降噪推送(纯 service,可被脚本/测试直调) | C5 | | `app/service/risk/agent_behavior_service.py` | RISK-008 滑动窗口聚合 audit_log + 行为链判定 + 出单(纯 service) | C6 | | `scripts/cron/escalation_scan.py` | 定时任务壳:sys.path 引导 → `new_trace()` → 调 escalation_service → 打印摘要(15min 周期挂系统 cron) | C5 | | `scripts/cron/agent_behavior_scan.py` | 同上(30min 周期) | C6 | | `tests/test_escalation_service.py` | C5 单测(含幂等/边界/降噪) | C5 | | `tests/test_agent_behavior_service.py` | C6 单测(A/B/C 边界 + 本人排除 + 去重) | C6 | ### 1.2 修改文件(10 个) | 文件 | 改动 | 所属任务 | | --- | --- | --- | | `app/config/settings.py` | 追加 11 个配置项(§5) | C4~C6 | | `.env.example` | 同步登记(含注释) | C4~C6 | | `app/service/risk/rules.py` | `RuleHit` 加 `alert_subtype` 字段;新增 `rule_concentration` 纯函数;`RiskThresholds` 加 `concentration_threshold` | C4 | | `app/service/risk/engine.py` | `process_trade_event` 追加持仓规则分支(RISK-001~005 之后、AML 之前) | C4 | | `app/service/risk/alert_service.py` | `record_trade_alerts` 维护 `payload.alert_subtype`;`_publish_alert` 支持覆盖 notify_role/附带 subtype | C4(C5/C6 复用) | | `app/repository/risk_repository.py` | 新增 `append_alert_subtypes` / `update_alert_escalation` / `list_pending_alerts_all` / `find_agent_behavior_alert` / `list_audit_events` | C5/C6 | | `app/repository/core_ro.py` | 新增 `concentration_profile(customer_id, limit=500)`(R4+R5 聚合一次 SQL 搞定,避免 Python 端 N+1) | C4 | | `app/api/deps.py` | `AGENT_ACCESS_MATRIX["risk"]["roles"]` 增补 `risk_manager` | C5 | | `app/api/risk.py` | `list_alerts_api` 加 risk_manager 全量只读分支;`chat` 相关不动 | C5 | | `app/api/chat.py` | risk 分支准入后加显式守卫:roles 含 `risk_manager` → deny(对话线维持仅 risk_officer) | C5 | | `app/service/risk/chat_tools.py` | `customer_context` 扩展 `concentration_ratio`;新增 `query_overdue_alerts` / `query_agent_behavior` 两个 Tool + 注册表条目 | C4/C5/C6 | | `app/service/tool_service.py` | `_INTENT_KEYWORDS["risk"]` 追加两条意图词组;`summarize` 加三个 Tool 摘要分支 | C4/C5/C6 | | `app/gateway/trade_gateway.py` | `submit_trade` 加可选参数 `actor_id`,透传到 `_audit`(C6 条件 A 数据源前提) | C6 | | `app/api/simulate.py` | 调 `submit_trade` 时传 `actor_id=auth.actor_id` | C6 | | `scripts/core/02-seed-base.sql` | 追加 STAFF-31001/31002(staff_type='risk_officer',roles='["risk_manager"]') | C5 前置 | | `tests/conftest.py` | 新增 `backdated_alert` / `backdated_audit_event` fixture(免等待测试基建) | C5/C6 | > `chat_tools.py` / `tool_service.py` 三个任务都碰,按 C4→C5→C6 顺序合并提交,避免冲突。 --- ## 2. C4 · RISK-006 集中度预警(FR-8) ### 2.1 数据流转 ```text 交易落库 → engine.process_trade_event ├─ run_rules(RISK-001~005) (现有,不动) ├─ [新增] core.concentration_profile(customer_id) │ → rule_concentration(profile, th) → RuleHit(alert_subtype="concentration") │ 命中 → 并入 hits,走既有 record_trade_alerts(聚合锚点 = find_pending_event_alert │ **排除 payload.alert_subtype 含 'agent_behavior' 的单**——客户维度事件单与 │ 代理人维度行为链单是两条出单线,聚合并存、互不影响,PRD 4A.0 #3; │ 评审 P0-1 收口,见 §2.5 锚点过滤) │ + L3:upsert_profile_l3(customer_id, "pattern", monitor_tags=["high_risk_concentration"]) │ + 审计:event_type='risk_concentration'(金额 DESENS-005 截断) └─ match_customer(AML) (现有,不动) ``` 关键点:RISK-006 **不新建出单路径**——并入 `record_trade_alerts` 既有聚合,`risk_score=60` 经 `max()` 语义不覆盖 RISK-001/002 的 70;`alert_type` 随最高分规则动态更新的现有逻辑不变。 ### 2.2 core_ro 新增:`concentration_profile`(挂账 #1 收口) ```python def concentration_profile(self, customer_id: str, limit: int = 500) -> dict[str, Any]: """持仓集中度画像(RISK-006 专用,一次 SQL 聚合,避免 N+1)。 返回 {r45_value: Decimal, total_value: Decimal, holdings_truncated: bool, rows: [...]}。 rows 仅保留 market_value 降序前 limit 条明细(进审计 input_summary 用摘要,不进 payload 全量)。 触及 limit → holdings_truncated=True(调用方按保守口径视同达标,PRD FR-8 截断防护)。 """ # SQL:SELECT h.market_value, p.min_risk_code FROM core_holding h # JOIN core_product p ... WHERE h.customer_id=:cid # ORDER BY h.market_value DESC LIMIT :lim+1 ← 多取 1 行判定截断 # Python 端按 min_risk_code in ("R4","R5") 分组求和(Decimal) ``` > 采用「SQL 取明细 + Python 聚合」而非 SUM(CASE):明细行还要进审计摘要,且 sqlite/MySQL DDL 差异下 Python 判 `min_risk_code` 最稳。limit+1 探测截断,避免二次 COUNT 查询。 > 注:PRD FR-8 字面为「调用 `core_ro.list_holdings`」;本方法为挂账 #1 的收口(聚合封装,内部同源查询,市值口径不变),偏离字面已在方法 docstring 注明。 ### 2.3 rules.py 改动 ```python @dataclass(frozen=True) class RuleHit: rule_id: str alert_type: str risk_score: int detail: str alert_subtype: str | None = None # 新增:RISK-006→"concentration",RISK-008→"agent_behavior",RISK-001~005 恒 None RULE_SCORES["RISK-006"] = 60 RULE_ALERT_TYPES["RISK-006"] = "pattern" @dataclass(frozen=True) class RiskThresholds: ... concentration_threshold: float = 0.80 # settings.risk_concentration_threshold # from_settings() 必须同步补读:concentration_threshold=s.risk_concentration_threshold # (评审 P1-2:引擎/网关测试全部走 from_settings 默认路径,漏读字段会导致 # autouse 回归 fixture 的 monkeypatch 失效) def rule_concentration(profile: dict[str, Any], th: RiskThresholds) -> RuleHit | None: """RISK-006 集中度:R4+R5 市值占比 ≥ 阈值(纯函数)。 total_value == 0(空仓)不触发;holdings_truncated=True 视同达标(保守告警)。 detail 含 r45_value/total_value/占比(Decimal→str)。 """ ``` > **不动 `run_rules` 签名**——RISK-001~005 吃 trades、RISK-006 吃持仓画像,输入域不同;引擎分别调用后合并 hits,比改 run_rules 兼容面小(现有 13 处 run_rules 调用/断言零改动)。 ### 2.4 engine.py 接入点 ```python # process_trade_event 内,hits = run_rules(...) 之后: profile = core.concentration_profile(trade["customer_id"]) # 截断防护在 rule 内判定 conc = rule_concentration(profile, th) if conc: hits = [*hits, conc] # 并入后走 record_trade_alerts 既有聚合/审计/推送 upsert_profile_l3(trade["customer_id"], "pattern", monitor_tags=["high_risk_concentration"], last_alert_id=... ) # 聚合出单后补(同现有 best 分支) _audit_concentration(repo, trade, profile, conc) # event_type='risk_concentration' ``` > 注意顺序:RISK-006 分支放在 `run_rules` 之后、`record_trade_alerts` 之前,使 `record_trade_alerts` 一次调用覆盖合并;L3 upsert 与专属审计放在出单之后(拿 alert_id)。命中 RISK-006 但 RISK-001~005 未命中时,`record_trade_alerts` 同样正常建单(hits 非空即走聚合分支),**无需**为「仅集中度命中」写独立分支。 ### 2.5 alert_service.py 改动 **① 聚合锚点排除代理人维度单(评审 P0-1 收口,必须做)** `RiskRepository.find_pending_event_alert` 改为「查候选 + Python 过滤」: ```python def find_pending_event_alert(self, customer_id: str, day_start: datetime) -> dict | None: """当日该客户的事件类 pending 单(聚合锚点)——**排除代理人维度行为链单**。 LIMIT 1 会取到最新 pattern 单;若它恰是 agent_behavior 单而更早还有客户维度 事件单,会误判无锚点导致当日出第二张客户维度单。故候选取 LIMIT 50 再过滤: SELECT * FROM risk_alert WHERE customer_id=:cid AND status='pending_review' AND alert_type IN (事件类) AND created_at>=:day_start ORDER BY created_at DESC LIMIT 50 → Python 过滤:'agent_behavior' not in (payload.get("alert_subtype") or []) → 取最新 """ ``` > 依据 PRD §4A.0 #3:FR-4「同客户同自然日仅一张」只对**客户维度**事件单成立,RISK-008 行为链单(同 pattern 但 `payload.actor_type='agent'`)是独立出单线,两者不得互相并入。C6 的 `find_agent_behavior_alert` 已按 subtype 过滤,天然安全(单向污染只在 C4 侧)。 **② alert_subtype 集合维护** ```python subtypes = sorted({h.alert_subtype for h in hits if h.alert_subtype}) # 新建单:仅当 subtypes 非空才写 payload["alert_subtype"] = subtypes(空集不注入该键,评审 P2-3) # 追加单(锚点命中):risk_repo.append_alert_event(..., extra_subtypes=subtypes) ``` **③ 追加时 alert_type 按合并后规则重算(评审 P1-3,既有缺陷顺带修正)** 现状缺陷:`append_alert_event` 用「本批 best.alert_type」更新——既有 large_amount(70) 单,本批仅命中 RISK-006(60) 时会被错误翻转为 pattern。修法(alert_service 追加分支): ```python merged_rules = sorted(set(pending["triggered_rules"]) | {h.rule_id for h in hits}) best_rule = max((r for r in merged_rules if r in RULE_SCORES), key=lambda r: RULE_SCORES[r]) risk_repo.append_alert_event(..., risk_score=max(old, 本批max), alert_type=RULE_ALERT_TYPES[best_rule]) ``` (`RULE_SCORES`/`RULE_ALERT_TYPES` 从 rules.py 导入——alert_service 已依赖 rules.RuleHit,无新增分层问题。编码后跑全量回归确认无现有断言依赖旧翻转行为。) - `RiskRepository.append_alert_event` 加可选参数 `extra_subtypes: list[str] | None = None`:合并进 `payload.setdefault("alert_subtype", [])`(sorted set 并集),老单无该字段时首次追加自动创建。**向后兼容**:不传时行为与现在完全一致,现有用例不受影响。 - `_publish_alert` 加可选参数 `notify_role: list[str] | None = None, extra: dict | None = None`:缺省保持 `["risk_officer"]`(+aml compliance) 现状;`extra` 合并进推送体(C5 放 `escalation_level`,C4/C6 放 `alert_subtype`)。`02-redis-keys.md` 同步增补字段说明(PRD §13 联动清单)。 ### 2.6 对话线:`customer_context` 扩展 `chat_tools.customer_context` 返回体追加: ```python "profile": {"concentration_ratio": 0.9, "r45_value": "...", "total_value": "...", "holdings_truncated": False} # 复用 core_ro.concentration_profile,脱敏后百分比 ``` 调用同一聚合函数,不重复实现口径。`tool_service.summarize` 的 customer_context 分支加一句「高风险持仓占比 90%(仅供参考)」。意图词无需新增(customer_context 词组已有)。 --- ## 3. C5 · RISK-007 时效升级(FR-9) ### 3.1 鉴权与角色联动(**C5 开工前置**,挂账 #6) 1. **种子**:`scripts/core/02-seed-base.sql` 追加: ```sql INSERT INTO core_staff VALUES ('STAFF-31001','风控经理甲','risk_officer','["risk_manager"]',1,...); INSERT INTO core_staff VALUES ('STAFF-31002','风控经理乙','risk_officer','["risk_manager"]',1,...); ``` 2. **`app/api/deps.py`**:`AGENT_ACCESS_MATRIX["risk"]["roles"]` → `("risk_officer", "risk_manager", "service_risk")`。理由:HTTP 台账请求经 `get_auth_context` 的交叉校验,不加则 manager 连 GET /api/risk/alerts 都过不去(401 通道 AUTH_403_AGENT_MISMATCH)。 3. **`app/api/chat.py`(对话线例外)**:risk 分支在 `assert_agent_access` 之后加守卫: ```python if "risk_manager" in auth.roles: deny(auth, "AUTH_403_ROLE", repo, message="对话线仅限 risk_officer,请走 HTTP 台账") ``` > 这是对 PRD 4A.1「AGENT_ACCESS_MATRIX 增补行」与「对话线不放行 risk_manager」两条的联合落地:矩阵放行解决 HTTP 通道,chat 层显式拒绝保住 FR-6 冻结口径。Tool 层 `assert_tool_access` 不动(manager 根本进不了对话线,天然 fail-closed)。 4. **`app/api/risk.py::list_alerts_api`**: ```python if auth.has_role("risk_officer"): ... # 现有全量 elif auth.has_role("risk_manager"): ... # 全量只读(同 risk_officer 查询路径,无处置入口) elif auth.has_role("compliance"): ... # 现有 aml 强制过滤 else: deny(...) ``` `handle_alert_api` 白名单**不动**(manager 无 risk_officer 角色 → 自然 403,PRD「零代码改动」)。JWT 手册 §5.3 增补为文档动作,随 C5 提交。 ### 3.2 仓储新增(risk_repository.py) ```python def list_pending_alerts_all(self, page_size: int = 1000) -> list[dict]: """全部 pending_review 单(C5 扫描输入;演示规模一次取回,Python 端判级)。 复用 _parse_alert。""" def update_alert_escalation(self, alert_id: str, escalation_level: int, escalated_at: datetime, trace_id: str) -> bool: """payload 升级标记独占写入(RISK-007;status/handler_* 列一律不碰)。 实现:**统一读改写**(评审 P1-1 定案)——SELECT payload → Python 合并 escalation_level/escalated_at/escalation_trace_id → UPDATE payload。 与现有 append_alert_event 同模式(跨 MySQL/sqlite 已验证可行),不引入 JSON_SET 方言分支;并发窗口由「单发定时脚本 + 幂等闸门(仅升不降)」兜底。 WHERE alert_id=:aid,返回 rowcount==1。""" ``` > 人工处置 API 走 `handle_alert_with_audit`(只写 status/handler_*/handled_at 列),与 `update_alert_escalation` **写集不相交**,互不覆盖——对应 A-11「payload 与人工处置字段互不覆盖」断言。 ### 3.3 escalation_service.py(核心职责) ```python """预警处置时效升级(C5 · PRD FR-9 / 规则表 v1.1 RISK-007 补充约束)。""" def scan_and_escalate( now: datetime | None = None, risk_repo: RiskRepository | None = None, thresholds: EscalationThresholds | None = None, ) -> dict[str, Any]: """单次扫描入口(定时脚本/测试/手动演示共用)。 EscalationThresholds(frozen dataclass, from_settings()): l1_hours=4, l2_hours=24, aml_l1_hours=1, aml_l2_hours=4 流程: 1. rows = repo.list_pending_alerts_all() 2. 逐单计算 computed_level: overdue_h = (now - created_at).total_seconds()/3600 aml 单用 aml_l1/l2,普通单用 l1/l2;level = 2 if overdue>=l2 else 1 if overdue>=l1 else 0 3. 幂等闸门:computed_level > (payload.escalation_level or 0) 才动作; 同级/降级一律跳过(每级别至多推送一次,A-11 幂等断言锚点) 4. 动作顺序(LEVEL_1 与 LEVEL_2 同机制,PRD 拍板「先持久化再推送」): a. repo.update_alert_escalation(alert_id, level, now, trace_id) b. 审计 event_type='alert_escalation'(input_summary 含 alert_id/已达级别/ 超时时长/升级原因;trace 贯通) 5. 降噪合并:按 (customer_id, level) 分组,同组多单合并为一次推送 (alert_ids 列表进推送体;审计仍逐单落,不合并) 6. 推送:redis_gateway.publish("risk:pub:alert", {..., "escalation_level": lvl, "notify_role": 通知链**累积**(评审 P1-4 修正): L1 → ["risk_officer", "risk_manager"] L2 → ["risk_officer", "risk_manager", "compliance"] ← manager 不因升到 L2 而移出 → 复用 alert_service._publish_alert 的 extra/notify_role 扩展 7. 任务级审计一条:decision='scan_completed',input_summary={scanned, escalated, merged, skipped} 返回 {"scanned": n, "escalated": [...], "merged_notices": m, "skipped": k} """ ``` **边界细则(编码时执行):** - `created_at` 为 naive datetime(repo 直返),`now` 缺省 `datetime.now()`,同口径相减。 - AML 判定按 `alert_type == "aml"`;suitability 单**同样纳入**扫描(属 pending_review,超时同样积压——PRD 未排除,从「时效监控」本义)。 - 升级后单据 status 保持 `pending_review`,仍在下轮扫描范围(level 已达 2,幂等闸门自动跳过);人工处置后自然退出(status ≠ pending_review)。 - `escalation_level` 读缺省 0:`int((alert["payload"] or {}).get("escalation_level") or 0)`。 ### 3.4 定时任务脚本 `scripts/cron/escalation_scan.py`(结构对齐 `scripts/demo/rebuild_alerts.py`): ```python """RISK-007 处置时效升级扫描(建议 15min 周期;RISK_ESCALATION_SCAN_MINUTES 可配)。 初期独立脚本 + 系统 cron;内嵌 lifespan 归 M4 评估(开发计划挂账 #2)。""" # ① sys.path 引导项目根(rebuild_alerts 先例) # ② from app.utils.trace import new_trace; new_trace() ← 显式生成 trace(无 HTTP 上下文) # ③ scan_and_escalate() → print JSON 摘要(供 cron 日志/演示走查) # 审计写库失败:service 内降级 + logger.error 留底(红线 5 口径,不阻塞下次扫描) ``` ### 3.5 对话 Tool:`query_overdue_alerts` ```python def query_overdue_alerts(customer_id, core_ro=None, risk_repo=None, **params): """超期 pending 单列表(只读;FR-9 API 变更)。 hours = params.get("hours") # 缺省取 settings L1 阈值 → repo.list_pending_alerts_all() → 过滤 overdue_h >= hours → [{alert_id, customer_id, alert_type, created_at, escalation_level(读 payload,缺省0), overdue_hours(round 1位), risk_score}],按 overdue_hours 降序(对话线查询置顶语义) """ # 注册表条目:requires_customer=False(risk_officer 无绑定客户可查全量,C2 守卫先例) # param_whitelist=("hours",), int_bounds={"hours": (1, 720)} ``` 意图词组:`("query_overdue_alerts", ("超期", "超时", "逾期", "多久没处理", "处置时效"))`——**置于 risk 组 alert_query 之前**(评审 P2-4:「超时预警」类话术不得误命中台账)。 --- ## 4. C6 · RISK-008 代理人行为链(FR-10) ### 4.0 前置改造:交易审计 actor 透传(条件 A 数据源前提) 现状:`trade_gateway._audit` 写 `actor_id="SYSTEM"`(B6 挂账注释),条件 A 无法归属发起人。 ```python # trade_gateway.submit_trade 加参数:actor_id: str | None = None # _audit(..., actor_id=actor_id or "SYSTEM") ← 缺省 SYSTEM,现有测试/脚本零改动 # app/api/simulate.py:submit_trade(req, ..., actor_id=auth.actor_id) ``` 演示期 risk_demo 账号(roles 含 risk_demo)扮演代理人发起 → audit_log.trade_request 行的 actor_id = STAFF-90001;客户本人发起时 actor_id == customer_id,扫描时**排除**(防误报,A-12 断言)。 ### 4.1 仓储新增 ```python def list_audit_events(self, event_type: str, since: datetime, actor_id: str | None = None, decision: str | None = None, limit: int = 2000) -> list[dict]: """audit_log 滑窗查询(RISK-008 数据源;input_summary 解析回 dict)。 WHERE event_type=:et AND created_at>=:since [AND actor_id=:aid] [AND decision=:d] ORDER BY created_at ASC。条件 B/C 用 actor_id+decision='forbidden'; 条件 A 用 event_type='trade_request' 全量(actor 过滤在 Python 端做,排除 SYSTEM)。""" def find_agent_behavior_alert(self, actor_id: str, day_start: datetime) -> dict | None: """同代理人当日 pending 行为链单(去重锚点)。 SQL:SELECT * FROM risk_alert WHERE status='pending_review' AND alert_type='pattern' AND created_at>=:day_start ORDER BY created_at DESC(演示规模可接受); Python 过滤 payload.alert_subtype=='agent_behavior' and payload.actor_id==actor_id。 不用 payload LIKE:JSON 键值歧义 + 序列化空格差异两坑(挂账 #9 登记性能,一期 Python 过滤最稳)。""" ``` ### 4.2 agent_behavior_service.py(核心职责) ```python """代理人异常行为链识别(C6 · PRD FR-10 / 规则表 v1.1 RISK-008 补充约束)。 数据源(定案口径,input_guard_log 不作为数据源): 条件 A:audit_log(event_type='trade_request') ∩ core_trade 时间线——**join 键 = trade_id** (评审 P2-2 补细):审计行 input_summary.trade_id 归属发起人 actor_id; 扫描先取 trade_request 行 → 过滤 actor_id=='SYSTEM'(历史无归属)与 actor_id==customer_id(本人交易)→ 以行内 trade_id 集合从 core_trade (list_trades_range 当日窗 + get_trade_by_id 补漏)取 traded_at/trade_type/ product_id 组装时间线 → 模式判定只在「有归属 actor 的 trade_id 子集」内做 条件 B:audit_log(event_type='authz', decision='forbidden') 且 input_summary.code=='AUTH_403_SCOPE' 条件 C:同 B 族,code IN ('AUTH_403_NOT_OWNER','AUTH_403_NOT_ASSIGNED') (对话线 blocked 经 record_authz_denial 双写 audit_log,天然同口径可查) """ @dataclass(frozen=True) class BehaviorThresholds: a_window_hours: int = 24; a_count: int = 3 b_window_hours: int = 72; b_count: int = 5 c_window_hours: int = 24; c_count: int = 10 # from_settings() 对齐 settings.risk_agent_behavior_* def detect_hits(now, core, repo, th) -> dict[str, dict[str, Any]]: """按代理人分组聚合三条件证据。返回 {actor_id: {"A": [...], "B": [...], "C": [...]}}, 每条证据含 {trace_id, at, customer_id, detail}。纯查询,无副作用(单测友好)。""" def scan_and_alert(now=None, core_ro=None, risk_repo=None, thresholds=None) -> dict: """扫描入口(脚本/测试共用;模式对齐 escalation_service)。 1. detect_hits 收集证据 2. 逐 actor: existing = repo.find_agent_behavior_alert(actor_id, day_start) ├─ 有 → 新子条件才 append_alert_event(triggered_rules 追加 "RISK-008", │ payload.events 追加证据、alert_subtype 集合并集——同日多子条件只一张单) └─ 无 → 出单: alert_type='pattern', triggered_rules=["RISK-008"], risk_score=75 customer_id = 证据条数众数客户(PRD:涉及客户语义,payload 显式声明) payload = {alert_subtype:'agent_behavior', actor_id, actor_type:'agent', subtypes_hit:['A'|'B'|'C'], timeline:[...], customers:[...], evidence_trace_ids:[...], note:'customer_id 为涉及客户非归属客户'} L3 不写(代理人画像本期仅 payload 承载,挂账 #3) 3. 推送 notify_role=["risk_officer","risk_manager"](_publish_alert extra 带 alert_subtype) 4. 审计 event_type='agent_behavior_detected'(每 actor 每次命中一条;显式 new_trace 由脚本入口给) 返回 {"scanned_windows": ..., "hits": [...], "created": [...], "appended": [...]} """ ``` **边界细则:** - 条件 A 模式判定:归属映射 audit 行 → (actor_id, customer_id, trade_id),trade_id 经 core_trade 取时间线;同 (actor, customer) 组内,对每个 redeem 事件找 2h 内该 customer 的 subscribe 且 product_id 不同 → 计一次诱导调仓;同一 redeem 不重复配对(消费制)。窗口按 core_trade.traded_at(事件时点,非扫描墙钟,与 RISK-004 rebuild 口径一致)。 - 窗口均为滑动窗口 `created_at >= now - window`,自然日去重键只用于出单(PRD 明确两套口径并存)。 - 行为链明细展示脱敏:payload 存 customer_id(内部键不脱敏);Tool 输出层对 display_name 走 `utils/desensitize` 二次脱敏(PRD 系统动作 6)。 ### 4.3 对话 Tool:`query_agent_behavior` ```python def query_agent_behavior(customer_id, core_ro=None, risk_repo=None, **params): """代理人行为链查询(只读)。agent_id = params.get("agent_id")(可缺省=全部)。 rows, _ = repo.list_alerts(alert_type='pattern', status=None, page_size=100) # 注意返回 tuple(评审 P2-1) → Python 过滤 payload.alert_subtype=='agent_behavior'(+ actor_id==agent_id 若传) → [{alert_id, actor_id, subtypes_hit, customers(脱敏展示), created_at, risk_score, status, evidence_trace_ids[:5]}] 可见性:对话线准入已限 risk_officer(manager 走 HTTP 台账),Tool 层天然 fail-closed。 """ # 注册表:requires_customer=False, param_whitelist=("agent_id",), int_bounds={} # 意图词组:("query_agent_behavior", ("代理人", "行为链", "异常行为", "诱导", "越权记录")) ``` > 意图词注意与现有 risk 组内 `alert_query`("预警")词序:`query_agent_behavior` 放 alert_query **之前**("代理人预警"应命中行为链而非台账);"预警"仍是 alert_query 专属词,无冲突。 ### 4.4 定时脚本 `scripts/cron/agent_behavior_scan.py`:结构同 escalation_scan.py(sys.path 引导 → `new_trace()` → `scan_and_alert()` → JSON 摘要打印)。 --- ## 5. 配置项清单(settings.py + .env.example,C4~C6 开工时一次性加齐) ```python # ===== 风控追加 v1.1(FR-8/9/10 · PRD §4A)===== risk_concentration_threshold: float = 0.80 # FR-8 R4+R5 占比阈值 risk_escalation_scan_minutes: int = 15 # FR-9 扫描周期(cron 侧参考) risk_escalation_l1_hours: int = 4 # FR-9 普通单 LEVEL_1 risk_escalation_l2_hours: int = 24 # FR-9 普通单 LEVEL_2 risk_escalation_aml_l1_hours: int = 1 # FR-9 AML 通道 L1 risk_escalation_aml_l2_hours: int = 4 # FR-9 AML 通道 L2 risk_agent_behavior_scan_minutes: int = 30 # FR-10 扫描周期 risk_agent_behavior_a_window_hours: int = 24 # FR-10 条件 A risk_agent_behavior_a_count: int = 3 risk_agent_behavior_b_window_hours: int = 72 # FR-10 条件 B risk_agent_behavior_b_count: int = 5 risk_agent_behavior_c_window_hours: int = 24 # FR-10 条件 C risk_agent_behavior_c_count: int = 10 ``` --- ## 6. 测试方案(A-10 / A-11 / A-12) ### 6.1 conftest 新增 fixture(免等待基建) ```python @pytest.fixture() def backdated_alert(sqlite_engine): """注入 created_at 回拨的 pending 单(C5);teardown 按 alert_id 精确删除。 用法:alert_id = backdated_alert(alert_id="ALT-TEST-...", customer_id=..., hours_ago=5, payload={...})""" @pytest.fixture() def backdated_audit_event(sqlite_engine): """注入 created_at 回拨的 audit_log 行(C6);trace_id 用 'TEST-TRACE-' 前缀, teardown DELETE WHERE trace_id LIKE 'TEST-TRACE-%'(不污染正常 trace 空间)。""" ``` > 真 MySQL 集成测试沿用 `risk_demo_env`;回拨行 teardown 注意 `_cleanup_test_rows` 按 `created_at >= started_at` 清不掉回拨行 → 两个 fixture **必须自带精确清理**(alert_id / trace 前缀),这是评审必查点。 ### 6.2 用例清单 | 文件 | 覆盖 | | --- | --- | | `tests/test_rules.py`(扩) | rule_concentration:90% R5 命中 / 100% R1 不触发 / 空仓不触发 / 截断视同达标 / 阈值边界(79.9% 不触发 / 80% 触发) | | `tests/test_engine.py`(扩) | 引擎合并:RISK-006 与 RISK-001 同单聚合、risk_score=max(60,70)=70、payload.alert_subtype 含 concentration、L3 tag 追加、risk_concentration 审计 | | `tests/test_alert_service.py`(扩) | append_alert_event extra_subtypes 合并(老单无字段首次创建);**P0-1 回归:当日已有 agent_behavior 单后再触发 RISK-006 → 出第二张客户维度单,两单 payload 互不污染**;alert_type 合并重算(large_amount 单 + 仅 RISK-006 追加 → 保持 large_amount 不翻转);现有断言回归 | | `tests/test_escalation_service.py`(新) | 3h59m 不升 / 4h 升 L1 / 24h 升 L2 / AML 1h 短通道(两条断言分开)· 同单重复扫描不重复推送(幂等)· 同客户多单同级别合并一次推送 · 阈值内处置不升级 · 处置后退出扫描 · status 全程 pending_review · notify_role 断言(L1 含 risk_manager / L2 含 compliance)· update_alert_escalation 不碰 handler 列 | | `tests/test_agent_behavior_service.py`(新) | 条件 A:2 次不触发 / 3 次触发 / 本人交易(actor==customer)不计数 / 跨产品 2h 窗口边界 / SYSTEM actor 排除;条件 B:4/5 边界 + code 精确匹配;条件 C:9/10 边界 + 双 code 并集;同日多子条件只一张单(subtypes_hit 并集);customer_id=众数;payload.actor_id 正确指向代理人 | | `tests/test_chat_tools.py`(扩) | customer_context 带 concentration_ratio;query_overdue_alerts 边界与缺省 hours;query_agent_behavior 过滤与 agent_id 匹配 | | `tests/test_risk_api.py`(扩) | risk_manager:GET alerts 200 全量 / handle 403 / **suitability/check 403 / aml/scan 403(评审 P2-5:防未来新增端点只依赖矩阵漏加角色校验的回归断言)** / chat risk 入口 deny(AUTH_403_ROLE)/ debug 头与 JWT 通道各一条 | | 集成(risk_demo_env) | A-10(TRD-TEST- 交易落库后 5s 内出单)/ A-11 / A-12 按验收表逐条 | **全量回归口径(PRD A-10 预声明 · 评审 P1-2 修订)**:conftest 增加自动生效的回归隔离 fixture(autouse,作用于全部单测): ```python @pytest.fixture(autouse=True) def _disable_concentration_rule(monkeypatch): """RISK-006 阈值推到不可达,现有用例断言零改动;RISK-006 专测内再手动改回真实阈值。""" monkeypatch.setattr(settings, "risk_concentration_threshold", 1.01) ``` 前置条件:`RiskThresholds.from_settings()` **必须补读** `concentration_threshold` 字段(引擎/网关测试全部走 `from_settings()` 默认路径,漏读则 monkeypatch 失效、425 绿被 RISK-006 打穿)。RISK-006 专属测试显式传 `RiskThresholds(...)` 或再次 monkeypatch 真实阈值 0.80。 --- ## 7. 实现优先级与执行顺序 | 步骤 | 内容 | 理由 | | --- | --- | --- | | 0 | C5 前置联动:seed 两个 manager 账号 + deps 矩阵 + chat 守卫 + risk.py 台账分支 + JWT 手册文档 | PRD「C5 开工前执行」;独立可测(不依赖 C4) | | 1 | **C4**(rules + core_ro + engine + alert_service + customer_context + settings) | 改动最小、复用度最高、无新基建;先拿下一个完整验收(A-10) | | 2 | **C5**(repo 方法 + escalation_service + cron 脚本 + Tool + conftest fixture) | 建立定时任务基建与 fixture 模式,C6 直接复用;P0 红线规则 | | 3 | **C6**(trade_gateway actor 透传 + repo 审计查询 + agent_behavior_service + Tool + 脚本) | 前置改造最多(actor 链路),基建已由 C5 备好 | | 4 | 演示 SOP 补章节(risk_demo 扮演代理人步骤)+ 02-redis-keys.md 增补 + MEMORY/TODO 收口 | 随各任务顺手,最后统一核对挂账 #1~#9 | 每个任务独立 commit(C4 → `feat: T-C4 ...`),全量 `python -m pytest` 绿后再进下一步;C4~C6 全部完成后按开发计划 M3 口径验收 A-10/A-11/A-12 并打 tag。 ### 红线自查(每步提交前过一遍) - 不改表结构(alert_type/status 复用 payload 承载;audit_log.event_type 为 VARCHAR 可直接扩)✓ - 升级标记仅定时任务可写(update_alert_escalation 只有 cron 链路调用;handle API 不触碰 payload)✓ - 人工处置唯一入口 risk_officer(manager 全链路无处置权)✓ - 四 Agent 不互调 LLM(C5/C6 定时任务无 LLM;新 Tool 全部只读)✓ - 审计只 INSERT(三个新 event_type 走 insert_audit_log)✓ - 客户 Agent 边界不涉及(三个需求全部收口在 risk Agent + 定时任务)✓ ### 已登记挂账(不属本期,编码时勿"顺手"处理) #2 定时任务内嵌 lifespan、#3 agent_profile_l3 建表、#4 escalated 状态改表、#5 alert_type ENUM 扩展、#7 customer_query 事件类型补建、#8 core_trade actor 字段、#9 payload JSON 检索性能。