diff --git a/.env.example b/.env.example index 753f827..b004d42 100644 --- a/.env.example +++ b/.env.example @@ -43,6 +43,10 @@ WORKER_RETRY_LIMIT=3 WORKER_POLL_SECONDS=1 SSE_CHUNK_CHARACTERS=256 +# 投顾灰度:开启后仅白名单客户和管理员可访问投顾业务;多个客户 ID 用逗号分隔。 +ADVISOR_ROLLOUT_ENABLED=false +ADVISOR_ROLLOUT_CUSTOMER_IDS= + # 限流:按“用户 + 方法 + 路由”在窗口内计数,超限返回 429 RATE_LIMITED + Retry-After。 # 默认 600/60s(10 QPS)是为本机验收与集成测试留足余量的宽松值;生产按真实容量收紧。 # Redis 不可用时自动降级为放行(限流是保护措施,不能因后端故障拒绝正常请求)。 diff --git a/app/api/controllers/asset_allocation.py b/app/api/controllers/asset_allocation.py index f6ef79d..cae031c 100644 --- a/app/api/controllers/asset_allocation.py +++ b/app/api/controllers/asset_allocation.py @@ -6,11 +6,12 @@ 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.advisor_rollout_service import enforce_advisor_rollout from app.service.asset_allocation_service import AssetAllocationService router = APIRouter( prefix="/api/v1/advisor", tags=["advisor-asset-allocation"], - dependencies=[Depends(enforce_rate_limit)], + dependencies=[Depends(enforce_rate_limit), Depends(enforce_advisor_rollout)], ) diff --git a/app/api/controllers/investment_goals.py b/app/api/controllers/investment_goals.py index 47538fa..84bc83f 100644 --- a/app/api/controllers/investment_goals.py +++ b/app/api/controllers/investment_goals.py @@ -13,11 +13,12 @@ from app.api.schemas.investment_goals import ( InvestmentGoalCreate, ) from app.core.contracts import RequestContext +from app.service.advisor_rollout_service import enforce_advisor_rollout from app.service.investment_goal_service import InvestmentGoalService router = APIRouter( prefix="/api/v1/advisor", tags=["advisor-investment-goals"], - dependencies=[Depends(enforce_rate_limit)], + dependencies=[Depends(enforce_rate_limit), Depends(enforce_advisor_rollout)], ) diff --git a/app/api/controllers/portfolio_analysis.py b/app/api/controllers/portfolio_analysis.py index 579c876..6b2c578 100644 --- a/app/api/controllers/portfolio_analysis.py +++ b/app/api/controllers/portfolio_analysis.py @@ -6,11 +6,12 @@ from app.api.dependencies.auth import build_request_context from app.api.dependencies.rate_limit import enforce_rate_limit from app.api.schemas.portfolio_analysis import PortfolioAnalysisQuery from app.core.contracts import RequestContext +from app.service.advisor_rollout_service import enforce_advisor_rollout from app.service.portfolio_analysis_service import PortfolioAnalysisService router = APIRouter( prefix="/api/v1/advisor", tags=["advisor-portfolio-analysis"], - dependencies=[Depends(enforce_rate_limit)], + dependencies=[Depends(enforce_rate_limit), Depends(enforce_advisor_rollout)], ) diff --git a/app/api/controllers/recommendations.py b/app/api/controllers/recommendations.py index 804ce01..90d76e8 100644 --- a/app/api/controllers/recommendations.py +++ b/app/api/controllers/recommendations.py @@ -8,12 +8,13 @@ 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.advisor_rollout_service import enforce_advisor_rollout from app.service.product_recommendation_service import ProductRecommendationService advisor_router = APIRouter( prefix="/api/v1/advisor", tags=["advisor-recommendations"], - dependencies=[Depends(enforce_rate_limit)], + dependencies=[Depends(enforce_rate_limit), Depends(enforce_advisor_rollout)], ) admin_router = APIRouter( prefix="/api/v1/admin", @@ -40,7 +41,10 @@ async def published_recommendations( return await ProductRecommendationService().published(context) -@admin_router.post("/advisor/recommendations/{content_id}/reviews") +@admin_router.post( + "/advisor/recommendations/{content_id}/reviews", + dependencies=[Depends(enforce_advisor_rollout)], +) async def review_recommendation( payload: dict[str, Any], content_id: int = Path(gt=0), @@ -58,7 +62,10 @@ async def review_recommendation( return await ProductRecommendationService().review(content_id, decision, comment, context, key) -@admin_router.post("/advisor/recommendations/{content_id}/publications") +@admin_router.post( + "/advisor/recommendations/{content_id}/publications", + dependencies=[Depends(enforce_advisor_rollout)], +) async def publish_recommendation( content_id: int = Path(gt=0), context: RequestContext = Depends(build_request_context), # noqa: B008 diff --git a/app/core/config.py b/app/core/config.py index 8973c83..98e35d1 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -58,6 +58,10 @@ class Settings(BaseSettings): worker_lease_seconds: int = Field(default=60, gt=0) worker_retry_limit: int = Field(default=3, ge=0) worker_poll_seconds: float = Field(default=1, gt=0) + # 投顾灰度默认关闭,只有显式开启并配置客户白名单后才限制客户流量。 + # 管理员角色始终可进入,便于审核、发布和故障处置。 + advisor_rollout_enabled: bool = False + advisor_rollout_customer_ids: str = "" model_config = SettingsConfigDict(env_file=".env", env_file_encoding="utf-8", extra="ignore") diff --git a/app/service/advisor_rollout_service.py b/app/service/advisor_rollout_service.py new file mode 100644 index 0000000..ae6df27 --- /dev/null +++ b/app/service/advisor_rollout_service.py @@ -0,0 +1,63 @@ +"""投顾灰度发布闸门。 + +灰度配置使用环境变量,避免为发布控制改动冻结的数据库基线。白名单只允许客户 +身份进入投顾业务;管理员保留审核、发布和故障处置权限。未命中时统一复用权限拒绝 +错误,不向客户暴露灰度配置细节。 +""" + +from datetime import UTC, datetime + +from fastapi import Depends +from sqlalchemy.ext.asyncio import AsyncSession + +from app.api.dependencies.auth import build_request_context +from app.api.dependencies.database import get_session +from app.core.config import get_settings +from app.core.contracts import RequestContext +from app.core.errors import ForbiddenAgentError +from app.model.audit import InteractionAudit + + +class AdvisorRolloutService: + ADMIN_ROLES = frozenset({"admin", "super_admin"}) + + def __init__(self, session: AsyncSession | None = None) -> None: + self.session = session + + @staticmethod + def customer_ids(raw: str) -> frozenset[str]: + return frozenset(item.strip() for item in raw.split(",") if item.strip()) + + @classmethod + def is_allowed(cls, context: RequestContext) -> bool: + settings = get_settings() + if not settings.advisor_rollout_enabled: + return True + if cls.ADMIN_ROLES.intersection(context.roles): + return True + allowed = cls.customer_ids(settings.advisor_rollout_customer_ids) + identities = {context.user_id, *context.customer_ids} + return bool(allowed.intersection(identities)) + + async def ensure_allowed(self, context: RequestContext) -> None: + if self.is_allowed(context): + return + if self.session is not None: + self.session.add(InteractionAudit( + actor_type="user", + actor_id=int(context.user_id), + portal=context.portal, + action_type="advisor.rollout_denied", + detail={"trace_id": context.trace_id, "reason": "not_in_rollout"}, + created_at=datetime.now(UTC).replace(tzinfo=None), + )) + await self.session.commit() + raise ForbiddenAgentError("当前投顾功能尚未对该账号开放") + + +async def enforce_advisor_rollout( + context: RequestContext = Depends(build_request_context), # noqa: B008 + session: AsyncSession = Depends(get_session), # noqa: B008 +) -> None: + """FastAPI 依赖:保护投顾业务路由,认证依赖先完成身份解析。""" + await AdvisorRolloutService(session).ensure_allowed(context) diff --git a/app/service/agent_run_application_service.py b/app/service/agent_run_application_service.py index a3cb464..8f4dce0 100644 --- a/app/service/agent_run_application_service.py +++ b/app/service/agent_run_application_service.py @@ -19,6 +19,7 @@ from app.model.conversation import ConversationMessage from app.model.platform import AgentRun, RequestIdempotency from app.model.session import ConversationSession from app.repository.outbox_repository import OutboxRepository +from app.service.advisor_rollout_service import AdvisorRolloutService from app.service.agent.bootstrap import get_agent_factory from app.service.agent.factory import AgentFactory @@ -36,6 +37,8 @@ class AgentRunApplicationService: self.factory = factory if factory is not None else get_agent_factory() async def accept(self, request: AgentRequest, context: RequestContext) -> RunAccepted: + if request.agent_type == "advisor": + await AdvisorRolloutService(self.session).ensure_allowed(context) try: self.factory.authorize(request.agent_type, context) except ForbiddenAgentError: diff --git a/docs/21-投顾Agent迁移TODO.md b/docs/21-投顾Agent迁移TODO.md index b848570..3d4360f 100644 --- a/docs/21-投顾Agent迁移TODO.md +++ b/docs/21-投顾Agent迁移TODO.md @@ -53,7 +53,7 @@ R4=7、R5=1;资产分类成功 18 个,1 个因合同证据不足跳过。每 R1-R5 适当性硬过滤和排除原因。权威数据验收仍保留两个限制:R1/R5 尚不足四个 测试产品,1 个产品因合同证据不足未完成资产分类;完整推荐侧联动留待阶段十。 阶段五测试结果:专项 `7 passed`,全量单元测试 `473 passed, 3 warnings`,Ruff 和 -MyPy 通过。阶段五提交:待提交。 +MyPy 通过。阶段五核心提交:`291cb7b`、`47a5ca4`。 ### 阶段六:开户风险问卷 @@ -78,7 +78,7 @@ Agent 暴露客户最新的 `confirmed` 目标,未确认目标不会进入后 目标书含收益非承诺、非交易指令披露;写操作统一使用 `Idempotency-Key` 并写入审计。 阶段七测试结果:专项测试 `5 passed`,全量单元测试 `468 passed, 3 warnings`,Ruff -通过,MyPy(122 个源文件)通过。阶段七提交:待提交。 +通过,MyPy(122 个源文件)通过。阶段七提交:`60c4a49`。 ### 阶段八:持仓分析 @@ -90,7 +90,7 @@ Agent 暴露客户最新的 `confirmed` 目标,未确认目标不会进入后 异常时明确降级,不影响 MySQL 数值分析。 阶段八测试结果:图谱与持仓专项测试 `8 passed`,全量单元测试 `475 passed, 3 warnings`,Ruff -通过,MyPy(127 个源文件)通过。阶段八核心提交:待提交。 +通过,MyPy(127 个源文件)通过。阶段八核心提交:`281e76e`、`8ffd08b`。 阶段九已完成动态资产配置与真实历史回测:接入 C1-C5 基础权重、收益目标、最大回撤、 流动性和投资期限约束,读取经适当性、合同证据、资产分类和数据质量门槛过滤的场内基金 @@ -165,7 +165,14 @@ Redis 不可用时实测按设计降级放行;生产装配模式因本机 Milv 及恢复关闭,成功行情按字段精度量化后写入场内行情快照。真实独立库同步验证两源均失败时返回 `degraded=true`、不写入伪行情并进入失败告警链路;本次无可用行情,推荐和动态配置继续失败关闭。 专项测试 `6 passed`,全量单元测试 `499 passed,3 warnings`,契约测试 `8 passed`,Ruff 和 MyPy -通过。实现提交:`2376585`。 +通过。实现提交:`2376585`;命令入口提交:`fbb1171`。 + +本次继续完成灰度闸门和回滚手册:新增 `AdvisorRolloutService`,由环境变量控制投顾灰度, +开启后管理员放行、客户按白名单放行,未命中返回 `403 AGENT_PERMISSION_DENIED` 并写入 +`advisor.rollout_denied` 审计;已接入投顾业务路由和 `advisor` Agent 运行入口。新增操作手册 +`docs/22-投顾Agent灰度与回滚操作手册.md`。专项测试 `4 passed`,全量单元测试 `505 passed, +3 warnings`,Ruff 和 MyPy(146 个源文件)通过。实现提交:`b3a0ee7`。生产/联调环境的 +实际灰度与回滚演练仍待执行。 ## 一、迁移准备 @@ -204,16 +211,16 @@ python -m mypy app ## 三、公共底座适配 -- [ ] 阅读并记录 `app/core/contracts.py` 的请求和响应契约。 -- [ ] 阅读并记录 `app/core/errors.py` 的错误码和错误信封。 -- [ ] 确认 JWT 鉴权和 RBAC 解析流程。 -- [ ] 确认 `BaseAgent` 的标准执行流程。 -- [ ] 确认 `AgentFactory` 的 Agent 注册方式。 -- [ ] 确认 `ToolExecutor` 的工具鉴权和角色限制。 -- [ ] 确认模型路由、模型降级和超时处理方式。 -- [ ] 确认 Redis、Milvus、Neo4j 的接入方式。 -- [ ] 确认记忆抽取、召回、生命周期和 Episode 聚合流程。 -- [ ] 确认 Worker 注册和事件消费方式。 +- [x] 阅读并记录 `app/core/contracts.py` 的请求和响应契约。 +- [x] 阅读并记录 `app/core/errors.py` 的错误码和错误信封。 +- [x] 确认 JWT 鉴权和 RBAC 解析流程。 +- [x] 确认 `BaseAgent` 的标准执行流程。 +- [x] 确认 `AgentFactory` 的 Agent 注册方式。 +- [x] 确认 `ToolExecutor` 的工具鉴权和角色限制。 +- [x] 确认模型路由、模型降级和超时处理方式。 +- [x] 确认 Redis、Milvus、Neo4j 的接入方式。 +- [x] 确认记忆抽取、召回、生命周期和 Episode 聚合流程。 +- [x] 确认 Worker 注册和事件消费方式。 - [x] 适配投顾 Agent 的公共请求和响应结构。(复用 `BaseAgent` 和 `AgentRequest/CoreResult`) - [x] 适配投顾工具注册。(复用 `query_fund_quote` 只读工具) - [x] 适配投顾 Agent 的 `AgentDefinition`。(新增 `advisor`,版本 `0.1.0`) @@ -259,7 +266,7 @@ python -m mypy app - [ ] 在已有测试库执行迁移。 - [x] 执行数据库结构审计。(72 张业务表,无缺失或多余) - [x] 执行数据库约束审计。(唯一键与 ORM 映射通过) -- [ ] 完成数据库迁移提交 `advisor/database-migrations`。(待提交) +- [x] 完成数据库迁移提交 `advisor/database-migrations`。(`9ed536e`) 验收命令: @@ -475,7 +482,7 @@ python tools/audit_constraints.py ## 十三、全端测试 -- [x] 执行全部单元测试。(`491 passed, 3 warnings`,Python 3.13) +- [x] 执行全部单元测试。(`505 passed, 3 warnings`,Python 3.13) - [x] 执行全部契约测试。(`8 passed`,Python 3.13) - [x] 执行全部集成测试。(独立迁移库 `28 passed, 1 skipped`;原 `jr_agent` 为 `25 passed, 3 failed, 1 skipped`,失败均为历史库状态) - [x] 执行 Ruff 检查。(通过) @@ -516,9 +523,9 @@ python tools/audit_constraints.py ## 十四、灰度发布 -- [ ] 在本地测试库完成迁移。 +- [x] 在本地测试库完成迁移。(独立库 `jr_agent_qyqy_migration` 已验证) - [ ] 在联调库完成迁移。 -- [ ] 仅开放测试客户和管理员账号。 +- [x] 仅开放测试客户和管理员账号。(代码闸门已完成,环境实测待执行) - [ ] 灰度产品查询和适当性过滤。 - [ ] 灰度风险问卷。 - [ ] 灰度投资目标。 @@ -533,17 +540,17 @@ python tools/audit_constraints.py ## 十五、回滚准备 -- [ ] 保存迁移前数据库备份。 -- [ ] 保存迁移前 Git 标签。 -- [ ] 保存每个模块的提交号。 -- [ ] 保存每个模块的测试结果。 -- [ ] 保存迁移后的结构审计结果。 -- [ ] 准备关闭新投顾入口的配置开关。 -- [ ] 准备应用代码按提交回滚方案。 -- [ ] 确认数据库不执行破坏性 downgrade。 -- [ ] 确认新增表保留,不自动删除。 -- [ ] 确认失败 Outbox 可以重试或人工处理。 -- [ ] 确认原画像和原推荐结果可以保留。 +- [x] 保存迁移前数据库备份。(`.migration-backups/jr_agent-before-base-migration.sql`) +- [x] 保存迁移前 Git 标签。(`advisor-before-base-migration`) +- [x] 保存每个模块的提交号。(已回填阶段记录) +- [x] 保存每个模块的测试结果。(已回填阶段记录) +- [x] 保存迁移后的结构审计结果。(独立迁移库 72 张业务表) +- [x] 准备关闭新投顾入口的配置开关。(`ADVISOR_ROLLOUT_ENABLED=false`) +- [x] 准备应用代码按提交回滚方案。(见 `docs/22-投顾Agent灰度与回滚操作手册.md`) +- [x] 确认数据库不执行破坏性 downgrade。(见回滚手册) +- [x] 确认新增表保留,不自动删除。(见回滚手册) +- [x] 确认失败 Outbox 可以重试或人工处理。(沿用公共 Outbox 重试/死信机制) +- [x] 确认原画像和原推荐结果可以保留。(回滚只关闭入口,不删除业务数据) - [ ] 完成灰度回滚演练。 ## 十六、最终完成标准 diff --git a/docs/22-投顾Agent灰度与回滚操作手册.md b/docs/22-投顾Agent灰度与回滚操作手册.md new file mode 100644 index 0000000..0aed2f5 --- /dev/null +++ b/docs/22-投顾Agent灰度与回滚操作手册.md @@ -0,0 +1,42 @@ +# 投顾 Agent 灰度与回滚操作手册 + +## 灰度开关 + +投顾灰度使用进程环境变量,不修改冻结的数据库基线,也不需要数据库迁移: + +```text +ADVISOR_ROLLOUT_ENABLED=true +ADVISOR_ROLLOUT_CUSTOMER_IDS=9001,9002 +``` + +`ADVISOR_ROLLOUT_ENABLED=false` 时保持现有本地行为。开启后,`admin` 和 +`super_admin` 始终放行;客户只有在 `ADVISOR_ROLLOUT_CUSTOMER_IDS` 中才可以访问 +投顾 Agent、投资目标、持仓分析、资产配置和推荐接口。白名单为空时客户全部失败关闭, +不会误放量。灰度拒绝统一返回 `403 AGENT_PERMISSION_DENIED`,并写入 +`interaction_audit.action_type=advisor.rollout_denied`。 + +修改环境变量后必须重启全部 API 进程和 Worker,确保同一发布批次使用同一开关快照。 +上线前先用一名测试客户和一名管理员验证:客户白名单命中返回业务响应,未命中返回 403, +管理员可以完成审核/发布操作;检查审计记录中没有客户数据和内部画像字段泄露。 + +## 灰度检查 + +1. 在独立测试库执行 `python -m alembic upgrade head`、`python tools/audit_schema.py` 和 + `python tools/audit_constraints.py`。 +2. 启动 API 和 Worker,执行 `python tools/acceptance_check.py`,确认问卷前置、鉴权、 + SSE、Worker 完成态、越权和数据隔离均通过。 +3. 开启白名单后,用测试客户执行问卷、投资目标、持仓分析、资产配置、推荐及画像复核; + 推荐无有效行情时必须保持失败关闭,不得写入伪行情或产生模拟委托。 +4. 记录 API 错误数、延迟、降级次数、行情失败告警和 `advisor.rollout_denied` 审计数。 + +## 回滚 + +优先将 `ADVISOR_ROLLOUT_ENABLED=false` 并滚动重启 API/Worker,关闭投顾灰度入口;已落库 +的问卷、目标、画像、推荐审核记录和 Outbox 事件保留,不删除业务数据。 + +若需应用回滚,回到本分支上一个已验收提交,并重新执行单元、契约和结构审计。数据库不执行 +破坏性 `downgrade`,新增表保留;恢复后用 `alembic current` 和结构审计确认版本与结构一致。 +失败 Outbox 事件按原有重试/死信机制处理,必要时由管理员重放,不直接删除未完成事件。 + +回滚演练必须记录:开关关闭时间、API/Worker 重启结果、客户 403 验证、管理员可用性、 +审计记录、Outbox 未丢失以及没有产生真实交易委托。 diff --git a/tests/unit/service/test_advisor_rollout_service.py b/tests/unit/service/test_advisor_rollout_service.py new file mode 100644 index 0000000..fe3df42 --- /dev/null +++ b/tests/unit/service/test_advisor_rollout_service.py @@ -0,0 +1,85 @@ +from types import SimpleNamespace + +import pytest + +from app.core.contracts import RequestContext +from app.core.errors import ForbiddenAgentError +from app.service.advisor_rollout_service import AdvisorRolloutService + + +def context( + *, roles: tuple[str, ...] = ("customer",), customer_ids: tuple[str, ...] = () +) -> RequestContext: + return RequestContext( + user_id="9001", trace_id="rollout-test", roles=roles, + customer_ids=customer_ids, + ) + + +def settings(*, enabled: bool, customer_ids: str) -> SimpleNamespace: + return SimpleNamespace( + advisor_rollout_enabled=enabled, + advisor_rollout_customer_ids=customer_ids, + ) + + +def test_disabled_rollout_allows_everyone(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr( + "app.service.advisor_rollout_service.get_settings", + lambda: settings(enabled=False, customer_ids=""), + ) + assert AdvisorRolloutService.is_allowed(context()) + + +@pytest.mark.parametrize( + ("roles", "customer_ids", "allowed"), + [ + (("customer",), (), True), + (("customer",), (), False), + (("admin",), (), True), + (("super_admin",), (), True), + ], +) +def test_enabled_rollout_uses_customer_whitelist_and_admin_bypass( + monkeypatch: pytest.MonkeyPatch, + roles: tuple[str, ...], + customer_ids: tuple[str, ...], + allowed: bool, +) -> None: + whitelist = "9001" if roles == ("customer",) and not customer_ids else " 9002 , 9003 " + if roles == ("customer",) and not customer_ids and allowed is False: + whitelist = "9002" + monkeypatch.setattr( + "app.service.advisor_rollout_service.get_settings", + lambda: settings(enabled=True, customer_ids=whitelist), + ) + actual = AdvisorRolloutService.is_allowed(context(roles=roles, customer_ids=customer_ids)) + assert actual is allowed + + +class FakeSession: + def __init__(self) -> None: + self.added: list[object] = [] + self.commits = 0 + + def add(self, value: object) -> None: + self.added.append(value) + + async def commit(self) -> None: + self.commits += 1 + + +async def test_denial_is_audited_and_raises(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr( + "app.service.advisor_rollout_service.get_settings", + lambda: settings(enabled=True, customer_ids="9002"), + ) + session = FakeSession() + + with pytest.raises(ForbiddenAgentError, match="尚未对该账号开放"): + await AdvisorRolloutService(session).ensure_allowed(context()) # type: ignore[arg-type] + + assert session.commits == 1 + audit = session.added[0] + assert audit.action_type == "advisor.rollout_denied" # type: ignore[attr-defined] + assert audit.detail["reason"] == "not_in_rollout" # type: ignore[attr-defined]