feat: add advisor rollout gate and rollback playbook

This commit is contained in:
Windows
2026-09-11 20:28:21 +08:00
parent fbb1171a45
commit f5dd5b8ba2
11 changed files with 253 additions and 35 deletions
+4
View File
@@ -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 不可用时自动降级为放行(限流是保护措施,不能因后端故障拒绝正常请求)。
+2 -1
View File
@@ -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)],
)
+2 -1
View File
@@ -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)],
)
+2 -1
View File
@@ -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)],
)
+10 -3
View File
@@ -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
+4
View File
@@ -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")
+63
View File
@@ -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)
@@ -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:
+36 -29
View File
@@ -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] 确认原画像和原推荐结果可以保留。(回滚只关闭入口,不删除业务数据)
- [ ] 完成灰度回滚演练。
## 十六、最终完成标准
@@ -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 未丢失以及没有产生真实交易委托。
@@ -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]