## 新增 tools/assign_customer_scope.py
`sys_customer_assignment` 是**全平台**的"谁负责哪个客户"口径,不是记忆私有:
记忆召回(app/core/memory_scope.py)、`memory:read:customer` /
`investment-goal:*:customer` 的 `own_customers` 授权判定、投顾"已发布方案"的可见客户
都读它。所以"给某人补归属"是一句**授权动作**,不该用一次性 SQL 随手写进库。
本工具把它做成幂等、可复跑、默认只预览(--apply 才写),并**当场用真实身份链路回验**
(IdentityRepository.load_context → customer_memory_scope),还会在目标账号缺少
`memory:read:customer` 时明确提示"归属关系不等于授权"。
内置两个已踩过的坑:
1. `assigned_at` 不能写"当前时间" —— MySQL DATETIME(0) 四舍五入到秒可能落进未来,
于是 `assigned_at <= now` 判定"尚未生效",归属拿不到**且不报错**;统一往前留 5 秒;
2. `employee_role` 是 NOT NULL 且参与唯一键 `(customer_id, employee_role)`,
故从 RBAC 反查员工真实角色码,不靠手填。
`id` 在 9901+ 段分配(避开业务号段)。另外修掉一个写入期 bug:校验查询已经开了隐式事务,
再 `async with session.begin()` 会抛 `InvalidRequestError: A transaction is already begun`。
用法:
python tools/assign_customer_scope.py # 预览:9002(风控) ← 9001
python tools/assign_customer_scope.py --apply
python tools/assign_customer_scope.py --show # 只读:每个员工的实际可见范围
## 订正 tools/seed_advisor_demo.py 的过期注释
该文件原写"`sys_customer_assignment.id` 没有自增(EXTRA 为空),必须自己给" ——
这条**已被 20260914_baseline_auto_increment 迁移改变**(实测 EXTRA 现为 auto_increment)。
注释改为记录现状,并说明仍然显式给 id 的理由(保持号段约定 + 让种子可复现,
自增会让同一份种子在不同环境落到不同 id)。
## 本轮实际执行
- 补入一行归属:员工 9002(risk_operator)→ 客户 9001(id=9902,幂等复跑正确跳过);
- 回验:风控 9002 的记忆可读范围从 `()` 变为 `(9001,)`,召回从 0 条变为 2 条;
投顾 9020 仍 2 条;
- 该动作**只影响记忆召回**:9002 其余权限都是 `data_scope='all'`,
`scope_from_context()` 见 all 直接放行、不看 customer_ids。
269 lines
11 KiB
Python
269 lines
11 KiB
Python
"""投顾演示数据准备:能配的配好,配不了的**说清卡在哪**。
|
||
|
||
## 背景
|
||
|
||
投顾那条线合并进来后,21 张 `advisor_*` 表全是空的,于是:
|
||
|
||
- `GET /api/v1/advisor/recommendations/published` 返回 `[]`
|
||
- 投顾登录后 `customer_ids` 为空(看不到任何客户)
|
||
- 客户画像的前置"完成开户风险测评问卷"未满足(9001 没有测评记录)
|
||
|
||
## 本脚本做什么
|
||
|
||
按依赖顺序准备**能自己造**的那几样:
|
||
|
||
1. **产品目录**:调 `tools/import_hq_test_products.py`(从行情源拉真实 ETF/LOF)。
|
||
这一步是**真实数据**,不是伪造。
|
||
2. **客户风险测评**(`fin_risk_assessment`):给演示客户 9001 补一条
|
||
**C5 激进型**、有效期一年。**这是演示数据**,脚本与 `answers` 字段里都写明。
|
||
3. **客户-投顾归属**(`sys_customer_assignment`):把 9001 指派给投顾 9020。
|
||
这决定投顾的 `customer_ids`,也决定 `data_scope` 下能看到谁。
|
||
|
||
## 本脚本**不**做什么,以及为什么
|
||
|
||
**"生成投顾方案"这一步配不了** —— 缺的是**证据来源**,不是技术。
|
||
|
||
`authoritative_tradable_products` 的文档写得很清楚:
|
||
|
||
Missing or unverified evidence excludes a product (fail closed).
|
||
|
||
而它要求的证据(`advisor_product_suitability_reference`)必须带上 **`source_url` 与
|
||
`document_sha256`** —— 即"这份适当性等级是从哪份文件来的"。要造这些行,就得有**真实的
|
||
销售适当性披露 / 基金合同文件**,也就是 `tools/import_product_governance_reference.py`
|
||
的 `--suitability` / `--contracts` 两个 CSV(它们**不在仓库里**,属环境数据)。
|
||
|
||
**本脚本不会去编 `source_url` / `document_sha256`**:那是金融项目里的证据链红线,
|
||
这道门挡住"没有依据的适当性等级"正是它该做的事。所以这里只**检查并报告**,
|
||
把真正的解法(拿到披露文件后跑哪条命令)打印出来。
|
||
|
||
## 用法
|
||
|
||
python tools/seed_advisor_demo.py --dry-run # 只看会写什么
|
||
python tools/seed_advisor_demo.py # 执行
|
||
|
||
# 方案那一步需要先准备好证据文件:
|
||
python tools/import_product_governance_reference.py --suitability <披露.csv> --contracts <合同.csv>
|
||
python tools/import_product_asset_classifications.py
|
||
python tools/sync_advisor_market_data.py
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import asyncio
|
||
import json
|
||
import subprocess
|
||
import sys
|
||
from datetime import UTC, datetime, timedelta
|
||
from pathlib import Path
|
||
|
||
from sqlalchemy import text
|
||
|
||
from app.infrastructure.db import SessionFactory
|
||
|
||
if hasattr(sys.stdout, "reconfigure"):
|
||
sys.stdout.reconfigure(errors="replace") # type: ignore[union-attr]
|
||
|
||
#: 演示客户与投顾(`sys_user`)。
|
||
DEMO_CUSTOMER_ID = 9001
|
||
DEMO_ADVISOR_ID = 9020
|
||
|
||
#: 演示测评的有效期。
|
||
ASSESSMENT_VALID_DAYS = 365
|
||
|
||
#: `sys_customer_assignment.id` 的自增**已于 2026-09-14 恢复**
|
||
#: (`alembic/versions/20260914_baseline_auto_increment.py`;此前 `information_schema`
|
||
#: 的 EXTRA 为空、不显式给 id 就报 `Field 'id' doesn't have a default value`)。
|
||
#: 这里仍然显式给 id:一是保持与既有行的号段约定(9901 起,避开业务号段),
|
||
#: 二是让脚本可复现 —— 自增会让同一份种子在不同环境落到不同 id 上。
|
||
ASSIGNMENT_ID = 9901
|
||
|
||
#: 与 `identity_repository` 的授权校验同源的坑:`assigned_at` 若用"当前时间"写,
|
||
#: MySQL 的 DATETIME(0) 会把它四舍五入到秒、可能落进未来,于是 `assigned_at <= now`
|
||
#: 判定"尚未生效",归属与角色都拿不到,且**不报错**。统一往前留 5 秒。
|
||
BACKDATE_SECONDS = 5
|
||
|
||
|
||
def _json(value: object) -> str:
|
||
return json.dumps(value, ensure_ascii=False)
|
||
|
||
|
||
def run_import(script: str, *, dry_run: bool) -> int:
|
||
"""跑投顾线自带的导入脚本(子进程,stdout/stderr 直接继承)。
|
||
|
||
刻意**不捕获输出**:一是保留它自己的日志格式,二是避免在受限沙箱里因为管道
|
||
被拒(`EPERM`)而把一个能跑的命令变成失败。
|
||
"""
|
||
args = [sys.executable, str(Path("tools") / script)]
|
||
if dry_run:
|
||
args.append("--dry-run")
|
||
print(f" $ {' '.join(args[1:])}")
|
||
return subprocess.run(args, check=False).returncode # noqa: S603 - 参数均为本模块常量
|
||
|
||
|
||
async def seed_assessment(*, dry_run: bool) -> None:
|
||
"""给演示客户补一条有效期内的风险测评。"""
|
||
now = datetime.now(UTC).replace(tzinfo=None)
|
||
assess = {
|
||
"customer_id": DEMO_CUSTOMER_ID,
|
||
"questionnaire_version": "demo-v1",
|
||
"total_score": 50,
|
||
"investor_type": "C5",
|
||
"assessed_at": now,
|
||
"valid_until": now + timedelta(days=ASSESSMENT_VALID_DAYS),
|
||
"created_at": now,
|
||
}
|
||
answers = {
|
||
"_note": "演示数据:由 tools/seed_advisor_demo.py 生成,不是真实测评结果",
|
||
"score": 50,
|
||
}
|
||
async with SessionFactory() as session, session.begin():
|
||
existing = await session.scalar(
|
||
text(
|
||
"SELECT COUNT(*) FROM fin_risk_assessment "
|
||
"WHERE customer_id = :cid AND questionnaire_version = :v"
|
||
),
|
||
{"cid": DEMO_CUSTOMER_ID, "v": assess["questionnaire_version"]},
|
||
)
|
||
if existing:
|
||
print(f" 测评:已存在(customer_id={DEMO_CUSTOMER_ID}),跳过")
|
||
return
|
||
if dry_run:
|
||
print(f" 测评:将插入 C5 / 有效期至 {assess['valid_until']:%Y-%m-%d}")
|
||
return
|
||
# `id` 无自增,由库内最大值推进。
|
||
next_id = int(
|
||
await session.scalar(
|
||
text("SELECT COALESCE(MAX(id), 0) + 1 FROM fin_risk_assessment")
|
||
)
|
||
or 1
|
||
)
|
||
await session.execute(
|
||
text(
|
||
"""
|
||
INSERT INTO fin_risk_assessment
|
||
(id, customer_id, questionnaire_version, answers, total_score,
|
||
investor_type, assessed_at, valid_until, created_at)
|
||
VALUES
|
||
(:id, :customer_id, :questionnaire_version, :answers, :total_score,
|
||
:investor_type, :assessed_at, :valid_until, :created_at)
|
||
"""
|
||
),
|
||
{"id": next_id, **assess, "answers": _json(answers)},
|
||
)
|
||
print(f" 测评:已插入 id={next_id}(C5,有效期 {ASSESSMENT_VALID_DAYS} 天)")
|
||
|
||
|
||
async def seed_assignment(*, dry_run: bool) -> None:
|
||
"""把演示客户指派给演示投顾。"""
|
||
now = datetime.now(UTC).replace(tzinfo=None)
|
||
assigned_at = now - timedelta(seconds=BACKDATE_SECONDS)
|
||
async with SessionFactory() as session, session.begin():
|
||
exists = await session.scalar(
|
||
text(
|
||
"SELECT COUNT(*) FROM sys_customer_assignment "
|
||
"WHERE customer_id = :cid AND employee_role = 'advisor' "
|
||
"AND unassigned_at IS NULL"
|
||
),
|
||
{"cid": DEMO_CUSTOMER_ID},
|
||
)
|
||
if exists:
|
||
print(f" 归属:已存在(客户 {DEMO_CUSTOMER_ID} 已指派给投顾),跳过")
|
||
return
|
||
if dry_run:
|
||
print(f" 归属:将把客户 {DEMO_CUSTOMER_ID} 指派给投顾 {DEMO_ADVISOR_ID}")
|
||
return
|
||
await session.execute(
|
||
text(
|
||
"""
|
||
INSERT INTO sys_customer_assignment
|
||
(id, customer_id, employee_id, employee_role, assigned_at, unassigned_at)
|
||
VALUES
|
||
(:id, :customer_id, :employee_id, 'advisor', :assigned_at, NULL)
|
||
"""
|
||
),
|
||
{"id": ASSIGNMENT_ID, "customer_id": DEMO_CUSTOMER_ID,
|
||
"employee_id": DEMO_ADVISOR_ID, "assigned_at": assigned_at},
|
||
)
|
||
print(
|
||
f" 归属:客户 {DEMO_CUSTOMER_ID} → 投顾 {DEMO_ADVISOR_ID}"
|
||
f"(assigned_at 往前留 {BACKDATE_SECONDS} 秒)"
|
||
)
|
||
|
||
|
||
async def evidence_gate_report() -> bool:
|
||
"""检查方案那一步的证据门是否已满足;True 表示可以继续。"""
|
||
async with SessionFactory() as session:
|
||
products = int(await session.scalar(text("SELECT COUNT(*) FROM fin_product")) or 0)
|
||
suitability = int(
|
||
await session.scalar(
|
||
text(
|
||
"SELECT COUNT(*) FROM advisor_product_suitability_reference "
|
||
"WHERE review_status='verified' AND source_url <> '' AND document_sha256 <> ''"
|
||
)
|
||
)
|
||
or 0
|
||
)
|
||
goals = int(
|
||
await session.scalar(text("SELECT COUNT(*) FROM advisor_investment_goal")) or 0
|
||
)
|
||
|
||
print("\n方案那一步的前置检查:")
|
||
print(f" 在售产品目录 fin_product : {products}")
|
||
print(f" 带来源证据的适当性参考(verified + 来源非空) : {suitability}")
|
||
print(f" 投资目标 advisor_investment_goal : {goals}")
|
||
|
||
if suitability == 0:
|
||
print(
|
||
"\n[结论] '生成投顾方案' 这一步**现在配不了**,缺的不是技术而是**证据来源**:\n"
|
||
" authoritative_tradable_products 是 fail closed 的,没有 verified 的\n"
|
||
" 适当性参考就会排除全部产品;而造那些行必须带上 source_url 与\n"
|
||
" document_sha256 —— 也就是真实的销售适当性披露 / 基金合同文件。\n"
|
||
" 本脚本**不会伪造**这两个字段(证据链红线)。\n"
|
||
"\n 拿到文件后按这个顺序跑:\n"
|
||
" python tools/import_product_governance_reference.py \\\n"
|
||
" --suitability <销售适当性披露.csv> --contracts <基金合同.csv>\n"
|
||
" python tools/import_product_asset_classifications.py\n"
|
||
" python tools/sync_advisor_market_data.py\n"
|
||
" 之后再跑本脚本,方案那条链路就能走通。"
|
||
)
|
||
return False
|
||
if goals == 0:
|
||
print(
|
||
"\n[结论] 证据门已满足,但还没有**投资目标**(advisor_investment_goal)。\n"
|
||
" generate 需要 goal.investment_horizon_months,先通过投顾的\n"
|
||
" 投资目标入口收集一份,再回来生成方案。"
|
||
)
|
||
return False
|
||
print("\n[结论] 证据门与投资目标都已就绪,可以生成投顾方案了。")
|
||
return True
|
||
|
||
|
||
async def main() -> int:
|
||
parser = argparse.ArgumentParser(description="投顾演示数据准备")
|
||
parser.add_argument("--dry-run", action="store_true", help="只看会写什么")
|
||
args = parser.parse_args()
|
||
|
||
print("[1/4] 产品目录(调用投顾自带的导入脚本,取真实行情数据)")
|
||
code = run_import("import_hq_test_products.py", dry_run=args.dry_run)
|
||
if code != 0:
|
||
print(f" [警告] 导入脚本返回 {code},产品目录可能没准备好(继续后续步骤)")
|
||
|
||
print("\n[2/4] 客户风险测评")
|
||
await seed_assessment(dry_run=args.dry_run)
|
||
|
||
print("\n[3/4] 客户-投顾归属")
|
||
await seed_assignment(dry_run=args.dry_run)
|
||
|
||
print("\n[4/4] 检查'生成投顾方案'的前置是否满足")
|
||
if args.dry_run:
|
||
print(" (dry-run 跳过:该检查依赖真实数据)")
|
||
return 0
|
||
ready = await evidence_gate_report()
|
||
print("\n完成。" + ("" if ready else " 方案那一步见上面的说明。"))
|
||
return 0
|
||
|
||
|
||
if __name__ == "__main__":
|
||
sys.exit(asyncio.run(main()))
|