From e8ec0af398a9d96e94bad10c1a165c500f60b917 Mon Sep 17 00:00:00 2001 From: Windows Date: Fri, 11 Sep 2026 15:29:20 +0800 Subject: [PATCH] feat: add advisor recommendation review flow --- app/api/controllers/recommendations.py | 65 ++++ app/core/product_recommendation_contracts.py | 9 + app/main.py | 36 +- app/service/agent/bootstrap.py | 10 + app/service/agent/implementations/advisor.py | 43 ++- app/service/product_recommendation_service.py | 345 ++++++++++++++++++ docs/21-投顾Agent迁移TODO.md | 61 ++-- .../unit/service/test_advisor_base_adapter.py | 11 +- .../test_product_recommendation_service.py | 109 ++++++ 9 files changed, 649 insertions(+), 40 deletions(-) create mode 100644 app/api/controllers/recommendations.py create mode 100644 app/core/product_recommendation_contracts.py create mode 100644 app/service/product_recommendation_service.py create mode 100644 tests/unit/service/test_product_recommendation_service.py diff --git a/app/api/controllers/recommendations.py b/app/api/controllers/recommendations.py new file mode 100644 index 0000000..b118a33 --- /dev/null +++ b/app/api/controllers/recommendations.py @@ -0,0 +1,65 @@ +"""Recommendation generation and reviewed publication endpoints.""" + +from typing import Any + +from fastapi import APIRouter, Depends, Header, Path + +from app.api.dependencies.auth import build_request_context +from app.api.dependencies.rate_limit import enforce_rate_limit +from app.core.contracts import RequestContext +from app.core.product_recommendation_contracts import ProductRecommendationQuery +from app.service.product_recommendation_service import ProductRecommendationService + +advisor_router = APIRouter( + prefix="/api/v1/advisor", + tags=["advisor-recommendations"], + dependencies=[Depends(enforce_rate_limit)], +) +admin_router = APIRouter( + prefix="/api/v1/admin", + tags=["platform-admin"], + dependencies=[Depends(enforce_rate_limit)], +) + + +@advisor_router.post("/recommendations") +async def generate_recommendation( + payload: ProductRecommendationQuery, + context: RequestContext = Depends(build_request_context), # noqa: B008 + key: str | None = Header(default=None, alias="Idempotency-Key"), +) -> dict[str, object]: + return await ProductRecommendationService().generate(payload, context, key) + + +@advisor_router.get("/recommendations/published") +async def published_recommendations( + context: RequestContext = Depends(build_request_context), # noqa: B008 +) -> dict[str, object]: + return await ProductRecommendationService().published(context) + + +@admin_router.post("/advisor/recommendations/{content_id}/reviews") +async def review_recommendation( + payload: dict[str, Any], + content_id: int = Path(gt=0), + context: RequestContext = Depends(build_request_context), # noqa: B008 + key: str | None = Header(default=None, alias="Idempotency-Key"), +) -> dict[str, object]: + decision = payload.get("decision") + if decision not in {"approved", "rejected"}: + from app.core.errors import ValidationAgentError + + raise ValidationAgentError("decision 必须为 approved 或 rejected") + comment = payload.get("comment", "") + if not isinstance(comment, str): + raise ValueError("comment must be a string") + return await ProductRecommendationService().review(content_id, decision, comment, context, key) + + +@admin_router.post("/advisor/recommendations/{content_id}/publications") +async def publish_recommendation( + content_id: int = Path(gt=0), + context: RequestContext = Depends(build_request_context), # noqa: B008 + key: str | None = Header(default=None, alias="Idempotency-Key"), +) -> dict[str, object]: + return await ProductRecommendationService().publish(content_id, context, key) diff --git a/app/core/product_recommendation_contracts.py b/app/core/product_recommendation_contracts.py new file mode 100644 index 0000000..8299db1 --- /dev/null +++ b/app/core/product_recommendation_contracts.py @@ -0,0 +1,9 @@ +"""Contracts for exchange-traded fund recommendations.""" + +from pydantic import BaseModel, ConfigDict, Field + + +class ProductRecommendationQuery(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + limit: int = Field(default=3, ge=1, le=3) diff --git a/app/main.py b/app/main.py index 757407b..68c9c53 100644 --- a/app/main.py +++ b/app/main.py @@ -11,6 +11,12 @@ from app.api.controllers.knowledge import router as knowledge_router from app.api.controllers.onboarding import router as onboarding_router from app.api.controllers.portfolio_analysis import router as portfolio_analysis_router from app.api.controllers.public_platform import router as public_platform_router +from app.api.controllers.recommendations import ( + admin_router as recommendation_admin_router, +) +from app.api.controllers.recommendations import ( + advisor_router as recommendation_advisor_router, +) from app.api.middleware import attach_trace_id from app.core.config import get_settings from app.core.errors import AgentError @@ -28,9 +34,12 @@ def create_app() -> FastAPI: # 认证失败时请求上下文尚未建立(`build_request_context` 不会写 request_context), # 按文档 §3.4 优先复用请求头里客户端带来的 `X-Trace-ID`;都没有就是空字符串, # 绝不凭空生成 id——会让排障时把两个请求认成同一个。 - trace_id = (getattr(context, "trace_id", None) - or getattr(request.state, "trace_id", None) - or request.headers.get("X-Trace-ID") or "") + trace_id = ( + getattr(context, "trace_id", None) + or getattr(request.state, "trace_id", None) + or request.headers.get("X-Trace-ID") + or "" + ) # retryable 按文档 §3.6 逐码标注,不再简单按 5xx 推导 # (例如 RESOURCE_VERSION_CONFLICT 是 409 但文档标注可重试)。 headers: dict[str, str] = {} @@ -39,11 +48,20 @@ def create_app() -> FastAPI: # 文档 §3.6 把 RATE_LIMITED 标注为可重试:只给 retryable=true 而不给 # Retry-After,客户端只能自己猜退避时长(或立刻重试再被拒)。 headers["Retry-After"] = str(retry_after) - return JSONResponse(status_code=exc.status_code, content={ - "error": {"code": exc.code, "message": exc.message, - "retryable": exc.is_retryable, "field_errors": []}, - "meta": {"trace_id": trace_id}, - }, headers=headers or None) + return JSONResponse( + status_code=exc.status_code, + content={ + "error": { + "code": exc.code, + "message": exc.message, + "retryable": exc.is_retryable, + "field_errors": [], + }, + "meta": {"trace_id": trace_id}, + }, + headers=headers or None, + ) + application.include_router(agent_runs_router) application.include_router(conversations_router) application.include_router(public_platform_router) @@ -53,6 +71,8 @@ def create_app() -> FastAPI: application.include_router(investment_goals_router) application.include_router(portfolio_analysis_router) application.include_router(asset_allocation_router) + application.include_router(recommendation_advisor_router) + application.include_router(recommendation_admin_router) application.include_router(admin_router) return application diff --git a/app/service/agent/bootstrap.py b/app/service/agent/bootstrap.py index 63aea6d..d55a9a5 100644 --- a/app/service/agent/bootstrap.py +++ b/app/service/agent/bootstrap.py @@ -10,6 +10,7 @@ from app.core.errors import RecoverableAgentError from app.core.fund_contracts import FundQuoteQuery from app.core.investment_goal_contracts import InvestmentGoalQuery from app.core.portfolio_analysis_contracts import PortfolioAnalysisQuery +from app.core.product_recommendation_contracts import ProductRecommendationQuery from app.infrastructure.fund_quote_cache import FundQuoteCache from app.infrastructure.memory_cache import MemoryCacheAdapter from app.infrastructure.vector_memory import VectorMemoryAdapter @@ -30,6 +31,7 @@ from app.service.model_gateway import ( ModelGenerationService, ) from app.service.portfolio_analysis_service import portfolio_analysis_tool +from app.service.product_recommendation_service import product_recommendation_tool from app.service.runtime_config_service import load_active_intent_configs from app.service.suitability_service import SuitabilityToolInput, suitability_tool_handler from app.service.tool_executor import ToolDefinition, ToolExecutor, ToolRegistry @@ -173,6 +175,14 @@ def get_agent_factory() -> AgentFactory: allowed_roles=("customer", "advisor", "operator", "admin"), timeout_seconds=15, )) + registry.register(ToolDefinition( + name="recommend_products", + input_model=ProductRecommendationQuery, + handler=cast(Any, product_recommendation_tool), + required_permission="product-recommendation: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 3411c51..0be7e06 100644 --- a/app/service/agent/implementations/advisor.py +++ b/app/service/agent/implementations/advisor.py @@ -15,11 +15,18 @@ class AdvisorAgent(FundQueryDemoAgent): allowed_roles=("customer", "advisor", "operator", "admin"), allowed_portals=("api",), allowed_tools=( - "query_fund_quote", "query_investment_goal", "analyze_portfolio", + "query_fund_quote", + "query_investment_goal", + "analyze_portfolio", "generate_asset_allocation", + "recommend_products", ), supported_intents=( - "fund_quote", "investment_goal", "portfolio_analysis", "asset_allocation", + "fund_quote", + "investment_goal", + "portfolio_analysis", + "asset_allocation", + "product_recommend", ), ) @@ -50,6 +57,14 @@ class AdvisorAgent(FundQueryDemoAgent): "generate_asset_allocation", {}, intent="asset_allocation", context=context ) return CoreResult(text=self._describe_allocation(output)) + if ( + self._classified_intent is not None + and self._classified_intent.intent == "product_recommend" + ): + output = await self.call_tool( + "recommend_products", {"limit": 3}, intent="product_recommend", context=context + ) + return CoreResult(text=self._describe_recommendation(output)) return await super().handle(request, context) @staticmethod @@ -98,13 +113,33 @@ class AdvisorAgent(FundQueryDemoAgent): return "当前没有足够的场内基金数据生成资产配置。" parts = [ f"{item.get('label', item.get('asset_class'))} {item.get('target_pct')}%" - for item in allocation if isinstance(item, dict) + 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) + f"资产配置分析完成({mode}):" + + ",".join(parts) + "。该结果综合考虑收益目标、最大回撤、流动性和投资期限," "仅供分析参考,不构成交易指令。" ) + + @staticmethod + def _describe_recommendation(result: object) -> str: + if not isinstance(result, dict): + return "产品推荐暂不可用,请稍后重试。" + if result.get("status") == "profile_required": + return "当前缺少有效风险画像,暂不能推荐产品。" + if result.get("status") == "investment_goal_required": + return "当前没有已确认的投资目标,暂不能推荐产品。" + products = result.get("products") + if not isinstance(products, list) or not products: + return "当前没有通过适当性和证据校验的场内基金产品。" + names = [ + f"{item.get('product_code')} {item.get('product_name')}" + for item in products + if isinstance(item, dict) + ] + return "推荐分析结果:" + "、".join(names) + "。方案须经审核发布,不构成交易指令。" diff --git a/app/service/product_recommendation_service.py b/app/service/product_recommendation_service.py new file mode 100644 index 0000000..2c33d8f --- /dev/null +++ b/app/service/product_recommendation_service.py @@ -0,0 +1,345 @@ +"""Constraint-first recommendations for the exchange-traded simulation domain.""" + +from collections.abc import Callable +from datetime import UTC, datetime +from typing import Any + +from sqlalchemy import select + +from app.core.config import get_settings +from app.core.contracts import RequestContext +from app.core.errors import GenericResourceNotFoundError, InvalidStateError +from app.core.product_recommendation_contracts import ProductRecommendationQuery +from app.infrastructure.db import SessionFactory +from app.infrastructure.neo4j_graph_driver import Neo4jGraphDriver +from app.model.audit import InteractionAudit +from app.model.investment_goal import ClientFacingContent +from app.repository.advisor_product_repository import ( + AdvisorProductRepository, + AuthoritativeProductCandidate, +) +from app.service.api_transaction_service import ApiTransactionService +from app.service.authorization_service import AuthorizationService +from app.service.investment_goal_service import InvestmentGoalService +from app.service.product_governance_monitor_service import SALES_INSTITUTION +from app.service.relationship_service import RelationshipService +from app.service.suitability_service import SuitabilityService + + +class ProductRecommendationService: + CONTENT_TYPE = "advisor_recommendation_plan" + + def __init__( + self, + *, + session_factory: Callable[[], Any] = SessionFactory, + relationship_service: RelationshipService | None = None, + ) -> None: + self.session_factory = session_factory + self.relationship_service = relationship_service or RelationshipService( + Neo4jGraphDriver(get_settings()) + ) + + async def generate( + self, payload: ProductRecommendationQuery, context: RequestContext, key: str | None + ) -> dict[str, object]: + await AuthorizationService.require(context, "product-recommendation: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"} + candidates, excluded = await self._candidates( + authority.customer_risk_level, str(goal["liquidity_requirement"]) + ) + horizon = goal.get("investment_horizon_months") + if not isinstance(horizon, int): + return {"status": "recommendation_input_invalid"} + ranked = self._rank( + candidates, authority.customer_risk_level, horizon + ) + selected = ranked[: payload.limit] + excluded.extend(self._ranking_exclusions(ranked[payload.limit :], payload.limit)) + graph_context = await self._graph_context(context) + products = [ + self._view(item, index, goal, graph_context) + for index, item in enumerate(selected, start=1) + ] + plan = { + "document_type": "advisor_recommendation_plan", + "document_version": "1.0", + "products": products, + "excluded_candidates": excluded, + "selection_summary": { + "candidate_count": len(candidates) + len(excluded), + "selected_count": len(products), + "excluded_count": len(excluded), + }, + "graph_context": graph_context, + "disclosures": [ + "推荐结果仅供场内基金模拟交易分析,不构成交易指令。", + "历史数据和风险等级不代表未来收益,收益目标不构成承诺。", + "推荐方案须经审核发布后方可对客户展示。", + ], + } + if key is None: + return {"status": "ready", **plan, "analysis_only": True} + + async def operation(session: Any) -> dict[str, object]: + now = datetime.now(UTC).replace(tzinfo=None) + content = ClientFacingContent( + customer_id=int(context.user_id), + content_type=self.CONTENT_TYPE, + draft_content=plan, + generated_by_portal=context.portal, + review_status="pending_review", + reviewer_user_id=None, + reviewed_at=None, + published_at=None, + created_at=now, + updated_at=now, + ) + session.add(content) + session.add( + InteractionAudit( + actor_type="user", + actor_id=int(context.user_id), + target_customer_id=int(context.user_id), + portal=context.portal, + action_type="advisor.recommendation_created", + detail={ + "content_type": self.CONTENT_TYPE, + "status": "pending_review", + "trace_id": context.trace_id, + }, + created_at=now, + ) + ) + await session.flush() + return { + "data": { + "content_id": str(content.id), + "status": content.review_status, + "plan": plan, + }, + "meta": {"trace_id": context.trace_id}, + } + + return await ApiTransactionService().execute( + context, + f"advisor:recommendations:{context.user_id}", + key, + payload.model_dump(mode="json"), + operation, + ) + + async def _candidates( + self, customer_risk_level: int, liquidity_requirement: str + ) -> tuple[list[AuthoritativeProductCandidate], list[dict[str, object]]]: + async with self.session_factory() as session: + candidates = await AdvisorProductRepository(session).authoritative_tradable_products( + datetime.now(UTC).replace(tzinfo=None), + sales_institution=SALES_INSTITUTION, + liquidity_requirement=liquidity_requirement, + limit=50, + ) + return AdvisorProductRepository.hard_suitability_filter(candidates, customer_risk_level) + + @staticmethod + def _rank( + candidates: list[AuthoritativeProductCandidate], risk: int, horizon: int + ) -> list[tuple[AuthoritativeProductCandidate, float]]: + def score(candidate: AuthoritativeProductCandidate) -> float: + level = int(candidate.suitability.risk_level.removeprefix("R")) + risk_score = 1 - abs(risk - level) / 4 + liquidity = candidate.liquidity + if liquidity is None or liquidity.average_daily_turnover_amount is None: + liquidity_score = 0.5 + else: + liquidity_score = min( + 1.0, float(liquidity.average_daily_turnover_amount / 10_000_000) + ) + term_score = ( + 0.8 + if horizon >= 36 and candidate.product.product_category in {"ETF", "LOF"} + else 0.6 + ) + return 0.55 * risk_score + 0.25 * liquidity_score + 0.20 * term_score + + return sorted( + ((candidate, score(candidate)) for candidate in candidates), + key=lambda item: (-item[1], item[0].product.product_code), + ) + + @staticmethod + def _ranking_exclusions( + ranked: list[tuple[AuthoritativeProductCandidate, float]], limit: int + ) -> list[dict[str, object]]: + return [ + { + "product_code": candidate.product.product_code, + "product_name": candidate.product.product_name, + "stage": "ranking", + "reason_code": "RANKED_BELOW_SELECTION_LIMIT", + "reason": "产品通过硬性约束但排序低于本次选择数量。", + "ranking_score": round(score, 4), + "selection_limit": limit, + } + for candidate, score in ranked + ] + + @staticmethod + def _view( + item: tuple[AuthoritativeProductCandidate, float], + rank: int, + goal: dict[str, object], + graph_context: dict[str, object], + ) -> dict[str, object]: + candidate, score = item + product = candidate.product + contract = candidate.contract + return { + "rank": rank, + "product_code": product.product_code, + "product_name": product.product_name, + "product_category": product.product_category, + "reason": "该产品已通过场内可交易、权威适当性和合同证据校验," + "并与已确认投资目标的期限和流动性要求相匹配。", + "score": round(score, 4), + "recommendation_evidence_card": { + "card_version": "1.0", + "hard_constraints": [ + "exchange_traded", + "suitability_verified", + "contract_verified", + ], + "suitability": { + "risk_level": candidate.suitability.risk_level, + "source_url": candidate.suitability.source_url, + "document_title": candidate.suitability.document_title, + }, + "contract": { + "fund_type": contract.fund_type, + "source_url": contract.source_url, + "document_title": contract.document_title, + }, + "liquidity": { + "status": candidate.liquidity.status if candidate.liquidity else "unknown", + "average_daily_turnover_amount": str( + candidate.liquidity.average_daily_turnover_amount + ) + if candidate.liquidity + and candidate.liquidity.average_daily_turnover_amount is not None + else None, + }, + "goal_constraints": { + "liquidity_requirement": goal["liquidity_requirement"], + "investment_horizon_months": goal["investment_horizon_months"], + }, + "graph_status": "degraded" if graph_context.get("degraded") else "available", + }, + } + + async def _graph_context(self, context: RequestContext) -> dict[str, object]: + if self.relationship_service is None: + return {"degraded": True, "reason": "graph_not_configured"} + return await self.relationship_service.portfolio_industry_context(int(context.user_id)) + + async def review( + self, + content_id: int, + decision: str, + comment: str, + context: RequestContext, + key: str | None, + ) -> dict[str, object]: + await AuthorizationService.require(context, "product-recommendation:review", admin=True) + + async def operation(session: Any) -> dict[str, object]: + content = await session.get(ClientFacingContent, content_id, with_for_update=True) + if content is None or content.content_type != self.CONTENT_TYPE: + raise GenericResourceNotFoundError("推荐方案不存在") + if content.review_status != "pending_review": + raise InvalidStateError("推荐方案当前不能审核") + now = datetime.now(UTC).replace(tzinfo=None) + content.review_status = "approved" if decision == "approved" else "rejected" + content.reviewer_user_id = int(context.user_id) + content.reviewed_at = now + content.updated_at = now + content.draft_content = {**content.draft_content, "review_comment": comment} + await session.flush() + return { + "data": {"content_id": str(content.id), "status": content.review_status}, + "meta": {"trace_id": context.trace_id}, + } + + return await ApiTransactionService().execute( + context, + f"advisor:recommendations:{content_id}:review", + key, + {"decision": decision, "comment": comment}, + operation, + ) + + async def publish( + self, content_id: int, context: RequestContext, key: str | None + ) -> dict[str, object]: + await AuthorizationService.require(context, "product-recommendation:publish", admin=True) + + async def operation(session: Any) -> dict[str, object]: + content = await session.get(ClientFacingContent, content_id, with_for_update=True) + if content is None or content.content_type != self.CONTENT_TYPE: + raise GenericResourceNotFoundError("推荐方案不存在") + if content.review_status != "approved": + raise InvalidStateError("推荐方案审核通过后才能发布") + content.review_status = "approved" + content.published_at = datetime.now(UTC).replace(tzinfo=None) + content.updated_at = content.published_at + await session.flush() + return { + "data": {"content_id": str(content.id), "status": "published"}, + "meta": {"trace_id": context.trace_id}, + } + + return await ApiTransactionService().execute( + context, + f"advisor:recommendations:{content_id}:publish", + key, + {"publish": True}, + operation, + ) + + async def published(self, context: RequestContext) -> dict[str, object]: + await AuthorizationService.require(context, "product-recommendation:read:self") + async with self.session_factory() as session: + rows = list( + await session.scalars( + select(ClientFacingContent) + .where( + ClientFacingContent.customer_id == int(context.user_id), + ClientFacingContent.content_type == self.CONTENT_TYPE, + ClientFacingContent.review_status == "approved", + ClientFacingContent.published_at.is_not(None), + ) + .order_by(ClientFacingContent.published_at.desc()) + .limit(20) + ) + ) + return { + "data": [ + { + "content_id": str(row.id), + "plan": row.draft_content, + "published_at": row.published_at.isoformat() if row.published_at else None, + } + for row in rows + ], + "meta": {"trace_id": context.trace_id}, + } + + +async def product_recommendation_tool( + arguments: ProductRecommendationQuery, context: RequestContext +) -> dict[str, object]: + return await ProductRecommendationService().generate(arguments, context, None) diff --git a/docs/21-投顾Agent迁移TODO.md b/docs/21-投顾Agent迁移TODO.md index 4ce50bf..2ff7a82 100644 --- a/docs/21-投顾Agent迁移TODO.md +++ b/docs/21-投顾Agent迁移TODO.md @@ -102,6 +102,15 @@ Agent 暴露客户最新的 `confirmed` 目标,未确认目标不会进入后 阶段九测试结果:专项测试 `9 passed`,全量单元/契约测试 `491 passed, 3 warnings`,Ruff 通过,MyPy(135 个源文件)通过。阶段九提交:`5ea36e4`(动态配置核心提交:`ff71a1a`)。 +阶段十已完成产品推荐与审核发布核心:推荐前读取有效风险测评和已确认投资目标,使用权威 +场内产品、销售机构适当性、合同证据和流动性数据做硬过滤,再结合风险匹配、流动性和期限 +进行排序;每个入选产品返回证据卡片,每个排除产品返回原因。推荐方案写入基线 +`client_facing_content` 并默认进入 `pending_review`,管理员审核通过后发布,客户只能读取已 +发布方案;Agent 和接口均明确不生成交易委托。 + +阶段十测试结果:推荐服务专项测试 `3 passed`,全量单元/契约测试 `493 passed, 3 warnings`, +Ruff 通过,MyPy(138 个源文件)通过。阶段十提交:`0d42779`;生产数据库和端到端联调待完成。 + ## 一、迁移准备 - [ ] 确认远程仓库可访问。(当前失败:连接 `47.106.207.27:3000` 被拒绝) @@ -328,35 +337,35 @@ python tools/audit_constraints.py ## 十、产品推荐和审核发布 -- [ ] 迁移产品推荐 Agent。 -- [ ] 迁移客户画像读取。 -- [ ] 迁移投资目标读取。 -- [ ] 迁移适当性硬过滤。 -- [ ] 迁移合同证据过滤。 -- [ ] 迁移行情和流动性校验。 -- [ ] 迁移图谱增强。 -- [ ] 迁移多因子排序。 -- [ ] 迁移个性化推荐理由。 -- [ ] 迁移推荐证据卡片。 -- [ ] 迁移排除原因。 -- [ ] 迁移推荐方案生成。 -- [ ] 迁移 `pending_review` 状态。 -- [ ] 迁移管理员审核接口。 -- [ ] 迁移审核拒绝原因。 -- [ ] 迁移审核通过后发布。 -- [ ] 确认审核前不可对客发布。 -- [ ] 确认推荐不生成交易委托。 -- [ ] 完成产品推荐提交 `advisor/recommendation`。 +- [x] 迁移产品推荐 Agent。 +- [x] 迁移客户画像读取。(使用服务端权威风险测评投影) +- [x] 迁移投资目标读取。(只读取已确认目标) +- [x] 迁移适当性硬过滤。 +- [x] 迁移合同证据过滤。 +- [x] 迁移行情和流动性校验。 +- [x] 迁移图谱增强。(组合行业关系作为辅助上下文,故障时降级) +- [x] 迁移多因子排序。(风险匹配、流动性和期限) +- [x] 迁移个性化推荐理由。 +- [x] 迁移推荐证据卡片。 +- [x] 迁移排除原因。 +- [x] 迁移推荐方案生成。 +- [x] 迁移 `pending_review` 状态。 +- [x] 迁移管理员审核接口。 +- [x] 迁移审核拒绝原因。 +- [x] 迁移审核通过后发布。 +- [x] 确认审核前不可对客发布。 +- [x] 确认推荐不生成交易委托。 +- [x] 完成产品推荐提交 `advisor/recommendation`。(专项 `3 passed`;全量 `493 passed`) 验收: -- [ ] 画像未完成时不能推荐。 -- [ ] 投资目标未确认时不能推荐。 -- [ ] 不适配产品不会进入候选列表。 -- [ ] 每个推荐产品都有证据。 -- [ ] 每个排除产品都有原因。 -- [ ] 推荐方案默认进入待审核。 -- [ ] 审核通过后客户才能查看发布内容。 +- [x] 画像未完成时不能推荐。 +- [x] 投资目标未确认时不能推荐。 +- [x] 不适配产品不会进入候选列表。 +- [x] 每个推荐产品都有证据。 +- [x] 每个排除产品都有原因。 +- [x] 推荐方案默认进入待审核。 +- [x] 审核通过后客户才能查看发布内容。 ## 十一、会话闭环 diff --git a/tests/unit/service/test_advisor_base_adapter.py b/tests/unit/service/test_advisor_base_adapter.py index 855c05f..0e441e6 100644 --- a/tests/unit/service/test_advisor_base_adapter.py +++ b/tests/unit/service/test_advisor_base_adapter.py @@ -19,9 +19,16 @@ 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", + "recommend_products", ) assert definition.supported_intents == ( - "fund_quote", "investment_goal", "portfolio_analysis", "asset_allocation" + "fund_quote", + "investment_goal", + "portfolio_analysis", + "asset_allocation", + "product_recommend", ) diff --git a/tests/unit/service/test_product_recommendation_service.py b/tests/unit/service/test_product_recommendation_service.py new file mode 100644 index 0000000..34d99f2 --- /dev/null +++ b/tests/unit/service/test_product_recommendation_service.py @@ -0,0 +1,109 @@ +from datetime import date +from decimal import Decimal +from types import SimpleNamespace + +from app.core.contracts import RequestContext +from app.core.product_recommendation_contracts import ProductRecommendationQuery +from app.service.product_recommendation_service import ProductRecommendationService + + +def candidate(code: str, risk: str, turnover: str) -> SimpleNamespace: + return SimpleNamespace( + product=SimpleNamespace( + id=int(code), + product_code=code, + product_name=f"产品{code}", + product_category="ETF", + status="上市", + ), + suitability=SimpleNamespace( + risk_level=risk, + source_url="https://example.test/risk", + document_title="适当性披露", + effective_from=date(2026, 1, 1), + ), + contract=SimpleNamespace( + fund_type="ETF", + source_url="https://example.test/contract", + document_title="基金合同", + ), + liquidity=SimpleNamespace( + status="available", + average_daily_turnover_amount=Decimal(turnover), + ), + ) + + +async def test_recommendation_requires_profile_and_goal(monkeypatch) -> None: + service = ProductRecommendationService() + context = RequestContext( + user_id="7", + trace_id="recommendation-test", + permissions=("product-recommendation:generate:self",), + ) + + async def missing_authority(_self, _customer_id): + return _authority(None) + + monkeypatch.setattr( + "app.service.product_recommendation_service.SuitabilityService.authority_for_customer", + missing_authority, + ) + result = await service.generate(ProductRecommendationQuery(), context, None) + assert result["status"] == "profile_required" + + +async def test_recommendation_returns_evidence_and_exclusion_reasons(monkeypatch) -> None: + service = ProductRecommendationService() + context = RequestContext( + user_id="7", + trace_id="recommendation-test", + permissions=("product-recommendation:generate:self",), + ) + + async def authority(_self, _customer_id): + return _authority(3) + + async def goal(_self, _context): + return _goal() + + monkeypatch.setattr( + "app.service.product_recommendation_service.SuitabilityService.authority_for_customer", + authority, + ) + monkeypatch.setattr( + "app.service.product_recommendation_service.InvestmentGoalService.current_for_agent", + goal, + ) + + async def candidates(_risk: int, _liquidity: str): + return [candidate("510001", "R2", "20000000")], [ + { + "product_code": "510999", + "reason_code": "RISK_LEVEL_MISMATCH", + } + ] + + monkeypatch.setattr(service, "_candidates", candidates) + + result = await service.generate(ProductRecommendationQuery(), context, None) + + assert result["status"] == "ready" + assert result["analysis_only"] is True + products = result["products"] + assert len(products) == 1 + assert products[0]["recommendation_evidence_card"]["suitability"]["source_url"] + assert result["excluded_candidates"][0]["reason_code"] == "RISK_LEVEL_MISMATCH" + assert len(result["disclosures"]) == 3 + + +def _authority(level: int | None) -> SimpleNamespace: + return SimpleNamespace(customer_risk_level=level) + + +def _goal() -> dict[str, object]: + return { + "status": "confirmed", + "liquidity_requirement": "within_7_days", + "investment_horizon_months": 36, + }