Files
group_fqcd_jr/app/service/suitability_service.py
T
lzf_0626 6516ccb385 feat: 第二版——接口契约对齐 docs/05,修复静默故障与数据库基线
相对第一版 46fc976 的完整变更。组员迁移对照表见 docs/20。

一、对外契约对齐 docs/05(破坏性,共 4 处,组员需按 docs/20 调整)
1) 配置发布端点改为文档规定的复数资源名:submit→validations、
   approve→reviews(需 body decision)、activate→activations、
   rollback→rollbacks;第一版这 4 个动词式路径 docs/05 从未定义过。
2) 错误码由 8 个笼统码改为 15 个具体语义码(FORBIDDEN→AGENT_PERMISSION_DENIED、
   UNAUTHORIZED→AUTHENTICATION_REQUIRED、CONFLICT→RESOURCE_VERSION_CONFLICT、
   RESOURCE_NOT_FOUND→RUN_NOT_FOUND/SESSION_NOT_FOUND 等),
   输入类错误状态码 400→422。
3) POST /api/v1/agent-runs 与 GET /api/v1/agent-runs/{run_id} 统一为
   {data, meta} 信封(data 内字段名与语义未变)。
4) 错误响应体统一为 {error:{code,message,retryable,field_errors}, meta:{trace_id}},
   不再返回 FastAPI 默认的 {"detail": ...}。

二、数据库基线与约束
新增 39 张表的基线迁移(链根)与联合唯一键纠偏(4 张表、删 8 增 4,幂等收敛);
撤下 config_release 的双人复核 CHECK(应用层已允许自审,审核节点保留,
自审如实写入 reviewer_id);记忆 active key 生成列与唯一键;
activate 开始记录 supersedes_release_id 使版本链可追溯。
docs/00 基线未修改,未重命名或删除任何表与字段。

三、修复会静默出错或无报错的缺陷
- 跑完集成测试后平台会静默失去生效配置:清理只删自己创建的版本,却没有恢复被它
  顶成 superseded 的原生效版本,且审计一并删除因而完全无痕,表现为所有工具被拒
  但没有任何报错。已修清理逻辑并加恢复。
- Worker 单轮异常导致进程退出;记忆抽取调用方的“事务已开始”异常;
  召回缓存丢失 degraded 标记;连接时区未生效导致 created_at/updated_at 差 8 小时;
  .env 与 os.getenv 密钥来源分裂导致“没有可用的已批准模型端点”。
- 记忆信号识别漏判与跨键误命中;SSE 未带 Accept 的协商行为。

四、功能补齐
记忆链路 P1/P2/P3(抽取、受控词表、召回与缓存、生命周期级联及投影事件)、
fin_* 场内交易只读 ORM 层、agent_intent_config 状态流转并在运行期真正生效、
限流(Redis 固定窗口、故障一律放行)、游标校验、trace_id 中间件、
示例业务 Agent fund_query_demo 与一键端到端验证脚本,以及审计/指纹/迁移状态工具。

五、文档与验证
新增 docs/19(业务 Agent 接入实操)、docs/20(第一版迁移指南)与 docs/evidence 证据;
docs/01/02/06/08/09/17 同步实现现状。

验证结果:ruff 通过、mypy 103 文件无错、unit+contract 447 passed、
integration 29 passed、acceptance_check --production 7 PASS、
demo_agent_e2e 9/9 PASS(含失败关闭反证)。
2026-09-10 15:55:54 +08:00

272 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""公共投顾适当性校验服务。
适当性是治理边界,不属于任何一个业务 Agent。该服务只读客户/产品风险信息,
返回不可变决定;拒绝决定必须留下审计记录,且不会修改交易或产品数据。
合规要点(B2 修复):
1. **客户风险等级只能来自服务端权威来源**:取 ``fin_risk_assessment`` 中该客户最新一条
测评(``assessed_at`` 最新),调用方传入的 ``customer_risk_level`` 一律不接受
(DTO ``extra="forbid"`` 直接拒绝伪造字段)。
2. **测评有效期以权威记录的 ``valid_until`` 为准**,调用方不能自报过期时间;测评缺失或
过期一律失败关闭,不做静默降级。
3. **专业投资者身份来自 ``sys_user``**(``is_professional_investor`` 且
``professional_investor_status='已认定'``),而非调用方参数;已认定的专业投资者可豁免
C/R 等级匹配,但仍强制风险揭示、确认与录音,且不能绕过测评有效期与审计。
"""
import re
from collections.abc import Callable, Mapping
from datetime import UTC, datetime
from typing import Any, Literal
from pydantic import BaseModel, ConfigDict, Field
from sqlalchemy import text
from app.core.contracts import RequestContext
from app.core.errors import ForbiddenAgentError
from app.infrastructure.db import SessionFactory
from app.model.audit import InteractionAudit
PROFESSIONAL_INVESTOR_CERTIFIED = "已认定"
RISK_LEVEL_SOURCE = "fin_risk_assessment"
PROFESSIONAL_INVESTOR_SOURCE = "sys_user"
CUSTOMER_SCOPE_EXEMPT_ROLES = frozenset({"admin", "super_admin"})
_INVESTOR_TYPE_PATTERN = re.compile(r"^C([1-5])$")
AuthorityReason = Literal[
"AUTHORITY_OK",
"CUSTOMER_NOT_FOUND",
"ASSESSMENT_MISSING",
"ASSESSMENT_EXPIRED",
"RISK_LEVEL_INVALID",
]
# 只读查询:客户专业投资者身份(sys_user)+ 最新一条风险测评(fin_risk_assessment)。
# 不读取 answers 问卷原文,避免把敏感测评明细带入服务层。
_AUTHORITY_SQL = text(
"""
SELECT u.is_professional_investor,
u.professional_investor_status,
a.investor_type,
a.assessed_at,
a.valid_until
FROM sys_user u
LEFT JOIN fin_risk_assessment a ON a.id = (
SELECT x.id FROM fin_risk_assessment x
WHERE x.customer_id = u.id
ORDER BY x.assessed_at DESC, x.id DESC
LIMIT 1
)
WHERE u.id = :customer_id
"""
)
class RiskAuthorityProfile(BaseModel):
"""服务端权威风险画像(只读汇总,不含测评问卷原文)。"""
model_config = ConfigDict(extra="forbid", frozen=True)
customer_id: str
customer_risk_level: int | None = None
professional_investor: bool = False
assessed_at: datetime | None = None
valid_until: datetime | None = None
authority_reason: AuthorityReason = "AUTHORITY_OK"
class SuitabilityToolInput(BaseModel):
"""ToolExecutor 使用的严格输入模型,避免业务 Agent 自行拼接规则。
调用方只能声明“给谁、买什么等级的产品、是否需要揭示/确认”,
风险等级与测评有效期一律由服务端权威来源解析。
"""
model_config = ConfigDict(extra="forbid", frozen=True)
customer_id: str = Field(min_length=1, max_length=20, pattern=r"^[0-9]+$")
product_risk_level: int = Field(ge=1, le=5)
product_requires_disclosure: bool = True
requires_confirmation: bool = False
class SuitabilityDecision(BaseModel):
model_config = ConfigDict(extra="forbid", frozen=True)
allowed: bool
reason_code: str
required_disclosure: bool
requires_confirmation: bool
requires_recording: bool
customer_risk_level: int | None = None
risk_level_source: str = RISK_LEVEL_SOURCE
professional_investor: bool = False
assessment_valid_until: datetime | None = None
def _as_utc(value: object) -> datetime | None:
"""库内 DATETIME 为 UTC naive,统一规范化为带时区,便于安全比较。"""
if not isinstance(value, datetime):
return None
return value if value.tzinfo is not None else value.replace(tzinfo=UTC)
def _risk_level_from_investor_type(investor_type: object) -> int | None:
if not isinstance(investor_type, str):
return None
matched = _INVESTOR_TYPE_PATTERN.match(investor_type.strip().upper())
return int(matched.group(1)) if matched is not None else None
class SuitabilityService:
"""执行 C1-C5/R1-R5 的公共、只读适当性规则。"""
def __init__(self, *, session_factory: Callable[[], Any] | None = None) -> None:
self._session_factory: Callable[[], Any] = session_factory or SessionFactory
async def evaluate(
self, request: SuitabilityToolInput, context: RequestContext, *, now: datetime | None = None
) -> SuitabilityDecision:
current = now or datetime.now(UTC)
self._assert_customer_scope(request.customer_id, context)
profile = await self._load_authority_profile(request.customer_id)
return self._decide(request, profile, current)
async def _load_authority_profile(self, customer_id: str) -> RiskAuthorityProfile:
async with self._session_factory() as session:
result = await session.execute(_AUTHORITY_SQL, {"customer_id": int(customer_id)})
row: Mapping[str, Any] | None = result.mappings().first()
if row is None:
return RiskAuthorityProfile(
customer_id=customer_id, authority_reason="CUSTOMER_NOT_FOUND"
)
level = _risk_level_from_investor_type(row["investor_type"])
valid_until = _as_utc(row["valid_until"])
if row["investor_type"] is None:
reason: AuthorityReason = "ASSESSMENT_MISSING"
elif level is None:
reason = "RISK_LEVEL_INVALID"
elif valid_until is None:
# 有效期缺失视为不可用,绝不按“长期有效”放行。
reason = "ASSESSMENT_EXPIRED"
else:
reason = "AUTHORITY_OK"
return RiskAuthorityProfile(
customer_id=customer_id,
customer_risk_level=level,
professional_investor=(
bool(row["is_professional_investor"])
and str(row["professional_investor_status"]) == PROFESSIONAL_INVESTOR_CERTIFIED
),
assessed_at=_as_utc(row["assessed_at"]),
valid_until=valid_until,
authority_reason=reason,
)
def _decide(
self, request: SuitabilityToolInput, profile: RiskAuthorityProfile, current: datetime
) -> SuitabilityDecision:
if profile.authority_reason == "CUSTOMER_NOT_FOUND":
return self._denied("CUSTOMER_NOT_FOUND", request, profile)
if profile.authority_reason == "ASSESSMENT_MISSING":
return self._denied("ASSESSMENT_MISSING", request, profile)
if profile.authority_reason == "RISK_LEVEL_INVALID":
return self._denied("RISK_LEVEL_INVALID", request, profile)
if profile.valid_until is None or profile.valid_until <= current:
# 测评过期即拒绝:专业投资者也不能绕过有效期。
return self._denied("ASSESSMENT_EXPIRED", request, profile)
if profile.customer_risk_level is None:
return self._denied("ASSESSMENT_MISSING", request, profile)
if profile.professional_investor:
# 已认定专业投资者可豁免等级匹配,但必须揭示、确认并录音留痕。
return SuitabilityDecision(
allowed=True,
reason_code="SUITABLE_PROFESSIONAL_INVESTOR",
required_disclosure=True,
requires_confirmation=True,
requires_recording=True,
customer_risk_level=profile.customer_risk_level,
professional_investor=True,
assessment_valid_until=profile.valid_until,
)
if profile.customer_risk_level < request.product_risk_level:
return self._denied("RISK_LEVEL_MISMATCH", request, profile)
required_disclosure = request.product_requires_disclosure
return SuitabilityDecision(
allowed=True,
reason_code="SUITABLE",
required_disclosure=required_disclosure,
requires_confirmation=request.requires_confirmation or required_disclosure,
requires_recording=required_disclosure or request.requires_confirmation,
customer_risk_level=profile.customer_risk_level,
professional_investor=False,
assessment_valid_until=profile.valid_until,
)
async def evaluate_and_audit(
self, request: SuitabilityToolInput, context: RequestContext, *, now: datetime | None = None
) -> SuitabilityDecision:
decision = await self.evaluate(request, context, now=now)
# 适当性决定是受监管业务决策,拒绝和通过都留痕;只记录权威来源摘要,
# 不保存测评问卷原文,也不保存调用方自报的任何等级。
async with self._session_factory() as session, session.begin():
actor_id = int(context.user_id) if context.user_id.isdecimal() else None
session.add(InteractionAudit(
actor_type="agent",
actor_id=actor_id,
portal=context.portal,
action_type="suitability.checked",
detail={
"trace_id": context.trace_id,
"status": "allowed" if decision.allowed else "denied",
"reason_code": decision.reason_code,
"customer_id": request.customer_id,
"customer_risk_level": decision.customer_risk_level,
"risk_level_source": decision.risk_level_source,
"professional_investor": decision.professional_investor,
"professional_investor_source": PROFESSIONAL_INVESTOR_SOURCE,
"assessment_valid_until": (
decision.assessment_valid_until.isoformat()
if decision.assessment_valid_until is not None
else None
),
"product_risk_level": request.product_risk_level,
},
created_at=datetime.now(UTC).replace(tzinfo=None),
))
return decision
@staticmethod
def _assert_customer_scope(customer_id: str, context: RequestContext) -> None:
"""公共鉴权:除管理员外不得查询他人风险测评。"""
if set(context.roles).intersection(CUSTOMER_SCOPE_EXEMPT_ROLES):
return
if customer_id == context.user_id or customer_id in context.customer_ids:
return
raise ForbiddenAgentError("不能查询该客户的风险测评")
@staticmethod
def _denied(
reason_code: str, request: SuitabilityToolInput, profile: RiskAuthorityProfile
) -> SuitabilityDecision:
return SuitabilityDecision(
allowed=False,
reason_code=reason_code,
required_disclosure=request.product_requires_disclosure,
requires_confirmation=True,
requires_recording=True,
customer_risk_level=profile.customer_risk_level,
professional_investor=profile.professional_investor,
assessment_valid_until=profile.valid_until,
)
async def suitability_tool_handler(
arguments: SuitabilityToolInput, context: RequestContext
) -> dict[str, Any]:
decision = await SuitabilityService().evaluate_and_audit(arguments, context)
return decision.model_dump(mode="json")