Files
group_xinghuo_jinrong/tests/test_demo_scripts.py
T
GaoYiYuan_0626 22f2a41192 基金转换 T-12:补偿脚本(详情 + 预警)
阶段二失败(或阶段 1.5 引擎失败)后,仅凭 Core 侧数据把「详情 + 预警」两件事补回来
(PRD §7.1 / 架构 §5.4)。

落地(改 3 + 新增 2 脚本 + 测试 2 文件)
- convert_service:新增公开 compensate_convert(group_id, ...) —— 补偿的服务端单点入口
  · 锁键 convert:rerun:{gid},与 convert_fund 幂等重试路径同一个键
  · 详情侧:非 completed 才补写,复用 _finalize_from_core(不另写第二份阶段二)
  · 预警侧:幂等锚点 = 转出端 out_trade_id,复用 find_alerts_by_trade;
    命中即 skipped,否则跑 process_convert_event
- risk_repository:has_engine_error_audit 加 decision 参数(默认值不变)
  · convert 线阶段 1.5 失败审计用 engine_error,普通交易用 risk_engine_error,不是同一个码
- rebuild_alerts.py:新增 --convert-group(与 location 参数 trade_ids 互斥),薄封装
- 新增 scripts/agent/cleanup_pending_convert.py(架构 §2 与开发计划 §9 指定路径):
  超 convert_compensate_sla_hours 的 pending 占位 → status='expired'(标记不硬删,S2)
- 新增 scripts/dev/verify_convert_compensate.py:真库验证脚本(MySQL 8.0.46)
- 测试 +12:test_convert_service +4(补写 / 幂等 / missing / locked)、
  test_demo_scripts +8(--convert-group 分派与接线 + cleanup 脚本)

state 四态与 CLI 退出码
- rebuilt(0) / skipped(0 幂等) / missing(1 零写入) / locked(3)
- 退出码 2 保留给 argparse 用法错误,故 locked 取 3

真库专属证据(sqlite 单测给不了的,本任务核心增量)
- status='expired' 在 MySQL ENUM 上被接受(sqlite 该列是 VARCHAR,写什么都收)
- created_at < cutoff 在 DATETIME(3) 上的时间边界正确(超时进候选 / 未超时不进 / 复跑幂等)
- input_summary 是真 JSON 列,而 has_engine_error_audit 用 LIKE 判定:脚本先断言
  information_schema 的 DATA_TYPE='json' 再验命中,并反向断言决策码不匹配则不命中

顺带收口(用户指示)
- core_ro.concentration_profile 补 h.qty > 0,与 list_holdings 真正同口径
  · ratio 不变(归零行市值为 0),但 rows 不再多出已清仓产品、不虚占截断判定位
  · 新用例含跨出口一致性断言;突变验证:去掉 qty > 0 → 精准 1 条红

验证
- pytest -q → 731 passed / 3 skipped(基线 719 加 12,零回归)
- 突变验证 3 组精准命中:去掉幂等锚点(1 红)/ 去掉 status 过滤(2 红)/ 补偿无视锁(1 红)
- 真库 verify_convert_compensate.py 34/34,隔离数据零残留
- 全套 7 个真库脚本复跑零回归:seed 全 PASS / apply 24 / service 35 / engine 31 / lots 20 / tools 14 / compensate 34
2026-09-10 18:59:42 +08:00

357 lines
14 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""B9a 演示/运维脚本单测(开发计划 B9a · sqlite)。
rebuild_alerts:补偿重放出单并推送、幂等跳过(防重复 append/aml 重复出单)、
missing 不落库;subscribe_alerts:payload → 单行可读文本。
**T-12 追加**:`rebuild_alerts --convert-group` 的分派与接线、
`cleanup_pending_convert` 超时占位清理(标记不硬删)。
脚本目录非包,动态入 sys.path 后按模块名导入。
"""
import sys
from datetime import datetime, timedelta
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"
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(DEMO_DIR))
sys.path.insert(0, str(ROOT / "scripts" / "agent"))
import rebuild_alerts # noqa: E402
from rebuild_alerts import main, rebuild_convert_group, rebuild_trade # noqa: E402
from subscribe_alerts import format_alert # noqa: E402
from cleanup_pending_convert import cleanup # noqa: E402
from app.repository.convert_repository import ConvertRepository # noqa: E402
from app.repository.core_ro import CoreReadOnlyRepository # noqa: E402
from app.repository.risk_repository import RiskRepository # noqa: E402
from app.service.risk import alert_service # noqa: E402
from app.service.risk.alert_service import handle_alert # noqa: E402
from app.service.risk.engine import process_trade_event # noqa: E402
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():
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"
assert _counts(engine, "audit_log", "decision='alert_created'") == 1 # P-05 留痕链
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_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(
{
"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
assert format_alert({"raw": "not-a-json-dict"}) == "[raw] not-a-json-dict"
# ---------- T-12 · rebuild_alerts --convert-group 分派与接线 ----------
def test_convert_group_delegates_to_service(monkeypatch):
"""`--convert-group` 只是薄封装:所有逻辑必须落在 compensate_convert(不复制实现)。"""
seen: dict = {}
def fake(group_id, **kwargs):
seen["gid"] = group_id
seen["kwargs"] = kwargs
return {"state": "skipped", "alert_ids": ["ALT-9"]}
monkeypatch.setattr(rebuild_alerts, "compensate_convert", fake)
out = rebuild_convert_group("CNV-T12-1", "CORE", "REPO", "CREPO")
assert out["state"] == "skipped"
assert seen["gid"] == "CNV-T12-1"
assert seen["kwargs"] == {"core_ro": "CORE", "risk_repo": "REPO", "convert_repo": "CREPO"}
def test_convert_group_missing_is_reported_without_writes(env):
"""未知 group(Core 侧无流水)→ missing,且不写任何东西。"""
core, repo, pub, engine = env
out = rebuild_convert_group("CNV-T12-NOPE", core, repo, ConvertRepository(engine=engine))
assert out["state"] == "missing" and out["trade_count"] == 0
assert _counts(engine, "risk_alert") == 0
assert _counts(engine, "audit_log") == 0
def test_cli_rejects_convert_group_with_trade_ids(monkeypatch):
"""`--convert-group` 与 trade_id 位置参数互斥(argparse 用法错误 → 退出码 2)。"""
monkeypatch.setattr(
sys, "argv", ["rebuild_alerts.py", "TRD-X", "--convert-group", "CNV-Y"]
)
with pytest.raises(SystemExit) as exc:
main()
assert exc.value.code == 2
def test_cli_requires_one_of_the_two_forms(monkeypatch):
"""两者都不给 → 用法错误(不是静默什么都不做)。"""
monkeypatch.setattr(sys, "argv", ["rebuild_alerts.py"])
with pytest.raises(SystemExit) as exc:
main()
assert exc.value.code == 2
# ---------- T-12 · cleanup_pending_convert:超时占位 → expired ----------
def _seed_convert_detail(engine, rows) -> None:
"""rows: (group_id, status, created_at)。"""
with engine.begin() as conn:
for gid, status, created in rows:
conn.execute(
text(
"INSERT INTO risk_convert_detail"
" (convert_group_id, status, estimated, created_at)"
" VALUES (:g, :s, 0, :c)"
),
{"g": gid, "s": status, "c": created},
)
def _status_map(engine) -> dict:
with engine.connect() as conn:
return {
r["convert_group_id"]: r["status"]
for r in conn.execute(
text("SELECT convert_group_id, status FROM risk_convert_detail")
).mappings()
}
def test_cleanup_expires_only_overdue_pending(sqlite_engine):
"""只动超时的 `pending`:fresh pending / completed / failed 一律不碰。"""
overdue = datetime.now() - timedelta(hours=30)
fresh = datetime.now() - timedelta(hours=1)
_seed_convert_detail(
sqlite_engine,
[
("CNV-OLD", "pending", overdue),
("CNV-NEW", "pending", fresh),
("CNV-DONE", "completed", overdue),
("CNV-FAIL", "failed", overdue),
],
)
summary = cleanup(ConvertRepository(engine=sqlite_engine), 24)
assert summary["sla_hours"] == 24 and summary["dry_run"] is False
assert summary["candidate_count"] == 1
assert summary["expired"] == ["CNV-OLD"]
assert _status_map(sqlite_engine) == {
"CNV-OLD": "expired",
"CNV-NEW": "pending",
"CNV-DONE": "completed",
"CNV-FAIL": "failed",
}
# S2:标记不硬删 —— 行仍在(4 行一行不少)
assert _counts(sqlite_engine, "risk_convert_detail") == 4
def test_cleanup_idempotent_second_run_finds_nothing(sqlite_engine):
"""复跑 → 候选为空(已 expired 的行不再进候选),行数不变。"""
_seed_convert_detail(
sqlite_engine, [("CNV-OLD", "pending", datetime.now() - timedelta(hours=30))]
)
repo = ConvertRepository(engine=sqlite_engine)
first = cleanup(repo, 24)
second = cleanup(repo, 24)
assert first["expired"] == ["CNV-OLD"]
assert second["candidate_count"] == 0 and second["expired"] == []
assert _status_map(sqlite_engine) == {"CNV-OLD": "expired"}
assert _counts(sqlite_engine, "risk_convert_detail") == 1
def test_cleanup_dry_run_writes_nothing(sqlite_engine):
"""--dry-run 只报告候选,状态不动。"""
_seed_convert_detail(
sqlite_engine, [("CNV-OLD", "pending", datetime.now() - timedelta(hours=30))]
)
summary = cleanup(ConvertRepository(engine=sqlite_engine), 24, dry_run=True)
assert summary["dry_run"] is True
assert summary["candidate_count"] == 1 and summary["expired_count"] == 0
assert summary["candidates"][0]["convert_group_id"] == "CNV-OLD"
assert _status_map(sqlite_engine) == {"CNV-OLD": "pending"}
def test_cleanup_honours_custom_sla(sqlite_engine):
"""--hours 覆盖 SLA:2h 前落下的占位在 1h 口径下不算超时、在 3h 口径下算。"""
created = datetime.now() - timedelta(hours=2)
_seed_convert_detail(sqlite_engine, [("CNV-MID", "pending", created)])
repo = ConvertRepository(engine=sqlite_engine)
assert cleanup(repo, 3, dry_run=True)["candidate_count"] == 0
assert cleanup(repo, 1, dry_run=True)["candidate_count"] == 1