diff --git a/app/repository/core_ro.py b/app/repository/core_ro.py index 7897806..418efcb 100644 --- a/app/repository/core_ro.py +++ b/app/repository/core_ro.py @@ -73,6 +73,13 @@ class CoreReadOnlyRepository: ).mappings() ] + def get_trade_by_id(self, trade_id: str) -> dict[str, Any] | None: + """按 trade_id 查单笔交易流水(rebuild_alerts 补偿重放用,仅 SELECT)。""" + sql = text("SELECT * FROM core_trade WHERE trade_id = :tid") + with self._engine.connect() as conn: + row = conn.execute(sql, {"tid": trade_id}).mappings().first() + return dict(row) if row else None + def sum_trades_on_date(self, customer_id: str, day: date) -> Decimal: """当日申赎合计金额(RISK-002 累计口径:仅 confirmed 的 subscribe/redeem)。 diff --git a/app/repository/risk_repository.py b/app/repository/risk_repository.py index 27625a6..96a5f3f 100644 --- a/app/repository/risk_repository.py +++ b/app/repository/risk_repository.py @@ -137,6 +137,27 @@ class RiskRepository: ).mappings().first() return self._parse_alert(dict(row)) if row else None + def find_alerts_by_trade(self, trade_id: str) -> list[dict]: + """trade_id 已入的预警单(rebuild_alerts 幂等检查;含已处置单——重放不得绕过处置结论)。 + + payload 为 JSON 列(MySQL)/TEXT JSON(sqlite 测试),LIKE 匹配 trade_id 字符串; + trade_id 形如 TRD-{date}-{uuid8} 唯一性强,无误匹配面。首笔交易的 trade_id + 恒写 payload.events[0].trade_id,LIKE 覆盖事件类/suitability/aml 全部出单路径。 + """ + sql = text( + """ + SELECT alert_id, alert_type, status + FROM risk_alert + WHERE payload LIKE :pat + ORDER BY created_at + """ + ) + with self._engine.connect() as conn: + return [ + dict(r) + for r in conn.execute(sql, {"pat": f"%{trade_id}%"}).mappings() + ] + def list_alerts( self, status: str | None = None, diff --git a/docs/memory/FLOW.md b/docs/memory/FLOW.md index ba28620..025a756 100644 --- a/docs/memory/FLOW.md +++ b/docs/memory/FLOW.md @@ -48,11 +48,11 @@ RBAC 联调账号:scripts/dev/rbac-seed-reference.md Client → Gateway(JWT/RBAC) → api/chat → agent_service(LangGraph) → Tools → 存储 → 响应 + audit_log ``` -当前:**风控事件线已集成 + 测试基建完成**(交易网关 → 规则引擎 → 预警/AML/L3 → 4 API,`main.py` 路由挂载 + trace 中间件 + lifespan 完成,B1~B8;B8 = conftest 演示数据校验 + 集成测试 A-1~A-5/A-7/A-9 + `tests/_ddl.py` DDL 单一事实源;B7/B8 均经独立 AI 评审闭环,196 测试绿);剩 B9a 脚本 → B9b 演示走查;对话链路待 T-01/T-07。 +当前:**风控事件线已集成 + 测试基建完成**(交易网关 → 规则引擎 → 预警/AML/L3 → 4 API,`main.py` 路由挂载 + trace 中间件 + lifespan 完成,B1~B8;B8 = conftest 演示数据校验 + 集成测试 A-1~A-5/A-7/A-9 + `tests/_ddl.py` DDL 单一事实源;B7/B8 均经独立 AI 评审闭环;B9a = subscribe_alerts/rebuild_alerts 演示与补偿脚本,201 测试绿);剩 B9b 演示走查;对话链路待 T-01/T-07。 -**本机已就位状态(2026-09-06 · 已完成上述 ①~⑤,无需重做):** `.env` 已配置(学习项目,`MYSQL_PASSWORD=123456`;NEO4J/DEEPSEEK 留空暂不影响风控);`jinrong_core` + `jinrong_agent` 已灌库(28 客户 / AML 名单 8 条 / 演示测评已刷新 / 归属同步 28 行);`python -m pytest` 196 绿(集成测试真连本机 MySQL;未灌库的机器自动 skip 集成模块,单测不受影响)。 +**本机已就位状态(2026-09-06 · 已完成上述 ①~⑤,无需重做):** `.env` 已配置(学习项目,`MYSQL_PASSWORD=123456`;NEO4J/DEEPSEEK 留空暂不影响风控);`jinrong_core` + `jinrong_agent` 已灌库(28 客户 / AML 名单 8 条 / 演示测评已刷新 / 归属同步 28 行);`python -m pytest` 201 绿(集成测试真连本机 MySQL;未灌库的机器自动 skip 集成模块,单测不受影响);本机 Redis 服务在跑(`redis://127.0.0.1:6379/0`,B9a 订阅脚本已验证)。 -**本机已知坑:** `mysql.exe` 不在 PATH(位于 `C:\Program Files\MySQL\MySQL Server 8.0\bin`);`reset.ps1` 的交互式 `-p` 在自动化执行时会卡死——脚本化重灌用 `MYSQL_PWD` 环境变量传密码(交互执行不受影响,不把密码写进仓库脚本)。 +**本机已知坑:** `mysql.exe` 不在 PATH(位于 `C:\Program Files\MySQL\MySQL Server 8.0\bin`);`reset.ps1` 的交互式 `-p` 在自动化执行时会卡死——脚本化重灌用 `MYSQL_PWD` 环境变量传密码(交互执行不受影响,不把密码写进仓库脚本);`redis` 包 requirements 有但初始 `pip install -r` 漏装(B9a 于 2026-09-07 补装,重装环境时留意)。 ------ diff --git a/docs/memory/MEMORY.md b/docs/memory/MEMORY.md index 5d48cfa..51f35cd 100644 --- a/docs/memory/MEMORY.md +++ b/docs/memory/MEMORY.md @@ -9,7 +9,7 @@ **项目是什么:** 金融四 Agent(客户财富 / 代理人 / 数据分析 / 风控)共用数据层与合规底座;**不**互调 LLM,跨 Agent 走 L1/L2/L3 画像与预警表。 -**当前进度:** 需求与表设计已定 · **风控模块已落地 B1~B8**(规则/预警聚合/L3/AML/引擎/交易网关/鉴权+4 API/适当性校验/main 集成+挂账①~⑦/**B8 conftest+集成测试**,B7/B8 均经独立 AI 评审闭环,**196 测试绿**)· 剩 **B9a 脚本 → B9b 演示走查** · **JWT(T-01) / 审计中间件(T-02) / LangGraph 对话线(T-07) 未做**(chat/knowledge/admin 仍空壳)。**开发在分支 `feature/risk`(未合入 main)。** +**当前进度:** 需求与表设计已定 · **风控模块已落地 B1~B8+B9a**(规则/预警聚合/L3/AML/引擎/交易网关/鉴权+4 API/适当性校验/main 集成+挂账①~⑦/B8 conftest+集成测试/**B9a 订阅+重放补偿脚本**,B7/B8 均经独立 AI 评审闭环,**201 测试绿**)· 剩 **B9b 演示走查** · **JWT(T-01) / 审计中间件(T-02) / LangGraph 对话线(T-07) 未做**(chat/knowledge/admin 仍空壳)。**开发在分支 `feature/risk`(未合入 main)。** **仓库地图:** @@ -26,7 +26,7 @@ | `app/utils/` | **基本就绪** | trace / desensitize / db(引擎工厂)/ response(统一错误体)/ exceptions(含 ApiError)已实现;logger 占位 | | `app/config/settings.py` | **已实现** | 双库 `jinrong_agent` + `jinrong_core` + risk_* 阈值 | | `scripts/core/*.sql` + `reset.ps1` | **已实现** | Core 模拟库 DDL + 种子 | -| `scripts/agent/` `scripts/demo/` | **已实现** | AML 名单种子 + 风控演示数据 | +| `scripts/agent/` `scripts/demo/` | **已实现** | AML 名单种子 + 风控演示数据 + subscribe_alerts/rebuild_alerts(B9a 订阅与幂等重放补偿) | | `scripts/sync/*.py` | **已实现** | 归属同步 + Neo4j 全图 | | `tests/` | **已实现** | 16 个测试文件 196 用例(sqlite 隔离;DDL 单一事实源 `_ddl.py`;`test_integration_risk.py` 走真 MySQL + TRD-TEST- 前缀隔离) | | `docs/需求拆解/` | 已定 | 场景 P0、矩阵、合规原文 | @@ -51,7 +51,7 @@ > **本机(2026-09-06)①~⑤已执行、`.env` 已配置,勿重做**;本机状态与已知坑(mysql.exe 路径 / reset.ps1 交互式 -p)见 `FLOW.md` §0 尾注。 -**下一步开发(见 TODO):** 风控 B9a 演示/运维脚本(subscribe_alerts.py / rebuild_alerts.py)→ B9b 演示走查(含 B6/B7 挂账核查单);随后 Wave 0 T-01 JWT / T-02 审计中间件 / T-06 / T-07 LangGraph。 +**下一步开发(见 TODO):** 风控 B9b 演示走查(reset → demo SQL → Swagger 过 A-1~A-5/A-7~A-9 + B6/B7/B8 挂账核查单①~⑥);随后 Wave 0 T-01 JWT / T-02 审计中间件 / T-06 / T-07 LangGraph。 **禁止(改代码前必记):** Core 正式 C1~C5 不可被画像覆盖 · 审计表只 INSERT · 代理人草稿不外发 · 仅 R-02 可阻断交易 · 四 Agent 不互调 LLM。 @@ -124,7 +124,8 @@ Core 模拟:scripts/core/reset.ps1 · 文档 docs/项目框架设计/Core模 种子:scripts/agent/seed-aml-list.sql(AML 名单)· scripts/demo/prepare_risk_demo.sql(reset 后重跑) 依赖:requirements.txt(LangGraph + langchain-core/openai + FastAPI + SQLAlchemy) 启动:uvicorn app.main:app --reload → GET /health -测试:python -m pytest(196 用例;集成测试需本机演示数据,未灌库时自动 skip) +测试:python -m pytest(201 用例;集成测试需本机演示数据,未灌库时自动 skip) +运维/演示脚本:scripts/demo/subscribe_alerts.py(订阅推送演示)· rebuild_alerts.py TRD-xxx(引擎异常补偿重放) 配置:.env(见 .env.example) RBAC 联调账号:scripts/dev/rbac-seed-reference.md ``` @@ -165,6 +166,6 @@ RBAC 联调账号:scripts/dev/rbac-seed-reference.md 2. 改动属于 api / service / tool / repository 哪一层? 3. 是否需 customer_id 归属与 JWT RBAC? 4. Core 是模拟库只读还是 agent 库读写? -5. 如何验证?(`python -m pytest` 全量(当前 196 绿)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) +5. 如何验证?(`python -m pytest` 全量(当前 201 绿)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) 大任务:FRAMEWORK/FLOW 与实现状态不符时先更新 memory 再编码(用户确认跳过除外)。 diff --git a/docs/memory/TODO.md b/docs/memory/TODO.md index 1f12d6e..81b9d94 100644 --- a/docs/memory/TODO.md +++ b/docs/memory/TODO.md @@ -5,11 +5,10 @@ ## 进行中 -- [ ] 风控 **B9a 演示/运维脚本**:`scripts/demo/subscribe_alerts.py`(订阅演示)+ `scripts/demo/rebuild_alerts.py`(按 trade_id 幂等重放补偿)——验收:手工执行验证;B9b 走查前先处理核查单⑥(B8 复审 P2-1 演示库耦合前置断言) +- [ ] 风控 **B9b 演示链路走查**(reset → demo SQL → Swagger 过 A-1~A-5/A-7~A-9 + 核查单①~⑥)——走查前先处理核查单⑥(B8 复审 P2-1 演示库耦合前置断言) ## 待办(推荐顺序) -- [ ] 风控 **B9b 演示链路走查**(reset → demo SQL → Swagger 过 A-1~A-5/A-7~A-9 + 核查单①~⑥) - [ ] 风控 M2 收尾:MEMORY/TODO 状态复核(开发计划 M4 前置) - [ ] T-01 Auth SDK / JWT 中间件(对照 `02-JWT-RBAC鉴权手册.md`) @@ -23,7 +22,7 @@ ### 风控模块(PRD v1.0 已冻结 · `docs/PRD/PRD-风控监测Agent.md`,事件驱动线不依赖 T-07 可先行) -- [ ] T-30 风控事件线:`app/gateway/` 交易网关 + `service/risk/` 规则引擎(RISK-001~005)+ 预警单聚合 + L3 最小写入 + `risk:pub:alert` 推送 + AML(含 `risk_aml_list` 种子)+ 演示数据脚本验收 A-1~A-5/A-9 —— **进度(2026-09-06):B1 规则纯函数 / B2 预警服务 / B3 L3 写入(risk_score 一期不写)/ B4 AML+引擎编排 / B5 交易网关 / B6 鉴权+4 API+处置编排 / B7 main 集成(路由挂载 + trace 中间件 + lifespan 启动期 debug 校验;挂账①锁公共化 ②L3 缓存 DEL ③死代码 ④统一错误体 ⑤启动校验 ⑥引擎工厂 ⑦处置原子事务均落地;⑧ input_guard_log 双写仍随 T-02)均完成并经独立 AI 评审闭环(B7 复审有条件通过→P1-1 dispose 已修闭环;183 测试绿;P2/P3 已登记开发计划 B7/B8/B9b/T-02);**B8 已完成(2026-09-06):conftest(演示数据校验/幂等代跑 prepare_risk_demo/TRD-TEST- teardown + L3 快照还原)+ 集成测试 11 例(A-1~A-5/A-7/A-9 + trace 一致性 + 审计 JSON)+ sqlite DDL 单一事实源 `_ddl.py`(10 文件收敛,localtime 时区收敛 B5 P3-4)+ L3 DEL 行为断言(P2-1)+ locks 文案(P3-2),196 绿;B8 复审有条件通过→P1-1 conftest skip 路径炸收集已修闭环(故障注入验证),P2-1 演示库耦合前置断言挂 B9b 核查单⑥**;剩 B9a 脚本、B9b 演示走查** +- [ ] T-30 风控事件线:`app/gateway/` 交易网关 + `service/risk/` 规则引擎(RISK-001~005)+ 预警单聚合 + L3 最小写入 + `risk:pub:alert` 推送 + AML(含 `risk_aml_list` 种子)+ 演示数据脚本验收 A-1~A-5/A-9 —— **进度(2026-09-06):B1 规则纯函数 / B2 预警服务 / B3 L3 写入(risk_score 一期不写)/ B4 AML+引擎编排 / B5 交易网关 / B6 鉴权+4 API+处置编排 / B7 main 集成(路由挂载 + trace 中间件 + lifespan 启动期 debug 校验;挂账①锁公共化 ②L3 缓存 DEL ③死代码 ④统一错误体 ⑤启动校验 ⑥引擎工厂 ⑦处置原子事务均落地;⑧ input_guard_log 双写仍随 T-02)均完成并经独立 AI 评审闭环(B7 复审有条件通过→P1-1 dispose 已修闭环;183 测试绿;P2/P3 已登记开发计划 B7/B8/B9b/T-02);**B8 已完成(2026-09-06):conftest(演示数据校验/幂等代跑 prepare_risk_demo/TRD-TEST- teardown + L3 快照还原)+ 集成测试 11 例(A-1~A-5/A-7/A-9 + trace 一致性 + 审计 JSON)+ sqlite DDL 单一事实源 `_ddl.py`(10 文件收敛,localtime 时区收敛 B5 P3-4)+ L3 DEL 行为断言(P2-1)+ locks 文案(P3-2);B8 复审有条件通过→P1-1 conftest skip 路径炸收集已修闭环(故障注入验证),P2-1 演示库耦合前置断言挂 B9b 核查单⑥**;**B9a 已完成(2026-09-07):`scripts/demo/subscribe_alerts.py`(risk:pub:alert 订阅演示:连接自检/--duration/心跳)+ `rebuild_alerts.py`(按 trade_id 幂等重放:`find_alerts_by_trade` 以 payload LIKE 查已入单(含 aml/已处置单)防重复出单与重复 append;core_ro 新增只读 `get_trade_by_id`)+ tests/test_demo_scripts.py 5 例(出单推送/幂等跳过/missing 不落库/aml 幂等/格式化),201 绿;真库手工验证通过(rebuilt→skipped→missing exit1;订阅端到端收到假消息+真交易推送且 trace 贯通),验证现场已清理**;剩 B9b 演示走查** - [ ] T-31 `service/suitability.py` 公共校验(SUIT-001~008)+ `POST /api/risk/suitability/check` + 单测验收 A-8 —— **代码已完成(A1~A4 · 77 测试绿;B6 起该 API 带鉴权+直调审计);MySQL 手工 SQL 对照挂账至 B9b 执行(阶段 A 评审 P2-8)** - [ ] T-32 预警台账与人工处置 API(`GET /alerts`、`POST /handle`,risk_officer/compliance 权限)+ 对话线(依赖 T-01/T-03/T-07)验收 A-6/A-7 —— **API 部分已由 B6 覆盖并复审通过(A-7/A-9 用例绿);对话线仍依赖 T-01/T-03/T-07** diff --git a/docs/项目框架设计/开发计划-风控模块.md b/docs/项目框架设计/开发计划-风控模块.md index f997c81..bdd791f 100644 --- a/docs/项目框架设计/开发计划-风控模块.md +++ b/docs/项目框架设计/开发计划-风控模块.md @@ -32,6 +32,8 @@ | B8 | **`tests/conftest.py`**:a) session fixture 启动校验演示数据就位(CUST-4001 测评 <365 天、risk_aml_list ≥8),缺失则中止并提示先跑 FLOW §0 ③④;b) fixture 幂等代跑 `prepare_risk_demo.sql`;c) teardown 按 `TRD-TEST-` 清 core_trade + 关联 risk_alert/risk_suitability_log/audit_log + 还原 L3 行。**顺手集中 sqlite 测试 DDL 为单一事实源(B4 评审 P3-12,各测试文件手写 DDL 收敛;CURRENT_TIMESTAMP 改 localtime 或 fixture 固定时间,防 UTC/本地日界错位——B5 评审 P3-4)**;集成测试交易统一走 `trade_id_factory` 注入 `TRD-TEST-` 前缀(trade_gateway 已留参数,B5 评审 P3-2)。集成测试:A-1~A-5、A-7(状态机/compliance 403/GET 强制 aml)、A-9 越权、**trace 一致性断言**;补 platform 审计 input_summary 的 JSON 解析断言(含引擎输出/阻断 reasons,B5 复审 L1);**B7 复审 P2-1:补 L3 DEL 钩子行为断言(fake 记录 deletes,断言 key=`profile:l3:{customer_id}` 与降级路径);P3-2 locks 降级文案中性化顺手改** | 测试套件 + fixture | `pytest` 全绿 | B7 | | B9a | 演示/运维脚本开发:`scripts/demo/subscribe_alerts.py`(订阅演示)+ `scripts/demo/rebuild_alerts.py`(按 trade_id 幂等重放补偿) | 2 个脚本 | 手工执行验证 | B2、B4(可与 B5~B8 并行) | +**B9a 完成(2026-09-07):** 两脚本 + repo 只读扩展(core_ro.get_trade_by_id / risk_repo.find_alerts_by_trade:payload LIKE 查已入单,含 aml/已处置单——aml 出单路径本身无去重,幂等检查是其防重放双单的唯一闸门)+ tests/test_demo_scripts.py 5 例,201 绿。真库手工验证:rebuilt(RISK-001/002 出单)→ 再跑 skipped → missing exit 1;订阅端到端收假消息 + 真交易推送(trace 贯通、notify_role 正确),验证现场已清理。本机补装 redis 包(requirements 有而环境漏装,见 FLOW §0 尾注)。 + **B8 复审(独立 AI 评审 · 2026-09-06):有条件通过 → 已闭环。** 开发计划 B8 行 a/b/c + DDL 单一事实源/localtime + TRD-TEST- 注入 + A-1~A-5/A-7/A-9/trace/审计 JSON 全部达成,196 绿实跑,B7 复审 P2-1/P3-2 落地。**P1-1 `conftest.py` `except pytest.SkipRequested` 引用不存在属性,演示数据缺失时 skip 路径反噬整个 pytest 收集(未灌库机器单测全跑不了)——已修(Skipped 继承 BaseException 直传,故障注入验证:错密码收集期 skip 不 error)**;P3 顺手项已落:teardown `l3_snapshot` 哨兵(防 setup 失败掩盖原始异常)、A-1 补 trace 环①响应头与环②suitability_log.trace_id。挂账:**P2-1 集成测试对演示库当日状态隐式耦合(A-4 依赖 CUST-9527 当日无历史交易、A-5 aml 单查询未按本测试 trade_id 定位、A-1 同日已有 pending 单会并入旧单)→ B9b 核查单⑥,演示与 pytest 同日交叉前必须处理**;P3-3 `_state` 模块级顺序耦合(单跑 A-7 等 KeyError,文件头已声明保序,可接受留痕);P3-4 阈值冒烟用例在 .env 覆盖阈值时会红(注释已声明前提);P3-5 teardown 时间窗隐含 MySQL 时钟==本机时钟(docstring 已声明单机约定)。 | B9b | 演示链路走查:`reset.ps1` → `prepare_risk_demo.sql` → agent 库建表 → `seed-aml-list.sql` → Swagger 逐条过 **A-1~A-5、A-7~A-9(A-6 归 M3)**。**B6 挂账核查单:① 预警类 API 响应体含固定 disclaimer「本预警由系统自动生成,最终判定需经风控专员人工审核」(PRD §6/规则表 §5,B6 复审 P3-7)② aml/scan 幂等防护(重复扫描同命中客户重复出单,演示点击即复现,B6 评审 P3-6)③ B6 时代码以 TestClient 独立挂 router 等价验证,本次补一次真 Swagger 手测(B6 复审遗漏⑤)④ analyst 台账只读权限扩展待 Wave 3 分析 Agent 接入时定权限矩阵(手册 §5.3 有 risk:alert:read,现 fail-closed 拒绝,B6 复审观察③)⑤ 生产/演示机 `app_env=development` 误配检查(B7 复审 P3-3:误配时 debug 头可达且启动校验放行,列入演示 SOP)⑥ B8 复审 P2-1:集成测试对演示库当日状态的隐式耦合前置断言(A-4 前断言 CUST-9527 当日事件计数==0、A-5 aml 单按本测试 trade_id 定位、A-1 同日 pending 单并入防护),演示与 pytest 同日交叉执行前必须处理** | 演示 SOP | 按 PRD §8 验收表逐条打勾 | B8、B9a | diff --git a/scripts/demo/rebuild_alerts.py b/scripts/demo/rebuild_alerts.py new file mode 100644 index 0000000..b3e1a3f --- /dev/null +++ b/scripts/demo/rebuild_alerts.py @@ -0,0 +1,87 @@ +#!/usr/bin/env python3 +"""按 trade_id 幂等重放风控引擎(B9a 补偿 · trade_gateway engine_error 场景)。 + +背景(PRD FR-1 / 架构 §5.3):网关 INSERT core_trade 成功后 process_trade_event +异常时交易已成立——响应带 engine_error=true、审计 decision='risk_engine_error', +预警缺失。用本脚本按 trade_id 补放引擎(也可用于人工补放任何一笔已落库交易)。 + +幂等:trade_id 已存在于任一 risk_alert.payload(含 aml 单与已处置单)即跳过, +不重复出单、不重复 append 事件。重放以 core_trade.traded_at 为引擎事件时点, +RISK-004 窗口判定可复现(engine.py 模块注释口径),不受脚本执行时刻影响。 + +用法: + python scripts/demo/rebuild_alerts.py TRD-20260907-AB12CD34 [更多 trade_id ...] +""" + +from __future__ import annotations + +import argparse +import sys +from pathlib import Path +from typing import Any + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402 +from app.repository.risk_repository import RiskRepository # noqa: E402 +from app.service.risk.engine import process_trade_event # noqa: E402 + + +def rebuild_trade( + trade_id: str, + core: CoreReadOnlyRepository, + repo: RiskRepository, +) -> dict[str, Any]: + """重放单笔已落库交易(CLI 与单测共用入口)。 + + 返回 {"state": "rebuilt"|"skipped"|"missing", "trade_id", "alert_ids", ...}: + rebuilt 附 triggered_rules/aml_hit;skipped 附已入的预警单 alert_ids; + missing 表示 core_trade 无此流水(未做任何写操作)。 + """ + trade = core.get_trade_by_id(trade_id) + if trade is None: + return {"state": "missing", "trade_id": trade_id, "alert_ids": []} + existing = repo.find_alerts_by_trade(trade_id) + if existing: + return { + "state": "skipped", + "trade_id": trade_id, + "alert_ids": [row["alert_id"] for row in existing], + } + result = process_trade_event(trade, core_ro=core, risk_repo=repo) + return { + "state": "rebuilt", + "trade_id": trade_id, + "alert_ids": result["alert_ids"], + "triggered_rules": result["triggered_rules"], + "aml_hit": result["aml_hit"], + } + + +def main() -> None: + parser = argparse.ArgumentParser(description="按 trade_id 幂等重放风控引擎(引擎异常补偿)") + parser.add_argument("trade_ids", nargs="+", help="core_trade 中的 trade_id,可一次多笔") + args = parser.parse_args() + + core = CoreReadOnlyRepository() + repo = RiskRepository() + + missing = 0 + for tid in args.trade_ids: + out = rebuild_trade(tid, core, repo) + if out["state"] == "rebuilt": + rules = ",".join(out["triggered_rules"]) or "-" + print(f"[rebuilt] {tid}: rules={rules} alerts={out['alert_ids']} aml={out['aml_hit']}") + elif out["state"] == "skipped": + print(f"[skipped] {tid}: 已入预警单 {out['alert_ids']}(幂等跳过)") + else: + print(f"[missing] {tid}: core_trade 无此流水", file=sys.stderr) + missing += 1 + + if missing: + raise SystemExit(f"{missing}/{len(args.trade_ids)} 笔 trade_id 未找到") + + +if __name__ == "__main__": + main() diff --git a/scripts/demo/subscribe_alerts.py b/scripts/demo/subscribe_alerts.py new file mode 100644 index 0000000..b714a0c --- /dev/null +++ b/scripts/demo/subscribe_alerts.py @@ -0,0 +1,95 @@ +#!/usr/bin/env python3 +"""risk:pub:alert 预警推送订阅演示(B9a · PRD FR-4 通知 / A-3 验收第三条)。 + +演示 SOP: + 终端 1 python scripts/demo/subscribe_alerts.py + 终端 2 uvicorn app.main:app --reload → Swagger 发演示交易 + (A-3:CUST-3001 申购 50 万 PROD-510300,RISK-001/002 命中) + 终端 1 应实时打印预警推送行;无推送时每 30s 打印心跳提示仍在监听。 + +payload 口径(PRD §4 FR-4 · 02-redis-keys.md §2.4): + {alert_id, alert_type, customer_id_mask, risk_score, trace_id, notify_role} + +--duration N:监听 N 秒后自动退出(默认 0 = 不限,Ctrl+C 退出)。 +""" + +from __future__ import annotations + +import argparse +import json +import sys +import time +from pathlib import Path +from typing import Any + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from app.config.settings import settings # noqa: E402 + +CHANNEL = "risk:pub:alert" + + +def format_alert(payload: dict[str, Any]) -> str: + """推送 payload → 单行可读文本(演示展示用,字段缺失容错)。""" + roles = ",".join(payload.get("notify_role") or []) + return ( + f"[{payload.get('alert_type', '?'):>12}] {payload.get('alert_id', '?')}" + f" score={payload.get('risk_score', '?')}" + f" customer={payload.get('customer_id_mask', '?')}" + f" trace={payload.get('trace_id', '?')}" + f" notify=[{roles}]" + ) + + +def run(duration: float) -> None: + import redis + + client = redis.Redis.from_url(settings.redis_url, decode_responses=True) + try: + client.ping() + except Exception as exc: + print(f"Redis 不可达({settings.redis_url}):{exc}", file=sys.stderr) + print("请先启动本机 Redis 服务(见 FLOW §0 本机状态)。", file=sys.stderr) + raise SystemExit(1) + + pubsub = client.pubsub(ignore_subscribe_messages=True) + pubsub.subscribe(CHANNEL) + print(f"已订阅 {CHANNEL} @ {settings.redis_url}(Ctrl+C 退出)") + print("等待预警推送…(另开终端发演示交易,SOP 见脚本头注释)") + + deadline = time.monotonic() + duration if duration > 0 else None + received = 0 + last_heartbeat = time.monotonic() + try: + while True: + msg = pubsub.get_message(timeout=1.0) + if msg and msg.get("type") == "message": + received += 1 + try: + payload = json.loads(msg["data"]) + except (TypeError, ValueError): + payload = {"raw": str(msg["data"])} + print(f"{time.strftime('%H:%M:%S')} #{received} {format_alert(payload)}") + elif deadline is not None and time.monotonic() >= deadline: + print(f"已监听 {duration:g}s,共收到 {received} 条推送,退出。") + return + elif deadline is None and time.monotonic() - last_heartbeat >= 30: + last_heartbeat = time.monotonic() + print(f"…监听中(已收到 {received} 条)") + except KeyboardInterrupt: + print(f"\n退出(共收到 {received} 条推送)。") + finally: + pubsub.close() + client.close() + + +def main() -> None: + parser = argparse.ArgumentParser(description="订阅 risk:pub:alert(风控预警推送演示)") + parser.add_argument("--duration", type=float, default=0.0, help="监听秒数;0=不限(Ctrl+C 退出)") + args = parser.parse_args() + run(args.duration) + + +if __name__ == "__main__": + main() diff --git a/tests/test_demo_scripts.py b/tests/test_demo_scripts.py new file mode 100644 index 0000000..ace1b0b --- /dev/null +++ b/tests/test_demo_scripts.py @@ -0,0 +1,152 @@ +"""B9a 演示/运维脚本单测(开发计划 B9a · sqlite)。 + +rebuild_alerts:补偿重放出单并推送、幂等跳过(防重复 append/aml 重复出单)、 +missing 不落库;subscribe_alerts:payload → 单行可读文本。 +脚本目录非包,动态入 sys.path 后按模块名导入。 +""" + +import sys +from datetime import datetime +from decimal import Decimal +from pathlib import Path + +import pytest +from sqlalchemy import text + +from _ddl import create_sqlite_engine + +DEMO_DIR = Path(__file__).resolve().parents[1] / "scripts" / "demo" +sys.path.insert(0, str(DEMO_DIR)) + +from rebuild_alerts import rebuild_trade # noqa: E402 +from subscribe_alerts import format_alert # noqa: E402 + +from app.repository.core_ro import CoreReadOnlyRepository +from app.repository.risk_repository import RiskRepository +from app.service.risk import alert_service + + +class FakePublisher: + def __init__(self): + self.messages = [] + + def publish(self, channel, payload): + self.messages.append((channel, payload)) + + +@pytest.fixture() +def env(): + engine = create_sqlite_engine() + with engine.begin() as conn: + conn.execute( + text( + "INSERT INTO core_customer (customer_id, display_name, age, is_active) VALUES" + " ('C1', '张某某', 40, 1), ('C2', '李四', 35, 1)" + ) + ) + conn.execute( + text( + "INSERT INTO core_product (product_id, product_name, min_risk_code, product_type)" + " VALUES ('P1', '测试混合基金', 'R3', 'mixed')" + ) + ) + conn.execute( + text( + "INSERT INTO risk_aml_list (list_id, list_type, full_name, match_threshold," + " source, list_version, is_active) VALUES ('PEP-1', 'pep', '李四', 0.85, 'mock', 'v1', 1)" + ) + ) + core = CoreReadOnlyRepository(engine=engine) + repo = RiskRepository(engine=engine) + pub = FakePublisher() + alert_service.set_publisher(pub) + yield core, repo, pub, engine + alert_service.set_publisher(None) + engine.dispose() + + +def _seed_trade(conn, trade_id, amount, customer="C1", at=datetime(2026, 9, 6, 14, 0, 0)): + """直接落 core_trade 不调引擎(engine_error 补偿场景:交易已成立、预警缺失)。""" + conn.execute( + text( + "INSERT INTO core_trade (trade_id, customer_id, product_id, trade_type, amount," + " trade_status, traded_at) VALUES (:tid, :cid, 'P1', 'subscribe', :amt, 'confirmed', :at)" + ), + {"tid": trade_id, "cid": customer, "amt": amount, "at": at}, + ) + + +def _counts(engine, table, where="1=1"): + with engine.connect() as conn: + return conn.execute(text(f"SELECT COUNT(*) FROM {table} WHERE {where}")).scalar_one() + + +def test_rebuild_creates_alert_and_publishes(env): + core, repo, pub, engine = env + with engine.begin() as conn: + _seed_trade(conn, "TRD-TEST-RB1", "600000") + out = rebuild_trade("TRD-TEST-RB1", core, repo) + assert out["state"] == "rebuilt" + assert out["triggered_rules"] == ["RISK-001", "RISK-002"] + assert len(out["alert_ids"]) == 1 and out["aml_hit"] is False + alert = repo.get_alert(out["alert_ids"][0]) + assert alert["status"] == "pending_review" and alert["risk_score"] == 70 + (channel, _), = pub.messages + assert channel == "risk:pub:alert" + + +def test_rebuild_idempotent_skips_second_run(env): + core, repo, pub, engine = env + with engine.begin() as conn: + _seed_trade(conn, "TRD-TEST-RB2", "600000") + first = rebuild_trade("TRD-TEST-RB2", core, repo) + second = rebuild_trade("TRD-TEST-RB2", core, repo) + assert first["state"] == "rebuilt" and second["state"] == "skipped" + assert second["alert_ids"] == first["alert_ids"] + assert _counts(engine, "risk_alert") == 1 + assert len(repo.get_alert(first["alert_ids"][0])["payload"]["events"]) == 1 + assert len(pub.messages) == 1 # 重放不重复推送 + + +def test_rebuild_missing_trade_touches_nothing(env): + core, repo, pub, engine = env + out = rebuild_trade("TRD-NO-SUCH", core, repo) + assert out["state"] == "missing" and out["alert_ids"] == [] + assert _counts(engine, "risk_alert") == 0 + assert _counts(engine, "audit_log") == 0 + assert pub.messages == [] + + +def test_rebuild_aml_hit_then_idempotent(env): + """aml 单幂等是 LIKE 检查的关键价值:record_aml_alert 本身无去重,重放防二次出单。""" + core, repo, pub, engine = env + with engine.begin() as conn: + _seed_trade(conn, "TRD-TEST-RB3", "1000", customer="C2") + first = rebuild_trade("TRD-TEST-RB3", core, repo) + assert first["state"] == "rebuilt" and first["aml_hit"] is True + aml_ids = [ + aid + for aid in first["alert_ids"] + if repo.get_alert(aid)["alert_type"] == "aml" + ] + assert aml_ids + second = rebuild_trade("TRD-TEST-RB3", core, repo) + assert second["state"] == "skipped" and second["alert_ids"] == first["alert_ids"] + assert _counts(engine, "risk_alert", "alert_type='aml'") == 1 + + +def test_format_alert_renders_payload_fields(): + line = format_alert( + { + "alert_id": "ALT-20260906-ABC", + "alert_type": "aml", + "customer_id_mask": "CUST-9**", + "risk_score": 95, + "trace_id": "tr-1", + "notify_role": ["risk_officer", "compliance"], + } + ) + for frag in ("aml", "ALT-20260906-ABC", "score=95", "CUST-9**", "tr-1", + "risk_officer,compliance"): + assert frag in line + assert "CUST-9527" not in line # payload 只有脱敏掩码,无原始 id