34 KiB
实现方案 · 风控追加需求 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 数据流转
交易落库 → 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 收口)
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 改动
@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 接入点
# 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 过滤」:
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 集合维护
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 追加分支):
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 返回体追加:
"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)
- 种子:
scripts/core/02-seed-base.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,...);
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)。app/api/chat.py(对话线例外):risk 分支在assert_agent_access之后加守卫:
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)。
app/api/risk.py::list_alerts_api:
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)
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(核心职责)
"""预警处置时效升级(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):
"""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
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 无法归属发起人。
# 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 仓储新增
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(核心职责)
"""代理人异常行为链识别(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
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 开工时一次性加齐)
# ===== 风控追加 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(免等待基建)
@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,作用于全部单测):
@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 检索性能。