From e62fb1d291cb745c1abd43a645bcc92d8cabc72b Mon Sep 17 00:00:00 2001 From: YUAN Date: Mon, 7 Sep 2026 01:02:12 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20B9a=20=E5=A4=8D=E5=AE=A1=E9=97=AD?= =?UTF-8?q?=E7=8E=AF=E2=80=94=E2=80=94P2-1=20skipped=20=E5=88=86=E6=94=AF?= =?UTF-8?q?=20has=5Fengine=5Ferror=5Faudit=20=E4=B8=AD=E9=97=B4=E6=80=81?= =?UTF-8?q?=E8=AD=A6=E7=A4=BA(input=5Fsummary=20LIKE)+docstring=20?= =?UTF-8?q?=E5=B1=80=E9=99=90=E5=A3=B0=E6=98=8E;=20P2-2=20=E9=9D=9E?= =?UTF-8?q?=E9=A6=96=E7=AC=94=20trade=20=E5=B9=82=E7=AD=89=E7=94=A8?= =?UTF-8?q?=E4=BE=8B;=20P3=20=E9=A1=BA=E6=89=8B:=20=E5=8F=8C=E8=84=9A?= =?UTF-8?q?=E6=9C=AC=E8=BF=9E=E6=8E=A5=E8=87=AA=E6=A3=80/ImportError=20?= =?UTF-8?q?=E5=88=86=E6=94=AF/=E9=9D=9E=20dict=20payload=20raw=20=E5=85=9C?= =?UTF-8?q?=E5=BA=95/get=5Ftrade=5Fby=5Fid=20confirmed=20=E8=BF=87?= =?UTF-8?q?=E6=BB=A4/=E5=8B=BF=E5=B9=B6=E8=A1=8C=E5=A3=B0=E6=98=8E,=208=20?= =?UTF-8?q?=E4=BE=8B=20204=20=E7=BB=BF,=20warning=20=E8=B7=AF=E5=BE=84?= =?UTF-8?q?=E7=9C=9F=E5=BA=93=E5=A4=8D=E9=AA=8C=E5=90=8E=E6=B8=85=E7=90=86?= =?UTF-8?q?;=20docs=20=E5=90=8C=E6=AD=A5=E5=A4=8D=E5=AE=A1=E7=BB=93?= =?UTF-8?q?=E8=AE=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/repository/core_ro.py | 12 ++++- app/repository/risk_repository.py | 16 +++++++ docs/memory/FLOW.md | 4 +- docs/memory/MEMORY.md | 6 +-- docs/memory/TODO.md | 2 +- docs/项目框架设计/开发计划-风控模块.md | 2 + scripts/demo/rebuild_alerts.py | 36 ++++++++++++--- scripts/demo/subscribe_alerts.py | 10 ++++- tests/test_demo_scripts.py | 61 ++++++++++++++++++++++++++ 9 files changed, 135 insertions(+), 14 deletions(-) diff --git a/app/repository/core_ro.py b/app/repository/core_ro.py index 418efcb..d01cd64 100644 --- a/app/repository/core_ro.py +++ b/app/repository/core_ro.py @@ -74,8 +74,16 @@ class CoreReadOnlyRepository: ] 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") + """按 trade_id 查单笔 confirmed 交易流水(rebuild_alerts 补偿重放用,仅 SELECT)。 + + 与引擎统计口径一致只取 confirmed——非 confirmed 流水不经引擎,重放无意义。 + """ + sql = text( + """ + SELECT * FROM core_trade + WHERE trade_id = :tid AND trade_status = 'confirmed' + """ + ) with self._engine.connect() as conn: row = conn.execute(sql, {"tid": trade_id}).mappings().first() return dict(row) if row else None diff --git a/app/repository/risk_repository.py b/app/repository/risk_repository.py index 96a5f3f..39710a4 100644 --- a/app/repository/risk_repository.py +++ b/app/repository/risk_repository.py @@ -158,6 +158,22 @@ class RiskRepository: for r in conn.execute(sql, {"pat": f"%{trade_id}%"}).mappings() ] + def has_engine_error_audit(self, trade_id: str) -> bool: + """存在 decision='risk_engine_error' 的该笔审计(rebuild_alerts 警示用)。 + + audit_log 无 trade_id 列,trade_id 在 input_summary(JSON 列/TEXT JSON) + 内,同 find_alerts_by_trade 的 LIKE 口径。 + """ + sql = text( + """ + SELECT 1 FROM audit_log + WHERE decision = 'risk_engine_error' AND input_summary LIKE :pat + LIMIT 1 + """ + ) + with self._engine.connect() as conn: + return conn.execute(sql, {"pat": f"%{trade_id}%"}).first() is not None + def list_alerts( self, status: str | None = None, diff --git a/docs/memory/FLOW.md b/docs/memory/FLOW.md index 025a756..3afb50f 100644 --- a/docs/memory/FLOW.md +++ b/docs/memory/FLOW.md @@ -48,9 +48,9 @@ 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 评审闭环;B9a = subscribe_alerts/rebuild_alerts 演示与补偿脚本,201 测试绿);剩 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 演示与补偿脚本,204 测试绿);剩 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` 201 绿(集成测试真连本机 MySQL;未灌库的机器自动 skip 集成模块,单测不受影响);本机 Redis 服务在跑(`redis://127.0.0.1:6379/0`,B9a 订阅脚本已验证)。 +**本机已就位状态(2026-09-06 · 已完成上述 ①~⑤,无需重做):** `.env` 已配置(学习项目,`MYSQL_PASSWORD=123456`;NEO4J/DEEPSEEK 留空暂不影响风控);`jinrong_core` + `jinrong_agent` 已灌库(28 客户 / AML 名单 8 条 / 演示测评已刷新 / 归属同步 28 行);`python -m pytest` 204 绿(集成测试真连本机 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` 环境变量传密码(交互执行不受影响,不把密码写进仓库脚本);`redis` 包 requirements 有但初始 `pip install -r` 漏装(B9a 于 2026-09-07 补装,重装环境时留意)。 diff --git a/docs/memory/MEMORY.md b/docs/memory/MEMORY.md index 51f35cd..919557b 100644 --- a/docs/memory/MEMORY.md +++ b/docs/memory/MEMORY.md @@ -9,7 +9,7 @@ **项目是什么:** 金融四 Agent(客户财富 / 代理人 / 数据分析 / 风控)共用数据层与合规底座;**不**互调 LLM,跨 Agent 走 L1/L2/L3 画像与预警表。 -**当前进度:** 需求与表设计已定 · **风控模块已落地 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)。** +**当前进度:** 需求与表设计已定 · **风控模块已落地 B1~B8+B9a**(规则/预警聚合/L3/AML/引擎/交易网关/鉴权+4 API/适当性校验/main 集成+挂账①~⑦/B8 conftest+集成测试/**B9a 订阅+重放补偿脚本**,B7/B8/B9a 均经独立 AI 评审闭环,**204 测试绿**)· 剩 **B9b 演示走查** · **JWT(T-01) / 审计中间件(T-02) / LangGraph 对话线(T-07) 未做**(chat/knowledge/admin 仍空壳)。**开发在分支 `feature/risk`(未合入 main)。** **仓库地图:** @@ -124,7 +124,7 @@ 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(201 用例;集成测试需本机演示数据,未灌库时自动 skip) +测试:python -m pytest(204 用例;集成测试需本机演示数据,未灌库时自动 skip) 运维/演示脚本:scripts/demo/subscribe_alerts.py(订阅推送演示)· rebuild_alerts.py TRD-xxx(引擎异常补偿重放) 配置:.env(见 .env.example) RBAC 联调账号:scripts/dev/rbac-seed-reference.md @@ -166,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` 全量(当前 201 绿)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) +5. 如何验证?(`python -m pytest` 全量(当前 204 绿)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) 大任务:FRAMEWORK/FLOW 与实现状态不符时先更新 memory 再编码(用户确认跳过除外)。 diff --git a/docs/memory/TODO.md b/docs/memory/TODO.md index 81b9d94..e393b37 100644 --- a/docs/memory/TODO.md +++ b/docs/memory/TODO.md @@ -22,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);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-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;engine_error 中断笔 skip 时 warning 提示人工核对;core_ro 新增只读 `get_trade_by_id`)+ tests/test_demo_scripts.py 8 例,204 绿;真库手工验证通过(rebuilt→skipped→missing exit1;订阅端到端收到假消息+真交易推送且 trace 贯通;中间态警示路径复验),验证现场已清理;独立 AI 评审有条件通过→P2-1 中间态警示/P2-2 非首笔幂等用例已闭环**;剩 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 bdd791f..4475d89 100644 --- a/docs/项目框架设计/开发计划-风控模块.md +++ b/docs/项目框架设计/开发计划-风控模块.md @@ -34,6 +34,8 @@ **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 尾注)。 +**B9a 复审(独立 AI 评审 · 2026-09-07):有条件通过 → 已闭环。** 幂等选型(payload LIKE)经穷举核实无误匹配/漏匹配面、三条出单路径全覆盖、已处置单覆盖、红线全守、风格与文档注记准确。闭环:P2-1 部分失败中间态补偿盲区——skipped 分支新增 has_engine_error_audit 检测(audit_log 无 trade_id 列,input_summary LIKE 同口径),命中输出 warning 提示人工核对 aml/审计/L3 缺口 + docstring 局限声明(聚合锚点按执行日/编排非原子);P2-2 补「同日第二笔 append 进同单后 rebuild 非首笔 trade 必 skip 且不重复 append」用例。P3 顺手落:连接自检(subscribe 的 ImportError 分支 + rebuild 双库探针)、非 dict payload 防护(raw 兜底)、get_trade_by_id 加 confirmed 过滤(对齐引擎统计口径)、docstring 勿并行声明。P3 留痕:find_alerts_by_trade 全表扫 + 跨进程并发重放无防护(演示规模可接受,生产化改 JSON_CONTAINS/events 明细表)。204 绿;warning 路径真库复验后现场清理。 + **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 index b3e1a3f..352fa72 100644 --- a/scripts/demo/rebuild_alerts.py +++ b/scripts/demo/rebuild_alerts.py @@ -6,8 +6,16 @@ 预警缺失。用本脚本按 trade_id 补放引擎(也可用于人工补放任何一笔已落库交易)。 幂等:trade_id 已存在于任一 risk_alert.payload(含 aml 单与已处置单)即跳过, -不重复出单、不重复 append 事件。重放以 core_trade.traded_at 为引擎事件时点, -RISK-004 窗口判定可复现(engine.py 模块注释口径),不受脚本执行时刻影响。 +不重复出单、不重复 append 事件。**勿并行运行多个本脚本实例**(聚合锁为进程内 +锁,跨进程并发重放同一 trade_id 无防护)。 + +可复现性与局限: +- RISK-004/005 规则窗口以 core_trade.traded_at 为事件时点,不受脚本执行时刻影响; + 聚合锚点(同日 pending 单查找)按执行日——跨日补放历史交易时事件并入执行日 + 活跃单(events 明细含原始 traded_at 可追溯),即时补偿(交易日=执行日)不受影响。 +- 引擎编排非原子:若当初 engine_error 发生在首张预警单落库之后(aml 单/审计/L3 + 部分缺失的中间态),该笔会被幂等跳过,脚本检测到 risk_engine_error 审计时输出 + warning 提示人工核对,不做选择性补齐。 用法: python scripts/demo/rebuild_alerts.py TRD-20260907-AB12CD34 [更多 trade_id ...] @@ -36,19 +44,26 @@ def rebuild_trade( """重放单笔已落库交易(CLI 与单测共用入口)。 返回 {"state": "rebuilt"|"skipped"|"missing", "trade_id", "alert_ids", ...}: - rebuilt 附 triggered_rules/aml_hit;skipped 附已入的预警单 alert_ids; - missing 表示 core_trade 无此流水(未做任何写操作)。 + rebuilt 附 triggered_rules/aml_hit;skipped 附已入的预警单 alert_ids,当初 + engine_error 中断过的笔附 warning(部分失败中间态,人工核对);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 { + out = { "state": "skipped", "trade_id": trade_id, "alert_ids": [row["alert_id"] for row in existing], } + if repo.has_engine_error_audit(trade_id): + out["warning"] = ( + "该笔曾引擎异常中断(risk_engine_error):已入单可能不完整" + "(aml 单/审计/L3 或缺失),请人工核对预警台账与 audit_log" + ) + return out result = process_trade_event(trade, core_ro=core, risk_repo=repo) return { "state": "rebuilt", @@ -67,6 +82,15 @@ def main() -> None: core = CoreReadOnlyRepository() repo = RiskRepository() + # 连接自检(惰性引擎构造不连库;零写入探针,兼验证双库可达) + try: + core.get_trade_by_id("__connectivity_probe__") + repo.find_alerts_by_trade("__connectivity_probe__") + except Exception as exc: + print(f"MySQL 不可达:{exc}", file=sys.stderr) + print("请确认本机 MySQL 服务已启动(见 FLOW §0 本机状态)。", file=sys.stderr) + raise SystemExit(1) + missing = 0 for tid in args.trade_ids: out = rebuild_trade(tid, core, repo) @@ -75,6 +99,8 @@ def main() -> None: 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']}(幂等跳过)") + if out.get("warning"): + print(f"[warning] {tid}: {out['warning']}", file=sys.stderr) else: print(f"[missing] {tid}: core_trade 无此流水", file=sys.stderr) missing += 1 diff --git a/scripts/demo/subscribe_alerts.py b/scripts/demo/subscribe_alerts.py index b714a0c..01bbd14 100644 --- a/scripts/demo/subscribe_alerts.py +++ b/scripts/demo/subscribe_alerts.py @@ -32,6 +32,8 @@ CHANNEL = "risk:pub:alert" def format_alert(payload: dict[str, Any]) -> str: """推送 payload → 单行可读文本(演示展示用,字段缺失容错)。""" + if "raw" in payload: + return f"[raw] {payload['raw']}" roles = ",".join(payload.get("notify_role") or []) return ( f"[{payload.get('alert_type', '?'):>12}] {payload.get('alert_id', '?')}" @@ -43,7 +45,11 @@ def format_alert(payload: dict[str, Any]) -> str: def run(duration: float) -> None: - import redis + try: + import redis + except ImportError: + print("缺少 redis 包:请执行 python -m pip install -r requirements.txt", file=sys.stderr) + raise SystemExit(1) client = redis.Redis.from_url(settings.redis_url, decode_responses=True) try: @@ -68,6 +74,8 @@ def run(duration: float) -> None: received += 1 try: payload = json.loads(msg["data"]) + if not isinstance(payload, dict): + payload = {"raw": str(msg["data"])} except (TypeError, ValueError): payload = {"raw": str(msg["data"])} print(f"{time.strftime('%H:%M:%S')} #{received} {format_alert(payload)}") diff --git a/tests/test_demo_scripts.py b/tests/test_demo_scripts.py index ace1b0b..7fbd138 100644 --- a/tests/test_demo_scripts.py +++ b/tests/test_demo_scripts.py @@ -24,15 +24,21 @@ 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 +from app.service.risk.alert_service import handle_alert +from app.service.risk.engine import process_trade_event class FakePublisher: def __init__(self): self.messages = [] + self.deletes = [] def publish(self, channel, payload): self.messages.append((channel, payload)) + def delete(self, *keys): + self.deletes.append(keys) + @pytest.fixture() def env(): @@ -93,6 +99,7 @@ def test_rebuild_creates_alert_and_publishes(env): assert alert["status"] == "pending_review" and alert["risk_score"] == 70 (channel, _), = pub.messages assert channel == "risk:pub:alert" + assert _counts(engine, "audit_log", "decision='alert_created'") == 1 # P-05 留痕链 def test_rebuild_idempotent_skips_second_run(env): @@ -135,6 +142,59 @@ def test_rebuild_aml_hit_then_idempotent(env): assert _counts(engine, "risk_alert", "alert_type='aml'") == 1 +def test_rebuild_non_first_trade_of_merged_alert_skips(env): + """B8 复审口径(P2-2):幂等保障依赖 LIKE 对 events[n](n≥1)命中—— + 同日第二笔 append 进同单后 rebuild,必须 skip 且不重复 append。""" + core, repo, pub, engine = env + with engine.begin() as conn: + _seed_trade(conn, "TRD-TEST-RB4", "600000") + _seed_trade(conn, "TRD-TEST-RB5", "600000", at=datetime(2026, 9, 6, 15, 0, 0)) + first = rebuild_trade("TRD-TEST-RB4", core, repo) + assert first["state"] == "rebuilt" + # 第二笔走正常引擎路径(等价网关同步调用)→ append 进同单 + process_trade_event(core.get_trade_by_id("TRD-TEST-RB5"), core_ro=core, risk_repo=repo) + merged_id = first["alert_ids"][0] + assert len(repo.get_alert(merged_id)["payload"]["events"]) == 2 + second = rebuild_trade("TRD-TEST-RB5", core, repo) + assert second["state"] == "skipped" and second["alert_ids"] == [merged_id] + assert len(repo.get_alert(merged_id)["payload"]["events"]) == 2 # 不重复 append + + +def test_rebuild_after_disposal_still_skips(env): + """已处置单仍 skip(重放不得绕过处置结论,评审 P3-7)。""" + core, repo, pub, engine = env + with engine.begin() as conn: + _seed_trade(conn, "TRD-TEST-RB6", "600000") + out = rebuild_trade("TRD-TEST-RB6", core, repo) + aid = out["alert_ids"][0] + handle_alert(aid, "confirmed_normal", "STAFF-R1", risk_repo=repo) + again = rebuild_trade("TRD-TEST-RB6", core, repo) + assert again["state"] == "skipped" and again["alert_ids"] == [aid] + assert _counts(engine, "risk_alert") == 1 + + +def test_rebuild_warns_on_engine_error_middle_state(env): + """P2-1:首张单落库后 engine_error 中断(部分失败中间态)→ skip + warning 提示人工核对。""" + core, repo, pub, engine = env + with engine.begin() as conn: + _seed_trade(conn, "TRD-TEST-RB7", "600000") + out = rebuild_trade("TRD-TEST-RB7", core, repo) + assert "warning" not in out # 正常补偿无警示 + # 模拟中断留痕:同一笔再走一遍"交易已成立但引擎异常"的审计(如 aml 单缺失场景) + repo.insert_audit_log( + { + "trace_id": "tr-x", "event_type": "trade_request", "agent_type": "platform", + "actor_id": "SYSTEM", "customer_id": "C1", "rule_id": None, + "input_summary": {"trade_id": "TRD-TEST-RB7", "error_stage": "process_trade_event"}, + "decision": "risk_engine_error", "risk_score": None, + "handler_id": None, "handler_result": None, "handler_comment": None, + } + ) + again = rebuild_trade("TRD-TEST-RB7", core, repo) + assert again["state"] == "skipped" + assert "人工核对" in again["warning"] + + def test_format_alert_renders_payload_fields(): line = format_alert( { @@ -150,3 +210,4 @@ def test_format_alert_renders_payload_fields(): "risk_officer,compliance"): assert frag in line assert "CUST-9527" not in line # payload 只有脱敏掩码,无原始 id + assert format_alert({"raw": "not-a-json-dict"}) == "[raw] not-a-json-dict"