Author SHA1 Message Date
wangjianlong_0626 1f62aca6f7 fix(seed): 当日行情重跑必须刷新 source_updated_at(修 503 FUND_QUOTE_UNAVAILABLE)
真机实测到的缺陷:`tools/seed_sim_account_demo.py` 的 `_upsert_market_price()`
在"当天已有行"时直接 `return`,于是同一天重跑种子**不刷新 `source_updated_at`**。
行情是时效数据,过期后 `FundQuoteService` 返回
`503 FUND_QUOTE_UNAVAILABLE:产品 510300 行情已过期`,
连带 `T001` 仪表盘与 `T006` 持仓一起不可用(`T010` 权益不查行情,仍 200)。

表现极具误导性:**刚灌完种子能用,过十几分钟就 503** ——
看起来像行情适配器或缓存故障,实际根因在种子脚本。归因过程:
`memory_sync` / Redis / Milvus 全部正常,`fin_market_price` 里当天那行
`source_updated_at` 停在首次灌入时刻。

修法:把列值抽成共享 `values` 字典,当天已有行时 `update` 刷新全部行情列
(含 `source_updated_at`),不存在才 `insert`。
**幂等的正确含义是"不产生重复行"(`(product_id, trade_date)` 唯一),
不是"不更新值"。** 账户/持仓的"已存在则跳过"保持不变 ——
那是业务数据,不该被种子覆盖。

验证(本机,测试客户 9001 `cust_t`):
- `python -X utf8 -m tools.seed_sim_account_demo --customer-id 9001` 正常
- `GET /api/v1/users/me/account/dashboard` → 200(修前 503)
- `GET /api/v1/users/me/holdings` → 200
- `GET /api/v1/users/me/entitlements` → 200

门禁:ruff `app tests tools alembic` 全过;`mypy app` 0 错(252 文件);
`audit_schema.py` 90 张业务表无差异;文档编号/端点编号/RBAC 种子一致性全过。
2026-09-12 17:32:24 +08:00
wangjianlong_0626 4766e3bd98 feat(benefit): 客户权益功能(T010)+ 修投顾迁移契约里写死 head 的脆弱断言
## 1. 新增客户权益(用户端)

`GET /api/v1/users/me/entitlements`(T010,权限 `benefit:read:self`):

- **层级**由 `fin_customer_profile.total_asset` **实时判定**
  (门槛来自 `knowledge/product/高净值客户服务规范.md`:
  金卡 50 万 / 白金 200 万 / 钻石 600 万 / 私行 1000 万;低于 50 万为普通客户);
- **权益按层级累积展开**(文档原文"含全部下级权益,新增以下"):
  金卡 9 条 / 白金 20 / 钻石 33 / 私行 54,各档已逐档实测;
- 返回**升级提示**(`next_tier`:下一层级与门槛),前端可直接渲染"再投 X 元升级"。

### 新增表 `fin_customer_benefit`(1 张)

层级 → 权益目录,54 条种子数据(`tools/seed_customer_benefits.py`,按 `benefit_code` 幂等)。

**基线合规证明**(规则 1/3/4):只新增这一张表;**未**重命名/删除任何已有表;
**未**重命名/删除/复用任何已有字段,**未**改任何已有字段的类型、可空性或业务含义;
未改 `docs/00`。
复核:`tools/audit_schema.py` → `90 business tables, no missing or unexpected tables`。

### 两条设计取舍

1. **不落"某客户享有哪些权益"**:层级可算,权益由层级推出,两者都不落库。
   与 `docs/00` L159(不保留 `net_worth_flag`,因为可算)同一取向。
2. **权益只存各层新增条目**,累积由服务层 `tier_chain()` 展开 ——
   否则改一条权益要改四处,漏一处就出现"白金没有金卡权益"。

### 数据来源与一处刻意省略

逐条照抄知识文档,不新增文档里没有的权益。**私行那条
「7×24小时私人银行专线:400-XXX-XXXX 转 8」不写号码** ——
文档里是占位符,而对客号码的唯一来源是 `customer_service_rules.CONTACT_PHONE`
(本线此前修过"同一客服给客户两个不同号码"的缺陷)。把占位符抄进库等于再造一份假号码。

## 2. 修投顾迁移契约里写死的断言

`tests/unit/test_advisor_migration_contract.py` 原先断言

```python
assert script.get_heads()[0] == "20260911_merge_adv_risk_heads"
```

那是"投顾迁移刚加完那一刻"的快照 —— 本 PR 一新增迁移(`20260912_customer_benefit`)
它就变红,**而红的原因与投顾链的对错无关**:断言测到的是时间,不是契约。

原意是"投顾链接在这条主链上、没另起分支"。改为断言**投顾链尾是当前 head 的祖先**
(链尾从 `ADVISOR_FILES[-1]` 派生,不写死),既保住原意又不受后续迁移影响。
`len(script.get_heads()) == 1`(链不分叉)与"投顾文件首尾相接"两条原样保留。

## 3. 顺带发现的既有缺口(**不在本次改动范围**)

`app/api/controllers/trading.py` 的 **T001–T009 未调用 `AuthorizationService.require`**:
`docs/05` §19 为它们登记了权限码(`account:read:self` / `trade:order:*` / `holding:read:self`),
但代码只做认证 + 开户测评门槛,**没有执行 RBAC 权限检查**。
对照:仓库里 **26 个 service** 都调了 `require`,`trade_service` 不在其中。

本线的 T010 **按正确做法实现**:`CustomerBenefitService.entitlements_for` 先鉴权再读数据,
且**鉴权在读取客户资产之前**(有测试断言"拒绝时未查库")。
T001–T009 如何补,需架构师定口径后另行处理。

## 4. 文档

- 新增 `docs/41-客户权益功能说明.md`:表登记 + 基线合规证明 + 分层口径 + 累积规则 +
  数据来源 + 权限 + 与仪表盘的关系 + 上述缺口
- `docs/05` §19 登记 T010,并**单独注明它引入了新表**(避免被误读为
  "T 段数据库零变更"的一部分)
- `AGENTS.md` 表数 89 → **90** 张业务表

## 验证

- `pytest tests/unit/service/test_customer_benefit_service.py` → **20 passed**
  (含边界:499999.99 不是金卡、500000 整是金卡、1000 万整是私行;累积条数;升级提示;
  鉴权先于读数据)
- 全量 `pytest tests` → `2 failed, 1469 passed, 1 skipped`
  (2 个失败为既有环境项:httpx 把中文序列化成 `\uXXXX`,非本次引入)
- `ruff check app tests tools alembic` → `All checks passed`
- `mypy app` → **0 错 / 252 文件**
- 真机:`GET /users/me/entitlements` → `200`;各档分层与累积条数逐档实测通过
- `audit_schema.py` → 90 张业务表无缺失/意外;文档守卫 55 份无编号冲突;
  端点编号无重复;RBAC 种子一致性通过
2026-09-12 17:24:37 +08:00
wangjianlong_0626 37926e1c88 Merge remote-tracking branch 'origin/qyqy_develop' into nl-merge-colleague 2026-09-12 16:39:39 +08:00
wangjianlong_0626 db4dae2fe6 chore: 追平主干(advisor 演示引导/风控权限配置 3 提交);同步测试基线 1424
主干 `c200a61..5634fdc` 3 个提交(advisor 演示环境引导、风控登录权限与模型能力配置),
与本线 5 个提交**无文件重叠**,合并**无冲突**(`tools/seed_test_rbac.py` 双方都改过,
Git 自动合上;合并后 `check_rbac_seed_consistency.py` 仍通过)。

追平后本线相对主干恢复 **BEHIND 0**,PR #9 可重新 fast-forward 合并。

## 验证(合并后重跑)

- `ruff check app tests tools alembic` → `All checks passed`
- `mypy app` → 0 错 / 245 文件
- `pytest tests`(全量)→ `2 failed, 1424 passed, 1 skipped`
  (2 个失败为既有环境项:httpx 把中文序列化成 `\uXXXX`,非本次引入)
- `check_authoritative_docs.py` → 53 份文档无编号冲突
- `check_rbac_seed_consistency.py` → 通过(种子 id 唯一、各 grant 脚本与种子逐条一致)
- `check_docs_endpoint_ids.py` → §19 端点编号无重复
- `audit_schema.py` → 89 张业务表无缺失/意外
2026-09-12 14:50:08 +08:00
wangjianlong_0626 aa3f9f9b3d Merge remote-tracking branch 'origin/qyqy_develop' into nl-merge-colleague 2026-09-12 14:48:45 +08:00
wangjianlong_0626 b2a2e4292b docs(37): 补「换到新环境必须先做」的前置清单(建集合/Milvus/embedding 端点/迁移) 2026-09-12 14:35:49 +08:00
wangjianlong_0626 c12f836d26 Merge remote-tracking branch 'origin/qyqy_develop' into nl-merge-colleague 2026-09-12 14:28:51 +08:00
wangjianlong_0626 86d1cf1ffc docs: 补回合并时丢失的 ruff 基线;合并主干 6 提交(含架构师对我疏漏的修正)
## 1. 补回被我丢失的门禁条目

主干 `AGENTS.md` 原本有一条 `ruff` 干净,但我在解决 `AGENTS.md` 冲突时**只保留了自己的
mypy/pytest 基线,把这条丢了** —— 合并时丢信息,和代码冲突一样是缺陷。

本次补回,并把命令写准(这点很重要):

- `ruff check app tests tools alembic` → **`All checks passed`(0 错)**
- 直接 `ruff check .`(全仓)会报 **40 个错,全部来自仓库根目录的 `hq.py` / `nl2sql_yc.py`**
  (袁聪线的演示脚本,不属本项目包结构)

所以"ruff 干净"**必须带范围**,否则会和别人的脚本混在一起、把一个健康状态误报成 40 个错。

## 2. 我的一个真实疏漏(已由架构师修正,本次合并带入)

架构师提交 `ffbcc22`:**`fix: 删掉 NL 合并后残留的未使用变量 role_ids(ruff F841)`**

该疏漏是我引入的:改 `tools/seed_test_rbac.py` 时把
`zip(user_ids, role_ids, strict=True)` 换成显式配对表 `USER_ROLES`,删掉了 `user_ids`
却**没删 `role_ids`**。根因是**我全程没跑过 ruff** —— 项目门禁里有它,
只跑 mypy 和 pytest 是不够的。

已在本机复跑 `ruff check app tests tools alembic` 确认:我改过的文件全部干净
(`role_ids` 已随主干修正进来)。

## 3. 合并主干 6 提交

`fedbf5a..c8cdc06`,含上条修正与投顾线的行情双源、验收归档等,**无冲突**。

## 验证

- `ruff check app tests tools alembic` → `All checks passed`
- `mypy app` → 0 错 / 245 文件
- `pytest tests`(全量)→ `2 failed, 1421 passed, 1 skipped`
  (2 个失败为既有环境项:httpx 把中文序列化成 `\uXXXX`,非本次引入)
2026-09-12 14:26:43 +08:00
wangjianlong_0626 8cb75992e6 Merge remote-tracking branch 'origin/qyqy_develop' into nl-merge-colleague 2026-09-12 14:21:34 +08:00
wangjianlong_0626 451aa4e915 fix(profile): ProfileAssemblyService 漏写 current_customer_id(唯一键失效的另一处)
## 现象(2026-09-12 端到端跑通后查数据时发现)

客户 9102 的画像快照状态**彼此矛盾**:

| version | is_current | current_customer_id |
|---|---|---|
| 2 | `0` | `9102` ← 清旧时没清空 |
| 3 | `1` | `NULL` ← 建新时没写入 |

## 根因

`ProfileAssemblyService._write_snapshot` 把 `current_customer_id` **当成了生成列**:

- 方法 docstring 原文写着"唯一键 `uk_profile_snapshot_current` 建立在**生成列** `current_customer_id` 上"
- 因此两处都不赋值(以为 DB 会自动填)

**但该列不是生成列** —— `alembic/baseline_generated.sql` 与真实库都是**普通可空列 + 唯一键**,
`app/model/profile.py` 的模块 docstring 第 2 条已明确:"当前版本必须由写入方**显式写入**客户 ID
(历史版本写 NULL),才能保证「每个客户最多一条当前快照」"。

后果:唯一键**形同虚设**(多个 NULL 不冲突)⇒ 不变式失效;且旧版本残留的值
一旦与新版本补上的值相同,就会**直接撞唯一键**。

> 这与本线先前修的 `CustomerProfileCandidateService._write_profile_snapshot` 是**同一个缺陷的另一处**
> —— 当时只找到一处,这次是靠真实链路跑出数据后核对才暴露出来。

## 改动

`app/service/profile_assembly_service.py`:

- 旧版本:`is_current = False` 的同时 `current_customer_id = None`
- 新版本:`is_current=True` 的同时 `current_customer_id=customer_id`
- 订正方法 docstring 的错误认知("生成列"→ 普通可空列 + 唯一键),并写明后果

## 已有数据订正

新增 `tools/fix_profile_snapshot_current.py`(**默认 dry-run**、幂等、`--apply` 才提交):

1. 先清空 `is_current=0` 却残留值的行
2. 再补写 `is_current=1` 却是 NULL 的行
3. **顺序要紧**:反过来的话第 2 步会与残留值撞唯一键

本机实测:清空 1 行、补写 1 行,复核两类异常均归零。

## 验证

- `mypy app` → 0 错 / 245 文件
- `pytest tests`(全量)→ `2 failed, 1418 passed, 1 skipped`
  (2 个失败为既有环境项:httpx 把中文序列化成 `\uXXXX`,非本次引入)
- 端到端:真实对话 → 记忆抽取 → `memory_unit` 落库已实测通过(客户 9102
  `preference:risk_level = "低风险"`,候选态)
2026-09-12 14:20:38 +08:00
16 changed files with 1042 additions and 25 deletions
+13 -6
View File
@@ -73,11 +73,12 @@
- 解释器:本机用 **`.\.venv\Scripts\python.exe`**;架构师环境用 `D:\conda\envs\jr_py313\python.exe`。
两者等价,**各用本机可用的那个**(`.venv` 被 `.gitignore` 忽略、不进仓库,不存在"需要统一"的问题)。
- 数据库现为 **90 张表**(含 `alembic_version`)= **89 张业务表** =
**场内 51 + 场外/推广 17 + 投顾 21**。
后 38 张(`offsite_*` / `promotion_*` / `advisor_*`)**不进 `docs/00` 基线**(规则 8):
- 数据库现为 **91 张表**(含 `alembic_version`)= **90 张业务表** =
**场内 51 + 场外/推广 17 + 投顾 21 + 客户权益 1**。
后 39 张(`offsite_*` / `promotion_*` / `advisor_*` / `fin_customer_benefit`)**不进 `docs/00` 基线**(规则 8):
场外/推广那 17 张逐表登记见 `docs/28-场外与推广域数据表登记.md`;
**投顾那 21 张的登记文档待补**(按同样口径另立一份)。
**投顾那 21 张的登记文档待补**(按同样口径另立一份);
客户权益 1 张见 `docs/41-客户权益功能说明.md`(含基线合规证明)。
核验命令:`python tools/audit_schema.py`(若报 `unexpected` 先分清是"库里多表"还是"迁移没进来")。
- 已注册业务 Agent(**7 个**,见 `app/service/agent/implementations/` 与 `app/service/agent/`):
`FundQueryDemoAgent`、`CustomerServiceAgent`、`RiskAgent`、`PlatformProbeAgent`、
@@ -117,13 +118,19 @@
**新增模型前先搜一遍 `__tablename__` 有没有被占用。**
- 测试基线(**2026-09-12 合并主干 PR #7 + 本线记忆投影链路之后实测**):
`mypy app` → **245 个文件 0 错**;
`pytest tests`(全量)→ `2 failed, 1415 passed, 2 skipped`;
`pytest tests`(全量)→ `2 failed, 1424 passed, 1 skipped`;
`pytest tests/integration` → `102 passed, 1 skipped`。
⚠️ **用例数会随开发增减,判断健康看"0 failed"而不是看绝对值**。
剩下那 2 个失败**都不是代码缺陷**,接手时不要"修"它们:
那 2 个失败**都不是代码缺陷**,接手时不要"修"它们:
`tests/unit/service/test_offsite_document_recognition_adapter.py` 的 2 个用例 —— **环境相关**:
它们断言请求体里是中文原文,而 httpx 会把中文序列化成 `\uXXXX`,字节序列自然不匹配。
功能无影响;若要修,正确做法是断言 `json.loads(body)` 后的字段值(字节级断言不该用来测 JSON)。
- **`ruff`:`ruff check app tests tools alembic` → 全部 `All checks passed`(0 错)。**
⚠️ **检查范围要写对**:直接 `ruff check .`(全仓)会报 40 个错,**全部来自仓库根目录的两个
散落脚本 `hq.py` / `nl2sql_yc.py`**(袁聪线的演示脚本,不属本项目包结构)。
所以门禁命令用上面那条;报"ruff 干净"时**必须带范围**,否则会和别人的脚本混在一起。
(`pyproject.toml` 的 `[tool.ruff.lint]` 选的是 `E,F,I,B,UP`、`line-length=100`;
`tools/*.py` 另配了 `E501/I001/B007` 的 per-file-ignores。)
- **mypy:`mypy app` → `Success: no issues found in 245 source files`(0 错)。**
⚠️ 曾在本机报 184 个错,**已查明是环境版本旧**,与代码质量无关 —— 复现矩阵:
@@ -0,0 +1,68 @@
"""add customer benefit catalog (tier → entitled benefits)
Compatibility proof: this revision creates only the additive
``fin_customer_benefit`` table. It does not alter, rename, delete, reuse, or
retype any baseline table or existing field.
为什么需要这张表(2026-09-12):
`docs/00` 基线里**已有**客户分层字段与承载分层规则的费率表:
- `sys_user.customer_tier VARCHAR(16)` —— 客户分层(仅客户使用)
- `fin_fee_rule.customer_tier` —— 费率规则按层级区分
但**"各层级享有哪些权益"没有任何载体**:知识库 `knowledge/product/高净值客户服务规范.md`
写全了四级分层(金卡/白金/钻石/私行)与各层权益(金融服务 + 非金融权益、累积制),
系统里却既无表也无接口。本表只补这一块 —— **层级的定义/来源不改**,
仍以 `sys_user.customer_tier` 与 `fin_customer_profile.total_asset` 为准
(基线 L159:「可由 `customer_tier` 或 `total_asset` 计算」)。
设计取舍:
- 只存**权益条目**(层级 → 有哪些权益),**不存"某客户享有什么"** ——
后者可由层级实时推出,落库会造成双份真相(与基线不保留 `net_worth_flag` 同一理由)。
- `customer_tier` 用英文码(`gold`/`platinum`/`diamond`/`private`),
与 `investor_type` 用 `C1-C5` 同一口径;中文名(金卡/白金/钻石/私行)由接口层映射,
避免把展示文案写进库。
- 权益按层级**累积**(白金含全部金卡权益,依此类推)——由**服务层展开**,
不在表里重复存父级条目,避免同一权益改两处。
"""
from alembic import op
revision = "20260912_customer_benefit"
down_revision = "20260911_merge_adv_risk_heads"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.execute(
"""
CREATE TABLE fin_customer_benefit (
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
benefit_code VARCHAR(64) NOT NULL,
customer_tier VARCHAR(16) NOT NULL,
category VARCHAR(16) NOT NULL,
name VARCHAR(128) NOT NULL,
description VARCHAR(512) NOT NULL,
display_order INT NOT NULL DEFAULT 0,
status VARCHAR(16) NOT NULL,
created_at DATETIME(6) NOT NULL,
updated_at DATETIME(6) NOT NULL,
PRIMARY KEY (id),
UNIQUE KEY uk_fin_customer_benefit_code (benefit_code),
KEY idx_fin_customer_benefit_tier (customer_tier, status, display_order),
CONSTRAINT chk_fin_customer_benefit_tier
CHECK (customer_tier IN ('gold', 'platinum', 'diamond', 'private')),
CONSTRAINT chk_fin_customer_benefit_category
CHECK (category IN ('financial', 'non_financial')),
CONSTRAINT chk_fin_customer_benefit_status
CHECK (status IN ('active', 'inactive'))
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci
"""
)
def downgrade() -> None:
raise RuntimeError("customer benefit catalog must not be dropped automatically")
+36
View File
@@ -0,0 +1,36 @@
"""客户权益 controller(`docs/05` §19 T 段)。
| § | 端点 | 权限码 | 摘要 |
|---|---|---|---|
| T010 | `GET /api/v1/users/me/entitlements` | `benefit:read:self` | 我的客户层级与应享权益 |
设计要点(与 `trading.py` 的 §T 端点保持一致):
- 走 `build_request_context`(数据范围 `self`),权限检查在 **Service 层**
(`CustomerBenefitService.entitlements_for` → `AuthorizationService.require`);
- 不走限流依赖 —— 低频只读查询,由底座网关层限流;
- 信封用 `envelope`(`docs/05` §3.3)。
"""
from __future__ import annotations
from fastapi import APIRouter, 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.api.views.envelope import envelope
from app.core.contracts import RequestContext
from app.service.customer_benefit_service import CustomerBenefitService
router = APIRouter(prefix="/api/v1/users/me", tags=["benefit"])
# T010 我的权益
@router.get("/entitlements")
async def get_my_entitlements(
context: RequestContext = Depends(build_request_context), # noqa: B008
session: AsyncSession = Depends(get_session), # noqa: B008
) -> dict[str, object]:
data = await CustomerBenefitService(session).entitlements_for(context)
return envelope(data, context)
+2
View File
@@ -10,6 +10,7 @@ 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.auth import router as auth_router
from app.api.controllers.benefit import router as benefit_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
@@ -138,6 +139,7 @@ def create_app() -> FastAPI:
application.include_router(recommendation_admin_router)
application.include_router(admin_router)
application.include_router(trading_router)
application.include_router(benefit_router)
application.mount(
"/customer-service-test",
StaticFiles(directory=Path(__file__).resolve().parent / "static", html=True),
+33
View File
@@ -0,0 +1,33 @@
"""ORM mapping for the additive customer benefit catalog.
只映射新增的 `fin_customer_benefit`(层级 → 权益条目)。
**不映射"某客户享有哪些权益"** —— 那可由客户层级实时推出,落库会造成双份真相。
客户层级本身仍以 `sys_user.customer_tier`(`app.model.fund`)与
`fin_customer_profile.total_asset` 为准,本模块不改它们。
"""
from datetime import datetime
from sqlalchemy import BigInteger, DateTime, Integer, String
from sqlalchemy.orm import Mapped, mapped_column
from app.model.base import Base
class CustomerBenefit(Base):
"""客户权益目录:一条 = 某个层级享有的一项权益。"""
__tablename__ = "fin_customer_benefit"
id: Mapped[int] = mapped_column(BigInteger, primary_key=True)
benefit_code: Mapped[str] = mapped_column(String(64), nullable=False)
#: 英文层级码(gold/platinum/diamond/private);中文名由接口层映射。
customer_tier: Mapped[str] = mapped_column(String(16), nullable=False)
#: financial / non_financial
category: Mapped[str] = mapped_column(String(16), nullable=False)
name: Mapped[str] = mapped_column(String(128), nullable=False)
description: Mapped[str] = mapped_column(String(512), nullable=False)
display_order: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
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)
+190
View File
@@ -0,0 +1,190 @@
"""客户权益:按可投资资产判定层级,并展开该层级(含以下各层)应享有的权益。
### 口径来源
`knowledge/product/高净值客户服务规范.md`(公司内部服务标准,四级分层 + 各层权益):
| 层级 | 名称 | 可投资资产门槛 |
|---|---|---|
| `gold` | 金卡 | 50 万 - 200 万 |
| `platinum` | 白金 | 200 万 - 600 万 |
| `diamond` | 钻石 | 600 万 - 1000 万 |
| `private` | 私行 | 1000 万以上 |
低于 50 万为**普通客户**(无层级)。文档写的是「以客户可投资金融资产(不含自住房产)
为主要分层依据」,本实现取 `fin_customer_profile.total_asset` —— 基线
(`docs/00` L213/L220)把它定为「风控研判所用资产快照」且「按统一口径计算」,
是系统里唯一可用的资产口径;**不另造口径**。
### 两条设计约束
1. **权益按层级累积**(文档原文"含全部金卡权益,新增以下"):
白金含金卡全部条目、依此类推。展开在**服务层**做,表里不重复存父级条目
—— 否则同一权益要改多处。
2. **不把"某客户享有哪些权益"落库**:它可由层级实时推出。
基线的同一取向见 `docs/00` L159(不保留 `net_worth_flag`,因为可算)。
"""
from dataclasses import dataclass
from decimal import Decimal
from typing import Any
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.contracts import RequestContext
from app.model.benefit import CustomerBenefit
from app.model.fund import FundCustomerProfile
from app.service.authorization_service import AuthorizationService
#: 本模块使用的权限码;定义源是 `tools/seed_test_rbac.py` 的 `PERMISSIONS`。
PERMISSION_READ_SELF = "benefit:read:self"
@dataclass(frozen=True)
class TierSpec:
"""一个层级的码、中文名与门槛(含)。"""
code: str
label: str
min_total_asset: Decimal
#: **从高到低**排列:判定时取第一个满足门槛的层级。
#: 顺序即累积顺序,`_tier_chain` 依赖它,不要随意重排。
TIERS: tuple[TierSpec, ...] = (
TierSpec("private", "私行", Decimal("10000000")),
TierSpec("diamond", "钻石", Decimal("6000000")),
TierSpec("platinum", "白金", Decimal("2000000")),
TierSpec("gold", "金卡", Decimal("500000")),
)
HEADLINE = "普通客户"
def _spec(code: str) -> TierSpec | None:
return next((t for t in TIERS if t.code == code), None)
def resolve_tier(total_asset: Decimal | int | float | None) -> TierSpec | None:
"""按可投资资产判定层级;低于最低门槛(50 万)或资产缺失时返回 None。"""
if total_asset is None:
return None
amount = Decimal(str(total_asset))
return next((t for t in TIERS if amount >= t.min_total_asset), None)
def tier_chain(spec: TierSpec) -> tuple[str, ...]:
"""该层级**及以下**所有层级码(用于累积展开)。
例:钻石 → `("gold", "platinum", "diamond")`。
`TIERS` 是从高到低,故取其尾部到该层级为止。
"""
codes_low_to_high = [t.code for t in reversed(TIERS)]
return tuple(codes_low_to_high[: codes_low_to_high.index(spec.code) + 1])
class CustomerBenefitService:
"""权益查询;**只读**,不写任何表(含不写 `sys_user.customer_tier`)。"""
def __init__(self, session: AsyncSession) -> None:
self.session = session
async def entitlements_for(self, context: RequestContext) -> dict[str, Any]:
"""端点入口:先鉴权,再按客户资产判定层级并展开权益。
⚠️ 权限检查放在 **Service 层**(项目惯例:26 个 service 都这么做,
见 `AuthorizationService.require`)—— `trading.py` 的 §T 端点目前**没有**
这一步,只有测评门槛,属既有缺口,不在本次改动范围内。
"""
await AuthorizationService.require(context, PERMISSION_READ_SELF)
total_asset = await self._total_asset(int(context.user_id))
return await self.entitlements(total_asset=total_asset)
async def _total_asset(self, customer_id: int) -> Decimal | None:
"""读客户画像的资产快照。
基线把 `fin_customer_profile.total_asset` 定为「风控研判所用资产快照」,
是系统里唯一的资产口径;客户未开户(无画像行)时返回 None,
由 `resolve_tier` 判为无层级,而不是抛错 —— "还没开户"不是异常。
"""
row = await self.session.scalar(
select(FundCustomerProfile.total_asset).where(
FundCustomerProfile.customer_id == customer_id
)
)
return row
async def entitlements(self, *, total_asset: Decimal | None) -> dict[str, Any]:
"""返回客户当前层级与应享权益。
`total_asset` 由调用方从 `fin_customer_profile` 读入并传入,
避免本服务再查一次客户主表(也在测试里便于给定资产直接断言分层)。
"""
spec = resolve_tier(total_asset)
if spec is None:
return {
"tier": None,
"tier_label": HEADLINE,
"min_total_asset": None,
"total_asset": self._amount(total_asset),
"next_tier": self._next_tier(None),
"benefits": [],
}
codes = tier_chain(spec)
rows = list(await self.session.scalars(
select(CustomerBenefit)
.where(
CustomerBenefit.customer_tier.in_(codes),
CustomerBenefit.status == "active",
)
.order_by(CustomerBenefit.display_order, CustomerBenefit.id)
))
return {
"tier": spec.code,
"tier_label": spec.label,
"min_total_asset": self._amount(spec.min_total_asset),
"total_asset": self._amount(total_asset),
"next_tier": self._next_tier(spec),
"benefits": [self._view(row) for row in rows],
}
@staticmethod
def _next_tier(spec: TierSpec | None) -> dict[str, Any] | None:
"""下一层级与还差多少 —— 前端可直接渲染"再投 X 元升级"。
已是最高的私行返回 None。
"""
if spec is None:
target = TIERS[-1] # 金卡
else:
higher = [t for t in TIERS if t.min_total_asset > spec.min_total_asset]
if not higher:
return None
target = min(higher, key=lambda t: t.min_total_asset)
return {
"tier": target.code,
"tier_label": target.label,
"min_total_asset": CustomerBenefitService._amount(target.min_total_asset),
}
@staticmethod
def _amount(value: Decimal | int | float | None) -> str | None:
"""金额统一用字符串输出(与项目其他接口一致,避免浮点误差)。"""
if value is None:
return None
return str(Decimal(str(value)).quantize(Decimal("0.01")))
@staticmethod
def _view(row: CustomerBenefit) -> dict[str, Any]:
spec = _spec(row.customer_tier)
return {
"benefit_code": row.benefit_code,
# 权益所属层级(用于前端按层级分组;累积展开后可能低于客户自身层级)
"tier": row.customer_tier,
"tier_label": spec.label if spec else None,
"category": row.category,
"name": row.name,
"description": row.description,
}
+13 -1
View File
@@ -226,8 +226,16 @@ class ProfileAssemblyService:
) -> None:
"""写入新版本快照并把旧版本置为非当前。
唯一键 `uk_profile_snapshot_current` 建立在生成列 `current_customer_id` 上,
唯一键 `uk_profile_snapshot_current` 建立在 `current_customer_id` 上,
保证「每个客户最多一条 current」;因此必须先清旧再写新,顺序不能反。
⚠️ **该列不是生成列**(订正 2026-09-12):它由 `alembic/baseline_generated.sql`
与真实库确认为**普通可空列 + 唯一键**,`app/model/profile.py` 的模块 docstring
第 2 条已写明"当前版本必须由写入方**显式写入**客户 ID(历史版本写 NULL)"。
此处原先按"生成列"理解而两处都不赋值,后果实测到:
旧版本 `is_current=0` 却仍留着 `current_customer_id`,新版本 `is_current=1`
却是 NULL —— **唯一键形同虚设**(多个 NULL 不冲突),
而且旧值一旦残留、新值再补写就会直接撞键。
"""
previous = list(await self.session.scalars(
select(ProfileSnapshot).where(
@@ -236,6 +244,8 @@ class ProfileAssemblyService:
))
for row in previous:
row.is_current = False
# 必须同时清空,否则与即将插入的新当前版本撞唯一键(见 docstring)。
row.current_customer_id = None
row.updated_at = now
await self.session.flush()
@@ -252,6 +262,8 @@ class ProfileAssemblyService:
generation_basis=basis,
snapshot_hash=sha256(payload.encode("utf-8")).hexdigest(),
is_current=True,
# 当前版本必须显式写入客户 ID(历史版本为 NULL),见 docstring。
current_customer_id=customer_id,
generated_at=now,
created_at=now,
updated_at=now,
+12
View File
@@ -1162,6 +1162,18 @@ GET /internal/metrics
| T007 | `GET /api/v1/users/me/transactions` | `trade:txn:read`(已登录) | 否 | `200` | 成交记录列表 |
| T008 | `GET /api/v1/users/me/transactions/{txn_no}` | `trade:txn:read`(资源所有者) | 否 | `200` | 成交详情 |
| T009 | `GET /api/v1/users/me/cash-ledger` | `account:read:self`(已登录) | 否 | `200` | 资金账本(按 id 倒序游标分页) |
| T010 | `GET /api/v1/users/me/entitlements` | `benefit:read:self`(已登录) | 否 | `200` | 我的客户层级与应享权益(含升级提示) |
> **T010 的两点说明**(与 T001–T009 **不同源**,避免混淆):
>
> - **它引入了新表**:`fin_customer_benefit`(层级 → 权益目录,54 条种子数据)。
> 下面的「T001 – T009 的四点说明」中"数据库零变更"**不覆盖 T010** ——
> 该表是**新增**的,未重命名/删除/修改任何基线表或既有字段(规则 1/3/4 均未触碰)。
> - **不落"某客户享有哪些权益"**:层级由 `fin_customer_profile.total_asset` **实时判定**
> (门槛见 `knowledge/product/高净值客户服务规范.md`:金卡 50 万 / 白金 200 万 /
> 钻石 600 万 / 私行 1000 万),权益按层级**累积**展开(白金含全部金卡条目,依此类推)。
> 与 `docs/00` L159 不保留 `net_worth_flag` 是同一取向:可算的不落库。
> `sys_user.customer_tier` 字段**本接口只读、不写**。
> **T001 – T009 的四点说明**:
>
+15
View File
@@ -82,6 +82,21 @@
建集合:`python tools/setup_milvus_profile_collection.py`
——**幂等**,集合已存在时只做结构比对报告、不覆盖不删重建(共享 Milvus 实例里还有别的项目的集合)。
> ## ⚠️ 换到新环境(例如最终验收在架构师机器上跑)必须先做这一步
>
> **集合是环境数据,不随代码合并。** 本线的代码里没有、也不该有"自动建集合"的逻辑
> (那会让"环境没准备好"被静默盖住)。所以在新机器上:
| 前置 | 命令 | 不做的后果 |
|---|---|---|
| **建画像向量集合** | `python -X utf8 tools/setup_milvus_profile_collection.py` | `load_collection` 抛 `RecoverableAgentError("画像向量集合不可用")` ⇒ `memory_sync_outbox` 的 `milvus` 行**失败重试到死信**。**这是环境问题,不是代码缺陷** |
| Milvus 可达 | `docker ps` 能看到 `milvus-standalone` 且 healthy | 同上;另注意 `MILVUS_LOCAL_URI` 若被设过会指向本地 Lite 文件、查的是另一个库 |
| embedding 端点 | 发布配置里有 `task_type=embedding` 的端点 | `_profile_embed` 抛"没有可用的 embedding 端点",事件同样重试到死信 |
| 数据库迁移 | `alembic upgrade heads` | `memory_sync_outbox` / `memory_unit` 等表不存在 |
> 建集合后可用 `python -X utf8 tools/setup_milvus_profile_collection.py` **再跑一次**确认输出为
> `exists`(幂等,不会重复建)。
### `memory_sources` 契约(适配器的输入)
```json
+182
View File
@@ -0,0 +1,182 @@
# 客户权益功能说明与数据表登记
**日期**:2026-09-12|**分支**:`NL_develop`|**端点**:`T010`|**权限码**:`benefit:read:self`
---
## 1. 一句话
客户可以查到**自己属于哪一层级、享有哪些权益、还差多少升级**。
```
GET /api/v1/users/me/entitlements (T010,已登录 + benefit:read:self)
```
响应(`docs/05` §3.3 信封,`data` 内):
```json
{
"tier": "platinum",
"tier_label": "白金",
"min_total_asset": "2000000.00",
"total_asset": "3000000.00",
"next_tier": {"tier": "diamond", "tier_label": "钻石", "min_total_asset": "6000000.00"},
"benefits": [
{"benefit_code": "tier:gold:01", "tier": "gold", "tier_label": "金卡",
"category": "financial", "name": "专属理财经理服务", "description": "..."}
]
}
```
---
## 2. 数据表登记(新增 1 张)
### `fin_customer_benefit` 客户权益目录
| 列 | 类型 | 说明 |
|---|---|---|
| `id` | BIGINT UNSIGNED | 主键 |
| `benefit_code` | VARCHAR(64) | 权益编号(唯一),种子按 `tier:<层级>:<序号>` 生成 |
| `customer_tier` | VARCHAR(16) | 适用层级:`gold` / `platinum` / `diamond` / `private` |
| `category` | VARCHAR(16) | `financial` / `non_financial` |
| `name` | VARCHAR(128) | 权益名称 |
| `description` | VARCHAR(512) | 权益说明 |
| `display_order` | INT | 展示顺序(按文档出现次序) |
| `status` | VARCHAR(16) | `active` / `inactive` |
| `created_at` / `updated_at` | DATETIME(6) | 审计时间 |
- 唯一键 `uk_fin_customer_benefit_code`;索引 `idx_fin_customer_benefit_tier`
- 三个 CHECK:层级、类别、状态取值受限
- 迁移:`alembic/versions/20260912_customer_benefit.py`(`down_revision = 20260911_merge_adv_risk_heads`)
### 基线合规证明(规则 1/3/4)
- **只新增**这一张表;**未**重命名/删除任何已有表(规则 3);
- **未**重命名/删除/复用任何已有字段,**未**改任何已有字段类型、可空性或业务含义(规则 4);
- **未改** `docs/00` 基线文档;
- 复核命令:`python -X utf8 tools/audit_schema.py` → 应显示业务表数 **+1**、且无 `missing`/`unexpected`。
---
## 3. 为什么不落"某客户享有哪些权益"
**层级由资产实时判定**,权益由层级推出 —— 两者都不落库。理由与 `docs/00` L159
(不保留 `net_worth_flag`,因为"可由 `customer_tier` 或 `total_asset` 计算,冗余存储会造成不一致")
完全一致:**能算的不要存**。
落库会立刻带来两个问题:资产变化后等级与已存权益脱节;以及同一事实两处可写(谁改都算对)。
> 需要"客户被人工特别授权某项权益"这类留痕需求时,再加一张**例外表**
> (`customer_id + benefit_code + 生效期 + 授权人`),而**不是**把全量权益快照落库。
---
## 4. 分层口径
来源:`knowledge/product/高净值客户服务规范.md`(公司内部服务标准)。
| 层级 | 名称 | 可投资资产门槛 |
|---|---|---|
| `gold` | 金卡 | 50 万 - 200 万 |
| `platinum` | 白金 | 200 万 - 600 万 |
| `diamond` | 钻石 | 600 万 - 1000 万 |
| `private` | 私行 | 1000 万以上 |
- **低于 50 万为"普通客户"**(`tier: null`,`tier_label: "普通客户"`),权益为空但**仍返回升级提示**;
- **门槛含等号**:恰好 50 万即金卡(有边界测试守着,防 off-by-one);
- 资产取 `fin_customer_profile.total_asset` —— 基线(`docs/00` L213/L220)把它定为
「风控研判所用资产快照」且「按统一口径计算」,是系统里唯一的资产口径,**不另造口径**;
文档写的是"可投资金融资产(不含自住房产)",与 `total_asset` 的口径差异**由该字段的维护方负责**,
权益模块不自行调整。
### 权益按层级累积
文档每层都写"含全部下级权益,新增以下"。表里**只存该层新增条目**,
累积展开由 `CustomerBenefitService.tier_chain()` 完成:
| 资产 | 层级 | 权益条数 |
|---|---|---|
| ¥60 万 | 金卡 | 9 |
| ¥300 万 | 白金 | 20(9+11) |
| ¥800 万 | 钻石 | 33(9+11+13) |
| ¥2000 万 | 私行 | 54(9+11+13+21) |
> 若把父级条目在每层重复存一遍,改一条权益要改四处、漏一处就出现"白金没有金卡权益"。
---
## 5. 权益数据的来源与一处刻意省略
**逐条照抄**《高净值客户服务规范》第二章,不新增文档里没有的权益。种子:
`tools/seed_customer_benefits.py`(54 条,按 `benefit_code` 幂等、已存在不覆盖)。
**刻意省略的一处**:文档私行条目原文是
> 7×24小时私人银行专线:**400-XXX-XXXX** 转 8
号码是**占位符**。项目已有明确口径:对客号码的唯一来源是
`app/core/customer_service_rules.py` 的 `CONTACT_PHONE`(本线此前修过
"同一客服给客户两个不同号码"的缺陷,见 `docs/37`)。把占位符抄进库等于又造一份假号码,
故库里只写权益名「7×24 小时私人银行专线」,号码一律走客服热线配置。
---
## 6. 权限
新权限码 `benefit:read:self`(种子 id **9066**,续 §T 的 9060-9065)。
- 定义源:`tools/seed_test_rbac.py` 的 `PERMISSIONS`(**唯一**定义源);
- 已挂入 `CUSTOMER_PERMISSIONS`(customer 角色自带);
- 一致性由 `python -X utf8 tools/check_rbac_seed_consistency.py` 守着;
- ⚠️ **`config_release` 是环境数据**:本权限走 RBAC(`sys_permission`),
不经 `agent_tools` 白名单,故**换环境重跑种子即可,无需重新发布配置**。
---
## 7. 与仪表盘的关系
`T001 /users/me/account/dashboard` 已返回账户、组合汇总与持仓(**接口完整**)。
本接口是**独立**的只读端点,前端可在仪表盘上以"我的等级 + 权益卡片"呈现:
- 仪表盘负责**资产与持仓**(`T001`)→ 本接口负责**等级与权益**(`T010`);
- 两者都归属"用户端(客户视角)",权限数据范围均为 `self`。
> 若后续希望一次请求拿全,可在 `T001` 响应里内联权益字段;**当前不这么做**,
> 因为权益数据的更新频率远低于资产(改权益是运营动作),内联会让每次刷新多查两表。
---
## 8. 已知缺口(**不在本次改动范围**,供架构师评估)
`app/api/controllers/trading.py` 的 **T001–T009 未调用 `AuthorizationService.require`** ——
`docs/05` §19 为它们登记了权限码(`account:read:self` / `trade:order:*` / `holding:read:self` 等),
但代码只做了认证 + 开户测评门槛,**没有执行 RBAC 权限检查**。
对照:仓库里 **26 个 service** 都调了 `AuthorizationService.require`,`trade_service` 不在其中。
本线的 T010 **按正确做法实现**(在 `CustomerBenefitService.entitlements_for` 里先鉴权再读数据,
且鉴权在读取客户资产**之前**,有测试守着)。T001–T009 的补法需架构师定:
是补 `require` 调用,还是明确"用户自助端点只靠认证 + 测评门槛"这一口径并同步 §19 的权限列。
---
## 9. 相关文件
**新增**
- `alembic/versions/20260912_customer_benefit.py`(建表)
- `app/model/benefit.py`(ORM)
- `app/service/customer_benefit_service.py`(分层 + 累积展开 + 鉴权)
- `app/api/controllers/benefit.py`(T010)
- `tools/seed_customer_benefits.py`(54 条权益种子)
- `tests/unit/service/test_customer_benefit_service.py`(20 例)
**修改**
- `app/main.py`(注册 `benefit_router`)
- `tools/seed_test_rbac.py`(加权限码 9066 + 挂 customer 角色)
- `docs/05-接口文档.md`(§19 登记 T010 并注明它引入新表)
- `AGENTS.md`(业务表数 89 → 90)
**验证**
- `pytest tests/unit/service/test_customer_benefit_service.py` → 20 passed
- 真机:`GET /users/me/entitlements` → `200`;各档分层与累积条数逐档实测通过
@@ -0,0 +1,182 @@
"""客户权益服务的定向测试:分层判定、累积展开、升级提示、鉴权。
权益条目由 `tools/seed_customer_benefits.py` 灌入(54 条)。本测试**不查库**,
用替身 session 直接给条目,保证分层与展开逻辑可独立验证。
"""
from decimal import Decimal
from typing import Any
import pytest
from app.core.contracts import RequestContext
from app.core.errors import ForbiddenAgentError
from app.model.benefit import CustomerBenefit
from app.service.customer_benefit_service import (
TIERS,
CustomerBenefitService,
resolve_tier,
tier_chain,
)
def benefit(code: str, tier: str, order: int) -> CustomerBenefit:
return CustomerBenefit(
id=order, benefit_code=code, customer_tier=tier, category="financial",
name=f"{tier}-{order}", description="d", display_order=order,
status="active", created_at=None, updated_at=None, # type: ignore[arg-type]
)
#: 与种子同一形状:每层条数不同,便于断言累积后的总数。
ALL: list[CustomerBenefit] = [
*[benefit(f"g{i}", "gold", i) for i in range(1, 10)], # 9
*[benefit(f"p{i}", "platinum", 100 + i) for i in range(1, 12)], # 11
*[benefit(f"d{i}", "diamond", 200 + i) for i in range(1, 14)], # 13
*[benefit(f"v{i}", "private", 300 + i) for i in range(1, 22)], # 21
]
class FakeSession:
"""只实现 `scalars`:把 `in_(codes)` 近似成"返回全部",由服务侧过滤数量断言。"""
def __init__(self, rows: list[CustomerBenefit]) -> None:
self.rows = rows
self.last_codes: tuple[str, ...] = ()
async def scalars(self, statement: Any) -> list[CustomerBenefit]:
# 从编译后的 SQL 参数里取层级码,保证"只取该层级及以下"确实生效。
params = statement.compile().params
codes = tuple(v for k, v in params.items() if "customer_tier" in str(k))
flat: list[str] = []
for c in codes:
if isinstance(c, str):
flat.append(c)
else:
flat.extend(str(x) for x in c)
self.last_codes = tuple(flat)
return [r for r in self.rows if r.customer_tier in self.last_codes]
# --- 分层判定 ---------------------------------------------------------------
@pytest.mark.parametrize(
("amount", "expected"),
[
(None, None),
("0", None),
("499999.99", None), # 差 1 分不到金卡
("500000", "gold"), # 门槛含等号
("1999999.99", "gold"),
("2000000", "platinum"),
("5999999", "platinum"),
("6000000", "diamond"),
("9999999", "diamond"),
("10000000", "private"), # 1000 万整
("99999999", "private"),
],
)
def test_resolve_tier_boundaries(amount: str | None, expected: str | None) -> None:
"""门槛含等号、边界不外溢 —— 这类 off-by-one 在金额分层里最容易错。"""
spec = resolve_tier(Decimal(amount) if amount is not None else None)
assert (spec.code if spec else None) == expected
def test_tiers_are_ordered_high_to_low() -> None:
"""`TIERS` 必须从高到低:`resolve_tier` 取第一个命中,`tier_chain` 依赖该顺序。"""
thresholds = [t.min_total_asset for t in TIERS]
assert thresholds == sorted(thresholds, reverse=True)
def test_tier_chain_is_accumulating_and_ordered_low_to_high() -> None:
assert tier_chain(next(t for t in TIERS if t.code == "gold")) == ("gold",)
assert tier_chain(next(t for t in TIERS if t.code == "platinum")) == ("gold", "platinum")
assert tier_chain(next(t for t in TIERS if t.code == "private")) == (
"gold", "platinum", "diamond", "private",
)
# --- 累积展开 ---------------------------------------------------------------
@pytest.mark.asyncio
@pytest.mark.parametrize(
("amount", "tier", "count"),
[
("500000", "gold", 9),
("2000000", "platinum", 20), # 9 + 11
("6000000", "diamond", 33), # 9 + 11 + 13
("10000000", "private", 54), # 9 + 11 + 13 + 21
],
)
async def test_benefits_accumulate_with_tier(amount: str, tier: str, count: int) -> None:
"""文档写明"含全部下级权益,新增以下" ⇒ 高等级必须拿到低等级的全部条目。"""
session = FakeSession(ALL)
data = await CustomerBenefitService(session).entitlements( # type: ignore[arg-type]
total_asset=Decimal(amount)
)
assert data["tier"] == tier
assert len(data["benefits"]) == count
# 低等级条目必须在场(累积的直接证据)
assert any(b["tier"] == "gold" for b in data["benefits"])
@pytest.mark.asyncio
async def test_below_lowest_threshold_gets_no_benefits_but_keeps_upgrade_hint() -> None:
"""低于 50 万是"普通客户",不是错误:权益空,但仍告诉他要多少才升级。"""
session = FakeSession(ALL)
data = await CustomerBenefitService(session).entitlements( # type: ignore[arg-type]
total_asset=Decimal("100")
)
assert data["tier"] is None
assert data["tier_label"] == "普通客户"
assert data["benefits"] == []
assert data["next_tier"] == {
"tier": "gold", "tier_label": "金卡", "min_total_asset": "500000.00",
}
# 未达门槛时不应去查权益表
assert session.last_codes == ()
@pytest.mark.asyncio
async def test_top_tier_has_no_next_tier() -> None:
data = await CustomerBenefitService(FakeSession(ALL)).entitlements( # type: ignore[arg-type]
total_asset=Decimal("50000000")
)
assert data["tier"] == "private"
assert data["next_tier"] is None
# --- 鉴权 -------------------------------------------------------------------
class FakeCtx:
def __init__(self, permissions: set[str]) -> None:
self.user_id = "9102"
self.permissions = permissions
self.roles: set[str] = set()
self.portal = "api"
@pytest.mark.asyncio
async def test_missing_permission_is_denied_before_touching_data(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""权限检查必须在读取任何客户数据**之前**(避免"拒绝请求却已经查了库")。"""
service = CustomerBenefitService(FakeSession(ALL)) # type: ignore[arg-type]
called = {"n": 0}
async def spy(customer_id: int) -> Decimal | None:
called["n"] += 1
return Decimal("10000000")
monkeypatch.setattr(service, "_total_asset", spy)
ctx = RequestContext.model_construct(
user_id=9102, permissions=frozenset(), roles=frozenset(),
portal="api", trace_id="t", request_id="r",
)
with pytest.raises(ForbiddenAgentError):
await service.entitlements_for(ctx)
assert called["n"] == 0
+17 -1
View File
@@ -36,7 +36,23 @@ def created_tables(path: Path) -> set[str]:
def test_advisor_migrations_form_one_chain_from_qyqy_head() -> None:
script = ScriptDirectory.from_config(Config(str(ROOT / "alembic.ini")))
assert len(script.get_heads()) == 1
assert script.get_heads()[0] == "20260911_merge_adv_risk_heads"
# ⚠️ 这里**不写死 head 名**。原先断言
# `script.get_heads()[0] == "20260911_merge_adv_risk_heads"`
# 那是"投顾迁移刚加完那一刻"的快照;之后任何人新增迁移(2026-09-12 的
# 客户权益迁移 `20260912_customer_benefit` 即是一例)都会让本用例变红,
# 而变红的原因与投顾链的对错**无关** —— 断言测到的是时间,不是契约。
#
# 原意是"投顾链确实接在这条主链上、没另起一条分支"。改为断言
# **投顾链尾是当前 head 的祖先**:既保住这个意思,又不受后续迁移影响。
head = script.get_heads()[0]
ancestors = {rev.revision for rev in script.iterate_revisions(head, "base")}
advisor_tail = re.search(
r'revision = "([^"]+)"',
(VERSIONS / ADVISOR_FILES[-1]).read_text(encoding="utf-8"),
)
assert advisor_tail is not None
assert advisor_tail.group(1) in ancestors
first = (VERSIONS / ADVISOR_FILES[0]).read_text(encoding="utf-8")
assert 'down_revision = "20260910_drop_review_separation"' in first
+88
View File
@@ -0,0 +1,88 @@
"""数据订正:`profile_snapshots.current_customer_id` 与 `is_current` 对齐。
**默认 dry-run**,加 `--apply` 才提交。幂等,可重复运行。
### 背景
`current_customer_id` 是**普通可空列 + 唯一键** `uk_profile_snapshot_current`(不是生成列),
契约是「当前版本写客户 ID、历史版本写 NULL」。但写入方曾把它误当生成列,
两处路径都未赋值(`ProfileAssemblyService._write_snapshot`、
`CustomerProfileCandidateService._write_profile_snapshot`),造成两类脏数据:
1. 历史版本 `is_current=0` 却**仍留着** `current_customer_id`;
2. 当前版本 `is_current=1` 却是 NULL。
后果:唯一键形同虚设(多个 NULL 不冲突)⇒「每个客户最多一条 current」失效;
而且 (1) 的值一旦与 (2) 补上的值相同就会**撞唯一键**。
### 顺序要紧
先清 (1) 再补 (2):反过来的话,第 2 步写入时旧值还在,会直接撞键。
用法:
.\\.venv\\Scripts\\python.exe tools\\fix_profile_snapshot_current.py # 先看
.\\.venv\\Scripts\\python.exe tools\\fix_profile_snapshot_current.py --apply # 再改
"""
import asyncio
import sys
from sqlalchemy import text
from app.infrastructure.db import engine
#: (1) 历史版本不该持有该列 → 清空。必须先做。
_CLEAR_STALE = text("""
UPDATE profile_snapshots
SET current_customer_id = NULL
WHERE is_current = 0 AND current_customer_id IS NOT NULL
""")
#: (2) 当前版本必须持有该列 → 补上。放在 (1) 之后,避免撞唯一键。
_FILL_CURRENT = text("""
UPDATE profile_snapshots
SET current_customer_id = customer_id
WHERE is_current = 1 AND current_customer_id IS NULL
""")
_COUNT_STALE = text(
"SELECT COUNT(*) FROM profile_snapshots "
"WHERE is_current = 0 AND current_customer_id IS NOT NULL"
)
_COUNT_MISSING = text(
"SELECT COUNT(*) FROM profile_snapshots "
"WHERE is_current = 1 AND current_customer_id IS NULL"
)
async def main(apply: bool) -> int:
async with engine.connect() as conn:
stale = await conn.scalar(_COUNT_STALE)
missing = await conn.scalar(_COUNT_MISSING)
print(f"待清空(is_current=0 却留着值): {stale}")
print(f"待补写(is_current=1 却是 NULL): {missing}")
if not apply:
print("dry-run:未提交。确认无误后加 --apply 再运行。")
await conn.rollback()
return 0
cleared = (await conn.execute(_CLEAR_STALE)).rowcount
filled = (await conn.execute(_FILL_CURRENT)).rowcount
await conn.commit()
print(f"已清空 {cleared} 行、已补写 {filled} 行")
# 复核:两类异常都应归零
print("复核 待清空:", await conn.scalar(_COUNT_STALE))
print("复核 待补写:", await conn.scalar(_COUNT_MISSING))
print("\n各客户的当前快照:")
for row in (await conn.execute(text(
"SELECT customer_id, version, is_current, current_customer_id "
"FROM profile_snapshots ORDER BY customer_id, version"
))).mappings().all():
print(" ", dict(row))
return 0
if __name__ == "__main__":
sys.exit(asyncio.run(main("--apply" in sys.argv)))
+148
View File
@@ -0,0 +1,148 @@
"""客户权益目录种子:把《高净值客户服务规范》的权益转成表数据。
数据源:`knowledge/product/高净值客户服务规范.md` 第二章「各层级专属权益」。
**逐条照抄文档**,不改写、不新增文档里没有的权益。
执行:
python -X utf8 -m tools.seed_customer_benefits
幂等:按 `benefit_code` 先查后插;已存在的**不覆盖**(避免把人工调整冲掉)。
## 两处刻意的处理
1. **私行那条"7×24 小时私人银行专线:400-XXX-XXXX 转 8" 不写号码**。
文档里是占位符,而项目已有明确口径:对客号码的唯一来源是
`app/core/customer_service_rules.py` 的 `CONTACT_PHONE`
(本线此前修过"同一客服给客户两个不同号码"的缺陷,见 `docs/37` §6.3 的 A1)。
把占位符抄进库,等于又造了第二份假号码。故只保留权益名称与说明。
2. **权益按层级累积分组**:文档每层都写"含全部下级权益,新增以下"。
表里**只存该层新增的条目**,累积展开由 `CustomerBenefitService.tier_chain()` 完成
—— 否则同一条权益要在多个层级重复存,改一处漏三处。
"""
from __future__ import annotations
import asyncio
import sys
from datetime import UTC, datetime
from sqlalchemy import select
from app.infrastructure.db import SessionFactory
from app.model.benefit import CustomerBenefit
#: (层级, 类别, 名称, 说明);display_order 按此列表顺序自动编号。
#: 层级顺序 gold → platinum → diamond → private = 累积顺序。
BENEFITS: tuple[tuple[str, str, str, str], ...] = (
# ---- 金卡(50 万+)----
("gold", "financial", "专属理财经理服务", "由理财经理提供专属服务(1:N,N≤300)"),
("gold", "financial", "基金申购费率 5 折优惠", "高于普通客户的 1 折优惠"),
("gold", "financial", "银行理财专属高收益产品", "较公开产品收益高 10-20BP"),
("gold", "financial", "每月 1 次免费资产配置报告", "每月可获取一次资产配置报告"),
("gold", "financial", "优先认购热门基金产品", "热门基金产品优先认购"),
("gold", "non_financial", "生日祝福礼遇", "精美礼品一份"),
("gold", "non_financial", "节日关怀", "春节、中秋礼品卡"),
("gold", "non_financial", "APP 金卡专属标识", "客户端展示金卡专属标识"),
("gold", "non_financial", "财富中心 VIP 区域使用", "可使用财富中心 VIP 区域"),
# ---- 白金(200 万+,含全部金卡权益)----
("platinum", "financial", "1 对 1 高级理财经理服务", "由高级理财经理提供 1 对 1 服务(N≤150)"),
("platinum", "financial", "基金申购费率 3 折优惠", "较金卡的 5 折进一步优惠"),
("platinum", "financial", "每季度 1 次投资策略会/市场研判会", "每季度参与资格一次"),
("platinum", "financial", "专属理财产品", "白金客户专享,年化收益较普通产品高 20-30BP"),
("platinum", "financial", "私募产品优先认购权", "私募产品优先认购"),
("platinum", "financial", "基金投顾服务费 8 折优惠", "投顾服务费 8 折"),
("platinum", "non_financial", "每年 2 次高端客户沙龙", "品酒、艺术品鉴赏等"),
("platinum", "non_financial", "三甲医院专家门诊预约", "每年 2 次"),
("platinum", "non_financial", "机场贵宾厅服务", "每年 6 次"),
("platinum", "non_financial", "高尔夫球场预约优惠", "合作球场 8 折"),
("platinum", "non_financial", "子女留学规划咨询", "合作机构免费 1 次"),
# ---- 钻石(600 万+,含全部白金权益)----
("diamond", "financial", "资深客户经理 1 对 1 专属服务", "由资深客户经理提供(N≤80)"),
("diamond", "financial", "基金申购费率 2 折优惠", "较白金的 3 折进一步优惠"),
("diamond", "financial", "家族办公室初步服务对接", "家族办公室服务初步对接"),
("diamond", "financial", "全球资产配置咨询", "全球范围资产配置咨询"),
("diamond", "financial", "私募产品优先配置权", "含稀缺额度"),
("diamond", "financial", "定制化投资报告", "月度/季度定制报告"),
("diamond", "financial", "税务筹划初步咨询", "每年 1 次"),
("diamond", "non_financial", "高端健康管理", "年度全面体检套餐"),
("diamond", "non_financial", "每年 4 次高端客户活动", "米其林晚宴、私人音乐会等"),
("diamond", "non_financial", "全球紧急救援服务", "全球范围紧急救援"),
("diamond", "non_financial", "机场专车接送服务", "每年 8 次"),
("diamond", "non_financial", "高端酒店会员权益", "合作五星级酒店 VIP 待遇"),
("diamond", "non_financial", "子女实习/就业推荐", "合作企业资源对接"),
# ---- 私行(1000 万+,含全部钻石权益)----
("private", "financial", "私人银行家 1 对 1 专属服务", "由私人银行家提供(N≤40)"),
("private", "financial", "基金申购费率 1 折优惠", "最低费率档"),
("private", "financial", "家族信托设立与管理服务", "家族信托全流程服务"),
("private", "financial", "家族办公室全方位服务", "家族办公室全方位服务"),
("private", "financial", "全球资产配置方案", "含海外置业、移民咨询"),
("private", "financial", "专属投委会成员定期沟通", "与投委会成员定期沟通"),
("private", "financial", "私募股权/创投基金认购权", "私募股权与创投基金认购"),
("private", "financial", "企业融资顾问服务", "免费提供"),
("private", "financial", "定制化资产配置白皮书", "年度"),
("private", "financial", "税务筹划与遗产规划", "CFA/CTA 专家服务"),
("private", "financial", "艺术品投资咨询", "艺术品投资咨询"),
("private", "financial", "专属理财产品定制", "单户可定制产品方案"),
# ⚠️ 文档原文为「7×24 小时私人银行专线:400-XXX-XXXX 转 8」——
# 号码是占位符,此处**只保留权益名**,号码一律走 customer_service_rules.CONTACT_PHONE。
("private", "non_financial", "7×24 小时私人银行专线", "全天候私人银行专线(号码统一由客服热线配置提供)"),
("private", "non_financial", "私人银行家上门服务", "每月至少 1 次"),
("private", "non_financial", "全球顶尖医疗资源对接", "全球医疗资源对接"),
("private", "non_financial", "机场贵宾厅及专车接送不限次", "每年不限次"),
("private", "non_financial", "私人飞机/游艇租赁服务", "合作供应商优惠价"),
("private", "non_financial", "高端社交圈层活动", "南方私行俱乐部年会、海外游学"),
("private", "non_financial", "家族传承规划", "法律、税务、治理综合方案"),
("private", "non_financial", "公益慈善顾问服务", "慈善顾问服务"),
("private", "non_financial", "奢侈品鉴赏", "珠宝、名表、红酒私人顾问"),
)
def _code(tier: str, index: int) -> str:
return f"tier:{tier}:{index:02d}"
async def seed() -> int:
now = datetime.now(UTC).replace(tzinfo=None)
inserted = 0
async with SessionFactory() as session:
existing = set(await session.scalars(select(CustomerBenefit.benefit_code)))
order_by_tier: dict[str, int] = {}
for tier, category, name, description in BENEFITS:
order_by_tier[tier] = order_by_tier.get(tier, 0) + 1
code = _code(tier, order_by_tier[tier])
if code in existing:
continue
session.add(CustomerBenefit(
benefit_code=code,
customer_tier=tier,
category=category,
name=name,
description=description,
display_order=order_by_tier[tier],
status="active",
created_at=now,
updated_at=now,
))
inserted += 1
await session.commit()
return inserted
async def main() -> None:
inserted = await seed()
async with SessionFactory() as session:
rows = (await session.execute(
select(CustomerBenefit.customer_tier, CustomerBenefit.category)
)).all()
per_tier: dict[str, int] = {}
for tier, _category in rows:
per_tier[tier] = per_tier.get(tier, 0) + 1
print(f"新增 {inserted} 条;库中现有 {len(rows)} 条:")
for tier, _label in (("gold", "金卡"), ("platinum", "白金"), ("diamond", "钻石"), ("private", "私行")):
print(f" {tier:9} {_label} {per_tier.get(tier, 0)} 条")
if __name__ == "__main__":
sys.stdout.reconfigure(encoding="utf-8")
asyncio.run(main())
+37 -17
View File
@@ -26,7 +26,7 @@ import sys
from datetime import UTC, datetime, timedelta
from decimal import Decimal
from sqlalchemy import func, insert, select
from sqlalchemy import func, insert, select, update
from sqlalchemy.orm import Session
from app.core.config import get_settings
@@ -118,7 +118,32 @@ async def _upsert_product(session: Session, spec: dict) -> int:
async def _upsert_market_price(session: Session, product_id: int, spec: dict) -> None:
"""写入/刷新当日行情。
⚠️ 当天已有行时**必须刷新**,不能直接 return(2026-09-12 实测到的缺陷):
行情是**时效数据**,`FundQuoteService` 会按 `source_updated_at` 判新鲜度,
过期即返回 `503 FUND_QUOTE_UNAVAILABLE`,连带 `T001` 仪表盘与 `T006` 持仓一起不可用。
原先"已存在就跳过"会让同一天重跑**不更新 `source_updated_at`** ——
表现为"刚灌完种子能用,过十几分钟仪表盘就 503",而这与种子无关、极难归因。
幂等的正确含义是**不产生重复行**(`(product_id, trade_date)` 唯一),
**不是"不更新值"**。账户/持仓的"已存在则跳过"是另一回事 ——
那是业务数据,不该被种子覆盖。
"""
today = datetime.now(UTC).date()
now = datetime.now(UTC).replace(tzinfo=None)
values = {
"open_price": spec["close_price"],
"high_price": spec["close_price"] + Decimal("0.05"),
"low_price": spec["close_price"] - Decimal("0.05"),
"close_price": spec["close_price"],
"volume": Decimal("1000000"),
"turnover_amount": spec["close_price"] * Decimal("1000000"),
"total_fund_shares": spec["total_fund_shares"],
"source": "eastmoney_demo_seed",
"source_updated_at": now,
}
existing = (
await session.execute(
select(FundMarketPrice.id).where(
@@ -128,25 +153,20 @@ async def _upsert_market_price(session: Session, product_id: int, spec: dict) ->
)
).scalar_one_or_none()
if existing is not None:
await session.execute(
update(FundMarketPrice).where(FundMarketPrice.id == existing).values(**values)
)
return
now = datetime.now(UTC).replace(tzinfo=None)
next_id = await _next_id(session, FundMarketPrice)
stmt = insert(FundMarketPrice).values(
id=next_id,
product_id=product_id,
trade_date=today,
open_price=spec["close_price"],
high_price=spec["close_price"] + Decimal("0.05"),
low_price=spec["close_price"] - Decimal("0.05"),
close_price=spec["close_price"],
volume=Decimal("1000000"),
turnover_amount=spec["close_price"] * Decimal("1000000"),
total_fund_shares=spec["total_fund_shares"],
source="eastmoney_demo_seed",
source_updated_at=now,
created_at=now,
await session.execute(
insert(FundMarketPrice).values(
id=next_id,
product_id=product_id,
trade_date=today,
created_at=now,
**values,
)
)
await session.execute(stmt)
async def _upsert_account(session: Session, customer_id: int) -> FundSimAccount:
+6
View File
@@ -149,6 +149,10 @@ PERMISSIONS: tuple[tuple[int, str, str, str, str], ...] = (
(9063, "trade:order:cancel", "trade", "order", "cancel"),
(9064, "holding:read:self", "holding", "read", "self"),
(9065, "trade:txn:read", "trade", "txn", "read"),
# ---- 9066:客户权益(本线新增,续 9065)----
# 只读"我的层级与应享权益";层级由 `fin_customer_profile.total_asset` 实时判定,
# 数据范围恒为 self,故用 `read:self` 而不是 `read:customer`。
(9066, "benefit:read:self", "benefit", "read", "self"),
)
# 客户:业务侧自助能力(自己的会话、反馈、转人工、自己的记忆画像)。
@@ -159,6 +163,8 @@ CUSTOMER_PERMISSIONS = (
9044,
# ZSY §T:账户看板与场内模拟交易(首版仅 customer 角色可用,留 admin 全量)
9060, 9061, 9062, 9063, 9064, 9065,
# 客户权益:客户看自己的层级与应享权益(本线新增)
9066,
)
# 风控专员:业务侧只读 + 跨客户记忆 + 审计只读,不含配置写权限。
RISK_PERMISSIONS = (