diff --git a/.env.example b/.env.example index 6941e88..69a1be7 100644 --- a/.env.example +++ b/.env.example @@ -48,3 +48,22 @@ RISK_PROBE_AMOUNT=400000 RISK_SMALL_AMOUNT=10000 RISK_SMALL_COUNT=3 RISK_AML_DEFAULT_THRESHOLD=0.85 + +# 风控追加 v1.1(FR-8/9/10 · PRD §4A · C4~C6 共用) +# FR-8 RISK-006 集中度:R4+R5 市值占比 ≥ 该阈值即命中(0~1) +RISK_CONCENTRATION_THRESHOLD=0.80 +# FR-9 RISK-007 时效升级:扫描周期(脚本参考)与两级超时小时数 +RISK_ESCALATION_SCAN_MINUTES=15 +RISK_ESCALATION_L1_HOURS=4 +RISK_ESCALATION_L2_HOURS=24 +# AML 单走短通道(1h/4h),与普通单分开计 +RISK_ESCALATION_AML_L1_HOURS=1 +RISK_ESCALATION_AML_L2_HOURS=4 +# FR-10 RISK-008 代理人行为链:扫描周期与 A/B/C 三条件窗口与次数 +RISK_AGENT_BEHAVIOR_SCAN_MINUTES=30 +RISK_AGENT_BEHAVIOR_A_WINDOW_HOURS=24 +RISK_AGENT_BEHAVIOR_A_COUNT=3 +RISK_AGENT_BEHAVIOR_B_WINDOW_HOURS=72 +RISK_AGENT_BEHAVIOR_B_COUNT=5 +RISK_AGENT_BEHAVIOR_C_WINDOW_HOURS=24 +RISK_AGENT_BEHAVIOR_C_COUNT=10 diff --git a/app/config/settings.py b/app/config/settings.py index c561011..afb4bc1 100644 --- a/app/config/settings.py +++ b/app/config/settings.py @@ -55,6 +55,25 @@ class Settings(BaseSettings): risk_small_count: int = 3 risk_aml_default_threshold: Decimal = Decimal("0.85") + # ===== 风控追加 v1.1(FR-8/9/10 · PRD §4A · C4~C6 共用,一次性加齐)===== + # FR-8 RISK-006 集中度:R4+R5 市值占比阈值(≥ 即命中) + risk_concentration_threshold: float = 0.80 + # FR-9 RISK-007 时效升级:扫描周期(脚本侧参考)与两级超时小时数 + risk_escalation_scan_minutes: int = 15 + risk_escalation_l1_hours: int = 4 + risk_escalation_l2_hours: int = 24 + # AML 单走短通道(1h/4h),与普通单分开计 + risk_escalation_aml_l1_hours: int = 1 + risk_escalation_aml_l2_hours: int = 4 + # FR-10 RISK-008 代理人行为链:扫描周期与 A/B/C 三条件窗口与次数 + risk_agent_behavior_scan_minutes: int = 30 + risk_agent_behavior_a_window_hours: int = 24 + risk_agent_behavior_a_count: int = 3 + risk_agent_behavior_b_window_hours: int = 72 + risk_agent_behavior_b_count: int = 5 + risk_agent_behavior_c_window_hours: int = 24 + risk_agent_behavior_c_count: int = 10 + # ===== 输入防护(T-03 · F-03)===== # 对话限流:actor 级固定窗口(拍板 2026-09-07:30 次/分钟,Redis 异常 fail-open) guard_rate_limit_max: int = 30 diff --git a/app/model/entities.py b/app/model/entities.py index 875d628..d2318bc 100644 --- a/app/model/entities.py +++ b/app/model/entities.py @@ -1 +1,835 @@ -"""SQLAlchemy ORM 实体(对齐 docs/项目框架设计/表设计 SQL)。""" +"""全库 ORM 实体定义(29 张表 · 代码内可读的表结构参照)。 + +**定位(先读这一段):** +- **权威表结构 = SQL 文件**:`docs/项目框架设计/表设计/01-mysql-共用底座.sql`、 + `02-mysql-agent专用.sql`(jinrong_agent 17 张)+ `scripts/core/01-ddl.sql` + (jinrong_core 12 张)。本文件与之**人工同步**,仅供 IDE 浏览 / 结构检索 / + 新人理解数据模型,**不用于建表**(建库走 `scripts/core/reset.ps1` + `mysql < *.sql`)。 +- **运行时读写不走 ORM**:repository 层统一用 SQLAlchemy `text()` 原生 SQL + + dict(见 `app/repository/*`、`app/gateway/gateway_repository.py`); + 禁止用 `Base.metadata.create_all()` 建表(会绕过 SQL 单一事实源)。 +- **jinrong_core 为只读库**:Core 正式 C1~C5 / 持仓 / 流水不可被画像覆盖; + 唯一例外 `app/gateway/gateway_repository.py` 仅可 INSERT `core_trade`(B5)。 +- 列注释、枚举值集、索引名 / 唯一键名均与 DDL **字面对齐**; + MySQL 专属精度(如 DATETIME(3)、UNSIGNED)在列 comment 里标注。 +- 每张表的 DDL 出处以类 docstring 标注,改动表结构时两处同步(改表需用户确认)。 +""" + +from __future__ import annotations + +from datetime import date, datetime +from decimal import Decimal +from typing import Any + +from sqlalchemy import ( + CHAR, + JSON, + BigInteger, + Boolean, + Date, + DateTime, + Enum, + ForeignKey, + Index, + Integer, + Numeric, + SmallInteger, + String, + Text, + UniqueConstraint, +) +from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column + +# ============================================================================= +# 两个库 → 两个独立 metadata +# ============================================================================= + + +class AgentBase(DeclarativeBase): + """jinrong_agent 库元数据(四 Agent 共用底座 + 各 Agent 专用表,17 张)。""" + + +class CoreBase(DeclarativeBase): + """jinrong_core 库元数据(Core 模拟底座,12 张,只读)。""" + + +# ============================================================================= +# jinrong_agent · 共用底座第一批(6 张) +# 权威 DDL:docs/项目框架设计/表设计/01-mysql-共用底座.sql +# ============================================================================= + + +class AgentSession(AgentBase): + """【共用】Agent 会话主表 · jinrong_agent.agent_session""" + + __tablename__ = "agent_session" + __table_args__ = ( + UniqueConstraint("session_id", name="uk_session_id"), + Index("idx_trace", "trace_id"), + Index("idx_actor", "agent_type", "actor_id", "created_at"), + Index("idx_customer", "customer_id", "created_at"), + Index("idx_advisor_customer", "advisor_id", "customer_id"), + {"comment": "【共用】Agent 会话主表", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True, comment="DDL: BIGINT UNSIGNED AUTO_INCREMENT") + session_id: Mapped[str] = mapped_column(String(64), comment="对外会话 UUID") + trace_id: Mapped[str] = mapped_column(String(64), comment="全链路追踪 ID") + agent_type: Mapped[str] = mapped_column( + Enum("customer", "advisor", "analyst", "risk"), comment="Agent 类型" + ) + actor_id: Mapped[str] = mapped_column(String(64), comment="操作者:customer_id / staff_id / SYSTEM") + actor_role: Mapped[str] = mapped_column(String(32), comment="customer/advisor/analyst/risk_officer/compliance") + customer_id: Mapped[str | None] = mapped_column(String(64), comment="会话关联客户") + advisor_id: Mapped[str | None] = mapped_column(String(64), comment="代理人归属校验用") + title: Mapped[str | None] = mapped_column(String(256)) + status: Mapped[str] = mapped_column( + Enum("active", "closed", "blocked"), comment="DDL: DEFAULT 'active'" + ) + # 列名 metadata 与 Declarative 基类属性冲突,Python 侧别名 metadata_ + metadata_: Mapped[dict[str, Any] | None] = mapped_column("metadata", JSON, nullable=True) + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + updated_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) ON UPDATE CURRENT_TIMESTAMP(3)") + closed_at: Mapped[datetime | None] = mapped_column(DateTime) + + +class AgentMessage(AgentBase): + """【共用】消息明细 · jinrong_agent.agent_message""" + + __tablename__ = "agent_message" + __table_args__ = ( + Index("idx_session_seq", "session_id", "seq_no"), + Index("idx_trace", "trace_id"), + Index("idx_created", "created_at"), + {"comment": "【共用】消息明细", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + session_id: Mapped[str] = mapped_column(String(64)) + trace_id: Mapped[str] = mapped_column(String(64)) + seq_no: Mapped[int] = mapped_column(Integer, comment="DDL: INT UNSIGNED;同会话内递增序号") + role: Mapped[str] = mapped_column(Enum("user", "assistant", "system", "tool")) + content: Mapped[str] = mapped_column(Text, comment="DDL: MEDIUMTEXT") + content_hash: Mapped[str | None] = mapped_column(CHAR(64)) + token_est: Mapped[int | None] = mapped_column(Integer, comment="DDL: INT UNSIGNED") + has_disclaimer: Mapped[bool] = mapped_column(Boolean, default=False, comment="DDL: TINYINT(1) DEFAULT 0") + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class AgentToolCall(AgentBase): + """【共用】Tool 调用审计(T-04 对话链路留痕)· jinrong_agent.agent_tool_call""" + + __tablename__ = "agent_tool_call" + __table_args__ = ( + Index("idx_session", "session_id", "created_at"), + Index("idx_trace", "trace_id"), + Index("idx_tool", "tool_name", "created_at"), + {"comment": "【共用】Tool 调用审计", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + session_id: Mapped[str] = mapped_column(String(64)) + trace_id: Mapped[str] = mapped_column(String(64)) + message_id: Mapped[int | None] = mapped_column(BigInteger, comment="一期 NULL:Tool 先于 LLM 执行(见 session_repository.insert_tool_call)") + tool_name: Mapped[str] = mapped_column(String(128)) + tool_input: Mapped[dict[str, Any]] = mapped_column(JSON, comment="JSON,调用方序列化传入") + tool_output: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + status: Mapped[str] = mapped_column(Enum("success", "error", "blocked", "timeout")) + error_code: Mapped[str | None] = mapped_column(String(64)) + latency_ms: Mapped[int | None] = mapped_column(Integer, comment="DDL: INT UNSIGNED") + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class AuditLog(AgentBase): + """【共用】审计总账(只 INSERT,禁止 UPDATE/DELETE)· jinrong_agent.audit_log""" + + __tablename__ = "audit_log" + __table_args__ = ( + Index("idx_trace", "trace_id"), + Index("idx_event_time", "event_type", "created_at"), + Index("idx_customer", "customer_id", "created_at"), + Index("idx_actor", "actor_id", "created_at"), + {"comment": "【共用】审计总账(只 INSERT)", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + trace_id: Mapped[str] = mapped_column(String(64)) + event_type: Mapped[str] = mapped_column(String(64), comment="非枚举:VARCHAR(64),事件类型可扩展") + agent_type: Mapped[str] = mapped_column(Enum("customer", "advisor", "analyst", "risk", "platform")) + actor_id: Mapped[str] = mapped_column(String(64)) + customer_id: Mapped[str | None] = mapped_column(String(64)) + rule_id: Mapped[str | None] = mapped_column(String(64)) + input_summary: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + decision: Mapped[str | None] = mapped_column(String(64)) + risk_score: Mapped[int | None] = mapped_column(SmallInteger, comment="DDL: SMALLINT UNSIGNED") + handler_id: Mapped[str | None] = mapped_column(String(64)) + handler_result: Mapped[str | None] = mapped_column(String(64)) + handler_comment: Mapped[str | None] = mapped_column(String(512)) + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class InputGuardLog(AgentBase): + """【共用】输入安全防护留痕(T-03 四类 guard)· jinrong_agent.input_guard_log""" + + __tablename__ = "input_guard_log" + __table_args__ = ( + Index("idx_session", "session_id"), + Index("idx_time", "created_at"), + {"comment": "【共用】输入安全防护", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + trace_id: Mapped[str] = mapped_column(String(64)) + session_id: Mapped[str | None] = mapped_column(String(64)) + agent_type: Mapped[str] = mapped_column(Enum("customer", "advisor", "analyst", "risk")) + actor_id: Mapped[str] = mapped_column(String(64)) + guard_type: Mapped[str] = mapped_column( + Enum("prompt_injection", "oversize", "illegal_param", "rate_limit"), + comment="ENUM 四值已用满(T-03)", + ) + raw_excerpt: Mapped[str | None] = mapped_column(String(1024)) + action: Mapped[str] = mapped_column(Enum("blocked", "sanitized", "passed")) + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class CustomerAdvisorRel(AgentBase): + """【共用】客户-代理人归属(Core 同步,sync_advisor_rel.py 灌)· jinrong_agent.customer_advisor_rel""" + + __tablename__ = "customer_advisor_rel" + __table_args__ = ( + UniqueConstraint("customer_id", "advisor_id", "effective_from", name="uk_customer_advisor"), + Index("idx_advisor", "advisor_id", "rel_status"), + {"comment": "【共用】客户-代理人归属(Core 同步)", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + customer_id: Mapped[str] = mapped_column(String(64)) + advisor_id: Mapped[str] = mapped_column(String(64)) + rel_status: Mapped[str] = mapped_column( + Enum("active", "transferred", "closed"), comment="DDL: DEFAULT 'active'" + ) + effective_from: Mapped[date] = mapped_column(Date) + effective_to: Mapped[date | None] = mapped_column(Date) + synced_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +# ============================================================================= +# jinrong_agent · 共用底座第二批:跨 Agent 交换(5 张) +# ============================================================================= + + +class CustomerProfileL1(AgentBase): + """【交换】L1 客户画像 · 客户 Agent 写(禁止覆盖 Core L0 正式测评)· jinrong_agent.customer_profile_l1""" + + __tablename__ = "customer_profile_l1" + __table_args__ = ({"comment": "【交换】L1 客户画像 · 客户 Agent 写", "mysql_engine": "InnoDB"},) + + customer_id: Mapped[str] = mapped_column(String(64), primary_key=True) + style_tags: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + style_questionnaire: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + allocation_plan: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + behavior_tags: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + attribution_pref: Mapped[str | None] = mapped_column(String(32)) + version: Mapped[int] = mapped_column(Integer, default=1, comment="DDL: INT UNSIGNED DEFAULT 1") + updated_by: Mapped[str] = mapped_column( + Enum("customer_agent", "system"), comment="DDL: DEFAULT 'customer_agent'" + ) + updated_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) ON UPDATE CURRENT_TIMESTAMP(3)") + + +class CustomerProfileL2(AgentBase): + """【交换】L2 服务画像 · 代理人 Agent 写 · jinrong_agent.customer_profile_l2""" + + __tablename__ = "customer_profile_l2" + __table_args__ = ( + UniqueConstraint("customer_id", "advisor_id", name="uk_customer_advisor"), + Index("idx_advisor", "advisor_id"), + {"comment": "【交换】L2 服务画像 · 代理人 Agent 写", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + customer_id: Mapped[str] = mapped_column(String(64)) + advisor_id: Mapped[str] = mapped_column(String(64)) + asset_snapshot: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + demands: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + follow_up_todos: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + service_tags: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + source_session_id: Mapped[str | None] = mapped_column(String(64)) + version: Mapped[int] = mapped_column(Integer, default=1, comment="DDL: INT UNSIGNED DEFAULT 1") + updated_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) ON UPDATE CURRENT_TIMESTAMP(3)") + + +class CustomerProfileL3(AgentBase): + """【交换】L3 监测画像 · 风控 Agent 写 · jinrong_agent.customer_profile_l3""" + + __tablename__ = "customer_profile_l3" + __table_args__ = ({"comment": "【交换】L3 监测画像 · 风控 Agent 写", "mysql_engine": "InnoDB"},) + + customer_id: Mapped[str] = mapped_column(String(64), primary_key=True) + monitor_tier: Mapped[str] = mapped_column( + Enum("normal", "watch", "high"), comment="DDL: DEFAULT 'normal'" + ) + risk_score: Mapped[int | None] = mapped_column(SmallInteger, comment="DDL: SMALLINT UNSIGNED") + score_dimensions: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + monitor_tags: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + last_alert_id: Mapped[str | None] = mapped_column(String(64)) + computed_at: Mapped[datetime] = mapped_column(DateTime) + updated_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) ON UPDATE CURRENT_TIMESTAMP(3)") + + +class RiskAlert(AgentBase): + """【交换】预警单 · 风控写 / 分析读(仅 R-02 可阻断交易)· jinrong_agent.risk_alert""" + + __tablename__ = "risk_alert" + __table_args__ = ( + Index("idx_status_time", "status", "created_at"), + Index("idx_customer", "customer_id", "created_at"), + Index("idx_type", "alert_type", "created_at"), + {"comment": "【交换】预警单 · 风控写 / 分析读", "mysql_engine": "InnoDB"}, + ) + + alert_id: Mapped[str] = mapped_column(String(64), primary_key=True) + trace_id: Mapped[str] = mapped_column(String(64)) + customer_id: Mapped[str] = mapped_column(String(64)) + trade_id: Mapped[str | None] = mapped_column(String(64)) + alert_type: Mapped[str] = mapped_column( + Enum("large_amount", "freq_trade", "suitability", "aml", "pattern"), + comment="C4 集中度事件复用 'pattern' + payload.alert_subtype(FR-8)", + ) + triggered_rules: Mapped[dict[str, Any]] = mapped_column(JSON) + risk_score: Mapped[int | None] = mapped_column(SmallInteger, comment="DDL: SMALLINT UNSIGNED") + status: Mapped[str] = mapped_column( + Enum("pending_review", "confirmed_normal", "confirmed_suspicious", "reported"), + comment="DDL: DEFAULT 'pending_review'", + ) + payload: Mapped[dict[str, Any]] = mapped_column(JSON, comment="C5 升级信息由 payload.escalation_level/escalated_at 承载(FR-9)") + handler_id: Mapped[str | None] = mapped_column(String(64)) + handler_result: Mapped[str | None] = mapped_column(String(64)) + handler_comment: Mapped[str | None] = mapped_column(String(512)) + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + handled_at: Mapped[datetime | None] = mapped_column(DateTime) + + +class RiskSuitabilityLog(AgentBase): + """【交换】适当性记录 · 风控写 / 客户+代理人读(AL-02 对齐 main 21 列契约; + 字段契约详见 docs/项目框架设计/表设计/07-risk_suitability_log说明.md) + · jinrong_agent.risk_suitability_log""" + + __tablename__ = "risk_suitability_log" + __table_args__ = ( + Index("idx_trace", "trace_id"), + Index("idx_customer", "customer_id", "created_at"), + Index("idx_product", "product_id", "created_at"), + Index("idx_blocked", "is_blocked", "created_at"), + Index("idx_match", "match_result", "created_at"), + Index("idx_mismatch", "mismatch_type", "created_at"), + {"comment": "【交换】适当性记录 · 风控写 / 客户+代理人读", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + trace_id: Mapped[str] = mapped_column(String(64)) + customer_id: Mapped[str] = mapped_column(String(64)) + product_id: Mapped[str] = mapped_column(String(64)) + product_name: Mapped[str | None] = mapped_column(String(128), comment="判定时产品名快照") + customer_risk_level: Mapped[str] = mapped_column(CHAR(2), comment="L0 C1~C5") + product_risk_level: Mapped[str] = mapped_column(CHAR(2), comment="产品最低 R1~R5") + investor_category: Mapped[str] = mapped_column( + Enum("ordinary", "professional", "professional_pending"), + comment="DDL: DEFAULT 'ordinary'", + ) + match_result: Mapped[str] = mapped_column( + Enum( + "allowed", "allowed_with_disclosure", "forbidden", + "risk_expired", "professional_exempt", + ), + comment="与 Core check_suitability 输出一致(五值)", + ) + mismatch_type: Mapped[str] = mapped_column( + Enum( + "none", "risk_level", "risk_expired", "age_branch_confirm", + "min_subscribe", "not_found", "professional_exempt", + ), + comment="阻断/特殊处理原因分类(七值),DDL: DEFAULT 'none'", + ) + is_matched: Mapped[bool] = mapped_column(Boolean) + is_blocked: Mapped[bool] = mapped_column(Boolean, default=False, comment="DDL: TINYINT(1) DEFAULT 0") + requires_disclosure: Mapped[bool] = mapped_column(Boolean, default=False, comment="DDL: TINYINT(1) DEFAULT 0") + needs_branch_confirm: Mapped[bool] = mapped_column(Boolean, default=False, comment="FM-01 年龄≥70 买 R3+;DDL: TINYINT(1) DEFAULT 0") + risk_was_expired: Mapped[bool] = mapped_column(Boolean, default=False, comment="判定时风评是否已过期;DDL: TINYINT(1) DEFAULT 0") + block_reason: Mapped[str | None] = mapped_column(String(512), comment="对人可读原因") + block_response_code: Mapped[str | None] = mapped_column(String(32), comment="API 机器码,如 SUIT_RISK_MISMATCH") + check_source: Mapped[str] = mapped_column( + Enum("r02_trade", "r02_chat", "c11_inquiry", "manual"), + comment="DDL: DEFAULT 'r02_trade'", + ) + actor_id: Mapped[str] = mapped_column(String(64), comment="发起者 customer_id / staff_id / svc-trade-suitability") + request_ref: Mapped[str | None] = mapped_column(String(64), comment="交易单号 / session_id") + profile_l1_version: Mapped[int | None] = mapped_column(Integer, comment="DDL: INT UNSIGNED") + rule_refs: Mapped[dict[str, Any] | None] = mapped_column(JSON, comment='依据规则,如 ["JR-AST-012","FM-03"]') + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +# ============================================================================= +# jinrong_agent · 各 Agent 专用(6 张) +# 权威 DDL:docs/项目框架设计/表设计/02-mysql-agent专用.sql +# ============================================================================= + + +class CustomerThresholdConfig(AgentBase): + """【客户专用】亏损阈值配置 · jinrong_agent.customer_threshold_config""" + + __tablename__ = "customer_threshold_config" + __table_args__ = ( + Index("idx_customer", "customer_id", "is_enabled"), + {"comment": "【客户专用】亏损阈值配置", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + customer_id: Mapped[str] = mapped_column(String(64)) + scope_type: Mapped[str] = mapped_column( + Enum("portfolio", "product"), comment="DDL: DEFAULT 'portfolio'" + ) + scope_ref: Mapped[str | None] = mapped_column(String(64)) + loss_threshold_pct: Mapped[Decimal] = mapped_column(Numeric(5, 2)) + notify_channel: Mapped[str] = mapped_column(String(32), comment="DDL: SET('app','sms','email') DEFAULT 'app'") + is_enabled: Mapped[bool] = mapped_column(Boolean, default=True, comment="DDL: TINYINT(1) DEFAULT 1") + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + updated_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) ON UPDATE CURRENT_TIMESTAMP(3)") + + +class CustomerNotifyLog(AgentBase): + """【客户专用】提醒留痕 · jinrong_agent.customer_notify_log""" + + __tablename__ = "customer_notify_log" + __table_args__ = ( + Index("idx_customer_time", "customer_id", "created_at"), + {"comment": "【客户专用】提醒留痕", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + customer_id: Mapped[str] = mapped_column(String(64)) + trace_id: Mapped[str] = mapped_column(String(64)) + notify_type: Mapped[str] = mapped_column(Enum("loss_threshold", "market_volatility")) + threshold_config_id: Mapped[int | None] = mapped_column(BigInteger) + payload: Mapped[dict[str, Any]] = mapped_column(JSON) + channel: Mapped[str] = mapped_column(String(16)) + send_status: Mapped[str] = mapped_column(Enum("sent", "failed")) + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class AdvisorDraft(AgentBase): + """【代理人专用】话术/跟进草稿(草稿不外发客户)· jinrong_agent.advisor_draft""" + + __tablename__ = "advisor_draft" + __table_args__ = ( + UniqueConstraint("draft_id", name="uk_draft_id"), + Index("idx_advisor_customer", "advisor_id", "customer_id", "created_at"), + Index("idx_review", "review_status", "created_at"), + {"comment": "【代理人专用】话术/跟进草稿", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + draft_id: Mapped[str] = mapped_column(String(64)) + session_id: Mapped[str] = mapped_column(String(64)) + trace_id: Mapped[str] = mapped_column(String(64)) + advisor_id: Mapped[str] = mapped_column(String(64)) + customer_id: Mapped[str] = mapped_column(String(64)) + draft_type: Mapped[str] = mapped_column(Enum("script", "follow_up")) + content: Mapped[str] = mapped_column(Text, comment="DDL: MEDIUMTEXT") + review_status: Mapped[str] = mapped_column( + Enum("pending", "approved", "rejected"), comment="DDL: DEFAULT 'pending'" + ) + reviewer_id: Mapped[str | None] = mapped_column(String(64)) + reviewed_at: Mapped[datetime | None] = mapped_column(DateTime) + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class ComplianceHitLog(AgentBase): + """【代理人专用】违规话术命中 · jinrong_agent.compliance_hit_log""" + + __tablename__ = "compliance_hit_log" + __table_args__ = ( + Index("idx_severity_time", "severity", "created_at"), + Index("idx_session", "session_id"), + {"comment": "【代理人专用】违规话术命中", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + session_id: Mapped[str] = mapped_column(String(64)) + trace_id: Mapped[str] = mapped_column(String(64)) + agent_type: Mapped[str] = mapped_column(Enum("customer", "advisor")) + actor_id: Mapped[str] = mapped_column(String(64)) + hit_category: Mapped[str] = mapped_column( + Enum("return_promise", "principal_guarantee", "buy_sell_guide", "product_recommend", "other") + ) + matched_terms: Mapped[dict[str, Any]] = mapped_column(JSON) + severity: Mapped[str] = mapped_column(Enum("low", "medium", "high")) + action_taken: Mapped[str] = mapped_column(Enum("flagged", "blocked", "alerted")) + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class AnalyticsQueryLog(AgentBase): + """【分析专用】查数 SQL 留痕(NL2SQL 审计)· jinrong_agent.analytics_query_log""" + + __tablename__ = "analytics_query_log" + __table_args__ = ( + Index("idx_staff_time", "staff_id", "created_at"), + Index("idx_trace", "trace_id"), + Index("idx_sql_hash", "sql_hash"), + {"comment": "【分析专用】查数 SQL 留痕", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + session_id: Mapped[str] = mapped_column(String(64)) + trace_id: Mapped[str] = mapped_column(String(64)) + staff_id: Mapped[str] = mapped_column(String(64)) + nl_question: Mapped[str] = mapped_column(Text) + generated_sql: Mapped[str] = mapped_column(Text) + sql_hash: Mapped[str] = mapped_column(CHAR(64)) + row_count: Mapped[int | None] = mapped_column(Integer, comment="DDL: INT UNSIGNED") + exec_status: Mapped[str] = mapped_column(Enum("success", "error", "blocked")) + exec_latency_ms: Mapped[int | None] = mapped_column(Integer, comment="DDL: INT UNSIGNED") + result_summary: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True) + has_disclaimer: Mapped[bool] = mapped_column(Boolean, default=False, comment="DDL: TINYINT(1) DEFAULT 0") + error_message: Mapped[str | None] = mapped_column(String(512)) + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class RiskAmlList(AgentBase): + """【风控专用】AML 名单本地镜像(seed-aml-list.sql 灌)· jinrong_agent.risk_aml_list""" + + __tablename__ = "risk_aml_list" + __table_args__ = ( + UniqueConstraint("list_id", name="uk_list_id"), + Index("idx_name", "full_name"), + Index("idx_active", "is_active"), + {"comment": "【风控专用】AML 名单本地镜像 · 依据 docs/PRD/PRD-风控监测Agent.md §6.2", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + list_id: Mapped[str] = mapped_column(String(64)) + list_type: Mapped[str] = mapped_column(Enum("sanction", "terror", "pep")) + full_name: Mapped[str] = mapped_column(String(128), comment="与 core_customer.display_name 同为脱敏展示名口径") + id_no: Mapped[str | None] = mapped_column(String(32), comment="预留:待 Core 提供证件数据后启用匹配") + bank_card_no: Mapped[str | None] = mapped_column(String(32), comment="预留:同上") + match_threshold: Mapped[Decimal] = mapped_column(Numeric(3, 2), default=Decimal("0.85")) + source: Mapped[str] = mapped_column(String(64)) + list_version: Mapped[str] = mapped_column(String(16)) + effective_date: Mapped[date] = mapped_column(Date) + is_active: Mapped[bool] = mapped_column(Boolean, default=True, comment="DDL: TINYINT(1) DEFAULT 1") + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +# ============================================================================= +# jinrong_core · Core 模拟底座(12 张,只读) +# 权威 DDL:scripts/core/01-ddl.sql(AL-01 对齐 main:is_hnw / 风评新列 / 矩阵表) +# ============================================================================= + + +class CoreRiskGrade(CoreBase): + """风险等级字典(客户 C1~C5 · 产品 R1~R5)· jinrong_core.core_risk_grade""" + + __tablename__ = "core_risk_grade" + __table_args__ = ({"comment": "风险等级字典", "mysql_engine": "InnoDB"},) + + code: Mapped[str] = mapped_column(CHAR(4), primary_key=True, comment="C1~C5 或 R1~R5") + grade_type: Mapped[str] = mapped_column(Enum("customer", "product")) + display_name: Mapped[str] = mapped_column(String(32)) + sort_order: Mapped[int] = mapped_column(SmallInteger, comment="DDL: TINYINT UNSIGNED") + + +class CoreIndustry(CoreBase): + """行业分类 · jinrong_core.core_industry""" + + __tablename__ = "core_industry" + __table_args__ = ({"comment": "行业分类", "mysql_engine": "InnoDB"},) + + industry_code: Mapped[str] = mapped_column(String(16), primary_key=True) + industry_name: Mapped[str] = mapped_column(String(64)) + + +class CoreSuitabilityRule(CoreBase): + """适当性匹配规则(L0 权威 · C×R 矩阵,《个人投资者适当性管理指南》第十二条) + · jinrong_core.core_suitability_rule""" + + __tablename__ = "core_suitability_rule" + __table_args__ = ( + UniqueConstraint("customer_risk_code", "product_risk_code", name="uk_cx_r"), + {"comment": "适当性匹配规则(L0 权威)", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(SmallInteger, primary_key=True, autoincrement=True, comment="DDL: TINYINT UNSIGNED AUTO_INCREMENT") + customer_risk_code: Mapped[str] = mapped_column( + CHAR(2), ForeignKey("core_risk_grade.code", name="fk_suit_cust_risk"), comment="C1~C5" + ) + product_risk_code: Mapped[str] = mapped_column( + CHAR(2), ForeignKey("core_risk_grade.code", name="fk_suit_prod_risk"), comment="R1~R5" + ) + match_result: Mapped[str] = mapped_column(Enum("allowed", "allowed_with_disclosure", "forbidden")) + rule_ref: Mapped[str] = mapped_column(String(32), comment="DDL: DEFAULT 'JR-AST-012'") + + +class CoreStaff(CoreBase): + """内部员工主档(模拟 IdP 账号源,含 RBAC 角色种子)· jinrong_core.core_staff""" + + __tablename__ = "core_staff" + __table_args__ = ({"comment": "内部员工主档(模拟 IdP 账号源)", "mysql_engine": "InnoDB"},) + + staff_id: Mapped[str] = mapped_column(String(64), primary_key=True) + display_name: Mapped[str] = mapped_column(String(64)) + staff_type: Mapped[str] = mapped_column( + Enum("advisor", "analyst", "risk_officer", "compliance", "ops") + ) + roles: Mapped[dict[str, Any]] = mapped_column(JSON, comment='JWT roles 数组,如 ["advisor"]') + tenant_id: Mapped[str] = mapped_column(String(32), comment="DDL: DEFAULT 'TENANT-001'") + is_active: Mapped[bool] = mapped_column(Boolean, default=True, comment="DDL: TINYINT(1) DEFAULT 1") + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class CoreCustomer(CoreBase): + """客户主档 L0 · KYC(对齐用户信息数据示例)· jinrong_core.core_customer""" + + __tablename__ = "core_customer" + __table_args__ = ( + Index("idx_service_tier", "service_tier"), + Index("idx_hnw", "is_hnw"), + Index("idx_aml", "aml_risk_level"), + {"comment": "客户主档 L0 · KYC", "mysql_engine": "InnoDB"}, + ) + + customer_id: Mapped[str] = mapped_column(String(64), primary_key=True) + display_name: Mapped[str] = mapped_column(String(64), comment="脱敏展示名") + gender: Mapped[str] = mapped_column(Enum("M", "F", "U"), comment="DDL: DEFAULT 'U'") + birth_date: Mapped[date | None] = mapped_column(Date) + age: Mapped[int | None] = mapped_column(SmallInteger, comment="DDL: TINYINT UNSIGNED") + id_no_mask: Mapped[str | None] = mapped_column(String(32), comment="如 310101199903XXXXXX") + occupation: Mapped[str | None] = mapped_column(String(64)) + employer: Mapped[str | None] = mapped_column(String(128)) + education: Mapped[str | None] = mapped_column( + Enum("high_school", "associate", "bachelor", "master", "doctor", "other") + ) + marital_status: Mapped[str | None] = mapped_column( + Enum("single", "married", "widowed", "divorced", "other") + ) + city: Mapped[str | None] = mapped_column(String(64)) + address_mask: Mapped[str | None] = mapped_column(String(256)) + phone_mask: Mapped[str | None] = mapped_column(String(16)) + email_mask: Mapped[str | None] = mapped_column(String(64)) + annual_income: Mapped[Decimal | None] = mapped_column(Numeric(18, 2), comment="家庭年收入(元)") + financial_asset: Mapped[Decimal | None] = mapped_column(Numeric(18, 2), comment="金融资产规模(不含房产)") + monthly_investable: Mapped[Decimal | None] = mapped_column(Numeric(18, 2), comment="月可投资金额") + is_hnw: Mapped[bool] = mapped_column(Boolean, default=False, comment="高净值客户;DDL: TINYINT(1) DEFAULT 0(AL-01)") + service_tier: Mapped[str] = mapped_column( + Enum("normal", "vip", "diamond"), comment="DDL: DEFAULT 'normal'" + ) + is_pep: Mapped[bool] = mapped_column(Boolean, default=False, comment="政治公众人物;DDL: TINYINT(1) DEFAULT 0") + aml_risk_level: Mapped[str] = mapped_column( + Enum("low", "medium", "high"), comment="DDL: DEFAULT 'low'" + ) + invest_experience_years: Mapped[int | None] = mapped_column(SmallInteger, comment="DDL: TINYINT UNSIGNED") + first_invest_date: Mapped[date | None] = mapped_column(Date) + tenant_id: Mapped[str] = mapped_column(String(32), comment="DDL: DEFAULT 'TENANT-001'") + open_date: Mapped[date] = mapped_column(Date) + is_active: Mapped[bool] = mapped_column(Boolean, default=True, comment="DDL: TINYINT(1) DEFAULT 1") + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class CoreCustomerRisk(CoreBase): + """客户正式风险测评 L0(对齐适当性指南 16 题问卷 + FM-03 风评过期) + · jinrong_core.core_customer_risk""" + + __tablename__ = "core_customer_risk" + __table_args__ = ( + UniqueConstraint("customer_id", name="uk_customer_current"), + Index("idx_risk", "risk_code"), + Index("idx_expires", "expires_at"), + Index("idx_investor_cat", "investor_category"), + {"comment": "客户正式风险测评 L0", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + customer_id: Mapped[str] = mapped_column( + String(64), ForeignKey("core_customer.customer_id", name="fk_cust_risk_customer") + ) + risk_code: Mapped[str] = mapped_column( + CHAR(2), ForeignKey("core_risk_grade.code", name="fk_cust_risk_code"), comment="C1~C5" + ) + questionnaire_score: Mapped[int | None] = mapped_column(SmallInteger, comment="问卷总分 20-100;DDL: SMALLINT UNSIGNED") + max_loss_tolerance_pct: Mapped[Decimal | None] = mapped_column(Numeric(5, 2), comment="最大可承受亏损比例") + investment_goal: Mapped[str | None] = mapped_column(String(128)) + investment_horizon: Mapped[str | None] = mapped_column( + Enum("short", "medium", "long", "flexible") + ) + investor_category: Mapped[str] = mapped_column( + Enum("ordinary", "professional", "professional_pending"), + comment="DDL: DEFAULT 'ordinary'(AL-01 新列)", + ) + professional_approved_at: Mapped[date | None] = mapped_column(Date, comment="AL-01 新列") + is_authoritative: Mapped[bool] = mapped_column(Boolean, default=True, comment="正式测评标记(画像不得覆盖);DDL: TINYINT(1) DEFAULT 1") + evaluated_at: Mapped[date] = mapped_column(Date) + expires_at: Mapped[date] = mapped_column(Date, comment="风评有效期(通常 evaluated_at+12 月,FM-03 过期判定依据;AL-01 新列)") + source: Mapped[str] = mapped_column(String(32), comment="DDL: DEFAULT 'risk_questionnaire'") + + +class CoreCustomerAdvisor(CoreBase): + """客户-代理人归属 · jinrong_core.core_customer_advisor""" + + __tablename__ = "core_customer_advisor" + __table_args__ = ( + UniqueConstraint("customer_id", "advisor_id", "effective_from", name="uk_cust_advisor_from"), + Index("idx_advisor", "advisor_id", "rel_status"), + {"comment": "客户-代理人归属", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + customer_id: Mapped[str] = mapped_column( + String(64), ForeignKey("core_customer.customer_id", name="fk_ca_customer") + ) + advisor_id: Mapped[str] = mapped_column( + String(64), ForeignKey("core_staff.staff_id", name="fk_ca_advisor"), comment="对应 core_staff.staff_id" + ) + rel_status: Mapped[str] = mapped_column( + Enum("active", "transferred", "closed"), comment="DDL: DEFAULT 'active'" + ) + effective_from: Mapped[date] = mapped_column(Date) + effective_to: Mapped[date | None] = mapped_column(Date) + + +class CoreProduct(CoreBase): + """产品主档(对齐 C-11:风险等级 + 期限 + 起购金额)· jinrong_core.core_product""" + + __tablename__ = "core_product" + __table_args__ = ( + Index("idx_min_risk", "min_risk_code"), + Index("idx_product_type", "product_type"), + {"comment": "产品主档", "mysql_engine": "InnoDB"}, + ) + + product_id: Mapped[str] = mapped_column(String(64), primary_key=True) + product_name: Mapped[str] = mapped_column(String(128)) + product_type: Mapped[str] = mapped_column( + Enum( + "money", "bond", "mixed", "stock", "index", + "wealth_mgmt", "private_fund", "insurance", "structured", + ) + ) + min_risk_code: Mapped[str] = mapped_column( + CHAR(2), ForeignKey("core_risk_grade.code", name="fk_product_risk"), comment="R1~R5 最低适配" + ) + min_subscribe_amount: Mapped[Decimal] = mapped_column(Numeric(18, 2), comment="起购金额(元);DDL: DEFAULT 1.00") + term_days: Mapped[int | None] = mapped_column(Integer, comment="产品期限(天),NULL=灵活开放;DDL: INT UNSIGNED(AL-01 新列)") + requires_disclosure: Mapped[bool] = mapped_column(Boolean, default=False, comment="购买前需签署风险揭示书;DDL: TINYINT(1) DEFAULT 0(AL-01 新列)") + industry_code: Mapped[str | None] = mapped_column( + String(16), ForeignKey("core_industry.industry_code", name="fk_product_industry") + ) + fee_rate: Mapped[Decimal | None] = mapped_column(Numeric(6, 4)) + is_open: Mapped[bool] = mapped_column(Boolean, default=True, comment="DDL: TINYINT(1) DEFAULT 1") + created_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3) DEFAULT CURRENT_TIMESTAMP(3)") + + +class CoreHolding(CoreBase): + """持仓快照 · jinrong_core.core_holding""" + + __tablename__ = "core_holding" + __table_args__ = ( + UniqueConstraint("customer_id", "product_id", name="uk_cust_product"), + Index("idx_customer", "customer_id"), + {"comment": "持仓快照", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + customer_id: Mapped[str] = mapped_column( + String(64), ForeignKey("core_customer.customer_id", name="fk_hold_customer") + ) + product_id: Mapped[str] = mapped_column( + String(64), ForeignKey("core_product.product_id", name="fk_hold_product") + ) + qty: Mapped[Decimal] = mapped_column(Numeric(18, 4)) + cost_amount: Mapped[Decimal] = mapped_column(Numeric(18, 2)) + market_value: Mapped[Decimal] = mapped_column(Numeric(18, 2)) + pnl_pct: Mapped[Decimal] = mapped_column(Numeric(8, 4), comment="盈亏比例") + as_of: Mapped[date] = mapped_column(Date) + + +class CoreTrade(CoreBase): + """交易流水(扩展 AML 字段 · 对齐反洗钱规则 RW-001~020) + 写入口唯一:app/gateway/gateway_repository.py(仅 INSERT,B5) + · jinrong_core.core_trade""" + + __tablename__ = "core_trade" + __table_args__ = ( + Index("idx_customer_time", "customer_id", "traded_at"), + Index("idx_amount", "amount", "traded_at"), + Index("idx_channel_time", "channel", "traded_at"), + {"comment": "交易流水", "mysql_engine": "InnoDB"}, + ) + + trade_id: Mapped[str] = mapped_column(String(64), primary_key=True) + customer_id: Mapped[str] = mapped_column( + String(64), ForeignKey("core_customer.customer_id", name="fk_trade_customer") + ) + product_id: Mapped[str] = mapped_column( + String(64), ForeignKey("core_product.product_id", name="fk_trade_product") + ) + trade_type: Mapped[str] = mapped_column(Enum("subscribe", "redeem", "convert")) + amount: Mapped[Decimal] = mapped_column(Numeric(18, 2)) + qty: Mapped[Decimal | None] = mapped_column(Numeric(18, 4)) + channel: Mapped[str] = mapped_column( + Enum("online", "mobile", "counter", "other"), comment="DDL: DEFAULT 'online'" + ) + counterparty_account_mask: Mapped[str | None] = mapped_column(String(32)) + counterparty_name: Mapped[str | None] = mapped_column(String(64)) + is_cash: Mapped[bool] = mapped_column(Boolean, default=False, comment="DDL: TINYINT(1) DEFAULT 0") + payer_name: Mapped[str | None] = mapped_column(String(64), comment="代付人(RW-014)") + is_third_party_pay: Mapped[bool] = mapped_column(Boolean, default=False, comment="DDL: TINYINT(1) DEFAULT 0") + trade_status: Mapped[str] = mapped_column( + Enum("confirmed", "pending", "cancelled"), comment="DDL: DEFAULT 'confirmed'" + ) + traded_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3)") + + +class CoreCashFlow(CoreBase): + """资金进出 · jinrong_core.core_cash_flow""" + + __tablename__ = "core_cash_flow" + __table_args__ = ( + Index("idx_customer", "customer_id", "occurred_at"), + Index("idx_flow_subtype", "flow_subtype", "occurred_at"), + {"comment": "资金进出", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + customer_id: Mapped[str] = mapped_column( + String(64), ForeignKey("core_customer.customer_id", name="fk_cf_customer") + ) + flow_type: Mapped[str] = mapped_column(Enum("in", "out")) + flow_subtype: Mapped[str] = mapped_column( + Enum("deposit", "withdraw", "transfer_in", "transfer_out", "subscribe", "redeem", "other"), + comment="DDL: DEFAULT 'other'", + ) + amount: Mapped[Decimal] = mapped_column(Numeric(18, 2)) + channel: Mapped[str] = mapped_column( + Enum("online", "mobile", "counter", "other"), comment="DDL: DEFAULT 'online'" + ) + counterparty_account_mask: Mapped[str | None] = mapped_column(String(32)) + counterparty_name: Mapped[str | None] = mapped_column(String(64)) + remark: Mapped[str | None] = mapped_column(String(128)) + occurred_at: Mapped[datetime] = mapped_column(DateTime, comment="DDL: DATETIME(3)") + + +class CoreProductNav(CoreBase): + """产品净值 · jinrong_core.core_product_nav""" + + __tablename__ = "core_product_nav" + __table_args__ = ( + UniqueConstraint("product_id", "nav_date", name="uk_product_date"), + {"comment": "产品净值", "mysql_engine": "InnoDB"}, + ) + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + product_id: Mapped[str] = mapped_column( + String(64), ForeignKey("core_product.product_id", name="fk_nav_product") + ) + nav: Mapped[Decimal] = mapped_column(Numeric(10, 4)) + daily_chg_pct: Mapped[Decimal] = mapped_column(Numeric(8, 4)) + nav_date: Mapped[date] = mapped_column(Date) diff --git a/app/repository/core_ro.py b/app/repository/core_ro.py index bf95e4c..6b51928 100644 --- a/app/repository/core_ro.py +++ b/app/repository/core_ro.py @@ -297,6 +297,61 @@ class CoreReadOnlyRepository: for r in conn.execute(sql, {"cid": customer_id, "lim": limit}).mappings() ] + def concentration_profile(self, customer_id: str, limit: int = 500) -> dict[str, Any]: + """持仓集中度画像(RISK-006 / FR-8 专用,一次 SQL 取明细 + Python 端聚合)。 + + PRD FR-8 字面为「调用 core_ro.list_holdings」;本方法是挂账 #1 的收口—— + 与 list_holdings 同源同口径(市值降序、min_risk_code 判定 R4/R5), + 只是把聚合收进仓储,避免调用方取回全量明细再自行求和(N+1)。 + + 采用「SQL 取明细 + Python 聚合」而非 SUM(CASE WHEN):明细行还要进审计 + 摘要,且 sqlite/MySQL 的 min_risk_code 类型差异(VARCHAR(8) vs CHAR(2)) + 下 Python 端判定最稳。limit+1 多取一行探测截断,省掉二次 COUNT 查询。 + + Returns: + r45_value: R4+R5 持仓市值合计(Decimal) + total_value: 全部持仓市值合计(Decimal) + ratio: R4+R5 占比(float;空仓为 0.0) + holdings_truncated: 明细是否触及 limit——调用方按保守口径视同达标 + rows: 市值降序前 limit 条明细(进审计摘要,不进 payload 全量) + """ + sql = text( + """ + SELECT h.market_value, p.min_risk_code + FROM core_holding h + JOIN core_product p ON p.product_id = h.product_id + WHERE h.customer_id = :cid + ORDER BY h.market_value DESC + LIMIT :lim + """ + ) + with self._engine.connect() as conn: + rows = [ + dict(r) + for r in conn.execute(sql, {"cid": customer_id, "lim": limit + 1}).mappings() + ] + + holdings_truncated = len(rows) > limit + rows = rows[:limit] + + total_value = Decimal(0) + r45_value = Decimal(0) + for row in rows: + value = row.get("market_value") + value = value if isinstance(value, Decimal) else Decimal(str(value or 0)) + total_value += value + risk_code = str(row.get("min_risk_code") or "").strip().upper() + if risk_code in ("R4", "R5"): + r45_value += value + + return { + "r45_value": r45_value, + "total_value": total_value, + "ratio": float(r45_value / total_value) if total_value else 0.0, + "holdings_truncated": holdings_truncated, + "rows": rows, + } + def list_trades( self, customer_id: str, since: date | None = None, limit: int = 50 ) -> list[dict[str, Any]]: diff --git a/app/repository/risk_repository.py b/app/repository/risk_repository.py index 5031eb1..48cf712 100644 --- a/app/repository/risk_repository.py +++ b/app/repository/risk_repository.py @@ -39,8 +39,42 @@ class RiskRepository: # ---------- risk_alert ---------- def find_pending_event_alert(self, customer_id: str, day_start: datetime) -> dict | None: - """当日该客户的事件类 pending 单(聚合锚点;PRD FR-4 同日一张)。""" - return self._find_pending(customer_id, day_start, types=EVENT_ALERT_TYPES) + """当日该客户的事件类 pending 单(聚合锚点;PRD FR-4 同日一张)——**排除代理人维度行为链单**。 + + PRD §4A.0 #3 / 实现方案 §2.5(评审 P0-1 收口):RISK-008 代理人行为链单 + 同为 `pattern` 类型,但属**代理人维度**的独立出单线,不能充当客户维度 + 事件单的聚合锚点——若 LIMIT 1 恰好取到行为链单,客户维度当日第二笔会被 + 误并进代理人单,或因误判"已有锚点可并入/锚点被占"而出重复单。 + + 故候选取最近 50 张,再按 `payload.alert_subtype` 过滤掉含 agent_behavior 的单, + 返回剩余中最新的一张(50 条为单客户单日事件单的保守上界)。 + """ + placeholders = ", ".join(f":t{i}" for i in range(len(EVENT_ALERT_TYPES))) + params: dict[str, Any] = {"cid": customer_id, "day_start": day_start} + for i, t in enumerate(EVENT_ALERT_TYPES): + params[f"t{i}"] = t + sql = text( + f""" + SELECT * FROM risk_alert + WHERE customer_id = :cid + AND status = 'pending_review' + AND alert_type IN ({placeholders}) + AND created_at >= :day_start + ORDER BY created_at DESC + LIMIT 50 + """ + ) + with self._engine.connect() as conn: + rows = [self._parse_alert(dict(r)) for r in conn.execute(sql, params).mappings()] + + for row in rows: + payload = row.get("payload") + if isinstance(payload, str): + payload = json.loads(payload) if payload else {} + subtypes = (payload or {}).get("alert_subtype") or [] + if "agent_behavior" not in subtypes: + return row + return None def find_pending_suitability_alert( self, customer_id: str, product_id: str, day_start: datetime @@ -111,8 +145,14 @@ class RiskRepository: triggered_rules: list[str], risk_score: int, alert_type: str, + extra_subtypes: list[str] | None = None, ) -> None: - """读改写:追加 payload.events、合并 triggered_rules、risk_score 取 max、alert_type 更新。""" + """读改写:追加 payload.events、合并 triggered_rules、risk_score 取 max、alert_type 更新。 + + extra_subtypes(C4 起):本批命中的出单子类型(如 RISK-006 的 concentration), + 合并进 `payload.alert_subtype`(sorted set 并集);老单无该字段时首次追加 + 自动创建。**向后兼容**:不传时行为与原先完全一致。 + """ with self._engine.begin() as conn: row = conn.execute( text("SELECT payload, triggered_rules, risk_score FROM risk_alert WHERE alert_id = :aid"), @@ -123,6 +163,9 @@ class RiskRepository: payload = json.loads(row["payload"]) if isinstance(row["payload"], str) else row["payload"] events = payload.setdefault("events", []) events.append(event) + if extra_subtypes: + merged_subtypes = sorted(set(payload.get("alert_subtype") or []) | set(extra_subtypes)) + payload["alert_subtype"] = merged_subtypes old_rules = ( json.loads(row["triggered_rules"]) if isinstance(row["triggered_rules"], str) diff --git a/app/service/risk/alert_service.py b/app/service/risk/alert_service.py index f2a3c42..188acc1 100644 --- a/app/service/risk/alert_service.py +++ b/app/service/risk/alert_service.py @@ -7,6 +7,7 @@ redis_gateway 单例(B7:连接由 lifespan 管理,publish 失败降级不 from __future__ import annotations +import json from datetime import date, datetime, time from typing import Any from uuid import uuid4 @@ -14,7 +15,7 @@ from uuid import uuid4 from app.repository.risk_repository import EVENT_ALERT_TYPES, RiskRepository from app.service.risk import redis_gateway from app.service.risk.locks import run_locked -from app.service.risk.rules import RuleHit +from app.service.risk.rules import RULE_ALERT_TYPES, RULE_SCORES, RuleHit from app.utils.exceptions import NotFoundError, StateConflict from app.utils.trace import current_trace, new_trace @@ -66,18 +67,30 @@ def _audit( ) -def _publish_alert(alert: dict[str, Any]) -> None: - redis_gateway.publish( - "risk:pub:alert", - { - "alert_id": alert["alert_id"], - "alert_type": alert["alert_type"], - "customer_id_mask": alert["customer_id"][:6] + "**", - "risk_score": alert["risk_score"], - "trace_id": alert["trace_id"], - "notify_role": ["risk_officer"] + (["compliance"] if alert["alert_type"] == "aml" else []), - }, - ) +def _publish_alert( + alert: dict[str, Any], + notify_role: list[str] | None = None, + extra: dict[str, Any] | None = None, +) -> None: + """预警推送(Redis PubSub)。 + + notify_role:覆盖接收角色(C5 升级单用,含 risk_manager/compliance); + 缺省保持现口径 risk_officer(+aml 追加 compliance)。 + extra:并入推送体的附加字段(C4/C6 放 alert_subtype,C5 放 escalation_level); + `02-redis-keys.md` 已同步登记。 + """ + body: dict[str, Any] = { + "alert_id": alert["alert_id"], + "alert_type": alert["alert_type"], + "customer_id_mask": alert["customer_id"][:6] + "**", + "risk_score": alert["risk_score"], + "trace_id": alert["trace_id"], + "notify_role": notify_role + or (["risk_officer"] + (["compliance"] if alert["alert_type"] == "aml" else [])), + } + if extra: + body.update(extra) + redis_gateway.publish("risk:pub:alert", body) def record_trade_alerts( @@ -100,6 +113,10 @@ def record_trade_alerts( best = max(hits, key=lambda h: h.risk_score) triggered_rules = sorted({h.rule_id for h in hits}) risk_score = max(h.risk_score for h in hits) + # C4 起:出单子类型集合(RISK-006→concentration、RISK-008→agent_behavior)。 + # 空集不注入 payload.alert_subtype(评审 P2-3:老单结构保持干净)。 + subtypes = sorted({h.alert_subtype for h in hits if h.alert_subtype}) + publish_extra = {"alert_subtype": subtypes} if subtypes else None event = { "trade_id": trade["trade_id"], "product_id": trade["product_id"], @@ -114,14 +131,28 @@ def record_trade_alerts( if locked: pending = risk_repo.find_pending_event_alert(trade["customer_id"], _day_start()) if pending: + # 评审 P1-3(既有缺陷顺带修正):alert_type 按「老单规则 ∪ 本批规则」 + # 重算,而不是取本批最高分规则——否则既有 large_amount(70) 单碰上 + # 本批仅命中 RISK-006(60) 时会被错误翻转为 pattern。 + old_rules = pending.get("triggered_rules") or [] + if isinstance(old_rules, str): + old_rules = json.loads(old_rules) if old_rules else [] + merged_rules = sorted(set(old_rules) | set(triggered_rules)) + best_rule = max( + (r for r in merged_rules if r in RULE_SCORES), + key=lambda r: RULE_SCORES[r], + default=None, + ) + merged_alert_type = RULE_ALERT_TYPES[best_rule] if best_rule else best.alert_type risk_repo.append_alert_event( - pending["alert_id"], event, triggered_rules, risk_score, best.alert_type + pending["alert_id"], event, triggered_rules, risk_score, merged_alert_type, + extra_subtypes=subtypes or None, ) updated = risk_repo.get_alert(pending["alert_id"]) - _audit(risk_repo, customer_id=trade["customer_id"], rule_id=",".join(triggered_rules), + _audit(risk_repo, customer_id=trade["customer_id"], rule_id=",".join(merged_rules), decision="alert_appended", risk_score=updated["risk_score"], input_summary=event) - _publish_alert(updated) + _publish_alert(updated, extra=publish_extra) return updated alert = { "alert_id": _new_alert_id(), @@ -138,10 +169,12 @@ def record_trade_alerts( "customer_context": customer_context or {}, }, } + if subtypes: + alert["payload"]["alert_subtype"] = subtypes risk_repo.insert_alert(alert) _audit(risk_repo, customer_id=trade["customer_id"], rule_id=",".join(triggered_rules), decision="alert_created", risk_score=risk_score, input_summary=event) - _publish_alert(alert) + _publish_alert(alert, extra=publish_extra) return alert return run_locked(f"agg:event:{trade['customer_id']}:{date.today()}", _agg) diff --git a/app/service/risk/chat_tools.py b/app/service/risk/chat_tools.py index 36c3164..bb5b3fd 100644 --- a/app/service/risk/chat_tools.py +++ b/app/service/risk/chat_tools.py @@ -109,6 +109,7 @@ def customer_context(customer_id: str, core_ro: CoreReadOnlyRepository | None = pending, pending_total = repo.list_alerts( customer_id=customer_id, status="pending_review", page_size=20 ) + profile = ro.concentration_profile(customer_id) return _jsonable( { "found": True, @@ -125,6 +126,14 @@ def customer_context(customer_id: str, core_ro: CoreReadOnlyRepository | None = }, "pending_alert_count": pending_total, "pending_alerts": [_alert_brief(r) for r in pending], + # FR-8:持仓集中度画像——复用 core_ro.concentration_profile, + # 与 RISK-006 判定同源同口径(不重复实现聚合逻辑)。 + "profile": { + "concentration_ratio": round(float(profile["ratio"]), 4), + "r45_value": str(profile["r45_value"]), + "total_value": str(profile["total_value"]), + "holdings_truncated": profile["holdings_truncated"], + }, } ) diff --git a/app/service/risk/engine.py b/app/service/risk/engine.py index de18fe4..b973ac1 100644 --- a/app/service/risk/engine.py +++ b/app/service/risk/engine.py @@ -23,8 +23,8 @@ from app.repository.risk_repository import RiskRepository from app.service.risk.alert_service import record_aml_alert, record_trade_alerts from app.service.risk.aml_service import match_customer from app.service.risk.profile_l3 import upsert_profile_l3 -from app.service.risk.rules import RiskThresholds, run_rules -from app.utils.trace import ensure_trace +from app.service.risk.rules import RiskThresholds, rule_concentration, run_rules +from app.utils.trace import current_trace, ensure_trace, new_trace def _as_datetime(value: Any) -> datetime: @@ -66,6 +66,51 @@ def _build_customer_context( } +def _audit_concentration( + repo: RiskRepository, + trade: dict[str, Any], + profile: dict[str, Any], + hit: Any, + alert: dict[str, Any] | None, +) -> None: + """RISK-006 专属审计(event_type='risk_concentration')。 + + 审计表只 INSERT(红线);金额按 DESENS-005 口径只落**合计与前 5 条摘要**, + 不把全量持仓明细写进审计(明细已在 payload 侧由事件承载)。 + """ + rows = profile.get("rows") or [] + top_holdings = [ + { + "market_value": str(r.get("market_value")), + "min_risk_code": r.get("min_risk_code"), + } + for r in rows[:5] + ] + repo.insert_audit_log( + { + "trace_id": current_trace() or new_trace(), + "event_type": "risk_concentration", + "agent_type": "risk", + "actor_id": "SYSTEM", + "customer_id": trade["customer_id"], + "rule_id": hit.rule_id, + "input_summary": { + "trade_id": trade.get("trade_id"), + "r45_value": str(profile.get("r45_value")), + "total_value": str(profile.get("total_value")), + "ratio": round(float(profile.get("ratio") or 0), 4), + "holdings_truncated": bool(profile.get("holdings_truncated")), + "top_holdings": top_holdings, + }, + "decision": "alert_created" if alert else "alert_appended", + "risk_score": hit.risk_score, + "handler_id": None, + "handler_result": None, + "handler_comment": None, + } + ) + + def process_trade_event( trade: dict[str, Any], core_ro: CoreReadOnlyRepository | None = None, @@ -92,6 +137,15 @@ def process_trade_event( result: dict[str, Any] = {"triggered_rules": [], "alert_ids": [], "aml_hit": False} hits = run_rules(trades, th, now=event_at) + + # FR-8 RISK-006 持仓集中度:输入是持仓画像而非流水,故不并入 run_rules + # (实现方案 §2.3/§2.4:避免改动现有 13 处 run_rules 调用与断言)。 + # 命中后并入 hits,走既有 record_trade_alerts 聚合/审计/推送,不新增出单路径。 + concentration_profile = core.concentration_profile(trade["customer_id"]) + conc = rule_concentration(concentration_profile, th) + if conc: + hits = [*hits, conc] + if hits: alert = record_trade_alerts( trade, @@ -106,9 +160,13 @@ def process_trade_event( upsert_profile_l3( trade["customer_id"], best.alert_type, + # FR-8:命中集中度即给 L3 打标签(后续画像/台账可筛高风险集中度客户) + monitor_tags=["high_risk_concentration"] if conc else None, last_alert_id=alert["alert_id"] if alert else None, risk_repo=repo, ) + if conc: + _audit_concentration(repo, trade, concentration_profile, conc, alert) else: # 未命中分支:pass 审计由 alert_service 统一落库(架构 §3.1 ④) record_trade_alerts(trade, [], risk_repo=repo) diff --git a/app/service/risk/rules.py b/app/service/risk/rules.py index 92f7eb8..8c5ea8f 100644 --- a/app/service/risk/rules.py +++ b/app/service/risk/rules.py @@ -21,6 +21,7 @@ RULE_SCORES: dict[str, int] = { "RISK-003": 50, "RISK-004": 80, "RISK-005": 80, + "RISK-006": 60, # FR-8 持仓集中度(PRD §4A) } RULE_ALERT_TYPES: dict[str, str] = { "RISK-001": "large_amount", @@ -28,6 +29,7 @@ RULE_ALERT_TYPES: dict[str, str] = { "RISK-003": "freq_trade", "RISK-004": "pattern", "RISK-005": "pattern", + "RISK-006": "pattern", } @@ -41,6 +43,10 @@ class RiskThresholds: probe_amount: Decimal = Decimal("400000") small_amount: Decimal = Decimal("10000") small_count: int = 3 + # FR-8 RISK-006:R4+R5 市值占比阈值(0~1)。 + # 评审 P1-2:from_settings 必须补读本字段——引擎/网关测试全走 from_settings + # 默认路径,漏读会让 conftest 的 autouse monkeypatch 失效,现有断言被 RISK-006 打穿。 + concentration_threshold: float = 0.80 @classmethod def from_settings(cls) -> "RiskThresholds": @@ -54,6 +60,7 @@ class RiskThresholds: probe_amount=s.risk_probe_amount, small_amount=s.risk_small_amount, small_count=s.risk_small_count, + concentration_threshold=s.risk_concentration_threshold, ) @@ -63,6 +70,9 @@ class RuleHit: alert_type: str risk_score: int detail: str # 触发说明(进预警单 payload / audit) + # 出单子类型:RISK-006→"concentration"、RISK-008→"agent_behavior"; + # RISK-001~005 恒 None(不写 payload.alert_subtype)。C4 起启用。 + alert_subtype: str | None = None def _eligible(trade: dict[str, Any]) -> bool: @@ -151,6 +161,43 @@ def rule_small_then_large(trades: list[dict[str, Any]], th: RiskThresholds) -> R return None +def rule_concentration(profile: dict[str, Any], th: RiskThresholds) -> RuleHit | None: + """RISK-006 持仓集中度:R4+R5 市值占比 ≥ 阈值即命中(FR-8 纯函数)。 + + 输入为 `core_ro.concentration_profile` 的返回值(持仓画像,非流水), + 故不并入 run_rules(输入域不同,见实现方案 §2.3:避免改动现有 13 处调用)。 + + - 空仓(total_value ≤ 0)不触发——无持仓无从谈集中度; + - `holdings_truncated=True` 视同达标(PRD FR-8 截断防护:明细被 limit 截断后 + 占比可能失真,按保守口径告警,宁可多报不可漏报)。 + """ + total = profile.get("total_value") or Decimal(0) + total = total if isinstance(total, Decimal) else Decimal(str(total)) + if total <= 0: + return None + + r45 = profile.get("r45_value") or Decimal(0) + r45 = r45 if isinstance(r45, Decimal) else Decimal(str(r45)) + ratio = r45 / total # Decimal 除法,避免 float 精度引发边界误判 + + truncated = bool(profile.get("holdings_truncated")) + if not truncated and ratio < Decimal(str(th.concentration_threshold)): + return None + + if truncated: + reason = f"持仓明细触及查询上限,按保守口径视同集中度达标(实际占比 {float(ratio):.1%})" + else: + reason = f"R4+R5 持仓占比 {float(ratio):.1%} ≥ 阈值 {float(th.concentration_threshold):.0%}" + + return RuleHit( + "RISK-006", + RULE_ALERT_TYPES["RISK-006"], + RULE_SCORES["RISK-006"], + f"{reason}(R4+R5 市值 {r45} 元 / 总市值 {total} 元)", + alert_subtype="concentration", + ) + + def run_rules( trades: list[dict[str, Any]], thresholds: RiskThresholds | None = None, diff --git a/app/service/tool_service.py b/app/service/tool_service.py index b0b8240..7099372 100644 --- a/app/service/tool_service.py +++ b/app/service/tool_service.py @@ -381,11 +381,19 @@ def summarize(record: dict[str, Any]) -> str: if not data.get("found"): return f"(客户风控上下文:未找到客户 {data.get('customer_id')})" l3 = data.get("l3") or {} - return ( + base = ( f"(客户 {data.get('customer_id')} {data.get('display_name', '')}:" f"测评 {data.get('risk_code')},L3 监测档 {l3.get('monitor_tier') or '无'}," - f"待审预警 {data.get('pending_alert_count', 0)} 条)" + f"待审预警 {data.get('pending_alert_count', 0)} 条" ) + # FR-8:高风险持仓占比(R4+R5);与 RISK-006 判定同源,仅供参考口径 + profile = data.get("profile") or {} + ratio = profile.get("concentration_ratio") + if isinstance(ratio, (int, float)): + base += f",高风险持仓占比 {ratio:.0%}(仅供参考)" + if profile.get("holdings_truncated"): + base += "(持仓明细较多,占比为截断口径保守值)" + return base + ")" if name == "suitability_check": if not data.get("found"): return f"(适当性校验:{data.get('error') or '未找到客户或产品'})" diff --git a/docs/项目框架设计/表设计/02-redis-keys.md b/docs/项目框架设计/表设计/02-redis-keys.md index c87046e..34f1f03 100644 --- a/docs/项目框架设计/表设计/02-redis-keys.md +++ b/docs/项目框架设计/表设计/02-redis-keys.md @@ -50,7 +50,7 @@ | Key | 类型 | TTL | 说明 | | --- | --- | --- | --- | -| `risk:pub:alert` | Pub/Sub | — | 新预警广播,payload=`{alert_id, alert_type, customer_id_mask, risk_score, trace_id, notify_role}`(口径以 `docs/PRD/PRD-风控监测Agent.md` FR-4 为准,2026-09-06 修订) | +| `risk:pub:alert` | Pub/Sub | — | 新预警广播,payload=`{alert_id, alert_type, customer_id_mask, risk_score, trace_id, notify_role}`,风控追加 v1.1 起**可带附加字段**:`alert_subtype`(C4/C6 出单子类型数组,如 `["concentration"]`/`["agent_behavior"]`)、`escalation_level`(C5 升级级别 `LEVEL_1`/`LEVEL_2`);两者仅在非空时下发(口径以 `docs/PRD/PRD-风控监测Agent.md` FR-4 + §4A 为准,2026-09-07 修订) | | `risk:dedup:{customer_id}:{rule_id}:{date}` | String | 24h | 同日同规则防重复预警风暴 | ### 2.5 输入防护与限流(F-03) diff --git a/tests/conftest.py b/tests/conftest.py index bb61078..778db47 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -23,6 +23,23 @@ ROOT = Path(__file__).resolve().parent.parent DEMO_SQL = ROOT / "scripts" / "demo" / "prepare_risk_demo.sql" +@pytest.fixture(autouse=True) +def _disable_concentration_rule(monkeypatch): + """RISK-006 回归隔离(C4 起):默认把集中度阈值推到不可达,现有用例断言零改动。 + + RISK-006 在引擎里对**每一笔**交易都会查持仓画像,若沿用默认 0.80, + 现有 436 条用例中断言「命中规则集合 / 单日一张单」的会被新规则打穿。 + 故全局默认禁用;RISK-006 专属测试内再显式传 `RiskThresholds(...)` 或 + 再次 monkeypatch 回真实阈值 0.80(实现方案 §6 全量回归口径,评审 P1-2)。 + + 前置条件:`RiskThresholds.from_settings()` 必须补读 concentration_threshold, + 否则此处 monkeypatch 不生效(rules.py 已补读并注明)。 + """ + from app.config.settings import settings + + monkeypatch.setattr(settings, "risk_concentration_threshold", 1.01) + + @pytest.fixture() def sqlite_engine(): """内存 sqlite 全表引擎(B8 前 DDL 散落各测试文件,收敛后统一走这里)。 diff --git a/tests/test_concentration_c4.py b/tests/test_concentration_c4.py new file mode 100644 index 0000000..98e0be7 --- /dev/null +++ b/tests/test_concentration_c4.py @@ -0,0 +1,312 @@ +"""C4 / FR-8 · RISK-006 持仓集中度:引擎接入 + 出单聚合 + 对话线专项测试。 + +覆盖实现方案 §6.2 中 C4 相关用例: + +- 引擎:RISK-006 与 RISK-001 同单聚合、risk_score 取 max、payload.alert_subtype + 含 concentration、L3 打 high_risk_concentration 标签、risk_concentration 审计; +- 出单:alert_subtype 集合维护(空集不注入 / 追加时合并)、 + **P0-1 回归**:当日已有 agent_behavior 单后再触发 RISK-006 应出第二张客户维度单、 + **P1-3 修正**:既有 large_amount 单追加仅 RISK-006 时 alert_type 不翻转; +- 对话线:customer_context 带 concentration_ratio。 + +阈值说明:conftest 的 autouse fixture 把 `risk_concentration_threshold` 推到 1.01 +(回归隔离),本模块内统一 monkeypatch 回真实阈值 0.80。 +""" + +from __future__ import annotations + +from datetime import datetime +from decimal import Decimal + +import pytest +from sqlalchemy import text + +from _ddl import create_sqlite_engine + +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.chat_tools import customer_context +from app.service.risk.engine import process_trade_event +from app.service.risk.rules import RiskThresholds, rule_concentration + +NOW = datetime(2026, 9, 6, 14, 0, 0) + + +class FakePublisher: + def __init__(self): + self.messages: list = [] + + def publish(self, channel, payload): + self.messages.append((channel, payload)) + + def delete(self, *keys): + pass + + +@pytest.fixture() +def env(monkeypatch): + """sqlite 环境:R3 产品(走交易)+ R5 产品(持仓主体,构成 90% 集中度)。 + + 持仓口径:P2(R5) 900000 + P1(R3) 100000 → R4+R5 占比 90% ≥ 0.80。 + """ + from app.config.settings import settings + + monkeypatch.setattr(settings, "risk_concentration_threshold", 0.80) + + 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)" + ) + ) + conn.execute( + text( + "INSERT INTO core_product (product_id, product_name, min_risk_code, product_type)" + " VALUES ('P1', '测试混合基金', 'R3', 'mixed')," + " ('P2', '测试股票基金', 'R5', 'equity')" + ) + ) + conn.execute( + text( + "INSERT INTO core_holding (customer_id, product_id, market_value, quantity)" + " VALUES ('C1', 'P2', 900000, 1000), ('C1', 'P1', 100000, 500)" + ) + ) + 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 _trade(trade_id, amount="600000", at=NOW): + return { + "trade_id": trade_id, + "customer_id": "C1", + "product_id": "P1", + "trade_type": "subscribe", + "amount": Decimal(amount), + "traded_at": at, + } + + +def _seed_trade(engine, trade): + 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, :pid, :tt, :amt, 'confirmed', :at)" + ), + { + "tid": trade["trade_id"], + "cid": trade["customer_id"], + "pid": trade["product_id"], + "tt": trade["trade_type"], + "amt": float(trade["amount"]), # sqlite 不支持绑定 Decimal,转 float + "at": trade["traded_at"], + }, + ) + + +def _insert_alert(engine, alert_id, alert_type, risk_score, rules, payload): + import json + + with engine.begin() as conn: + conn.execute( + text( + "INSERT INTO risk_alert (alert_id, trace_id, customer_id, trade_id," + " alert_type, triggered_rules, risk_score, status, payload)" + " VALUES (:aid, 'TRACE-TEST', 'C1', 'TRD-TEST-0', :atype, :rules," + " :score, 'pending_review', :payload)" + ), + { + "aid": alert_id, + "atype": alert_type, + "rules": json.dumps(rules), + "score": risk_score, + "payload": json.dumps(payload, ensure_ascii=False), + }, + ) + + +# ---------- 仓储聚合:concentration_profile ---------- + + +def test_concentration_profile_aggregates_r4_r5(env): + core, repo, pub, engine = env + profile = core.concentration_profile("C1") + assert profile["r45_value"] == Decimal(900000) + assert profile["total_value"] == Decimal(1000000) + assert profile["ratio"] == 0.9 + assert profile["holdings_truncated"] is False + + +def test_concentration_profile_empty_customer(env): + core, repo, pub, engine = env + profile = core.concentration_profile("NOT-EXIST") + assert profile["total_value"] == Decimal(0) + assert profile["ratio"] == 0.0 + assert rule_concentration(profile, RiskThresholds()) is None + + +# ---------- 引擎接入 ---------- + + +def test_engine_merges_concentration_with_large_amount(env): + """RISK-001(70) + RISK-006(60) 同单聚合:score 取 max=70,类型随最高分规则。""" + core, repo, pub, engine = env + trade = _trade("TRD-TEST-1") + _seed_trade(engine, trade) + + result = process_trade_event(trade, core_ro=core, risk_repo=repo) + + # 600000 同时触发 RISK-002(≥ 单日累计 500000),故断言"包含"而非全等 + assert "RISK-001" in result["triggered_rules"] + assert "RISK-006" in result["triggered_rules"] + assert len(result["alert_ids"]) == 1 # 聚合成一张单 + + alert = repo.get_alert(result["alert_ids"][0]) + assert alert["risk_score"] == 70 + assert alert["alert_type"] == "large_amount" # 不被 RISK-006 翻转 + assert alert["payload"]["alert_subtype"] == ["concentration"] + + +def test_engine_writes_concentration_audit_and_l3_tag(env): + core, repo, pub, engine = env + trade = _trade("TRD-TEST-2") + _seed_trade(engine, trade) + + result = process_trade_event(trade, core_ro=core, risk_repo=repo) + + with engine.connect() as conn: + audit = conn.execute( + text( + "SELECT event_type, rule_id, risk_score FROM audit_log" + " WHERE event_type = 'risk_concentration'" + ) + ).mappings().all() + assert len(audit) == 1 + assert audit[0]["rule_id"] == "RISK-006" + + l3 = repo.get_l3("C1") + assert "high_risk_concentration" in (l3.get("monitor_tags") or []) + + +def test_engine_concentration_only_still_creates_alert(env): + """仅命中集中度(未达大额)也要出单——无需为「仅 RISK-006」写独立分支。""" + core, repo, pub, engine = env + trade = _trade("TRD-TEST-3", amount="10000") + _seed_trade(engine, trade) + + result = process_trade_event(trade, core_ro=core, risk_repo=repo) + + assert result["triggered_rules"] == ["RISK-006"] + alert = repo.get_alert(result["alert_ids"][0]) + assert alert["risk_score"] == 60 + assert alert["alert_type"] == "pattern" + + +def test_engine_disabled_when_threshold_unreachable(env, monkeypatch): + """回归隔离口径:阈值推到 1.01 后 RISK-006 不触发(conftest autouse 同款行为)。""" + from app.config.settings import settings + + monkeypatch.setattr(settings, "risk_concentration_threshold", 1.01) + core, repo, pub, engine = env + trade = _trade("TRD-TEST-4", amount="10000") + _seed_trade(engine, trade) + + result = process_trade_event(trade, core_ro=core, risk_repo=repo) + assert result["triggered_rules"] == [] + + +# ---------- 出单:alert_subtype 与聚合锚点 ---------- + + +def test_append_merges_subtypes(env): + """追加到老单时,extra_subtypes 合并进 payload.alert_subtype(老单原本无该字段)。""" + core, repo, pub, engine = env + _insert_alert( + engine, "ALT-TEST-OLD", "large_amount", 70, ["RISK-001"], {"product_id": "P1", "events": []} + ) + trade = _trade("TRD-TEST-5", amount="10000") + _seed_trade(engine, trade) + + result = process_trade_event(trade, core_ro=core, risk_repo=repo) + + # 老单被追加,不新建 + assert result["alert_ids"] == ["ALT-TEST-OLD"] + alert = repo.get_alert("ALT-TEST-OLD") + assert alert["payload"]["alert_subtype"] == ["concentration"] + assert alert["risk_score"] == 70 # max(70, 60) + + +def test_alert_type_not_flipped_by_lower_score_rule(env): + """评审 P1-3:既有 large_amount(70) 单追加仅 RISK-006(60) 时,类型保持 large_amount。""" + core, repo, pub, engine = env + _insert_alert( + engine, "ALT-TEST-KEEP", "large_amount", 70, ["RISK-001"], {"product_id": "P1", "events": []} + ) + from app.service.risk.alert_service import record_trade_alerts + + profile = core.concentration_profile("C1") + hit = rule_concentration(profile, RiskThresholds.from_settings()) + assert hit is not None + + trade = _trade("TRD-TEST-6", amount="10000") + updated = record_trade_alerts(trade, [hit], risk_repo=repo) + assert updated["alert_type"] == "large_amount" + + +def test_p0_1_agent_behavior_alert_is_not_anchor(env): + """评审 P0-1 回归:当日已有 agent_behavior 单 → RISK-006 应新建客户维度单,不并入。""" + core, repo, pub, engine = env + _insert_alert( + engine, + "ALT-TEST-AGENT", + "pattern", + 70, + ["RISK-008"], + {"product_id": "P1", "events": [], "alert_subtype": ["agent_behavior"]}, + ) + trade = _trade("TRD-TEST-7", amount="10000") + _seed_trade(engine, trade) + + result = process_trade_event(trade, core_ro=core, risk_repo=repo) + + assert result["alert_ids"] and result["alert_ids"][0] != "ALT-TEST-AGENT" + new_alert = repo.get_alert(result["alert_ids"][0]) + assert new_alert["payload"]["alert_subtype"] == ["concentration"] + + agent_alert = repo.get_alert("ALT-TEST-AGENT") + assert agent_alert["payload"]["alert_subtype"] == ["agent_behavior"] # 未被污染 + + +# ---------- 对话线 ---------- + + +def test_customer_context_includes_concentration_ratio(env): + core, repo, pub, engine = env + data = customer_context("C1", core_ro=core, risk_repo=repo) + assert data["found"] is True + assert data["profile"]["concentration_ratio"] == 0.9 + assert data["profile"]["holdings_truncated"] is False + + +def test_customer_context_zero_holdings(env): + core, repo, pub, engine = env + with engine.begin() as conn: + conn.execute( + text( + "INSERT INTO core_customer (customer_id, display_name, age, is_active)" + " VALUES ('C9', '空仓客户', 30, 1)" + ) + ) + data = customer_context("C9", core_ro=core, risk_repo=repo) + assert data["found"] is True + assert data["profile"]["concentration_ratio"] == 0.0 diff --git a/tests/test_risk_rules.py b/tests/test_risk_rules.py index e1d3acc..08a3039 100644 --- a/tests/test_risk_rules.py +++ b/tests/test_risk_rules.py @@ -3,7 +3,7 @@ from datetime import datetime from decimal import Decimal -from app.service.risk.rules import RiskThresholds, run_rules +from app.service.risk.rules import RiskThresholds, rule_concentration, run_rules NOW = datetime(2026, 9, 6, 14, 0, 0) TH = RiskThresholds() # 默认冻结阈值 @@ -154,3 +154,50 @@ class TestDefenseAndAggregation: hits = run_rules(trades, TH, NOW) assert _ids(hits) == ["RISK-001", "RISK-002", "RISK-003", "RISK-005"] assert max(h.risk_score for h in hits) == 80 # pattern + + +class TestRISK006Concentration: + """FR-8 RISK-006 持仓集中度(纯函数;输入为 core_ro.concentration_profile 画像)。 + + 注意:run_rules 不含 RISK-006(输入域不同),由引擎单独调用后并入 hits。 + """ + + @staticmethod + def _profile(r45, total, truncated=False): + return { + "r45_value": Decimal(str(r45)), + "total_value": Decimal(str(total)), + "ratio": float(Decimal(str(r45)) / Decimal(str(total))) if total else 0.0, + "holdings_truncated": truncated, + "rows": [], + } + + def test_hit_90_percent_r5(self): + hit = rule_concentration(self._profile(900000, 1000000), TH) + assert hit is not None + assert hit.rule_id == "RISK-006" + assert hit.alert_type == "pattern" and hit.risk_score == 60 + assert hit.alert_subtype == "concentration" + + def test_miss_when_all_r1(self): + """全是 R1 持仓(R4+R5 为 0)→ 不触发。""" + assert rule_concentration(self._profile(0, 1000000), TH) is None + + def test_miss_on_empty_holdings(self): + """空仓(total=0)不触发——无持仓无从谈集中度。""" + assert rule_concentration(self._profile(0, 0), TH) is None + + def test_truncated_counts_as_hit(self): + """明细被截断时占比仅 10%,仍按保守口径视同达标(PRD FR-8 截断防护)。""" + hit = rule_concentration(self._profile(100000, 1000000, truncated=True), TH) + assert hit is not None + assert "保守" in hit.detail + + def test_threshold_boundary(self): + """阈值边界:79.9% 不触发,80%(等于阈值)触发。""" + assert rule_concentration(self._profile(799, 1000), TH) is None + assert rule_concentration(self._profile(800, 1000), TH) is not None + + def test_detail_carries_values(self): + hit = rule_concentration(self._profile(900000, 1000000), TH) + assert "900000" in hit.detail and "1000000" in hit.detail