From e00c4ad51d15e8e0bd4685918b9f47af024cfa8b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=8D=BF=E4=BA=91=E7=A7=8B=E6=9C=88?= <15273589815@163.com> Date: Mon, 14 Sep 2026 20:48:18 +0800 Subject: [PATCH] =?UTF-8?q?=E9=A3=8E=E6=8E=A7=E6=BC=94=E7=A4=BA=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E8=84=9A=E6=9C=AC=EF=BC=9A=E9=80=A0=E4=B8=8A=E6=B8=B8?= =?UTF-8?q?=E4=B8=9A=E5=8A=A1=E6=95=B0=E6=8D=AE=EF=BC=8C=E8=AE=A9=E9=A3=8E?= =?UTF-8?q?=E6=8E=A7=E6=89=AB=E6=8F=8F=E5=99=A8=E8=87=AA=E5=B7=B1=E7=94=9F?= =?UTF-8?q?=E6=88=90=E9=A2=84=E8=AD=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 来源 组员通过腾讯会议发来的 `seed_risk_demo_data.py`,按约定放入 `tools/`。 **脚本内容是组员写的**,我只做了两件事:放进仓库、把 docstring 里的用法示例 统一成实际文件名(组员原文写的是 `seed_risk_scan_demo_data.py`,与文件名不符, 已改成 `tools/seed_risk_demo_data.py`),**未改动任何逻辑**。 ## 它做什么 固定号段:客户 **12001–12005**、产品 **159991–159995**(独立号段,不覆盖人工数据 与其他演示数据)。写入客户 / 画像 / 产品 / 账户 / 持仓 / 交易 / 资金流水 / 登录 / 工单,**刻意不写 `fin_risk_alert` 与 `fin_risk_notification`** —— 预警必须由风控 扫描器按规则生成,这样演示才能真实覆盖规则、证据、通知三条链路。 五个场景各自瞄准一条规则: | 客户 | 场景 | 预期规则 | |---|---|---| | 12001 | 大额快进快出(3 天内入金 80 万后赎回 75 万,比例 93.75%) | RW-003 | | 12002 | 高风险适当性错配(缺风险揭示 + 二次确认) | RW-007 | | 12003 | 中风险适当性错配(等级差 1 且缺风险揭示) | RW-007 | | 12004 | 72 岁客户赎回 75 万,达历史均值 3 倍且用非常用设备 | RW-012 | | 12005 | 凌晨 2:15 小额自动定投(关联有效工单,预期合并预警) | RW-015、RW-018 | 模式:`--check`(只读核对)、默认 `--dry-run`(只打印计划、**不连数据库**)、 `--apply`(写入)、`--dry-run-scan`(跑规则引擎后回滚)。 ## 我跑了什么、核到什么 按约定的命令 `--apply --dry-run-scan`: [命中] 客户 12001:RW-003 (高) [命中] 客户 12002:RW-007 (高) [命中] 客户 12003:RW-007 (中) [命中] 客户 12004:RW-012 (高) [命中] 客户 12005:RW-015、RW-018 (低) 规则干跑完成,共命中 5 条演示预警,已回滚,不会写库。 **"已回滚"这句话我没有只信它的打印**,而是直接查了库: - 上游数据确实写入:客户 5 / 产品 5 / 账户 5 / 持仓 5 / 委托 6 / 成交 6 / 资金流水 1 / 登录 1 / 工单 3,演示产品行情 5 行(`source=custom_seed`); - **`fin_risk_alert` = 0、`fin_risk_notification` = 0** —— 回滚确实生效。 跑之前还先确认了它承诺回滚是**成立**的:`RiskRuleEngine.refresh_alerts()` (`risk_scan_service.py:101`)到 565 行之间**没有 `commit()`**,全文件仅有的 `commit/rollback` 在 `:566-568`,属于另一个方法(正式扫描 `scan()`)。 否则 `dry_run_scan` 的 `session.rollback()` 根本回滚不掉,那就是"假干跑"。 ## 复用与兼容 脚本 import 了 `tools/seed_custom_holdings.py` 的 `ensure_user` / `ensure_account` / `ensure_holding` / `ensure_product` / `ensure_nav_history` / `ensure_market_price`, 我逐个核对了签名,**全部匹配**(组员看过那个脚本)。 它自己的 `_next_id` 用 `SELECT MAX(id)+1` 发号:单线程脚本里没问题,且 `sys_user` 与 `fin_product` 本来就没有 AUTO_INCREMENT(被外键引用,迁移 `20260914_baseline_auto_increment` 已说明为何排除)。**这不是缺陷**,不改。 ## 实测 - `ruff check tools/seed_risk_demo_data.py` -> All checks passed - `--apply --dry-run-scan` -> 5/5 场景命中预期规则,退出码 0 - `pytest tests/unit tests/contract` -> **1428 passed, 2 skipped, 1 failed** (那个失败是投顾工作台页面被替换所致,与本次无关) --- tools/seed_risk_demo_data.py | 862 +++++++++++++++++++++++++++++++++++ 1 file changed, 862 insertions(+) create mode 100644 tools/seed_risk_demo_data.py diff --git a/tools/seed_risk_demo_data.py b/tools/seed_risk_demo_data.py new file mode 100644 index 0000000..0514ad1 --- /dev/null +++ b/tools/seed_risk_demo_data.py @@ -0,0 +1,862 @@ +"""生成风控演示所需的上游业务数据。 + +本脚本只写入客户、画像、产品、账户、持仓、交易、资金流水、登录和工单, +不直接写 `fin_risk_alert` 或 `fin_risk_notification`。预警必须由风控扫描器 +按规则生成,确保演示内容能真实覆盖规则、证据和通知链路。 + +固定演示客户为 12001 至 12005,产品编码为 159991 至 159995。 +重复执行按业务编号更新,不重复新增同一批演示记录。 + +用法: + python tools/seed_risk_demo_data.py + python tools/seed_risk_demo_data.py --check +""" + +from __future__ import annotations + +import argparse +import asyncio +import sys +from dataclasses import dataclass +from datetime import UTC, date, datetime, time, timedelta +from decimal import Decimal +from pathlib import Path + +from sqlalchemy import func, select, update +from sqlalchemy.ext.asyncio import AsyncSession + +PROJECT_ROOT = Path(__file__).resolve().parents[1] +if str(PROJECT_ROOT) not in sys.path: + sys.path.insert(0, str(PROJECT_ROOT)) + +from app.core.timeutil import local_zone, to_utc_naive # noqa: E402 +from app.infrastructure.db import SessionFactory # noqa: E402 +from app.model.fund import ( # noqa: E402 + FundCapitalFlow, + FundCustomerProfile, + FundProduct, + FundRiskAssessment, + FundSimAccount, + FundSimOrder, + FundTransaction, +) +from app.model.risk import RiskLoginRecord, RiskUser, RiskWorkOrder # noqa: E402 +from app.service.risk_scan_service import RiskRuleEngine # noqa: E402 +from tools.seed_custom_holdings import ( # noqa: E402 + ensure_account, + ensure_holding, + ensure_market_price, + ensure_nav_history, + ensure_product, + ensure_user, +) + +if hasattr(sys.stdout, "reconfigure"): + sys.stdout.reconfigure(errors="replace") # type: ignore[union-attr] + + +DEMO_PASSWORD = "risk12345" +DEMO_SOURCE = "risk_demo_seed" +DEMO_QUESTIONNAIRE_VERSION = "risk-demo-v1" + +#: 客户编号和产品编号都固定在独立号段,避免覆盖人工数据和其他演示数据。 +DEMO_CUSTOMERS = (12001, 12002, 12003, 12004, 12005) +DEMO_PRODUCTS = ("159991", "159992", "159993", "159994", "159995") + + +@dataclass(frozen=True) +class CustomerSpec: + customer_id: int + name: str + age: int + investor_type: str + total_asset: Decimal + behavior_score: int = 20 + + +@dataclass(frozen=True) +class ProductSpec: + code: str + name: str + risk_level: str + nav: Decimal + risk_disclosure_required: int = 0 + second_confirmation_required: int = 0 + recording_required: int = 0 + + +@dataclass(frozen=True) +class Scenario: + customer_id: int + label: str + expected_rules: tuple[str, ...] + + +CUSTOMER_SPECS: dict[int, CustomerSpec] = { + 12001: CustomerSpec(12001, "演示客户一", 42, "C4", Decimal("2000000.00")), + 12002: CustomerSpec(12002, "演示客户二", 38, "C1", Decimal("1500000.00")), + 12003: CustomerSpec(12003, "演示客户三", 46, "C3", Decimal("1800000.00")), + 12004: CustomerSpec(12004, "演示客户四", 72, "C3", Decimal("5000000.00")), + 12005: CustomerSpec(12005, "演示客户五", 51, "C4", Decimal("1200000.00")), +} + +PRODUCT_SPECS: dict[str, ProductSpec] = { + "159991": ProductSpec("159991", "风控演示稳健ETF", "R3", Decimal("1.000000")), + "159992": ProductSpec( + "159992", + "风控演示高风险ETF", + "R4", + Decimal("1.000000"), + risk_disclosure_required=1, + second_confirmation_required=1, + ), + "159993": ProductSpec( + "159993", + "风控演示进取ETF", + "R4", + Decimal("1.000000"), + risk_disclosure_required=1, + ), + "159994": ProductSpec("159994", "风控演示低风险ETF", "R2", Decimal("1.000000")), + "159995": ProductSpec("159995", "风控演示平衡ETF", "R3", Decimal("1.000000")), +} + +SCENARIOS: tuple[Scenario, ...] = ( + Scenario(12001, "大额快进快出", ("RW-003",)), + Scenario(12002, "高风险适当性错配", ("RW-007",)), + Scenario(12003, "中风险适当性错配", ("RW-007",)), + Scenario(12004, "高龄客户异常赎回", ("RW-012",)), + Scenario(12005, "凌晨小额自动定投", ("RW-015", "RW-018")), +) + + +def _now() -> datetime: + return datetime.now(UTC).replace(tzinfo=None) + + +def _birth_date(age: int) -> date: + today = date.today() + return date(today.year - age, 1, 1) + + +def _local_timestamp(hour: int, minute: int = 0) -> datetime: + zone = local_zone() + now = datetime.now(zone) + target = datetime.combine(now.date(), time(hour, minute), tzinfo=zone) + if target > now: + target -= timedelta(days=1) + return to_utc_naive(target) + + +async def _next_id(session: AsyncSession, model: type) -> int: + pk = model.__table__.primary_key.columns[0] # type: ignore[attr-defined] + value = await session.scalar(select(func.coalesce(func.max(pk), 0))) + return int(value or 0) + 1 + + +async def _ensure_assessment( + session: AsyncSession, + customer_id: int, + investor_type: str, + now: datetime, +) -> None: + existing = await session.scalar( + select(FundRiskAssessment.id).where( + FundRiskAssessment.customer_id == customer_id, + FundRiskAssessment.questionnaire_version == DEMO_QUESTIONNAIRE_VERSION, + ) + ) + answers = {f"q{index}": 1 for index in range(1, 14)} + values = { + "customer_id": customer_id, + "questionnaire_version": DEMO_QUESTIONNAIRE_VERSION, + "answers": answers, + "total_score": 30, + "investor_type": investor_type, + "assessed_at": now, + "valid_until": now + timedelta(days=365), + "created_at": now, + } + if existing is None: + session.add(FundRiskAssessment(id=await _next_id(session, FundRiskAssessment), **values)) + else: + await session.execute( + update(FundRiskAssessment).where(FundRiskAssessment.id == int(existing)).values(**values) + ) + + +async def _ensure_customer( + session: AsyncSession, + spec: CustomerSpec, + *, + now: datetime, +) -> tuple[str, int]: + await ensure_user( + session, + spec.customer_id, + password=DEMO_PASSWORD, + password_changed=True, + ) + account_no, _, _ = await ensure_account(session, spec.customer_id) + account_id = await session.scalar( + select(FundSimAccount.id).where(FundSimAccount.customer_id == spec.customer_id) + ) + if account_id is None: + raise RuntimeError(f"客户 {spec.customer_id} 的模拟账户未创建") + + await session.execute( + update(RiskUser) + .where(RiskUser.id == spec.customer_id) + .values( + investor_type=spec.investor_type, + investor_type_assessed_at=now, + updated_at=now, + ) + ) + + profile = await session.get(FundCustomerProfile, spec.customer_id) + profile_values = { + "trade_account": account_no, + "real_name": spec.name, + "birth_date": _birth_date(spec.age), + "occupation": "演示职业", + "mobile_masked": "138****0000", + "investor_type": spec.investor_type, + "investment_horizon": "3至5年", + "preferred_asset_class": ["固定收益类", "权益类"], + "trading_frequency": "中", + "last_active_at": now, + "total_asset": spec.total_asset, + "behavior_score": spec.behavior_score, + "risk_tags": [], + "opened_at": now - timedelta(days=365), + "updated_at": now, + } + if profile is None: + session.add(FundCustomerProfile(customer_id=spec.customer_id, **profile_values)) + else: + await session.execute( + update(FundCustomerProfile) + .where(FundCustomerProfile.customer_id == spec.customer_id) + .values(**profile_values) + ) + + await _ensure_assessment(session, spec.customer_id, spec.investor_type, now) + return account_no, int(account_id) + + +async def _ensure_product( + session: AsyncSession, + spec: ProductSpec, + *, + now: datetime, +) -> int: + product_id, _ = await ensure_product( + session, + spec.code, + name=spec.name, + exchange="SSE", + category="ETF", + risk_level=spec.risk_level, + nav=spec.nav, + ) + await session.execute( + update(FundProduct) + .where(FundProduct.id == product_id) + .values( + risk_level=spec.risk_level, + risk_disclosure_required=spec.risk_disclosure_required, + second_confirmation_required=spec.second_confirmation_required, + recording_required=spec.recording_required, + status="上市", + updated_at=now, + ) + ) + await ensure_nav_history(session, product_id, spec.nav, 120) + await ensure_market_price(session, product_id, spec.nav) + return int(product_id) + + +async def _ensure_work_order( + session: AsyncSession, + *, + work_order_no: str, + customer_id: int, + product_id: int, + amount: Decimal, + channel: str, + risk_disclosure_ack_at: datetime | None, + second_confirmation_at: datetime | None, + now: datetime, +) -> int: + values = { + "customer_id": customer_id, + "order_type": "风险演示", + "product_id": product_id, + "amount": amount, + "channel": channel, + "risk_rule_hits": None, + "risk_disclosure_ack_at": risk_disclosure_ack_at, + "second_confirmation_at": second_confirmation_at, + "recording_reference": None, + "submitted_at": now, + "status": "已完成", + "work_order_type": "演示工单", + "submitter_id": customer_id, + "priority": "普通", + "request_detail": {"source": DEMO_SOURCE}, + "created_at": now, + "updated_at": now, + } + existing = await session.scalar( + select(RiskWorkOrder.id).where(RiskWorkOrder.work_order_no == work_order_no) + ) + if existing is None: + work_order_id = await _next_id(session, RiskWorkOrder) + session.add(RiskWorkOrder(id=work_order_id, work_order_no=work_order_no, **values)) + return int(work_order_id) + await session.execute( + update(RiskWorkOrder).where(RiskWorkOrder.id == int(existing)).values(**values) + ) + return int(existing) + + +async def _ensure_order( + session: AsyncSession, + *, + order_no: str, + customer_id: int, + account_id: int, + product_id: int, + transaction_type: str, + amount: Decimal, + nav: Decimal, + confirmed_at: datetime, +) -> int: + side = "buy" if transaction_type == "申购" else "sell" + quantity = (amount / nav).quantize(Decimal("0.0001")) + values = { + "customer_id": customer_id, + "account_id": account_id, + "product_id": product_id, + "order_side": side, + "price_type": "market", + "quantity": quantity, + "limit_price": None, + "quote_price": nav, + "quote_at": confirmed_at, + "quote_source": DEMO_SOURCE, + "channel": "risk_demo", + "advisor_id": None, + "filled_quantity": quantity, + "average_executed_price": nav, + "status": "已成交", + "risk_rule_hits": None, + "risk_disclosure_ack_at": None, + "second_confirmation_at": None, + "recording_reference": None, + "submitted_at": confirmed_at, + "cancelled_at": None, + "created_at": confirmed_at, + "updated_at": confirmed_at, + } + existing = await session.scalar( + select(FundSimOrder.id).where(FundSimOrder.order_no == order_no) + ) + if existing is None: + order_id = await _next_id(session, FundSimOrder) + session.add(FundSimOrder(id=order_id, order_no=order_no, **values)) + return int(order_id) + await session.execute( + update(FundSimOrder).where(FundSimOrder.id == int(existing)).values(**values) + ) + return int(existing) + + +async def _ensure_transaction( + session: AsyncSession, + *, + transaction_no: str, + order_id: int, + work_order_id: int | None, + customer_id: int, + account_id: int, + product_id: int, + transaction_type: str, + amount: Decimal, + nav: Decimal, + confirmed_at: datetime, +) -> int: + side = "buy" if transaction_type == "申购" else "sell" + quantity = (amount / nav).quantize(Decimal("0.0001")) + values = { + "order_id": order_id, + "work_order_id": work_order_id, + "customer_id": customer_id, + "account_id": account_id, + "product_id": product_id, + "order_side": side, + "transaction_type": transaction_type, + "executed_price": nav, + "nav": nav, + "executed_quantity": quantity, + "shares": quantity, + "gross_amount": amount, + "amount": amount, + "fee_rule_id": None, + "fee_rate_snapshot": Decimal("0.000000"), + "fee_amount": Decimal("0.00"), + "fee": Decimal("0.00"), + "net_amount": amount, + "quote_at": confirmed_at, + "quote_source": DEMO_SOURCE, + "executed_at": confirmed_at, + "confirmed_at": confirmed_at, + "confirmed_by": None, + "auto_confirmed": 1, + "created_at": confirmed_at, + } + existing = await session.scalar( + select(FundTransaction.id).where(FundTransaction.transaction_no == transaction_no) + ) + if existing is None: + transaction_id = await _next_id(session, FundTransaction) + session.add(FundTransaction(id=transaction_id, transaction_no=transaction_no, **values)) + return int(transaction_id) + await session.execute( + update(FundTransaction).where(FundTransaction.id == int(existing)).values(**values) + ) + return int(existing) + + +async def _ensure_capital_flow( + session: AsyncSession, + *, + flow_no: str, + customer_id: int, + account_id: int, + amount: Decimal, + settled_at: datetime, + now: datetime, +) -> None: + values = { + "customer_id": customer_id, + "account_id": account_id, + "transaction_id": None, + "flow_type": "入金", + "amount": amount, + "balance_after": amount, + "status": "成功", + "settled_at": settled_at, + "occurred_at": settled_at, + "payer_name": "演示资金方", + "source_type": "BANK", + "related_work_order_id": None, + "match_status": "已匹配", + "created_at": now, + "updated_at": now, + } + existing = await session.scalar( + select(FundCapitalFlow.id).where(FundCapitalFlow.flow_no == flow_no) + ) + if existing is None: + session.add( + FundCapitalFlow( + id=await _next_id(session, FundCapitalFlow), + flow_no=flow_no, + **values, + ) + ) + else: + await session.execute( + update(FundCapitalFlow).where(FundCapitalFlow.id == int(existing)).values(**values) + ) + + +async def _ensure_login( + session: AsyncSession, + *, + customer_id: int, + device_id: str, + login_at: datetime, + is_common_device: bool, +) -> None: + values = { + "user_id": customer_id, + "login_at": login_at, + "login_result": "成功", + "ip_region": "上海", + "device_id": device_id, + "is_common_device": int(is_common_device), + "failure_reason": None, + "created_at": login_at, + } + existing = await session.scalar( + select(RiskLoginRecord.id).where( + RiskLoginRecord.user_id == customer_id, + RiskLoginRecord.device_id == device_id, + ) + ) + if existing is None: + session.add( + RiskLoginRecord( + id=await _next_id(session, RiskLoginRecord), + **values, + ) + ) + else: + await session.execute( + update(RiskLoginRecord).where(RiskLoginRecord.id == int(existing)).values(**values) + ) + + +async def seed(session: AsyncSession) -> None: + now = _now() + customer_data: dict[int, tuple[str, int]] = {} + product_ids: dict[str, int] = {} + + for spec in CUSTOMER_SPECS.values(): + customer_data[spec.customer_id] = await _ensure_customer(session, spec, now=now) + for spec in PRODUCT_SPECS.values(): + product_ids[spec.code] = await _ensure_product(session, spec, now=now) + + # 每个客户都有一笔持仓,保证客户证据和持仓证据不是空表。 + for customer_id, (account_no, _) in customer_data.items(): + product_code = list(PRODUCT_SPECS)[(customer_id - 12001) % len(PRODUCT_SPECS)] + spec = PRODUCT_SPECS[product_code] + await ensure_holding( + session, + customer_id, + account_no, + product_ids[product_code], + quantity=Decimal("10000"), + nav=spec.nav, + ) + + # RW-003:3 天内入金 80 万,随后赎回 75 万,赎回比例 93.75%。 + _, account_id = customer_data[12001] + product_id = product_ids["159991"] + await _ensure_capital_flow( + session, + flow_no="RISKDEMO-12001-FLOW-001", + customer_id=12001, + account_id=account_id, + amount=Decimal("800000.00"), + settled_at=now - timedelta(hours=24), + now=now, + ) + await _ensure_transaction( + session, + transaction_no="RISKDEMO-12001-TXN-001", + order_id=await _ensure_order( + session, + order_no="RISKDEMO-12001-ORDER-001", + customer_id=12001, + account_id=account_id, + product_id=product_id, + transaction_type="赎回", + amount=Decimal("750000.00"), + nav=PRODUCT_SPECS["159991"].nav, + confirmed_at=now, + ), + work_order_id=None, + customer_id=12001, + account_id=account_id, + product_id=product_id, + transaction_type="赎回", + amount=Decimal("750000.00"), + nav=PRODUCT_SPECS["159991"].nav, + confirmed_at=now, + ) + + # RW-007:高风险适当性错配,缺失风险揭示和二次确认。 + _, account_id = customer_data[12002] + work_order_id = await _ensure_work_order( + session, + work_order_no="RISKDEMO-12002-WO-001", + customer_id=12002, + product_id=product_ids["159992"], + amount=Decimal("300000.00"), + channel="APP", + risk_disclosure_ack_at=None, + second_confirmation_at=None, + now=now, + ) + await _ensure_transaction( + session, + transaction_no="RISKDEMO-12002-TXN-001", + order_id=await _ensure_order( + session, + order_no="RISKDEMO-12002-ORDER-001", + customer_id=12002, + account_id=account_id, + product_id=product_ids["159992"], + transaction_type="申购", + amount=Decimal("300000.00"), + nav=PRODUCT_SPECS["159992"].nav, + confirmed_at=now - timedelta(hours=2), + ), + work_order_id=work_order_id, + customer_id=12002, + account_id=account_id, + product_id=product_ids["159992"], + transaction_type="申购", + amount=Decimal("300000.00"), + nav=PRODUCT_SPECS["159992"].nav, + confirmed_at=now - timedelta(hours=2), + ) + + # RW-007:中风险适当性错配,等级差 1 且缺失风险揭示。 + _, account_id = customer_data[12003] + work_order_id = await _ensure_work_order( + session, + work_order_no="RISKDEMO-12003-WO-001", + customer_id=12003, + product_id=product_ids["159993"], + amount=Decimal("200000.00"), + channel="APP", + risk_disclosure_ack_at=None, + second_confirmation_at=None, + now=now, + ) + await _ensure_transaction( + session, + transaction_no="RISKDEMO-12003-TXN-001", + order_id=await _ensure_order( + session, + order_no="RISKDEMO-12003-ORDER-001", + customer_id=12003, + account_id=account_id, + product_id=product_ids["159993"], + transaction_type="申购", + amount=Decimal("200000.00"), + nav=PRODUCT_SPECS["159993"].nav, + confirmed_at=now - timedelta(hours=3), + ), + work_order_id=work_order_id, + customer_id=12003, + account_id=account_id, + product_id=product_ids["159993"], + transaction_type="申购", + amount=Decimal("200000.00"), + nav=PRODUCT_SPECS["159993"].nav, + confirmed_at=now - timedelta(hours=3), + ) + + # RW-012:72 岁客户赎回 75 万,金额达到历史均值 3 倍且使用非常用设备。 + _, account_id = customer_data[12004] + history_confirmed_at = now - timedelta(days=60) + await _ensure_transaction( + session, + transaction_no="RISKDEMO-12004-HIST-001", + order_id=await _ensure_order( + session, + order_no="RISKDEMO-12004-HIST-ORDER-001", + customer_id=12004, + account_id=account_id, + product_id=product_ids["159995"], + transaction_type="申购", + amount=Decimal("100000.00"), + nav=PRODUCT_SPECS["159995"].nav, + confirmed_at=history_confirmed_at, + ), + work_order_id=None, + customer_id=12004, + account_id=account_id, + product_id=product_ids["159995"], + transaction_type="申购", + amount=Decimal("100000.00"), + nav=PRODUCT_SPECS["159995"].nav, + confirmed_at=history_confirmed_at, + ) + await _ensure_login( + session, + customer_id=12004, + device_id="RISKDEMO-DEVICE-12004", + login_at=now - timedelta(hours=2), + is_common_device=False, + ) + await _ensure_transaction( + session, + transaction_no="RISKDEMO-12004-TXN-001", + order_id=await _ensure_order( + session, + order_no="RISKDEMO-12004-ORDER-001", + customer_id=12004, + account_id=account_id, + product_id=product_ids["159995"], + transaction_type="赎回", + amount=Decimal("750000.00"), + nav=PRODUCT_SPECS["159995"].nav, + confirmed_at=now, + ), + work_order_id=None, + customer_id=12004, + account_id=account_id, + product_id=product_ids["159995"], + transaction_type="赎回", + amount=Decimal("750000.00"), + nav=PRODUCT_SPECS["159995"].nav, + confirmed_at=now, + ) + + # RW-015 + RW-018:凌晨小额交易,关联有效自动定投工单,预期合并预警。 + _, account_id = customer_data[12005] + night_at = _local_timestamp(2, 15) + work_order_id = await _ensure_work_order( + session, + work_order_no="RISKDEMO-12005-WO-001", + customer_id=12005, + product_id=product_ids["159994"], + amount=Decimal("5000.00"), + channel="自动定投", + risk_disclosure_ack_at=None, + second_confirmation_at=None, + now=now, + ) + await _ensure_transaction( + session, + transaction_no="RISKDEMO-12005-TXN-001", + order_id=await _ensure_order( + session, + order_no="RISKDEMO-12005-ORDER-001", + customer_id=12005, + account_id=account_id, + product_id=product_ids["159994"], + transaction_type="申购", + amount=Decimal("5000.00"), + nav=PRODUCT_SPECS["159994"].nav, + confirmed_at=night_at, + ), + work_order_id=work_order_id, + customer_id=12005, + account_id=account_id, + product_id=product_ids["159994"], + transaction_type="申购", + amount=Decimal("5000.00"), + nav=PRODUCT_SPECS["159994"].nav, + confirmed_at=night_at, + ) + + +async def check(session: AsyncSession) -> int: + """只读检查演示数据是否完整,返回非零表示存在缺口。""" + missing: list[str] = [] + for spec in CUSTOMER_SPECS.values(): + user = await session.get(RiskUser, spec.customer_id) + profile = await session.get(FundCustomerProfile, spec.customer_id) + account = await session.scalar( + select(FundSimAccount.id).where(FundSimAccount.customer_id == spec.customer_id) + ) + if user is None or profile is None or account is None: + missing.append(f"客户 {spec.customer_id} 的账号、画像或账户缺失") + continue + if user.investor_type != spec.investor_type or profile.investor_type != spec.investor_type: + missing.append(f"客户 {spec.customer_id} 的风险等级与 sys_user / 画像不一致") + + for scenario in SCENARIOS: + count = await session.scalar( + select(func.count(FundTransaction.id)).where( + FundTransaction.customer_id == scenario.customer_id, + FundTransaction.transaction_no.like(f"RISKDEMO-{scenario.customer_id}-%"), + ) + ) + if not count: + missing.append( + f"场景 {scenario.label} 缺少交易数据,预期规则 {'、'.join(scenario.expected_rules)}" + ) + + if missing: + for item in missing: + print(f"[缺失] {item}") + return 1 + + print("风控演示数据检查通过。") + for scenario in SCENARIOS: + print( + f" 客户 {scenario.customer_id}:{scenario.label}," + f"预期规则 {'、'.join(scenario.expected_rules)}" + ) + return 0 + + +async def dry_run_scan(session: AsyncSession) -> int: + """在事务内运行规则引擎,打印结果后回滚,用于校验场景是否真正命中。""" + alerts = await RiskRuleEngine(session).refresh_alerts() + demo_alerts = [alert for alert in alerts if alert.customer_id in DEMO_CUSTOMERS] + for alert in demo_alerts: + print( + f" [命中] 客户 {alert.customer_id}:{'、'.join(alert.trigger_rule_codes)} " + f"({alert.alert_level})" + ) + await session.rollback() + if not demo_alerts: + print("没有命中新的演示预警。") + return 1 + print(f"规则干跑完成,共命中 {len(demo_alerts)} 条演示预警,已回滚,不会写库。") + return 0 + + +def print_plan() -> None: + """打印默认演示计划,默认模式不连接数据库。""" + print("风控演示数据脚本,当前为 dry-run 模式,不会写入数据库。") + print("固定演示号段:客户 12001-12005,产品 159991-159995。") + for scenario in SCENARIOS: + print( + f" 客户 {scenario.customer_id}:{scenario.label}," + f"预期规则 {'、'.join(scenario.expected_rules)}" + ) + print("确认无误后使用 --apply 正式写入。") + + +async def run(args: argparse.Namespace) -> int: + if args.check: + async with SessionFactory() as session: + return await check(session) + + if args.apply: + async with SessionFactory() as session: + async with session.begin(): + await seed(session) + status = await check(session) + if status != 0: + return status + if args.dry_run_scan: + return await dry_run_scan(session) + return 0 + + if args.dry_run_scan: + async with SessionFactory() as session: + status = await check(session) + if status != 0: + return status + return await dry_run_scan(session) + + print_plan() + return 0 + + +def main() -> int: + parser = argparse.ArgumentParser(description="生成风控扫描演示上游数据") + parser.add_argument("--check", action="store_true", help="只检查现有数据,不写入") + parser.add_argument( + "--dry-run", + action="store_true", + help="只打印生成计划,不连接数据库(默认行为)", + ) + parser.add_argument( + "--apply", + action="store_true", + help="正式写入上游演示数据", + ) + parser.add_argument( + "--dry-run-scan", + action="store_true", + help="运行规则干跑并在结束后回滚,不生成正式预警", + ) + args = parser.parse_args() + return asyncio.run(run(args)) + + +if __name__ == "__main__": + sys.exit(main())