Files
group_fqcd_jr/tools/seed_advisor_demo.py
lzf_0626 292acd2e2c 新增「补客户归属」工具,并订正一处因迁移而过期的注释
## 新增 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。
2026-09-14 21:46:29 +08:00

269 lines
11 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""投顾演示数据准备:能配的配好,配不了的**说清卡在哪**。
## 背景
投顾那条线合并进来后,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()))