diff --git a/app/api/controllers/asset_allocation.py b/app/api/controllers/asset_allocation.py new file mode 100644 index 0000000..642b620 --- /dev/null +++ b/app/api/controllers/asset_allocation.py @@ -0,0 +1,22 @@ +"""Analysis-only dynamic asset allocation endpoint.""" + +from fastapi import APIRouter, Depends + +from app.api.dependencies.auth import build_request_context +from app.api.dependencies.rate_limit import enforce_rate_limit +from app.api.schemas.asset_allocation import AssetAllocationQuery +from app.core.contracts import RequestContext +from app.service.asset_allocation_service import AssetAllocationService + +router = APIRouter( + prefix="/api/v1/advisor", tags=["advisor-asset-allocation"], + dependencies=[Depends(enforce_rate_limit)], +) + + +@router.post("/asset-allocation") +async def generate_asset_allocation( + payload: AssetAllocationQuery, + context: RequestContext = Depends(build_request_context), # noqa: B008 +) -> dict[str, object]: + return await AssetAllocationService().generate_for_agent(payload, context) diff --git a/app/api/schemas/asset_allocation.py b/app/api/schemas/asset_allocation.py new file mode 100644 index 0000000..3d2bb7f --- /dev/null +++ b/app/api/schemas/asset_allocation.py @@ -0,0 +1,3 @@ +from app.core.advisor_allocation_contracts import AssetAllocationQuery + +__all__ = ["AssetAllocationQuery"] diff --git a/app/core/advisor_allocation_contracts.py b/app/core/advisor_allocation_contracts.py new file mode 100644 index 0000000..56625e2 --- /dev/null +++ b/app/core/advisor_allocation_contracts.py @@ -0,0 +1,7 @@ +"""Contracts for analysis-only asset allocation.""" + +from pydantic import BaseModel, ConfigDict + + +class AssetAllocationQuery(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) diff --git a/app/main.py b/app/main.py index 018eac8..757407b 100644 --- a/app/main.py +++ b/app/main.py @@ -3,6 +3,7 @@ from fastapi.responses import JSONResponse from app.api.controllers.admin import router as admin_router from app.api.controllers.agent_runs import router as agent_runs_router +from app.api.controllers.asset_allocation import router as asset_allocation_router from app.api.controllers.conversations import router as conversations_router from app.api.controllers.health import router as health_router from app.api.controllers.investment_goals import router as investment_goals_router @@ -51,6 +52,7 @@ def create_app() -> FastAPI: application.include_router(onboarding_router) application.include_router(investment_goals_router) application.include_router(portfolio_analysis_router) + application.include_router(asset_allocation_router) application.include_router(admin_router) return application diff --git a/app/model/advisor_product.py b/app/model/advisor_product.py index 49eb1b2..b61dcf2 100644 --- a/app/model/advisor_product.py +++ b/app/model/advisor_product.py @@ -75,6 +75,60 @@ class AdvisorProductMetricSnapshot(Base): created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False) +class AdvisorProductAssetClassification(Base): + __tablename__ = "advisor_product_asset_classification" + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True) + product_id: Mapped[int] = mapped_column(BigInteger, nullable=False) + as_of_date: Mapped[date] = mapped_column(Date, nullable=False) + asset_class: Mapped[str] = mapped_column(String(32), nullable=False) + source: Mapped[str] = mapped_column(String(64), nullable=False) + status: Mapped[str] = mapped_column(String(16), nullable=False) + created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False) + updated_at: Mapped[datetime] = mapped_column(DateTime, nullable=False) + + +class AdvisorProductDataQualitySnapshot(Base): + __tablename__ = "advisor_product_data_quality_snapshot" + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True) + product_id: Mapped[int] = mapped_column(BigInteger, nullable=False) + as_of_date: Mapped[date] = mapped_column(Date, nullable=False) + observation_count: Mapped[int] = mapped_column(BigInteger, nullable=False) + expected_trading_days: Mapped[int] = mapped_column(BigInteger, nullable=False) + price_coverage_pct: Mapped[Decimal] = mapped_column(Numeric(7, 4), nullable=False) + turnover_coverage_pct: Mapped[Decimal] = mapped_column(Numeric(7, 4), nullable=False) + max_abs_daily_return_pct: Mapped[Decimal | None] = mapped_column(Numeric(10, 4)) + status: Mapped[str] = mapped_column(String(16), nullable=False) + reason_codes: Mapped[list[str]] = mapped_column(JSON, nullable=False) + rule_version: Mapped[str] = mapped_column(String(16), nullable=False) + created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False) + + +class AdvisorAllocationBacktestRun(Base): + __tablename__ = "advisor_allocation_backtest_run" + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True) + backtest_no: Mapped[str] = mapped_column(String(36), unique=True, nullable=False) + started_on: Mapped[date] = mapped_column(Date, nullable=False) + ended_on: Mapped[date] = mapped_column(Date, nullable=False) + profile_risk_level: Mapped[str] = mapped_column(String(8), nullable=False) + return_target_lower_pct: Mapped[Decimal] = mapped_column(Numeric(7, 4), nullable=False) + max_drawdown_pct: Mapped[Decimal] = mapped_column(Numeric(7, 4), nullable=False) + liquidity_requirement: Mapped[str] = mapped_column(String(32), nullable=False) + status: Mapped[str] = mapped_column(String(32), nullable=False) + observation_count: Mapped[int] = mapped_column(BigInteger, nullable=False) + static_total_return_pct: Mapped[Decimal | None] = mapped_column(Numeric(12, 4)) + dynamic_total_return_pct: Mapped[Decimal | None] = mapped_column(Numeric(12, 4)) + static_max_drawdown_pct: Mapped[Decimal | None] = mapped_column(Numeric(12, 4)) + dynamic_max_drawdown_pct: Mapped[Decimal | None] = mapped_column(Numeric(12, 4)) + dynamic_rebalance_count: Mapped[int] = mapped_column(BigInteger, nullable=False) + liquidity_history_coverage_pct: Mapped[Decimal] = mapped_column(Numeric(7, 4), nullable=False) + limitations: Mapped[list[str]] = mapped_column(JSON, nullable=False) + strategy_version: Mapped[str] = mapped_column(String(32), nullable=False) + created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False) + + class AdvisorProductSuitabilityReference(Base): __tablename__ = "advisor_product_suitability_reference" diff --git a/app/repository/portfolio_analysis_repository.py b/app/repository/portfolio_analysis_repository.py index 92d97b2..a150048 100644 --- a/app/repository/portfolio_analysis_repository.py +++ b/app/repository/portfolio_analysis_repository.py @@ -5,7 +5,12 @@ from datetime import date from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession -from app.model.advisor_product import AdvisorProductIndustryExposure, AdvisorProductMetricSnapshot +from app.model.advisor_product import ( + AdvisorProductAssetClassification, + AdvisorProductDataQualitySnapshot, + AdvisorProductIndustryExposure, + AdvisorProductMetricSnapshot, +) from app.model.fund import FundHolding, FundProduct @@ -73,3 +78,42 @@ class PortfolioAnalysisRepository: for row in rows: selected.setdefault(row.product_id, row) return selected + + async def latest_asset_classifications( + self, product_ids: tuple[int, ...], as_of_date: date + ) -> dict[int, AdvisorProductAssetClassification]: + if not product_ids: + return {} + rows = await self.session.scalars( + select(AdvisorProductAssetClassification).where( + AdvisorProductAssetClassification.product_id.in_(product_ids), + AdvisorProductAssetClassification.as_of_date <= as_of_date, + AdvisorProductAssetClassification.status == "active", + ).order_by( + AdvisorProductAssetClassification.product_id, + AdvisorProductAssetClassification.as_of_date.desc(), + ) + ) + selected: dict[int, AdvisorProductAssetClassification] = {} + for row in rows: + selected.setdefault(row.product_id, row) + return selected + + async def latest_quality( + self, product_ids: tuple[int, ...], as_of_date: date + ) -> dict[int, AdvisorProductDataQualitySnapshot]: + if not product_ids: + return {} + rows = await self.session.scalars( + select(AdvisorProductDataQualitySnapshot).where( + AdvisorProductDataQualitySnapshot.product_id.in_(product_ids), + AdvisorProductDataQualitySnapshot.as_of_date <= as_of_date, + ).order_by( + AdvisorProductDataQualitySnapshot.product_id, + AdvisorProductDataQualitySnapshot.as_of_date.desc(), + ) + ) + selected: dict[int, AdvisorProductDataQualitySnapshot] = {} + for row in rows: + selected.setdefault(row.product_id, row) + return selected diff --git a/app/service/agent/bootstrap.py b/app/service/agent/bootstrap.py index a6b2919..63aea6d 100644 --- a/app/service/agent/bootstrap.py +++ b/app/service/agent/bootstrap.py @@ -4,6 +4,7 @@ from typing import Any, cast from sqlalchemy.ext.asyncio import AsyncSession +from app.core.advisor_allocation_contracts import AssetAllocationQuery from app.core.config import get_settings from app.core.errors import RecoverableAgentError from app.core.fund_contracts import FundQuoteQuery @@ -16,6 +17,7 @@ from app.service.agent.factory import AgentFactory from app.service.agent.governance import PlatformGovernance from app.service.agent.implementations.advisor import AdvisorAgent from app.service.agent.implementations.fund_query_demo import FundQueryDemoAgent +from app.service.asset_allocation_service import asset_allocation_tool from app.service.fund_quote_service import query_fund_quote_tool from app.service.intent_classifier import IntentClassifier from app.service.investment_goal_service import investment_goal_query_tool @@ -163,6 +165,14 @@ def get_agent_factory() -> AgentFactory: allowed_roles=("customer", "advisor", "operator", "admin"), timeout_seconds=10, )) + registry.register(ToolDefinition( + name="generate_asset_allocation", + input_model=AssetAllocationQuery, + handler=cast(Any, asset_allocation_tool), + required_permission="asset-allocation:generate:self", + allowed_roles=("customer", "advisor", "operator", "admin"), + timeout_seconds=15, + )) model_service = get_model_service() endpoint_resolver = DatabaseModelEndpointResolver() factory = AgentFactory( diff --git a/app/service/agent/implementations/advisor.py b/app/service/agent/implementations/advisor.py index 114f2ad..3411c51 100644 --- a/app/service/agent/implementations/advisor.py +++ b/app/service/agent/implementations/advisor.py @@ -14,8 +14,13 @@ class AdvisorAgent(FundQueryDemoAgent): version="0.1.0", allowed_roles=("customer", "advisor", "operator", "admin"), allowed_portals=("api",), - allowed_tools=("query_fund_quote", "query_investment_goal", "analyze_portfolio"), - supported_intents=("fund_quote", "investment_goal", "portfolio_analysis"), + allowed_tools=( + "query_fund_quote", "query_investment_goal", "analyze_portfolio", + "generate_asset_allocation", + ), + supported_intents=( + "fund_quote", "investment_goal", "portfolio_analysis", "asset_allocation", + ), ) async def handle(self, request: AgentRequest, context: RequestContext) -> CoreResult: @@ -37,6 +42,14 @@ class AdvisorAgent(FundQueryDemoAgent): "analyze_portfolio", {}, intent="portfolio_analysis", context=context ) return CoreResult(text=self._describe_portfolio(output)) + if ( + self._classified_intent is not None + and self._classified_intent.intent == "asset_allocation" + ): + output = await self.call_tool( + "generate_asset_allocation", {}, intent="asset_allocation", context=context + ) + return CoreResult(text=self._describe_allocation(output)) return await super().handle(request, context) @staticmethod @@ -68,3 +81,30 @@ class AdvisorAgent(FundQueryDemoAgent): f"总市值 {summary.get('total_market_value')},产品集中度 HHI 为 {hhi}。" "分析结果仅供参考,不生成交易指令。" ) + + @staticmethod + def _describe_allocation(result: object) -> str: + if not isinstance(result, dict): + return "资产配置分析暂不可用,请稍后重试。" + status = result.get("status") + if status == "profile_required": + return "当前缺少有效风险画像,暂不能生成资产配置。" + if status == "investment_goal_required": + return "当前没有已确认的投资目标,暂不能生成资产配置。" + if status != "ready": + return "资产配置分析数据不完整,请稍后重试。" + allocation = result.get("allocation") + if not isinstance(allocation, list) or not allocation: + return "当前没有足够的场内基金数据生成资产配置。" + parts = [ + f"{item.get('label', item.get('asset_class'))} {item.get('target_pct')}%" + for item in allocation if isinstance(item, dict) + ] + optimization = result.get("optimization") + dynamic = isinstance(optimization, dict) and bool(optimization.get("dynamic")) + mode = "动态历史因子优化" if dynamic else "静态配置(历史数据覆盖不足)" + return ( + f"资产配置分析完成({mode}):" + ",".join(parts) + + "。该结果综合考虑收益目标、最大回撤、流动性和投资期限," + "仅供分析参考,不构成交易指令。" + ) diff --git a/app/service/asset_allocation_service.py b/app/service/asset_allocation_service.py new file mode 100644 index 0000000..9ef381c --- /dev/null +++ b/app/service/asset_allocation_service.py @@ -0,0 +1,183 @@ +"""Dynamic, analysis-only asset allocation from confirmed goals and history.""" + +from collections import defaultdict +from collections.abc import Callable +from datetime import UTC, date, datetime +from decimal import Decimal +from typing import Any + +from app.core.advisor_allocation_contracts import AssetAllocationQuery +from app.core.contracts import RequestContext +from app.infrastructure.db import SessionFactory +from app.repository.advisor_product_repository import AdvisorProductRepository +from app.repository.portfolio_analysis_repository import PortfolioAnalysisRepository +from app.service.authorization_service import AuthorizationService +from app.service.dynamic_allocation_optimizer import ( + ASSET_CLASSES, + AssetClassMarketMetric, + DynamicAllocationOptimizer, +) +from app.service.investment_goal_service import InvestmentGoalService +from app.service.product_governance_monitor_service import SALES_INSTITUTION +from app.service.suitability_service import SuitabilityService + +BASE_ALLOCATIONS = { + "C1": {"cash_management_etf": 50, "bond_etf": 40, "equity_etf": 10}, + "C2": {"cash_management_etf": 30, "bond_etf": 50, "equity_etf": 20}, + "C3": {"cash_management_etf": 15, "bond_etf": 45, "equity_etf": 40}, + "C4": {"cash_management_etf": 10, "bond_etf": 25, "equity_etf": 65}, + "C5": {"cash_management_etf": 5, "bond_etf": 15, "equity_etf": 80}, +} +ASSET_LABELS = { + "cash_management_etf": "现金管理类场内基金", + "bond_etf": "债券类场内基金", + "equity_etf": "权益类场内基金", +} + + +class AssetAllocationService: + def __init__(self, *, session_factory: Callable[[], Any] = SessionFactory) -> None: + self.session_factory = session_factory + + async def generate_for_agent( + self, _arguments: AssetAllocationQuery, context: RequestContext + ) -> dict[str, object]: + await AuthorizationService.require(context, "asset-allocation:generate:self") + authority = await SuitabilityService().authority_for_customer(int(context.user_id)) + if authority.customer_risk_level is None: + return {"status": "profile_required"} + goal = await InvestmentGoalService().current_for_agent(context) + if goal is None: + return {"status": "investment_goal_required"} + horizon = goal.get("investment_horizon_months") + if not isinstance(horizon, int): + return {"status": "investment_goal_invalid"} + risk = f"C{authority.customer_risk_level}" + liquidity = str(goal["liquidity_requirement"]) + strategic = self._strategic_weights( + risk, horizon, liquidity, Decimal(str(goal["max_drawdown_pct"])) + ) + metrics = await self._market_metrics(authority.customer_risk_level) + optimized = DynamicAllocationOptimizer.optimize( + strategic, + metrics, + return_target_lower_pct=Decimal(str(goal["annualized_return_lower_pct"])), + max_drawdown_pct=Decimal(str(goal["max_drawdown_pct"])), + liquidity_requirement=liquidity, + ) + return { + "status": "ready", + "allocation": [ + {"asset_class": key, "label": ASSET_LABELS[key], "target_pct": value} + for key, value in optimized.weights.items() + ], + "optimization": { + "method": "constrained_historical_multi_factor_v1", + "dynamic": optimized.dynamic, + "metric_coverage_pct": str(optimized.metric_coverage_pct.quantize(Decimal("0.01"))), + "strategic_allocation": strategic, + "factor_evidence": optimized.factors, + }, + "constraints": { + "annualized_return_lower_pct": goal["annualized_return_lower_pct"], + "max_drawdown_pct": goal["max_drawdown_pct"], + "liquidity_requirement": liquidity, + "investment_horizon_months": horizon, + }, + "analysis_only": True, + } + + async def _market_metrics(self, customer_risk_level: int) -> list[AssetClassMarketMetric]: + now = datetime.now(UTC).replace(tzinfo=None) + async with self.session_factory() as session: + candidates = await AdvisorProductRepository(session).authoritative_tradable_products( + now, sales_institution=SALES_INSTITUTION, limit=50 + ) + candidates, _excluded = AdvisorProductRepository.hard_suitability_filter( + candidates, customer_risk_level + ) + ids = tuple(item.product.id for item in candidates) + repository = PortfolioAnalysisRepository(session) + classifications = await repository.latest_asset_classifications(ids, date.today()) + snapshots = await repository.latest_metrics(ids, date.today()) + qualities = await repository.latest_quality(ids, date.today()) + grouped: dict[str, list[AssetClassMarketMetric]] = defaultdict(list) + for product_id in ids: + classification = classifications.get(product_id) + metric = snapshots.get(product_id) + quality = qualities.get(product_id) + if ( + classification is None + or classification.asset_class not in ASSET_CLASSES + or quality is None + or quality.status != "accepted" + or metric is None + or metric.trailing_120d_return_pct is None + or metric.max_drawdown_pct is None + or metric.average_daily_turnover_amount is None + or metric.observation_count < 20 + ): + continue + grouped[classification.asset_class].append( + AssetClassMarketMetric( + asset_class=classification.asset_class, + trailing_120d_return_pct=metric.trailing_120d_return_pct, + max_drawdown_pct=metric.max_drawdown_pct, + average_daily_turnover_amount=metric.average_daily_turnover_amount, + product_count=1, + ) + ) + return [ + AssetClassMarketMetric( + asset_class=key, + trailing_120d_return_pct=sum( + (item.trailing_120d_return_pct for item in rows), Decimal() + ) + / len(rows), + max_drawdown_pct=( + sum((item.max_drawdown_pct for item in rows), Decimal()) / len(rows) + ), + average_daily_turnover_amount=sum( + (item.average_daily_turnover_amount for item in rows), Decimal() + ) + / len(rows), + product_count=len(rows), + ) + for key, rows in grouped.items() + ] + + @staticmethod + def _strategic_weights( + risk: str, horizon: int, liquidity: str, max_drawdown: Decimal + ) -> dict[str, int]: + weights = dict(BASE_ALLOCATIONS[risk]) + if horizon <= 12: + AssetAllocationService._move(weights, "equity_etf", "cash_management_etf", 10) + elif horizon >= 60 and risk != "C1": + AssetAllocationService._move(weights, "bond_etf", "equity_etf", 5) + if liquidity == "daily": + AssetAllocationService._move(weights, "equity_etf", "cash_management_etf", 10) + AssetAllocationService._move(weights, "bond_etf", "cash_management_etf", 5) + if max_drawdown <= 10: + AssetAllocationService._cap_equity(weights, 20) + elif max_drawdown <= 20: + AssetAllocationService._cap_equity(weights, 40) + return weights + + @staticmethod + def _move(weights: dict[str, int], source: str, target: str, amount: int) -> None: + moved = min(weights[source], amount) + weights[source] -= moved + weights[target] += moved + + @staticmethod + def _cap_equity(weights: dict[str, int], cap: int) -> None: + excess = max(0, weights["equity_etf"] - cap) + weights["equity_etf"] -= excess + weights["bond_etf"] += excess + + +async def asset_allocation_tool( + arguments: AssetAllocationQuery, context: RequestContext +) -> dict[str, object]: + return await AssetAllocationService().generate_for_agent(arguments, context) diff --git a/app/service/dynamic_allocation_optimizer.py b/app/service/dynamic_allocation_optimizer.py new file mode 100644 index 0000000..d99446c --- /dev/null +++ b/app/service/dynamic_allocation_optimizer.py @@ -0,0 +1,160 @@ +"""Explainable constraint-first allocation optimizer.""" + +from dataclasses import dataclass +from decimal import ROUND_HALF_UP, Decimal + +HUNDRED = Decimal("100") +ASSET_CLASSES = ("cash_management_etf", "bond_etf", "equity_etf") + + +@dataclass(frozen=True) +class AssetClassMarketMetric: + asset_class: str + trailing_120d_return_pct: Decimal + max_drawdown_pct: Decimal + average_daily_turnover_amount: Decimal + product_count: int + + +@dataclass(frozen=True) +class DynamicAllocationResult: + weights: dict[str, int] + dynamic: bool + metric_coverage_pct: Decimal + factors: list[dict[str, object]] + + +class DynamicAllocationOptimizer: + @classmethod + def optimize( + cls, + strategic_weights: dict[str, int], + metrics: list[AssetClassMarketMetric], + *, + return_target_lower_pct: Decimal, + max_drawdown_pct: Decimal, + liquidity_requirement: str, + ) -> DynamicAllocationResult: + available = { + item.asset_class: item for item in metrics if item.asset_class in ASSET_CLASSES + } + coverage = Decimal(len(available)) / Decimal(len(ASSET_CLASSES)) * HUNDRED + if len(available) < 2: + return DynamicAllocationResult( + dict(strategic_weights), False, coverage, cls._factors(metrics, {}) + ) + scores = { + key: cls._score( + value, metrics, return_target_lower_pct, max_drawdown_pct, liquidity_requirement + ) + for key, value in available.items() + } + weights = {key: Decimal(value) for key, value in strategic_weights.items()} + mean = sum(scores.values(), Decimal()) / Decimal(len(scores)) + winners = {key: score - mean for key, score in scores.items() if score > mean} + losers = {key: mean - score for key, score in scores.items() if score < mean} + if winners and losers: + tilt = min(Decimal("15"), sum(weights.get(key, Decimal()) for key in losers)) + for key, value in winners.items(): + weights[key] = weights.get(key, Decimal()) + tilt * value / sum(winners.values()) + for key, value in losers.items(): + weights[key] = max( + Decimal(), weights.get(key, Decimal()) - tilt * value / sum(losers.values()) + ) + equity_cap = ( + Decimal("20") + if max_drawdown_pct <= 10 + else Decimal("40") + if max_drawdown_pct <= 20 + else Decimal("85") + ) + excess = max(Decimal(), weights["equity_etf"] - equity_cap) + weights["equity_etf"] -= excess + weights["bond_etf"] += excess + cash_min = { + "daily": Decimal("25"), + "within_7_days": Decimal("15"), + "within_30_days": Decimal("5"), + "over_30_days": Decimal(), + }[liquidity_requirement] + shortfall = max(Decimal(), cash_min - weights["cash_management_etf"]) + for key in ("bond_etf", "equity_etf"): + moved = min(shortfall, weights[key]) + weights[key] -= moved + weights["cash_management_etf"] += moved + shortfall -= moved + rounded = { + key: int(value.quantize(Decimal("1"), rounding=ROUND_HALF_UP)) + for key, value in weights.items() + } + largest = max(rounded, key=lambda key: rounded[key]) + rounded[largest] += 100 - sum(rounded.values()) + return DynamicAllocationResult(rounded, True, coverage, cls._factors(metrics, scores)) + + @staticmethod + def _score( + metric: AssetClassMarketMetric, + peers: list[AssetClassMarketMetric], + target: Decimal, + drawdown_limit: Decimal, + liquidity_requirement: str, + ) -> Decimal: + returns = [item.trailing_120d_return_pct for item in peers] + liquidities = [item.average_daily_turnover_amount for item in peers] + return_fit = ( + min( + Decimal("1"), + max(Decimal(), metric.trailing_120d_return_pct * Decimal("2.1") / target), + ) + if target > 0 + else Decimal("0.5") + ) + return_rank = DynamicAllocationOptimizer._rank(metric.trailing_120d_return_pct, returns) + drawdown = abs(metric.max_drawdown_pct) + drawdown_fit = ( + min(Decimal("1"), drawdown_limit / drawdown) if drawdown > 0 else Decimal("1") + ) + liquidity_floor = { + "daily": Decimal("10000000"), + "within_7_days": Decimal("3000000"), + "within_30_days": Decimal("500000"), + "over_30_days": Decimal(), + }[liquidity_requirement] + liquidity_fit = ( + min(Decimal("1"), metric.average_daily_turnover_amount / liquidity_floor) + if liquidity_floor + else Decimal("0.5") + ) + liquidity_rank = DynamicAllocationOptimizer._rank( + metric.average_daily_turnover_amount, liquidities + ) + return ( + Decimal("0.45") * (return_fit + return_rank) / 2 + + Decimal("0.35") * drawdown_fit + + Decimal("0.20") * (liquidity_fit + liquidity_rank) / 2 + ) + + @staticmethod + def _rank(value: Decimal, peers: list[Decimal]) -> Decimal: + low, high = min(peers), max(peers) + return Decimal("0.5") if low == high else (value - low) / (high - low) + + @staticmethod + def _factors( + metrics: list[AssetClassMarketMetric], scores: dict[str, Decimal] + ) -> list[dict[str, object]]: + return [ + { + "asset_class": item.asset_class, + "trailing_120d_return_pct": str(item.trailing_120d_return_pct), + "max_drawdown_pct": str(item.max_drawdown_pct), + "average_daily_turnover_amount": str(item.average_daily_turnover_amount), + "product_count": item.product_count, + "composite_score": ( + str(scores[item.asset_class].quantize(Decimal("0.0001"))) + if item.asset_class in scores + else None + ), + } + for item in sorted(metrics, key=lambda value: value.asset_class) + ] diff --git a/app/service/suitability_service.py b/app/service/suitability_service.py index 984fa5c..66d3456 100644 --- a/app/service/suitability_service.py +++ b/app/service/suitability_service.py @@ -134,6 +134,10 @@ class SuitabilityService: profile = await self._load_authority_profile(request.customer_id) return self._decide(request, profile, current) + async def authority_for_customer(self, customer_id: int) -> RiskAuthorityProfile: + """Expose the read-only risk authority for other advisory services.""" + return await self._load_authority_profile(str(customer_id)) + 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)}) diff --git a/docs/21-投顾Agent迁移TODO.md b/docs/21-投顾Agent迁移TODO.md index 0159420..f692923 100644 --- a/docs/21-投顾Agent迁移TODO.md +++ b/docs/21-投顾Agent迁移TODO.md @@ -92,6 +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` 接口。 +真实历史回测计算、回测证据落库和查询接口仍未完成。 + +阶段九测试结果:专项测试 `6 passed`,全量单元/契约测试 `488 passed, 3 warnings`,Ruff +通过,MyPy(133 个源文件)通过。阶段九提交:待提交。 + ## 一、迁移准备 - [ ] 确认远程仓库可访问。(当前失败:连接 `47.106.207.27:3000` 被拒绝) @@ -293,28 +302,28 @@ python tools/audit_constraints.py ## 九、动态资产配置 -- [ ] 迁移 C1-C5 基础配置。 -- [ ] 迁移投资期限约束。 -- [ ] 迁移流动性约束。 -- [ ] 迁移最大回撤约束。 -- [ ] 迁移收益目标约束。 -- [ ] 迁移历史行情指标读取。 -- [ ] 迁移动态权重优化器。 -- [ ] 迁移流动性覆盖率。 +- [x] 迁移 C1-C5 基础配置。 +- [x] 迁移投资期限约束。 +- [x] 迁移流动性约束。 +- [x] 迁移最大回撤约束。 +- [x] 迁移收益目标约束。 +- [x] 迁移历史行情指标读取。 +- [x] 迁移动态权重优化器。 +- [x] 迁移流动性覆盖率。(以历史指标覆盖率和平均日成交额作为证据) - [ ] 迁移动态配置回测。 - [ ] 保存配置回测证据。 -- [ ] 实现数据覆盖不足时的降级状态。 -- [ ] 确认输出配置比例而不是买卖指令。 -- [ ] 完成资产配置提交 `advisor/asset-allocation`。 +- [x] 实现数据覆盖不足时的降级状态。(覆盖少于两个资产类别时返回静态配置) +- [x] 确认输出配置比例而不是买卖指令。 +- [x] 完成资产配置提交 `advisor/asset-allocation`。(专项 `6 passed`;全量 `488 passed`) 验收: -- [ ] 收益目标参与优化。 -- [ ] 最大回撤参与优化。 -- [ ] 流动性要求参与优化。 -- [ ] 投资期限参与优化。 -- [ ] 动态配置与静态配置可以对比。 -- [ ] 回测结果包含限制条件和数据覆盖率。 +- [x] 收益目标参与优化。 +- [x] 最大回撤参与优化。 +- [x] 流动性要求参与优化。 +- [x] 投资期限参与优化。 +- [x] 动态配置与静态配置可以对比。(输出 `dynamic` 和 `strategic_allocation`) +- [ ] 回测结果包含限制条件和数据覆盖率。(待真实回测模块) ## 十、产品推荐和审核发布 diff --git a/tests/unit/service/test_advisor_asset_allocation.py b/tests/unit/service/test_advisor_asset_allocation.py new file mode 100644 index 0000000..0111e41 --- /dev/null +++ b/tests/unit/service/test_advisor_asset_allocation.py @@ -0,0 +1,26 @@ +from app.service.agent.implementations.advisor import AdvisorAgent + + +def test_asset_allocation_description_is_analysis_only() -> None: + text = AdvisorAgent._describe_allocation( + { + "status": "ready", + "allocation": [ + {"asset_class": "bond_etf", "label": "债券类场内基金", "target_pct": 60}, + {"asset_class": "equity_etf", "label": "权益类场内基金", "target_pct": 40}, + ], + "optimization": {"dynamic": True}, + } + ) + + assert "债券类场内基金 60%" in text + assert "动态历史因子优化" in text + assert "不构成交易指令" in text + assert "下单" not in text + + +def test_asset_allocation_description_exposes_missing_prerequisites() -> None: + assert "风险画像" in AdvisorAgent._describe_allocation({"status": "profile_required"}) + assert "已确认的投资目标" in AdvisorAgent._describe_allocation( + {"status": "investment_goal_required"} + ) diff --git a/tests/unit/service/test_advisor_base_adapter.py b/tests/unit/service/test_advisor_base_adapter.py index a9d3388..855c05f 100644 --- a/tests/unit/service/test_advisor_base_adapter.py +++ b/tests/unit/service/test_advisor_base_adapter.py @@ -19,8 +19,9 @@ def test_advisor_is_registered_through_the_new_base_factory() -> None: assert isinstance(agent, AdvisorAgent) assert agent.definition == definition assert definition.allowed_tools == ( - "query_fund_quote", "query_investment_goal", "analyze_portfolio" + "query_fund_quote", "query_investment_goal", "analyze_portfolio", + "generate_asset_allocation", ) assert definition.supported_intents == ( - "fund_quote", "investment_goal", "portfolio_analysis" + "fund_quote", "investment_goal", "portfolio_analysis", "asset_allocation" ) diff --git a/tests/unit/service/test_dynamic_allocation_optimizer.py b/tests/unit/service/test_dynamic_allocation_optimizer.py new file mode 100644 index 0000000..c7cb931 --- /dev/null +++ b/tests/unit/service/test_dynamic_allocation_optimizer.py @@ -0,0 +1,65 @@ +from decimal import Decimal + +from app.service.asset_allocation_service import AssetAllocationService +from app.service.dynamic_allocation_optimizer import ( + AssetClassMarketMetric, + DynamicAllocationOptimizer, +) + + +def metric( + asset_class: str, return_pct: str, drawdown: str, turnover: str +) -> AssetClassMarketMetric: + return AssetClassMarketMetric( + asset_class=asset_class, + trailing_120d_return_pct=Decimal(return_pct), + max_drawdown_pct=Decimal(drawdown), + average_daily_turnover_amount=Decimal(turnover), + product_count=2, + ) + + +def test_strategic_weights_apply_horizon_liquidity_and_drawdown_constraints() -> None: + weights = AssetAllocationService._strategic_weights("C1", 6, "daily", Decimal("10")) + assert weights == {"cash_management_etf": 65, "bond_etf": 35, "equity_etf": 0} + + weights = AssetAllocationService._strategic_weights("C5", 72, "over_30_days", Decimal("30")) + assert weights == {"cash_management_etf": 5, "bond_etf": 10, "equity_etf": 85} + + +def test_optimizer_uses_return_drawdown_and_liquidity_evidence() -> None: + metrics = [ + metric("cash_management_etf", "2", "1", "20000000"), + metric("bond_etf", "6", "8", "5000000"), + metric("equity_etf", "12", "25", "1000000"), + ] + result = DynamicAllocationOptimizer.optimize( + {"cash_management_etf": 15, "bond_etf": 45, "equity_etf": 40}, + metrics, + return_target_lower_pct=Decimal("6"), + max_drawdown_pct=Decimal("15"), + liquidity_requirement="within_7_days", + ) + + assert result.dynamic is True + assert sum(result.weights.values()) == 100 + assert result.weights["equity_etf"] <= 40 + assert result.metric_coverage_pct == Decimal("100") + evidence = {item["asset_class"]: item for item in result.factors} + assert evidence["equity_etf"]["composite_score"] is not None + + +def test_optimizer_falls_back_to_static_weights_when_coverage_is_insufficient() -> None: + strategic = {"cash_management_etf": 30, "bond_etf": 50, "equity_etf": 20} + result = DynamicAllocationOptimizer.optimize( + strategic, + [metric("bond_etf", "6", "8", "5000000")], + return_target_lower_pct=Decimal("6"), + max_drawdown_pct=Decimal("15"), + liquidity_requirement="within_7_days", + ) + + assert result.dynamic is False + assert result.weights == strategic + assert result.metric_coverage_pct == Decimal("33.33333333333333333333333333") + assert result.factors[0]["composite_score"] is None