feat: add allocation backtest evidence
This commit is contained in:
@@ -17,35 +17,57 @@ from app.api.schemas.admin import (
|
|||||||
ReviewPayload,
|
ReviewPayload,
|
||||||
RoutingPayload,
|
RoutingPayload,
|
||||||
)
|
)
|
||||||
|
from app.core.advisor_backtest_contracts import AllocationBacktestQuery
|
||||||
from app.core.contracts import RequestContext
|
from app.core.contracts import RequestContext
|
||||||
from app.service.admin_service import AdminService
|
from app.service.admin_service import AdminService
|
||||||
|
from app.service.allocation_backtest_service import AllocationBacktestService
|
||||||
|
|
||||||
router = APIRouter(prefix="/api/v1/admin", tags=["platform-admin"],
|
router = APIRouter(
|
||||||
dependencies=[Depends(enforce_rate_limit)])
|
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(
|
def register_resource(
|
||||||
resource: str, schema: type[BaseModel], id_name: str, *, scoped: bool = False,
|
resource: str,
|
||||||
update: bool = True, detail: bool = True,
|
schema: type[BaseModel],
|
||||||
|
id_name: str,
|
||||||
|
*,
|
||||||
|
scoped: bool = False,
|
||||||
|
update: bool = True,
|
||||||
|
detail: bool = True,
|
||||||
) -> None:
|
) -> None:
|
||||||
prefix = f"/config-releases/{{release_id}}/{resource}" if scoped else f"/{resource}"
|
prefix = f"/config-releases/{{release_id}}/{resource}" if scoped else f"/{resource}"
|
||||||
|
|
||||||
async def create(
|
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
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
||||||
key: str | None = Header(default=None, alias="Idempotency-Key"),
|
key: str | None = Header(default=None, alias="Idempotency-Key"),
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
result = await AdminService().mutate(resource, context, payload.model_dump(mode="json"),
|
result = await AdminService().mutate(
|
||||||
key, None, release_id=release_id)
|
resource, context, payload.model_dump(mode="json"), key, None, release_id=release_id
|
||||||
|
)
|
||||||
response.headers["ETag"] = f'"{result["meta"]["etag"]}"'
|
response.headers["ETag"] = f'"{result["meta"]["etag"]}"'
|
||||||
return result
|
return result
|
||||||
|
|
||||||
create.__annotations__["payload"] = schema
|
create.__annotations__["payload"] = schema
|
||||||
router.add_api_route(prefix, create, methods=["POST"], status_code=201,
|
router.add_api_route(
|
||||||
operation_id=f"create_{resource}")
|
prefix, create, methods=["POST"], status_code=201, operation_id=f"create_{resource}"
|
||||||
|
)
|
||||||
|
|
||||||
async def list_rows(
|
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),
|
cursor: str | None = Query(default=None),
|
||||||
context: RequestContext = Depends(build_request_context), # noqa: B008
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
@@ -57,7 +79,8 @@ def register_resource(
|
|||||||
router.add_api_route(prefix, list_rows, methods=["GET"], operation_id=f"list_{resource}")
|
router.add_api_route(prefix, list_rows, methods=["GET"], operation_id=f"list_{resource}")
|
||||||
|
|
||||||
async def get(
|
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
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
result = await AdminService().query(resource, context, row_id=row_id)
|
result = await AdminService().query(resource, context, row_id=row_id)
|
||||||
@@ -65,43 +88,67 @@ def register_resource(
|
|||||||
return result
|
return result
|
||||||
|
|
||||||
if detail:
|
if detail:
|
||||||
router.add_api_route(f"{prefix}/{{{id_name}}}", get, methods=["GET"],
|
router.add_api_route(
|
||||||
operation_id=f"get_{resource}")
|
f"{prefix}/{{{id_name}}}", get, methods=["GET"], operation_id=f"get_{resource}"
|
||||||
|
)
|
||||||
|
|
||||||
async def put(
|
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,
|
release_id: int | None = None,
|
||||||
context: RequestContext = Depends(build_request_context), # noqa: B008
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
||||||
key: str | None = Header(default=None, alias="Idempotency-Key"),
|
key: str | None = Header(default=None, alias="Idempotency-Key"),
|
||||||
if_match: str | None = Header(default=None, alias="If-Match"),
|
if_match: str | None = Header(default=None, alias="If-Match"),
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
result = await AdminService().mutate(resource, context, payload.model_dump(mode="json"),
|
result = await AdminService().mutate(
|
||||||
key, if_match, row_id=row_id, release_id=release_id)
|
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"]}"'
|
response.headers["ETag"] = f'"{result["meta"]["etag"]}"'
|
||||||
return result
|
return result
|
||||||
|
|
||||||
put.__annotations__["payload"] = schema
|
put.__annotations__["payload"] = schema
|
||||||
if update:
|
if update:
|
||||||
router.add_api_route(f"{prefix}/{{{id_name}}}", put, methods=["PUT"],
|
router.add_api_route(
|
||||||
operation_id=f"update_{resource}")
|
f"{prefix}/{{{id_name}}}", put, methods=["PUT"], operation_id=f"update_{resource}"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def register_transition(resource: str, id_name: str, action: str) -> None:
|
def register_transition(resource: str, id_name: str, action: str) -> None:
|
||||||
async def transition(
|
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
|
context: RequestContext = Depends(build_request_context), # noqa: B008
|
||||||
key: str | None = Header(default=None, alias="Idempotency-Key"),
|
key: str | None = Header(default=None, alias="Idempotency-Key"),
|
||||||
if_match: str | None = Header(default=None, alias="If-Match"),
|
if_match: str | None = Header(default=None, alias="If-Match"),
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
result = await AdminService().mutate(resource, context, payload.model_dump(mode="json"),
|
result = await AdminService().mutate(
|
||||||
key, if_match, row_id=row_id, action=action)
|
resource,
|
||||||
|
context,
|
||||||
|
payload.model_dump(mode="json"),
|
||||||
|
key,
|
||||||
|
if_match,
|
||||||
|
row_id=row_id,
|
||||||
|
action=action,
|
||||||
|
)
|
||||||
response.headers["ETag"] = f'"{result["meta"]["etag"]}"'
|
response.headers["ETag"] = f'"{result["meta"]["etag"]}"'
|
||||||
return result
|
return result
|
||||||
|
|
||||||
transition.__annotations__["payload"] = ReviewPayload if action == "reviews" else EmptyPayload
|
transition.__annotations__["payload"] = ReviewPayload if action == "reviews" else EmptyPayload
|
||||||
router.add_api_route(f"/{resource}/{{{id_name}}}/{action}", transition, methods=["POST"],
|
router.add_api_route(
|
||||||
status_code=201 if action == "rollbacks" else 200,
|
f"/{resource}/{{{id_name}}}/{action}",
|
||||||
operation_id=f"{action}_{resource}")
|
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)
|
register_resource("config-releases", ReleasePayload, "release_id", update=False)
|
||||||
|
|||||||
@@ -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
|
||||||
@@ -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,
|
||||||
|
}
|
||||||
+12
-11
@@ -92,14 +92,15 @@ Agent 暴露客户最新的 `confirmed` 目标,未确认目标不会进入后
|
|||||||
阶段八测试结果:图谱与持仓专项测试 `8 passed`,全量单元测试 `475 passed, 3 warnings`,Ruff
|
阶段八测试结果:图谱与持仓专项测试 `8 passed`,全量单元测试 `475 passed, 3 warnings`,Ruff
|
||||||
通过,MyPy(127 个源文件)通过。阶段八核心提交:待提交。
|
通过,MyPy(127 个源文件)通过。阶段八核心提交:待提交。
|
||||||
|
|
||||||
阶段九已完成动态资产配置核心:接入 C1-C5 基础权重、收益目标、最大回撤、流动性和投资
|
阶段九已完成动态资产配置与真实历史回测:接入 C1-C5 基础权重、收益目标、最大回撤、
|
||||||
期限约束,读取经适当性和证据门槛过滤的场内基金历史指标,输出动态权重、静态基线、数据
|
流动性和投资期限约束,读取经适当性、合同证据、资产分类和数据质量门槛过滤的场内基金
|
||||||
覆盖率和因子证据;覆盖不足时静态降级。已接入 `advisor` Agent 的 `asset_allocation`
|
历史指标,输出动态权重、静态基线、数据覆盖率和因子证据;覆盖不足时静态降级。已接入
|
||||||
意图、`generate_asset_allocation` 只读工具和 `POST /api/v1/advisor/asset-allocation` 接口。
|
`advisor` Agent 的 `asset_allocation` 意图、`generate_asset_allocation` 只读工具和
|
||||||
真实历史回测计算、回测证据落库和查询接口仍未完成。
|
`POST /api/v1/advisor/asset-allocation` 接口;管理员可通过回测接口生成静态/动态收益、
|
||||||
|
最大回撤、再平衡次数和限制条件,并写入新增回测证据表。
|
||||||
|
|
||||||
阶段九测试结果:专项测试 `6 passed`,全量单元/契约测试 `488 passed, 3 warnings`,Ruff
|
阶段九测试结果:专项测试 `9 passed`,全量单元/契约测试 `491 passed, 3 warnings`,Ruff
|
||||||
通过,MyPy(133 个源文件)通过。阶段九提交:待提交。
|
通过,MyPy(135 个源文件)通过。阶段九提交:`5ea36e4`(动态配置核心提交:`ff71a1a`)。
|
||||||
|
|
||||||
## 一、迁移准备
|
## 一、迁移准备
|
||||||
|
|
||||||
@@ -310,11 +311,11 @@ python tools/audit_constraints.py
|
|||||||
- [x] 迁移历史行情指标读取。
|
- [x] 迁移历史行情指标读取。
|
||||||
- [x] 迁移动态权重优化器。
|
- [x] 迁移动态权重优化器。
|
||||||
- [x] 迁移流动性覆盖率。(以历史指标覆盖率和平均日成交额作为证据)
|
- [x] 迁移流动性覆盖率。(以历史指标覆盖率和平均日成交额作为证据)
|
||||||
- [ ] 迁移动态配置回测。
|
- [x] 迁移动态配置回测。(滚动 120 个交易日、每 20 个交易日动态再平衡)
|
||||||
- [ ] 保存配置回测证据。
|
- [x] 保存配置回测证据。(`advisor_allocation_backtest_run`)
|
||||||
- [x] 实现数据覆盖不足时的降级状态。(覆盖少于两个资产类别时返回静态配置)
|
- [x] 实现数据覆盖不足时的降级状态。(覆盖少于两个资产类别时返回静态配置)
|
||||||
- [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] 投资期限参与优化。
|
- [x] 投资期限参与优化。
|
||||||
- [x] 动态配置与静态配置可以对比。(输出 `dynamic` 和 `strategic_allocation`)
|
- [x] 动态配置与静态配置可以对比。(输出 `dynamic` 和 `strategic_allocation`)
|
||||||
- [ ] 回测结果包含限制条件和数据覆盖率。(待真实回测模块)
|
- [x] 回测结果包含限制条件和数据覆盖率。
|
||||||
|
|
||||||
## 十、产品推荐和审核发布
|
## 十、产品推荐和审核发布
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
||||||
Reference in New Issue
Block a user