From c6ad0885d34283076e7ef64b32a3bfce1bbd1726 Mon Sep 17 00:00:00 2001 From: YUAN Date: Sun, 6 Sep 2026 23:37:03 +0800 Subject: [PATCH] =?UTF-8?q?test:=20B8=20=E6=94=B6=E6=95=9B=E2=80=94?= =?UTF-8?q?=E2=80=94sqlite=20=E6=B5=8B=E8=AF=95=20DDL=20=E5=8D=95=E4=B8=80?= =?UTF-8?q?=E4=BA=8B=E5=AE=9E=E6=BA=90=20=5Fddl.py(10=20=E6=96=87=E4=BB=B6?= =?UTF-8?q?=2045=20=E5=A4=84=E6=89=8B=E5=86=99=20DDL=20=E6=B8=85=E9=9B=B6)?= =?UTF-8?q?+localtime=20=E6=97=B6=E5=8C=BA=E6=94=B6=E6=95=9B(B5=20?= =?UTF-8?q?=E8=AF=84=E5=AE=A1=20P3-4,=20utcnow=20=E8=BE=B9=E7=95=8C?= =?UTF-8?q?=E4=BF=AE=E5=A4=8D)+=E6=B5=8B=E8=AF=95=E5=86=85=E9=81=97?= =?UTF-8?q?=E7=95=99=E5=BB=BA=E8=A1=A8=E6=B8=85=E7=90=86,=20196=20?= =?UTF-8?q?=E7=BB=BF;=20docs=20=E5=90=8C=E6=AD=A5=20B8=20=E5=AE=8C?= =?UTF-8?q?=E6=88=90=E7=8A=B6=E6=80=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/memory/MEMORY.md | 10 ++-- docs/memory/TODO.md | 2 +- tests/test_alert_service.py | 40 ++----------- tests/test_aml_service.py | 87 ++------------------------- tests/test_core_ro_sum.py | 49 +++++++-------- tests/test_main.py | 20 ++----- tests/test_profile_l3.py | 30 ++-------- tests/test_risk_api.py | 53 ++--------------- tests/test_risk_engine.py | 108 ++-------------------------------- tests/test_risk_repository.py | 71 ++-------------------- tests/test_suitability.py | 54 ++--------------- tests/test_trade_gateway.py | 53 +---------------- 12 files changed, 62 insertions(+), 515 deletions(-) diff --git a/docs/memory/MEMORY.md b/docs/memory/MEMORY.md index 518973a..820135c 100644 --- a/docs/memory/MEMORY.md +++ b/docs/memory/MEMORY.md @@ -9,7 +9,7 @@ **项目是什么:** 金融四 Agent(客户财富 / 代理人 / 数据分析 / 风控)共用数据层与合规底座;**不**互调 LLM,跨 Agent 走 L1/L2/L3 画像与预警表。 -**当前进度:** 需求与表设计已定 · **风控模块已落地 B1~B7**(规则/预警聚合/L3/AML/引擎/交易网关/鉴权+4 API/适当性校验/**main 集成+挂账①~⑦**,B7 复审闭环,**183 测试绿**)· 剩 **B8 conftest+集成测试 → B9a 脚本 → B9b 演示走查** · **JWT(T-01) / 审计中间件(T-02) / LangGraph 对话线(T-07) 未做**(chat/knowledge/admin 仍空壳)。**开发在分支 `feature/risk`(未合入 main)。** +**当前进度:** 需求与表设计已定 · **风控模块已落地 B1~B8**(规则/预警聚合/L3/AML/引擎/交易网关/鉴权+4 API/适当性校验/main 集成+挂账①~⑦/**B8 conftest+集成测试**,B7 复审闭环,**196 测试绿**)· 剩 **B9a 脚本 → B9b 演示走查** · **JWT(T-01) / 审计中间件(T-02) / LangGraph 对话线(T-07) 未做**(chat/knowledge/admin 仍空壳)。**开发在分支 `feature/risk`(未合入 main)。** **仓库地图:** @@ -28,7 +28,7 @@ | `scripts/core/*.sql` + `reset.ps1` | **已实现** | Core 模拟库 DDL + 种子 | | `scripts/agent/` `scripts/demo/` | **已实现** | AML 名单种子 + 风控演示数据 | | `scripts/sync/*.py` | **已实现** | 归属同步 + Neo4j 全图 | -| `tests/` | **已实现** | 15 个测试文件 183 用例(sqlite 隔离) | +| `tests/` | **已实现** | 16 个测试文件 196 用例(sqlite 隔离;DDL 单一事实源 `_ddl.py`;`test_integration_risk.py` 走真 MySQL + TRD-TEST- 前缀隔离) | | `docs/需求拆解/` | 已定 | 场景 P0、矩阵、合规原文 | | `docs/PRD/PRD-风控监测Agent.md` | **已冻结** | 风控 PRD v1.0 + 规则表附录 | | `docs/项目框架设计/表设计/` | 已定 | Agent 共用 11 表 + agent 专用 SQL | @@ -49,7 +49,7 @@ 7. uvicorn app.main:app --reload → GET /health;python -m pytest(180 绿) ``` -**下一步开发(见 TODO):** 风控 B8 conftest+集成测试(A-1~A-5/A-7~A-9 + trace 一致性 + sqlite DDL 收敛)→ B9a 演示/运维脚本 → B9b 演示走查;随后 Wave 0 T-01 JWT / T-02 审计中间件 / T-06 / T-07 LangGraph。 +**下一步开发(见 TODO):** 风控 B9a 演示/运维脚本(subscribe_alerts.py / rebuild_alerts.py)→ B9b 演示走查(含 B6/B7 挂账核查单);随后 Wave 0 T-01 JWT / T-02 审计中间件 / T-06 / T-07 LangGraph。 **禁止(改代码前必记):** Core 正式 C1~C5 不可被画像覆盖 · 审计表只 INSERT · 代理人草稿不外发 · 仅 R-02 可阻断交易 · 四 Agent 不互调 LLM。 @@ -122,7 +122,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(183 用例,sqlite 隔离) +测试:python -m pytest(196 用例;集成测试需本机演示数据,未灌库时自动 skip) 配置:.env(见 .env.example) RBAC 联调账号:scripts/dev/rbac-seed-reference.md ``` @@ -163,6 +163,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` 全量(当前 183 绿)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) +5. 如何验证?(`python -m pytest` 全量(当前 196 绿)· uvicorn 启动 + /health · SQL / sync 脚本 · 对照 REQUIREMENTS 验收列) 大任务:FRAMEWORK/FLOW 与实现状态不符时先更新 memory 再编码(用户确认跳过除外)。 diff --git a/docs/memory/TODO.md b/docs/memory/TODO.md index f3831a2..ed0fc81 100644 --- a/docs/memory/TODO.md +++ b/docs/memory/TODO.md @@ -20,7 +20,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 conftest+集成测试、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),196 绿;B8 复审待做**;剩 B9a 脚本、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/tests/test_alert_service.py b/tests/test_alert_service.py index c42e926..3fbd3e4 100644 --- a/tests/test_alert_service.py +++ b/tests/test_alert_service.py @@ -9,8 +9,9 @@ from decimal import Decimal from threading import Thread import pytest -from sqlalchemy import create_engine, text -from sqlalchemy.pool import StaticPool +from sqlalchemy import text + +from _ddl import create_sqlite_engine from app.repository.risk_repository import RiskRepository from app.service.risk import alert_service @@ -40,40 +41,7 @@ class FakePublisher: @pytest.fixture() def env(): - engine = create_engine( - "sqlite://", - poolclass=StaticPool, - connect_args={"check_same_thread": False}, - ) - with engine.begin() as conn: - conn.execute( - text( - """ - CREATE TABLE risk_alert ( - alert_id VARCHAR(64) PRIMARY KEY, trace_id VARCHAR(64), - customer_id VARCHAR(64), trade_id VARCHAR(64), alert_type VARCHAR(16), - triggered_rules TEXT, risk_score INTEGER, - status VARCHAR(24) DEFAULT 'pending_review', payload TEXT, - handler_id VARCHAR(64), handler_result VARCHAR(64), - handler_comment VARCHAR(512), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, handled_at TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE audit_log ( - id INTEGER PRIMARY KEY, trace_id VARCHAR(64), event_type VARCHAR(64), - agent_type VARCHAR(16), actor_id VARCHAR(64), customer_id VARCHAR(64), - rule_id VARCHAR(64), input_summary TEXT, decision VARCHAR(64), - risk_score INTEGER, handler_id VARCHAR(64), handler_result VARCHAR(64), - handler_comment VARCHAR(512), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) repo = RiskRepository(engine=engine) pub = FakePublisher() set_publisher(pub) diff --git a/tests/test_aml_service.py b/tests/test_aml_service.py index b98edcc..eb400b7 100644 --- a/tests/test_aml_service.py +++ b/tests/test_aml_service.py @@ -7,8 +7,9 @@ sqlite StaticPool 内存库;名单阈值边界用可控 threshold 值驱动( from datetime import datetime import pytest -from sqlalchemy import create_engine, text -from sqlalchemy.pool import StaticPool +from sqlalchemy import text + +from _ddl import create_sqlite_engine from app.repository.core_ro import CoreReadOnlyRepository from app.repository.risk_repository import RiskRepository @@ -32,88 +33,8 @@ class FakePublisher: @pytest.fixture() def env(): - engine = create_engine( - "sqlite://", - poolclass=StaticPool, - connect_args={"check_same_thread": False}, - ) + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) with engine.begin() as conn: - conn.execute( - text( - """ - CREATE TABLE risk_aml_list ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - list_id VARCHAR(64), list_type VARCHAR(16), full_name VARCHAR(128), - match_threshold REAL, source VARCHAR(64), list_version VARCHAR(16), - effective_date DATE, is_active TINYINT DEFAULT 1, - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE core_customer ( - customer_id VARCHAR(64) PRIMARY KEY, display_name VARCHAR(128), - age INTEGER, occupation VARCHAR(64), open_date DATE, is_active TINYINT DEFAULT 1 - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE core_customer_risk ( - customer_id VARCHAR(64), risk_code VARCHAR(8), evaluated_at TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE risk_alert ( - alert_id VARCHAR(64) PRIMARY KEY, trace_id VARCHAR(64), - customer_id VARCHAR(64), trade_id VARCHAR(64), alert_type VARCHAR(16), - triggered_rules TEXT, risk_score INTEGER, - status VARCHAR(24) DEFAULT 'pending_review', payload TEXT, - handler_id VARCHAR(64), handler_result VARCHAR(64), - handler_comment VARCHAR(512), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, handled_at TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE audit_log ( - id INTEGER PRIMARY KEY, trace_id VARCHAR(64), event_type VARCHAR(64), - agent_type VARCHAR(16), actor_id VARCHAR(64), customer_id VARCHAR(64), - rule_id VARCHAR(64), input_summary TEXT, decision VARCHAR(64), - risk_score INTEGER, handler_id VARCHAR(64), handler_result VARCHAR(64), - handler_comment VARCHAR(512), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE customer_profile_l3 ( - customer_id VARCHAR(64) PRIMARY KEY, - monitor_tier VARCHAR(16) NOT NULL, - risk_score INTEGER, - score_dimensions TEXT, - monitor_tags TEXT, - last_alert_id VARCHAR(64), - computed_at TIMESTAMP NOT NULL, - updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) # 名单:默认阈值 0.85 两行、严格 1.0 一行、行阈值覆盖(0.5/0.99 同名对)一停用行 conn.execute( text( diff --git a/tests/test_core_ro_sum.py b/tests/test_core_ro_sum.py index 8df8b74..5d299b0 100644 --- a/tests/test_core_ro_sum.py +++ b/tests/test_core_ro_sum.py @@ -4,43 +4,34 @@ from datetime import date, datetime from decimal import Decimal import pytest -from sqlalchemy import create_engine, text +from sqlalchemy import text + +from _ddl import create_sqlite_engine from app.repository.core_ro import CoreReadOnlyRepository @pytest.fixture() def repo(): - engine = create_engine("sqlite:///:memory:") - with engine.begin() as conn: - conn.execute( - text( - """ - CREATE TABLE core_trade ( - trade_id VARCHAR(64) PRIMARY KEY, - customer_id VARCHAR(64), - product_id VARCHAR(64), - trade_type VARCHAR(16), - amount NUMERIC(18,2), - trade_status VARCHAR(16) DEFAULT 'confirmed', - traded_at TIMESTAMP - ) - """ - ) - ) + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) def insert(trade_id, amount, traded_at, trade_type="subscribe", status="confirmed", cid="C1"): - conn = engine.connect() - conn.execute( - text( - "INSERT INTO core_trade (trade_id, customer_id, product_id, trade_type," - " amount, trade_status, traded_at)" - " VALUES (:tid, :cid, 'P1', :tt, :amt, :st, :at)" - ), - {"tid": trade_id, "cid": cid, "tt": trade_type, "amt": amount, "st": status, "at": traded_at}, - ) - conn.commit() - conn.close() + with engine.begin() as conn: + conn.execute( + text( + "INSERT INTO core_trade (trade_id, customer_id, product_id, trade_type," + " amount, trade_status, traded_at)" + " VALUES (:tid, :cid, 'P1', :tt, :amt, :st, :at)" + ), + { + "tid": trade_id, + "cid": cid, + "tt": trade_type, + "amt": amount, + "st": status, + "at": traded_at, + }, + ) yield CoreReadOnlyRepository(engine=engine), insert engine.dispose() diff --git a/tests/test_main.py b/tests/test_main.py index cc0ccdf..c09e9ca 100644 --- a/tests/test_main.py +++ b/tests/test_main.py @@ -9,8 +9,9 @@ import re import pytest from fastapi.testclient import TestClient -from sqlalchemy import create_engine, text -from sqlalchemy.pool import StaticPool +from sqlalchemy import text + +from _ddl import create_sqlite_engine from app.config.settings import settings from app.main import app @@ -34,20 +35,7 @@ class FakeGateway: @pytest.fixture() def client(monkeypatch): - engine = create_engine( - "sqlite://", poolclass=StaticPool, connect_args={"check_same_thread": False} - ) - with engine.begin() as conn: - conn.execute( - text( - """CREATE TABLE audit_log ( - id INTEGER PRIMARY KEY, trace_id VARCHAR(64), event_type VARCHAR(64), - agent_type VARCHAR(16), actor_id VARCHAR(64), customer_id VARCHAR(64), - rule_id VARCHAR(64), input_summary TEXT, decision VARCHAR(64), risk_score INTEGER, - handler_id VARCHAR(64), handler_result VARCHAR(64), handler_comment VARCHAR(512), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)""" - ) - ) + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) repo = RiskRepository(engine=engine) from app.api import deps as deps_mod from app.api import risk as risk_api diff --git a/tests/test_profile_l3.py b/tests/test_profile_l3.py index 712d0ef..558c2e7 100644 --- a/tests/test_profile_l3.py +++ b/tests/test_profile_l3.py @@ -8,8 +8,9 @@ risk_score 口径:一期不写(恒 NULL,归 R-05 评分模型首写,评 from datetime import datetime import pytest -from sqlalchemy import create_engine, text -from sqlalchemy.pool import StaticPool +from sqlalchemy import text + +from _ddl import create_sqlite_engine from app.repository.risk_repository import RiskRepository from app.service.risk import redis_gateway @@ -46,30 +47,7 @@ def _fake_redis(monkeypatch): @pytest.fixture() def env(): - engine = create_engine( - "sqlite://", - poolclass=StaticPool, - connect_args={"check_same_thread": False}, - ) - # 口径(评审 P3-3):以 VARCHAR/TEXT 近似真实 DDL 的 ENUM/SMALLINT/JSON, - # ENUM 档位防线不在测试 DB 层,由 tier_of 白名单 + highest_tier 保证。 - with engine.begin() as conn: - conn.execute( - text( - """ - CREATE TABLE customer_profile_l3 ( - customer_id VARCHAR(64) PRIMARY KEY, - monitor_tier VARCHAR(16) NOT NULL, - risk_score INTEGER, - score_dimensions TEXT, - monitor_tags TEXT, - last_alert_id VARCHAR(64), - computed_at TIMESTAMP NOT NULL, - updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) repo = RiskRepository(engine=engine) yield repo, engine engine.dispose() diff --git a/tests/test_risk_api.py b/tests/test_risk_api.py index 7d06a4e..71d5efe 100644 --- a/tests/test_risk_api.py +++ b/tests/test_risk_api.py @@ -8,9 +8,9 @@ from datetime import datetime, timedelta import pytest from fastapi import FastAPI from fastapi.testclient import TestClient -from sqlalchemy import create_engine, text -from sqlalchemy.pool import StaticPool +from sqlalchemy import text +from _ddl import create_sqlite_engine from app.api import risk as risk_api from app.api.deps import assert_customer_access from app.api.risk import router as risk_router @@ -51,51 +51,8 @@ def _alert(alert_id, customer, atype, status="pending_review", score=70): @pytest.fixture() def env(): - engine = create_engine( - "sqlite://", poolclass=StaticPool, connect_args={"check_same_thread": False} - ) + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) with engine.begin() as conn: - for ddl in [ - """CREATE TABLE core_customer ( - customer_id VARCHAR(64) PRIMARY KEY, display_name VARCHAR(128), age INTEGER, - occupation VARCHAR(64), open_date DATE, is_active TINYINT DEFAULT 1)""", - """CREATE TABLE core_customer_risk ( - customer_id VARCHAR(64), risk_code VARCHAR(8), evaluated_at TIMESTAMP)""", - """CREATE TABLE core_customer_advisor ( - advisor_id VARCHAR(64), customer_id VARCHAR(64), rel_status VARCHAR(16))""", - """CREATE TABLE core_product ( - product_id VARCHAR(64) PRIMARY KEY, product_name VARCHAR(128), - min_risk_code VARCHAR(8), product_type VARCHAR(32))""", - """CREATE TABLE risk_aml_list ( - id INTEGER PRIMARY KEY AUTOINCREMENT, list_id VARCHAR(64), list_type VARCHAR(16), - full_name VARCHAR(128), match_threshold REAL, source VARCHAR(64), - list_version VARCHAR(16), effective_date DATE, is_active TINYINT DEFAULT 1, - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)""", - """CREATE TABLE risk_alert ( - alert_id VARCHAR(64) PRIMARY KEY, trace_id VARCHAR(64), customer_id VARCHAR(64), - trade_id VARCHAR(64), alert_type VARCHAR(16), triggered_rules TEXT, - risk_score INTEGER, status VARCHAR(24) DEFAULT 'pending_review', payload TEXT, - handler_id VARCHAR(64), handler_result VARCHAR(64), handler_comment VARCHAR(512), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, handled_at TIMESTAMP)""", - """CREATE TABLE audit_log ( - id INTEGER PRIMARY KEY, trace_id VARCHAR(64), event_type VARCHAR(64), - agent_type VARCHAR(16), actor_id VARCHAR(64), customer_id VARCHAR(64), - rule_id VARCHAR(64), input_summary TEXT, decision VARCHAR(64), risk_score INTEGER, - handler_id VARCHAR(64), handler_result VARCHAR(64), handler_comment VARCHAR(512), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)""", - """CREATE TABLE risk_suitability_log ( - id INTEGER PRIMARY KEY AUTOINCREMENT, trace_id VARCHAR(64), customer_id VARCHAR(64), - product_id VARCHAR(64), customer_risk_level VARCHAR(8), product_risk_level VARCHAR(8), - is_matched TINYINT, is_blocked TINYINT, block_reason VARCHAR(512), - request_ref VARCHAR(64), profile_l1_version VARCHAR(32), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)""", - """CREATE TABLE customer_profile_l3 ( - customer_id VARCHAR(64) PRIMARY KEY, monitor_tier VARCHAR(16) NOT NULL, - risk_score INTEGER, score_dimensions TEXT, monitor_tags TEXT, - last_alert_id VARCHAR(64), computed_at TIMESTAMP NOT NULL, - updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)""", - ]: - conn.execute(text(ddl)) conn.execute( text( "INSERT INTO core_customer (customer_id, display_name, age, is_active) VALUES" @@ -193,9 +150,9 @@ def test_other_roles_cannot_list_alerts(client): def test_alerts_date_filter_and_pagination(client): """复审 P3:start_date/end_date 过滤与分页边界回归保护。 - sqlite CURRENT_TIMESTAMP 为 UTC,窗口用 utcnow 构造(localtime 收敛挂账 B8)。 + DDL 收敛后 created_at 为 localtime(_ddl.py,B5 评审 P3-4),窗口用本地时间构造。 """ - now = datetime.utcnow() + now = datetime.now() r = client.get("/api/risk/alerts", params={"start_date": (now - timedelta(hours=1)).isoformat()}, headers=OFFICER) assert r.json()["total"] == 3 # 窗口内全命中 r = client.get("/api/risk/alerts", params={"start_date": (now + timedelta(hours=1)).isoformat()}, headers=OFFICER) diff --git a/tests/test_risk_engine.py b/tests/test_risk_engine.py index 2431b4b..2057238 100644 --- a/tests/test_risk_engine.py +++ b/tests/test_risk_engine.py @@ -8,8 +8,9 @@ from datetime import datetime from decimal import Decimal import pytest -from sqlalchemy import create_engine, text -from sqlalchemy.pool import StaticPool +from sqlalchemy import text + +from _ddl import create_sqlite_engine from app.repository.core_ro import CoreReadOnlyRepository from app.repository.risk_repository import RiskRepository @@ -33,109 +34,8 @@ class FakePublisher: @pytest.fixture() def env(): - engine = create_engine( - "sqlite://", - poolclass=StaticPool, - connect_args={"check_same_thread": False}, - ) + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) with engine.begin() as conn: - conn.execute( - text( - """ - CREATE TABLE core_customer ( - customer_id VARCHAR(64) PRIMARY KEY, display_name VARCHAR(128), - age INTEGER, occupation VARCHAR(64), open_date DATE, is_active TINYINT DEFAULT 1 - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE core_customer_risk ( - customer_id VARCHAR(64), risk_code VARCHAR(8), evaluated_at TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE core_product ( - product_id VARCHAR(64) PRIMARY KEY, product_name VARCHAR(128), - min_risk_code VARCHAR(8), product_type VARCHAR(32) - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE core_trade ( - trade_id VARCHAR(64) PRIMARY KEY, customer_id VARCHAR(64), - product_id VARCHAR(64), trade_type VARCHAR(16), amount DECIMAL, - trade_status VARCHAR(16), traded_at TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE risk_aml_list ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - list_id VARCHAR(64), list_type VARCHAR(16), full_name VARCHAR(128), - match_threshold REAL, source VARCHAR(64), list_version VARCHAR(16), - effective_date DATE, is_active TINYINT DEFAULT 1, - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE risk_alert ( - alert_id VARCHAR(64) PRIMARY KEY, trace_id VARCHAR(64), - customer_id VARCHAR(64), trade_id VARCHAR(64), alert_type VARCHAR(16), - triggered_rules TEXT, risk_score INTEGER, - status VARCHAR(24) DEFAULT 'pending_review', payload TEXT, - handler_id VARCHAR(64), handler_result VARCHAR(64), - handler_comment VARCHAR(512), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, handled_at TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE audit_log ( - id INTEGER PRIMARY KEY, trace_id VARCHAR(64), event_type VARCHAR(64), - agent_type VARCHAR(16), actor_id VARCHAR(64), customer_id VARCHAR(64), - rule_id VARCHAR(64), input_summary TEXT, decision VARCHAR(64), - risk_score INTEGER, handler_id VARCHAR(64), handler_result VARCHAR(64), - handler_comment VARCHAR(512), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE customer_profile_l3 ( - customer_id VARCHAR(64) PRIMARY KEY, - monitor_tier VARCHAR(16) NOT NULL, - risk_score INTEGER, - score_dimensions TEXT, - monitor_tags TEXT, - last_alert_id VARCHAR(64), - computed_at TIMESTAMP NOT NULL, - updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) conn.execute( text( """ diff --git a/tests/test_risk_repository.py b/tests/test_risk_repository.py index b6a5658..c02c644 100644 --- a/tests/test_risk_repository.py +++ b/tests/test_risk_repository.py @@ -8,68 +8,16 @@ from datetime import datetime from decimal import Decimal import pytest -from sqlalchemy import create_engine, text +from sqlalchemy import text + +from _ddl import create_sqlite_engine from app.repository.risk_repository import RiskRepository @pytest.fixture() def repo(): - engine = create_engine("sqlite:///:memory:") - with engine.begin() as conn: - conn.execute( - text( - """ - CREATE TABLE risk_alert ( - alert_id VARCHAR(64) PRIMARY KEY, - trace_id VARCHAR(64), - customer_id VARCHAR(64), - trade_id VARCHAR(64), - alert_type VARCHAR(16), - triggered_rules TEXT, - risk_score INTEGER, - status VARCHAR(24) DEFAULT 'pending_review', - payload TEXT, - handler_id VARCHAR(64), - handler_result VARCHAR(64), - handler_comment VARCHAR(512), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, - handled_at TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE customer_profile_l3 ( - customer_id VARCHAR(64) PRIMARY KEY, - monitor_tier VARCHAR(8) DEFAULT 'normal', - risk_score INTEGER, - score_dimensions TEXT, - monitor_tags TEXT, - last_alert_id VARCHAR(64), - computed_at TIMESTAMP, - updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE risk_aml_list ( - id INTEGER PRIMARY KEY, - list_id VARCHAR(64), - list_type VARCHAR(8), - full_name VARCHAR(128), - match_threshold NUMERIC(3,2) DEFAULT 0.85, - is_active INTEGER DEFAULT 1 - ) - """ - ) - ) - + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) def make_alert(alert_id, alert_type="large_amount", cid="C1", rules=None, score=70): return { "alert_id": alert_id, @@ -204,17 +152,6 @@ def test_append_alert_event_empty_rules(repo): def test_insert_suitability_log(repo): r, engine, _ = repo - with engine.begin() as conn: - conn.execute( - text( - "CREATE TABLE risk_suitability_log (" - "id INTEGER PRIMARY KEY, trace_id VARCHAR(64), customer_id VARCHAR(64)," - " product_id VARCHAR(64), customer_risk_level VARCHAR(2)," - " product_risk_level VARCHAR(2), is_matched INTEGER, is_blocked INTEGER," - " block_reason VARCHAR(512), request_ref VARCHAR(64), profile_l1_version INTEGER," - " created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)" - ) - ) log = { "trace_id": "trc-x", "customer_id": "C1", diff --git a/tests/test_suitability.py b/tests/test_suitability.py index c53cc8f..c66df01 100644 --- a/tests/test_suitability.py +++ b/tests/test_suitability.py @@ -3,7 +3,9 @@ from datetime import date, datetime, timedelta import pytest -from sqlalchemy import create_engine, text +from sqlalchemy import text + +from _ddl import create_sqlite_engine from app.repository.core_ro import CoreReadOnlyRepository from app.repository.risk_repository import RiskRepository @@ -125,56 +127,8 @@ class TestCheckService: @pytest.fixture() def repos(self): - engine = create_engine("sqlite:///:memory:") + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) with engine.begin() as conn: - conn.execute( - text( - """ - CREATE TABLE core_customer ( - customer_id VARCHAR(64) PRIMARY KEY, display_name VARCHAR(64), - age INTEGER, occupation VARCHAR(64), phone_mask VARCHAR(16), - tenant_id VARCHAR(32), open_date DATE, is_active INTEGER DEFAULT 1, - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE core_customer_risk ( - id INTEGER PRIMARY KEY, customer_id VARCHAR(64), risk_code VARCHAR(2), - is_authoritative INTEGER DEFAULT 1, evaluated_at DATE, - source VARCHAR(32) DEFAULT 'risk_questionnaire' - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE core_product ( - product_id VARCHAR(64) PRIMARY KEY, product_name VARCHAR(128), - product_type VARCHAR(8), min_risk_code VARCHAR(2), - industry_code VARCHAR(16), fee_rate NUMERIC(6,4), - is_open INTEGER DEFAULT 1, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) - conn.execute( - text( - """ - CREATE TABLE risk_suitability_log ( - id INTEGER PRIMARY KEY, trace_id VARCHAR(64), customer_id VARCHAR(64), - product_id VARCHAR(64), customer_risk_level VARCHAR(2), - product_risk_level VARCHAR(2), is_matched INTEGER, is_blocked INTEGER, - block_reason VARCHAR(512), request_ref VARCHAR(64), - profile_l1_version INTEGER, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP - ) - """ - ) - ) conn.execute( text( "INSERT INTO core_customer (customer_id, display_name, age, open_date)" diff --git a/tests/test_trade_gateway.py b/tests/test_trade_gateway.py index b837aca..da94320 100644 --- a/tests/test_trade_gateway.py +++ b/tests/test_trade_gateway.py @@ -11,9 +11,9 @@ import pytest from fastapi import FastAPI from fastapi.testclient import TestClient -from sqlalchemy import create_engine, text -from sqlalchemy.pool import StaticPool +from sqlalchemy import text +from _ddl import create_sqlite_engine from app.api.simulate import router as simulate_router from app.gateway import trade_gateway as tg from app.gateway.gateway_repository import GatewayRepository @@ -37,57 +37,10 @@ class FakePublisher: self.deletes.append(keys) -DDL = [ - """CREATE TABLE core_customer ( - customer_id VARCHAR(64) PRIMARY KEY, display_name VARCHAR(128), age INTEGER, - occupation VARCHAR(64), open_date DATE, is_active TINYINT DEFAULT 1)""", - """CREATE TABLE core_customer_risk ( - customer_id VARCHAR(64), risk_code VARCHAR(8), evaluated_at TIMESTAMP)""", - """CREATE TABLE core_product ( - product_id VARCHAR(64) PRIMARY KEY, product_name VARCHAR(128), - min_risk_code VARCHAR(8), product_type VARCHAR(32))""", - """CREATE TABLE core_trade ( - trade_id VARCHAR(64) PRIMARY KEY, customer_id VARCHAR(64), product_id VARCHAR(64), - trade_type VARCHAR(16), amount DECIMAL, trade_status VARCHAR(16), traded_at TIMESTAMP)""", - """CREATE TABLE risk_aml_list ( - id INTEGER PRIMARY KEY AUTOINCREMENT, list_id VARCHAR(64), list_type VARCHAR(16), - full_name VARCHAR(128), match_threshold REAL, source VARCHAR(64), - list_version VARCHAR(16), effective_date DATE, is_active TINYINT DEFAULT 1, - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)""", - """CREATE TABLE risk_alert ( - alert_id VARCHAR(64) PRIMARY KEY, trace_id VARCHAR(64), customer_id VARCHAR(64), - trade_id VARCHAR(64), alert_type VARCHAR(16), triggered_rules TEXT, - risk_score INTEGER, status VARCHAR(24) DEFAULT 'pending_review', payload TEXT, - handler_id VARCHAR(64), handler_result VARCHAR(64), handler_comment VARCHAR(512), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, handled_at TIMESTAMP)""", - """CREATE TABLE audit_log ( - id INTEGER PRIMARY KEY, trace_id VARCHAR(64), event_type VARCHAR(64), - agent_type VARCHAR(16), actor_id VARCHAR(64), customer_id VARCHAR(64), - rule_id VARCHAR(64), input_summary TEXT, decision VARCHAR(64), risk_score INTEGER, - handler_id VARCHAR(64), handler_result VARCHAR(64), handler_comment VARCHAR(512), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)""", - """CREATE TABLE risk_suitability_log ( - id INTEGER PRIMARY KEY AUTOINCREMENT, trace_id VARCHAR(64), customer_id VARCHAR(64), - product_id VARCHAR(64), customer_risk_level VARCHAR(8), product_risk_level VARCHAR(8), - is_matched TINYINT, is_blocked TINYINT, block_reason VARCHAR(512), - request_ref VARCHAR(64), profile_l1_version VARCHAR(32), - created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)""", - """CREATE TABLE customer_profile_l3 ( - customer_id VARCHAR(64) PRIMARY KEY, monitor_tier VARCHAR(16) NOT NULL, - risk_score INTEGER, score_dimensions TEXT, monitor_tags TEXT, - last_alert_id VARCHAR(64), computed_at TIMESTAMP NOT NULL, - updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP)""", -] - - @pytest.fixture() def env(): - engine = create_engine( - "sqlite://", poolclass=StaticPool, connect_args={"check_same_thread": False} - ) + engine = create_sqlite_engine() # DDL 单一事实源(B4 评审 P3-12) with engine.begin() as conn: - for ddl in DDL: - conn.execute(text(ddl)) conn.execute( text( "INSERT INTO core_customer (customer_id, display_name, age, is_active) VALUES"