544 lines
34 KiB
Markdown
544 lines
34 KiB
Markdown
# 实现方案 · 风控追加需求 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 检索性能。
|