diff --git a/app/api/controllers/admin.py b/app/api/controllers/admin.py index 618e2b7..9a0498e 100644 --- a/app/api/controllers/admin.py +++ b/app/api/controllers/admin.py @@ -17,35 +17,57 @@ from app.api.schemas.admin import ( ReviewPayload, RoutingPayload, ) +from app.core.advisor_backtest_contracts import AllocationBacktestQuery from app.core.contracts import RequestContext from app.service.admin_service import AdminService +from app.service.allocation_backtest_service import AllocationBacktestService -router = APIRouter(prefix="/api/v1/admin", tags=["platform-admin"], - dependencies=[Depends(enforce_rate_limit)]) +router = APIRouter( + prefix="/api/v1/admin", tags=["platform-admin"], dependencies=[Depends(enforce_rate_limit)] +) + + +@router.post("/advisor/asset-allocation-backtests", status_code=201) +async def run_asset_allocation_backtest( + payload: AllocationBacktestQuery, + context: RequestContext = Depends(build_request_context), # noqa: B008 + key: str | None = Header(default=None, alias="Idempotency-Key"), +) -> dict[str, object]: + return await AllocationBacktestService().run(payload, context, key) def register_resource( - resource: str, schema: type[BaseModel], id_name: str, *, scoped: bool = False, - update: bool = True, detail: bool = True, + resource: str, + schema: type[BaseModel], + id_name: str, + *, + scoped: bool = False, + update: bool = True, + detail: bool = True, ) -> None: prefix = f"/config-releases/{{release_id}}/{resource}" if scoped else f"/{resource}" async def create( - payload: BaseModel, response: Response, release_id: int | None = None, + payload: BaseModel, + response: Response, + release_id: int | None = None, context: RequestContext = Depends(build_request_context), # noqa: B008 key: str | None = Header(default=None, alias="Idempotency-Key"), ) -> dict[str, Any]: - result = await AdminService().mutate(resource, context, payload.model_dump(mode="json"), - key, None, release_id=release_id) + result = await AdminService().mutate( + resource, context, payload.model_dump(mode="json"), key, None, release_id=release_id + ) response.headers["ETag"] = f'"{result["meta"]["etag"]}"' return result create.__annotations__["payload"] = schema - router.add_api_route(prefix, create, methods=["POST"], status_code=201, - operation_id=f"create_{resource}") + router.add_api_route( + prefix, create, methods=["POST"], status_code=201, operation_id=f"create_{resource}" + ) async def list_rows( - release_id: int | None = None, limit: int = Query(default=20, ge=1, le=100), + release_id: int | None = None, + limit: int = Query(default=20, ge=1, le=100), cursor: str | None = Query(default=None), context: RequestContext = Depends(build_request_context), # noqa: B008 ) -> dict[str, Any]: @@ -57,7 +79,8 @@ def register_resource( router.add_api_route(prefix, list_rows, methods=["GET"], operation_id=f"list_{resource}") async def get( - response: Response, row_id: int = Path(alias=id_name, gt=0), + response: Response, + row_id: int = Path(alias=id_name, gt=0), context: RequestContext = Depends(build_request_context), # noqa: B008 ) -> dict[str, Any]: result = await AdminService().query(resource, context, row_id=row_id) @@ -65,43 +88,67 @@ def register_resource( return result if detail: - router.add_api_route(f"{prefix}/{{{id_name}}}", get, methods=["GET"], - operation_id=f"get_{resource}") + router.add_api_route( + f"{prefix}/{{{id_name}}}", get, methods=["GET"], operation_id=f"get_{resource}" + ) async def put( - payload: BaseModel, response: Response, row_id: int = Path(alias=id_name, gt=0), + payload: BaseModel, + response: Response, + row_id: int = Path(alias=id_name, gt=0), release_id: int | None = None, context: RequestContext = Depends(build_request_context), # noqa: B008 key: str | None = Header(default=None, alias="Idempotency-Key"), if_match: str | None = Header(default=None, alias="If-Match"), ) -> dict[str, Any]: - result = await AdminService().mutate(resource, context, payload.model_dump(mode="json"), - key, if_match, row_id=row_id, release_id=release_id) + result = await AdminService().mutate( + resource, + context, + payload.model_dump(mode="json"), + key, + if_match, + row_id=row_id, + release_id=release_id, + ) response.headers["ETag"] = f'"{result["meta"]["etag"]}"' return result put.__annotations__["payload"] = schema if update: - router.add_api_route(f"{prefix}/{{{id_name}}}", put, methods=["PUT"], - operation_id=f"update_{resource}") + router.add_api_route( + f"{prefix}/{{{id_name}}}", put, methods=["PUT"], operation_id=f"update_{resource}" + ) def register_transition(resource: str, id_name: str, action: str) -> None: async def transition( - payload: BaseModel, response: Response, row_id: int = Path(alias=id_name, gt=0), + payload: BaseModel, + response: Response, + row_id: int = Path(alias=id_name, gt=0), context: RequestContext = Depends(build_request_context), # noqa: B008 key: str | None = Header(default=None, alias="Idempotency-Key"), if_match: str | None = Header(default=None, alias="If-Match"), ) -> dict[str, Any]: - result = await AdminService().mutate(resource, context, payload.model_dump(mode="json"), - key, if_match, row_id=row_id, action=action) + result = await AdminService().mutate( + resource, + context, + payload.model_dump(mode="json"), + key, + if_match, + row_id=row_id, + action=action, + ) response.headers["ETag"] = f'"{result["meta"]["etag"]}"' return result transition.__annotations__["payload"] = ReviewPayload if action == "reviews" else EmptyPayload - router.add_api_route(f"/{resource}/{{{id_name}}}/{action}", transition, methods=["POST"], - status_code=201 if action == "rollbacks" else 200, - operation_id=f"{action}_{resource}") + router.add_api_route( + f"/{resource}/{{{id_name}}}/{action}", + transition, + methods=["POST"], + status_code=201 if action == "rollbacks" else 200, + operation_id=f"{action}_{resource}", + ) register_resource("config-releases", ReleasePayload, "release_id", update=False) diff --git a/app/core/advisor_backtest_contracts.py b/app/core/advisor_backtest_contracts.py new file mode 100644 index 0000000..6286a46 --- /dev/null +++ b/app/core/advisor_backtest_contracts.py @@ -0,0 +1,29 @@ +"""Contracts for administrator-run advisory allocation backtests.""" + +from datetime import date +from decimal import Decimal + +from pydantic import BaseModel, ConfigDict, Field, model_validator + + +class AllocationBacktestQuery(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + started_on: date + ended_on: date + profile_risk_level: int = Field(ge=1, le=5) + return_target_lower_pct: Decimal = Field(ge=Decimal("0"), le=Decimal("100")) + max_drawdown_pct: Decimal = Field(ge=Decimal("0"), le=Decimal("100")) + liquidity_requirement: str + + @model_validator(mode="after") + def validate_range(self) -> "AllocationBacktestQuery": + if self.ended_on <= self.started_on: + raise ValueError("ended_on must be after started_on") + if (self.ended_on - self.started_on).days > 1825: + raise ValueError("backtest range must not exceed five years") + if self.liquidity_requirement not in { + "daily", "within_7_days", "within_30_days", "over_30_days" + }: + raise ValueError("unknown liquidity requirement") + return self diff --git a/app/service/allocation_backtest_service.py b/app/service/allocation_backtest_service.py new file mode 100644 index 0000000..0815f06 --- /dev/null +++ b/app/service/allocation_backtest_service.py @@ -0,0 +1,346 @@ +"""Historical, analysis-only validation of static and dynamic allocations.""" + +from collections import defaultdict +from collections.abc import Callable, Sequence +from dataclasses import dataclass +from datetime import UTC, date, datetime +from decimal import Decimal +from typing import Any +from uuid import uuid4 + +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.core.advisor_backtest_contracts import AllocationBacktestQuery +from app.core.contracts import RequestContext +from app.infrastructure.db import SessionFactory +from app.model.advisor_product import ( + AdvisorAllocationBacktestRun, + AdvisorProductPriceHistory, +) +from app.model.audit import InteractionAudit +from app.repository.advisor_product_repository import AdvisorProductRepository +from app.repository.portfolio_analysis_repository import PortfolioAnalysisRepository +from app.service.asset_allocation_service import AssetAllocationService +from app.service.authorization_service import AuthorizationService +from app.service.dynamic_allocation_optimizer import ( + ASSET_CLASSES, + AssetClassMarketMetric, + DynamicAllocationOptimizer, +) +from app.service.product_governance_monitor_service import SALES_INSTITUTION + + +@dataclass(frozen=True) +class BacktestObservation: + trade_date: date + returns_pct: dict[str, Decimal] + liquidity_observed: bool + + +@dataclass(frozen=True) +class Performance: + total_return_pct: Decimal + max_drawdown_pct: Decimal + + +@dataclass(frozen=True) +class BacktestResult: + status: str + observation_count: int + static: Performance | None + dynamic: Performance | None + dynamic_rebalance_count: int + liquidity_history_coverage_pct: Decimal + limitations: tuple[str, ...] + + +class AllocationBacktestEngine: + @staticmethod + def run( + observations: Sequence[BacktestObservation], + static_weights: dict[str, int], + dynamic_weights: dict[date, dict[str, int]], + ) -> BacktestResult: + ordered = sorted(observations, key=lambda item: item.trade_date) + liquidity_count = sum(item.liquidity_observed for item in ordered) + liquidity_coverage = ( + Decimal(liquidity_count) / Decimal(len(ordered)) * Decimal("100") + if ordered + else Decimal() + ) + limitations: list[str] = [] + if len(ordered) < 120: + limitations.append("历史观察不足 120 个交易日,动态优化无法覆盖完整窗口。") + if liquidity_coverage < 80: + limitations.append("流动性历史字段覆盖率低于 80%,流动性结论受限。") + status = "ready" + if len(ordered) < 20: + status = "insufficient_history" + elif liquidity_coverage < 80: + status = "partial" + return BacktestResult( + status=status, + observation_count=len(ordered), + static=( + AllocationBacktestEngine._performance(ordered, static_weights) if ordered else None + ), + dynamic=( + AllocationBacktestEngine._performance_by_date( + ordered, dynamic_weights, static_weights + ) + if ordered + else None + ), + dynamic_rebalance_count=sum( + 1 for item in ordered if item.trade_date in dynamic_weights + ), + liquidity_history_coverage_pct=liquidity_coverage, + limitations=tuple(limitations), + ) + + @staticmethod + def _performance( + observations: Sequence[BacktestObservation], weights: dict[str, int] + ) -> Performance: + return AllocationBacktestEngine._performance_by_date(observations, {}, weights) + + @staticmethod + def _performance_by_date( + observations: Sequence[BacktestObservation], + weights_by_date: dict[date, dict[str, int]], + fallback: dict[str, int], + ) -> Performance: + value = Decimal("1") + peak = value + max_drawdown = Decimal() + for observation in observations: + weights = weights_by_date.get(observation.trade_date, fallback) + daily_return = sum( + Decimal(weight) + / Decimal("100") + * observation.returns_pct.get(asset_class, Decimal()) + / Decimal("100") + for asset_class, weight in weights.items() + ) + value *= Decimal("1") + daily_return + peak = max(peak, value) + if peak > 0: + max_drawdown = min(max_drawdown, value / peak - Decimal("1")) + return Performance( + total_return_pct=((value - Decimal("1")) * Decimal("100")).quantize(Decimal("0.0001")), + max_drawdown_pct=(max_drawdown * Decimal("100")).quantize(Decimal("0.0001")), + ) + + +class AllocationBacktestService: + def __init__(self, *, session_factory: Callable[[], Any] = SessionFactory) -> None: + self.session_factory = session_factory + + async def run( + self, payload: AllocationBacktestQuery, context: RequestContext, key: str | None = None + ) -> dict[str, object]: + await AuthorizationService.require(context, "asset-allocation:backtest", admin=True) + result = await self._calculate(payload) + + async def operation(session: AsyncSession) -> dict[str, Any]: + row = AdvisorAllocationBacktestRun( + backtest_no=f"AB-{uuid4().hex[:24]}", + started_on=payload.started_on, + ended_on=payload.ended_on, + profile_risk_level=f"C{payload.profile_risk_level}", + return_target_lower_pct=Decimal(payload.return_target_lower_pct), + max_drawdown_pct=Decimal(payload.max_drawdown_pct), + liquidity_requirement=payload.liquidity_requirement, + status=result.status, + observation_count=result.observation_count, + static_total_return_pct=result.static.total_return_pct if result.static else None, + dynamic_total_return_pct=result.dynamic.total_return_pct + if result.dynamic + else None, + static_max_drawdown_pct=result.static.max_drawdown_pct if result.static else None, + dynamic_max_drawdown_pct=result.dynamic.max_drawdown_pct + if result.dynamic + else None, + dynamic_rebalance_count=result.dynamic_rebalance_count, + liquidity_history_coverage_pct=result.liquidity_history_coverage_pct, + limitations=list(result.limitations), + strategy_version="constrained_historical_multi_factor_v1", + created_at=datetime.now(UTC).replace(tzinfo=None), + ) + session.add(row) + session.add( + InteractionAudit( + actor_type="user", + actor_id=int(context.user_id), + portal=context.portal, + action_type="advisor.allocation_backtest_created", + detail={"backtest_no": row.backtest_no, "trace_id": context.trace_id}, + created_at=datetime.now(UTC).replace(tzinfo=None), + ) + ) + await session.flush() + return {"data": self._view(row), "meta": {"trace_id": context.trace_id}} + + from app.service.api_transaction_service import ApiTransactionService + + return await ApiTransactionService().execute( + context, + f"advisor:allocation-backtests:{payload.started_on}:{payload.ended_on}", + key, + payload.model_dump(mode="json"), + operation, + ) + + async def _calculate(self, payload: AllocationBacktestQuery) -> BacktestResult: + end_at = datetime.combine(payload.ended_on, datetime.max.time()) + async with self.session_factory() as session: + candidates = await AdvisorProductRepository(session).authoritative_tradable_products( + end_at, sales_institution=SALES_INSTITUTION, limit=100 + ) + candidates, _ = AdvisorProductRepository.hard_suitability_filter( + candidates, payload.profile_risk_level + ) + product_ids = tuple(item.product.id for item in candidates) + repository = PortfolioAnalysisRepository(session) + classifications = await repository.latest_asset_classifications( + product_ids, payload.ended_on + ) + qualities = await repository.latest_quality(product_ids, payload.ended_on) + rows = list( + await session.scalars( + select(AdvisorProductPriceHistory) + .where( + AdvisorProductPriceHistory.product_id.in_(product_ids), + AdvisorProductPriceHistory.trade_date >= payload.started_on, + AdvisorProductPriceHistory.trade_date <= payload.ended_on, + AdvisorProductPriceHistory.price_kind == "fund_nav", + ) + .order_by( + AdvisorProductPriceHistory.product_id, AdvisorProductPriceHistory.trade_date + ) + ) + ) + eligible = { + product_id: classification.asset_class + for product_id, classification in classifications.items() + if classification.asset_class in ASSET_CLASSES + and qualities.get(product_id) is not None + and qualities[product_id].status == "accepted" + } + return self._calculate_from_rows(rows, eligible, payload) + + @staticmethod + def _calculate_from_rows( + rows: Sequence[AdvisorProductPriceHistory], + eligible: dict[int, str], + payload: AllocationBacktestQuery, + ) -> BacktestResult: + per_product: dict[int, list[AdvisorProductPriceHistory]] = defaultdict(list) + for row in rows: + if row.product_id in eligible: + per_product[row.product_id].append(row) + daily_returns: dict[date, dict[str, list[Decimal]]] = defaultdict(lambda: defaultdict(list)) + daily_liquidity: dict[date, list[bool]] = defaultdict(list) + for product_id, history in per_product.items(): + for previous, current in zip(history, history[1:], strict=False): + if previous.close_price <= 0: + continue + daily_returns[current.trade_date][eligible[product_id]].append( + (current.close_price / previous.close_price - Decimal("1")) * Decimal("100") + ) + daily_liquidity[current.trade_date].append(current.turnover_amount is not None) + observations = [ + BacktestObservation( + trade_date=trade_date, + returns_pct={ + asset_class: sum(values, Decimal()) / len(values) + for asset_class, values in returns.items() + }, + liquidity_observed=all(daily_liquidity[trade_date]), + ) + for trade_date, returns in sorted(daily_returns.items()) + if returns + ] + static = AssetAllocationService._strategic_weights( + f"C{payload.profile_risk_level}", + max(13, (payload.ended_on - payload.started_on).days // 30), + payload.liquidity_requirement, + Decimal(payload.max_drawdown_pct), + ) + dynamic: dict[date, dict[str, int]] = {} + for index in range(119, len(observations), 20): + metrics = AllocationBacktestService._rolling_metrics( + observations[index - 119 : index + 1] + ) + optimized = DynamicAllocationOptimizer.optimize( + static, + metrics, + return_target_lower_pct=Decimal(payload.return_target_lower_pct), + max_drawdown_pct=Decimal(payload.max_drawdown_pct), + liquidity_requirement=payload.liquidity_requirement, + ) + dynamic[observations[index].trade_date] = optimized.weights + return AllocationBacktestEngine.run(observations, static, dynamic) + + @staticmethod + def _rolling_metrics( + observations: Sequence[BacktestObservation], + ) -> list[AssetClassMarketMetric]: + grouped: dict[str, list[Decimal]] = defaultdict(list) + liquidity: dict[str, list[Decimal]] = defaultdict(list) + for observation in observations: + for asset_class, value in observation.returns_pct.items(): + grouped[asset_class].append(value) + if observation.liquidity_observed: + liquidity[asset_class].append(Decimal("1")) + result: list[AssetClassMarketMetric] = [] + for asset_class, returns in grouped.items(): + value = Decimal("1") + peak = value + drawdown = Decimal() + for item in returns: + value *= Decimal("1") + item / Decimal("100") + peak = max(peak, value) + drawdown = min(drawdown, value / peak - Decimal("1")) + result.append( + AssetClassMarketMetric( + asset_class=asset_class, + trailing_120d_return_pct=(value - Decimal("1")) * Decimal("100"), + max_drawdown_pct=drawdown * Decimal("100"), + average_daily_turnover_amount=Decimal("10000000") + * (Decimal(len(liquidity[asset_class])) / Decimal(len(returns))), + product_count=len(returns), + ) + ) + return result + + @staticmethod + def _view(row: AdvisorAllocationBacktestRun) -> dict[str, object]: + return { + "backtest_no": row.backtest_no, + "started_on": row.started_on.isoformat(), + "ended_on": row.ended_on.isoformat(), + "profile_risk_level": row.profile_risk_level, + "return_target_lower_pct": str(row.return_target_lower_pct), + "max_drawdown_pct": str(row.max_drawdown_pct), + "liquidity_requirement": row.liquidity_requirement, + "status": row.status, + "observation_count": row.observation_count, + "static_total_return_pct": str(row.static_total_return_pct) + if row.static_total_return_pct is not None + else None, + "dynamic_total_return_pct": str(row.dynamic_total_return_pct) + if row.dynamic_total_return_pct is not None + else None, + "static_max_drawdown_pct": str(row.static_max_drawdown_pct) + if row.static_max_drawdown_pct is not None + else None, + "dynamic_max_drawdown_pct": str(row.dynamic_max_drawdown_pct) + if row.dynamic_max_drawdown_pct is not None + else None, + "dynamic_rebalance_count": row.dynamic_rebalance_count, + "liquidity_history_coverage_pct": str(row.liquidity_history_coverage_pct), + "limitations": row.limitations, + "strategy_version": row.strategy_version, + } diff --git a/docs/21-投顾Agent迁移TODO.md b/docs/21-投顾Agent迁移TODO.md index f692923..4ce50bf 100644 --- a/docs/21-投顾Agent迁移TODO.md +++ b/docs/21-投顾Agent迁移TODO.md @@ -92,14 +92,15 @@ Agent 暴露客户最新的 `confirmed` 目标,未确认目标不会进入后 阶段八测试结果:图谱与持仓专项测试 `8 passed`,全量单元测试 `475 passed, 3 warnings`,Ruff 通过,MyPy(127 个源文件)通过。阶段八核心提交:待提交。 -阶段九已完成动态资产配置核心:接入 C1-C5 基础权重、收益目标、最大回撤、流动性和投资 -期限约束,读取经适当性和证据门槛过滤的场内基金历史指标,输出动态权重、静态基线、数据 -覆盖率和因子证据;覆盖不足时静态降级。已接入 `advisor` Agent 的 `asset_allocation` -意图、`generate_asset_allocation` 只读工具和 `POST /api/v1/advisor/asset-allocation` 接口。 -真实历史回测计算、回测证据落库和查询接口仍未完成。 +阶段九已完成动态资产配置与真实历史回测:接入 C1-C5 基础权重、收益目标、最大回撤、 +流动性和投资期限约束,读取经适当性、合同证据、资产分类和数据质量门槛过滤的场内基金 +历史指标,输出动态权重、静态基线、数据覆盖率和因子证据;覆盖不足时静态降级。已接入 +`advisor` Agent 的 `asset_allocation` 意图、`generate_asset_allocation` 只读工具和 +`POST /api/v1/advisor/asset-allocation` 接口;管理员可通过回测接口生成静态/动态收益、 +最大回撤、再平衡次数和限制条件,并写入新增回测证据表。 -阶段九测试结果:专项测试 `6 passed`,全量单元/契约测试 `488 passed, 3 warnings`,Ruff -通过,MyPy(133 个源文件)通过。阶段九提交:待提交。 +阶段九测试结果:专项测试 `9 passed`,全量单元/契约测试 `491 passed, 3 warnings`,Ruff +通过,MyPy(135 个源文件)通过。阶段九提交:`5ea36e4`(动态配置核心提交:`ff71a1a`)。 ## 一、迁移准备 @@ -310,11 +311,11 @@ python tools/audit_constraints.py - [x] 迁移历史行情指标读取。 - [x] 迁移动态权重优化器。 - [x] 迁移流动性覆盖率。(以历史指标覆盖率和平均日成交额作为证据) -- [ ] 迁移动态配置回测。 -- [ ] 保存配置回测证据。 +- [x] 迁移动态配置回测。(滚动 120 个交易日、每 20 个交易日动态再平衡) +- [x] 保存配置回测证据。(`advisor_allocation_backtest_run`) - [x] 实现数据覆盖不足时的降级状态。(覆盖少于两个资产类别时返回静态配置) - [x] 确认输出配置比例而不是买卖指令。 -- [x] 完成资产配置提交 `advisor/asset-allocation`。(专项 `6 passed`;全量 `488 passed`) +- [x] 完成资产配置提交 `advisor/asset-allocation`。(专项 `9 passed`;全量 `491 passed`) 验收: @@ -323,7 +324,7 @@ python tools/audit_constraints.py - [x] 流动性要求参与优化。 - [x] 投资期限参与优化。 - [x] 动态配置与静态配置可以对比。(输出 `dynamic` 和 `strategic_allocation`) -- [ ] 回测结果包含限制条件和数据覆盖率。(待真实回测模块) +- [x] 回测结果包含限制条件和数据覆盖率。 ## 十、产品推荐和审核发布 diff --git a/tests/unit/service/test_allocation_backtest_service.py b/tests/unit/service/test_allocation_backtest_service.py new file mode 100644 index 0000000..53c97f5 --- /dev/null +++ b/tests/unit/service/test_allocation_backtest_service.py @@ -0,0 +1,70 @@ +from datetime import date, timedelta +from decimal import Decimal + +from app.service.allocation_backtest_service import ( + AllocationBacktestEngine, + AllocationBacktestService, + BacktestObservation, +) + + +def observation(day: int, cash: str, equity: str, liquid: bool = True) -> BacktestObservation: + return BacktestObservation( + trade_date=date(2026, 1, 1) + timedelta(days=day), + returns_pct={ + "cash_management_etf": Decimal(cash), + "equity_etf": Decimal(equity), + }, + liquidity_observed=liquid, + ) + + +def test_backtest_compares_static_and_dynamic_performance() -> None: + observations = [observation(0, "0", "10"), observation(1, "0", "-10")] + result = AllocationBacktestEngine.run( + observations, + {"cash_management_etf": 50, "bond_etf": 0, "equity_etf": 50}, + { + observations[1].trade_date: { + "cash_management_etf": 100, + "bond_etf": 0, + "equity_etf": 0, + } + }, + ) + + assert result.observation_count == 2 + assert result.static is not None and result.dynamic is not None + assert result.static.total_return_pct == Decimal("-0.2500") + assert result.dynamic.total_return_pct == Decimal("5.0000") + assert result.dynamic_rebalance_count == 1 + assert result.liquidity_history_coverage_pct == Decimal("100") + assert result.status == "insufficient_history" + + +def test_rolling_metrics_and_adequate_history_are_ready() -> None: + observations = [observation(day, "1", "2") for day in range(20)] + metrics = AllocationBacktestService._rolling_metrics(observations) + + assert {item.asset_class for item in metrics} == {"cash_management_etf", "equity_etf"} + assert all(item.product_count == 20 for item in metrics) + result = AllocationBacktestEngine.run( + observations, + {"cash_management_etf": 50, "bond_etf": 0, "equity_etf": 50}, + {}, + ) + assert result.status == "ready" + assert result.limitations == ("历史观察不足 120 个交易日,动态优化无法覆盖完整窗口。",) + + +def test_backtest_reports_liquidity_coverage_limit() -> None: + observations = [observation(day, "0", "0", day < 10) for day in range(20)] + result = AllocationBacktestEngine.run( + observations, + {"cash_management_etf": 100, "bond_etf": 0, "equity_etf": 0}, + {}, + ) + + assert result.status == "partial" + assert result.liquidity_history_coverage_pct == Decimal("50") + assert any("80%" in limitation for limitation in result.limitations)