diff --git a/app/repository/advisor_product_repository.py b/app/repository/advisor_product_repository.py index d421acd..a3aa651 100644 --- a/app/repository/advisor_product_repository.py +++ b/app/repository/advisor_product_repository.py @@ -13,6 +13,7 @@ from app.model.advisor_product import ( AdvisorProductContractSnapshot, AdvisorProductGovernanceCandidate, AdvisorProductMarketQuoteSnapshot, + AdvisorProductMetricSnapshot, AdvisorProductReferenceSnapshot, AdvisorProductSuitabilityReference, ) @@ -45,6 +46,13 @@ class ProductContractEvidence: document_published_at: date | None +@dataclass(frozen=True, slots=True) +class ProductLiquidityEvidence: + average_daily_turnover_amount: Decimal | None + latest_quote_observed_at: datetime | None + status: str + + @dataclass(frozen=True, slots=True) class AuthoritativeProductCandidate: product: FundProduct @@ -52,6 +60,7 @@ class AuthoritativeProductCandidate: contract: ProductContractEvidence asset_scale_billion: Decimal | None = None market_quote: AdvisorProductMarketQuoteSnapshot | None = None + liquidity: ProductLiquidityEvidence | None = None class _ProductRow(Protocol): @@ -96,6 +105,7 @@ class AdvisorProductRepository: fund_manager: str | None = None, quote_max_age_seconds: int = 300, min_asset_scale_billion: Decimal | None = None, + liquidity_requirement: str | None = None, limit: int = 50, ) -> list[AuthoritativeProductCandidate]: """Apply evidence gates before any product reaches an Agent. @@ -151,6 +161,15 @@ class AdvisorProductRepository: AdvisorProductReferenceSnapshot.product_id, AdvisorProductReferenceSnapshot.as_of_date.desc(), ))) + metric_rows = list(await self.session.scalars(select( + AdvisorProductMetricSnapshot + ).where( + AdvisorProductMetricSnapshot.product_id.in_(ids), + AdvisorProductMetricSnapshot.as_of_date <= as_of, + ).order_by( + AdvisorProductMetricSnapshot.product_id, + AdvisorProductMetricSnapshot.as_of_date.desc(), + ))) pending_ids = set(await self.session.scalars(select( AdvisorProductGovernanceCandidate.product_id ).where( @@ -169,6 +188,7 @@ class AdvisorProductRepository: suitability_by_product = self._latest_by_product(suitability_rows) contract_by_product = self._latest_by_product(contract_rows) reference_by_product = self._latest_by_product(reference_rows) + metric_by_product = self._latest_by_product(metric_rows) quote_cutoff = now - timedelta(seconds=quote_max_age_seconds) quote_by_product: dict[int, AdvisorProductMarketQuoteSnapshot] = {} for quote in quote_rows: @@ -189,6 +209,11 @@ class AdvisorProductRepository: and (scale is None or scale < min_asset_scale_billion) ): continue + metric = metric_by_product.get(product.id) + latest_quote = quote_by_product.get(product.id) + liquidity = self._liquidity_evidence(metric, latest_quote, liquidity_requirement) + if liquidity is not None and liquidity.status == "insufficient": + continue result.append(AuthoritativeProductCandidate( product=product, suitability=ProductSuitabilityEvidence( @@ -214,10 +239,64 @@ class AdvisorProductRepository: document_published_at=contract.document_published_at, ), asset_scale_billion=scale, - market_quote=quote_by_product.get(product.id), + market_quote=latest_quote, + liquidity=liquidity, )) return result + @staticmethod + def _liquidity_evidence( + metric: AdvisorProductMetricSnapshot | None, + quote: AdvisorProductMarketQuoteSnapshot | None, + requirement: str | None, + ) -> ProductLiquidityEvidence | None: + if metric is None and quote is None and requirement is None: + return None + thresholds = { + "daily": Decimal("10000000"), + "within_7_days": Decimal("1000000"), + "within_30_days": Decimal("100000"), + "over_30_days": Decimal("0"), + } + if requirement is not None and requirement not in thresholds: + raise ValueError("unknown liquidity requirement") + average = metric.average_daily_turnover_amount if metric is not None else None + status = "unknown" if average is None else "available" + if requirement is not None and (average is None or average < thresholds[requirement]): + status = "insufficient" + return ProductLiquidityEvidence( + average_daily_turnover_amount=average, + latest_quote_observed_at=quote.observed_at if quote is not None else None, + status=status, + ) + + @staticmethod + def hard_suitability_filter( + candidates: Sequence[AuthoritativeProductCandidate], customer_risk_level: int + ) -> tuple[list[AuthoritativeProductCandidate], list[dict[str, object]]]: + """Apply R-level hard filtering before ranking or model generation.""" + if not 1 <= customer_risk_level <= 5: + raise ValueError("customer risk level must be between 1 and 5") + selected: list[AuthoritativeProductCandidate] = [] + excluded: list[dict[str, object]] = [] + for candidate in candidates: + product_level = candidate.suitability.risk_level.upper().removeprefix("R") + if not product_level.isdigit() or not 1 <= int(product_level) <= 5: + excluded.append({ + "product_code": candidate.product.product_code, + "reason_code": "PRODUCT_RISK_LEVEL_INVALID", + "reason": "产品权威适当性等级无法解析,已失败关闭。", + }) + elif int(product_level) > customer_risk_level: + excluded.append({ + "product_code": candidate.product.product_code, + "reason_code": "RISK_LEVEL_MISMATCH", + "reason": "产品风险等级高于客户风险承受等级,已排除。", + }) + else: + selected.append(candidate) + return selected, excluded + @staticmethod def _latest_by_product(rows: Sequence[ProductRow]) -> dict[int, ProductRow]: result: dict[int, ProductRow] = {} diff --git a/docs/21-投顾Agent迁移TODO.md b/docs/21-投顾Agent迁移TODO.md index 6d1f169..892b30f 100644 --- a/docs/21-投顾Agent迁移TODO.md +++ b/docs/21-投顾Agent迁移TODO.md @@ -49,7 +49,11 @@ R4=7、R5=1;资产分类成功 18 个,1 个因合同证据不足跳过。每 阶段五已完成部分测试:产品证据 Repository `2 passed`,官网治理解析 `3 passed`, 治理导入/分类导入 `5 passed`,合计专项 `10 passed`;Ruff、MyPy、数据库结构审计和 -约束审计通过。阶段五尚未完成流动性指标细化和完整推荐侧适当性联动,暂不提交阶段完成标记。 +约束审计通过。阶段五后续补齐了历史平均成交额流动性证据、流动性门槛失败关闭、 +R1-R5 适当性硬过滤和排除原因。权威数据验收仍保留两个限制:R1/R5 尚不足四个 +测试产品,1 个产品因合同证据不足未完成资产分类;完整推荐侧联动留待阶段十。 +阶段五测试结果:专项 `7 passed`,全量单元测试 `473 passed, 3 warnings`,Ruff 和 +MyPy 通过。阶段五提交:待提交。 ### 阶段六:开户风险问卷 @@ -201,21 +205,21 @@ python tools/audit_constraints.py - [x] 接入产品治理变更监控。(`ProductGovernanceMonitorService` + 官网同步工具) - [x] 接入产品资产规模指标。(`advisor_product_reference_snapshot`) - [x] 接入产品历史行情指标。(增量 NAV 同步 + `ProductMetricService`) -- [ ] 接入产品流动性指标。 +- [x] 接入产品流动性指标。(平均日成交额、最新行情时间、流动性状态) - [x] 实现场内基金过滤。(Repository 强制 `SSE/SZSE`) -- [ ] 实现适当性硬过滤。 +- [x] 实现适当性硬过滤。(R1-R5 不匹配失败关闭并返回排除原因) - [x] 实现合同证据过滤。(verified + 有来源 URL/文档摘要) - [x] 实现来源缺失时的失败关闭。(缺失权威证据不进入候选) -- [ ] 完成产品数据提交 `advisor/product-data`。(阶段五进行中) +- [x] 完成产品数据提交 `advisor/product-data`。(专项 `7 passed`;完整推荐侧联动留待阶段十) 验收: -- [ ] R1、R2、R3、R4、R5 均有可查询产品。 +- [x] R1、R2、R3、R4、R5 均有可查询产品。(19 个南方场内产品,五级均有) - [ ] 每个风险等级至少有四个测试产品。 -- [ ] 各基金类型分类正确。 -- [ ] 场外基金不会进入场内交易表。 -- [ ] 不适配产品会被排除并说明原因。 -- [ ] 产品来源和合同证据可追溯。 +- [ ] 各基金类型分类正确。(1 个产品因合同证据不足跳过分类) +- [x] 场外基金不会进入场内交易表。 +- [x] 不适配产品会被排除并说明原因。(硬过滤返回 `reason_code`) +- [x] 产品来源和合同证据可追溯。 ## 六、开户风险问卷 diff --git a/tests/unit/repository/test_advisor_product_repository.py b/tests/unit/repository/test_advisor_product_repository.py index 80ac264..d529d6c 100644 --- a/tests/unit/repository/test_advisor_product_repository.py +++ b/tests/unit/repository/test_advisor_product_repository.py @@ -80,6 +80,7 @@ async def test_authoritative_candidates_fail_closed_and_keep_r1_without_scale_li ], [], [], + [], ]) result = await AdvisorProductRepository(session).authoritative_tradable_products( now, sales_institution="南方基金管理股份有限公司直销", @@ -95,10 +96,52 @@ async def test_authoritative_candidates_fail_closed_and_keep_r1_without_scale_li async def test_pending_governance_and_missing_evidence_are_excluded() -> None: now = datetime(2026, 9, 11, 10, 0) session = FakeSession([ - [product(1)], [suitability(1)], [contract(1)], [], [1], [], + [product(1)], [suitability(1)], [contract(1)], [], [], [1], [], ]) result = await AdvisorProductRepository(session).authoritative_tradable_products( now, sales_institution="南方基金管理股份有限公司直销", ) assert result == [] assert "exchange_code IN" in str(session.statements[0]) + + +def test_hard_suitability_filter_returns_exclusion_reasons() -> None: + from app.repository.advisor_product_repository import ( + AdvisorProductRepository, + AuthoritativeProductCandidate, + ProductContractEvidence, + ProductSuitabilityEvidence, + ) + + def candidate(code: str, risk: str) -> AuthoritativeProductCandidate: + item = product(1) + item.product_code = code + return AuthoritativeProductCandidate( + product=item, + suitability=ProductSuitabilityEvidence( + risk_level=risk, sales_institution="南方基金管理股份有限公司直销", + source_url="https://example.test/risk.pdf", document_title="风险等级披露", + document_sha256="a" * 64, effective_from=date(2026, 1, 1), + ), + contract=ProductContractEvidence( + fund_type="股票型ETF", investment_scope="场内基金", performance_benchmark=None, + risk_return_characteristics="净值波动", custodian_name=None, + management_fee_rate_pct=None, custodian_fee_rate_pct=None, inception_date=None, + source_url="https://example.test/contract.pdf", document_title="基金合同", + document_sha256="b" * 64, document_published_at=None, + ), + ) + + selected, excluded = AdvisorProductRepository.hard_suitability_filter( + [candidate("R1", "R1"), candidate("R4", "R4")], 2 + ) + assert [item.product.product_code for item in selected] == ["R1"] + assert excluded[0]["reason_code"] == "RISK_LEVEL_MISMATCH" + + +def test_liquidity_requirement_is_failed_closed_without_metrics() -> None: + from app.repository.advisor_product_repository import AdvisorProductRepository + + evidence = AdvisorProductRepository._liquidity_evidence(None, None, "within_7_days") + assert evidence is not None + assert evidence.status == "insufficient"