Files
group_fqcd_jr/tools/assign_customer_scope.py
T
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

207 lines
9.6 KiB
Python
Raw 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.
"""给员工补一条「客户归属」(`sys_customer_assignment`),并当场验证它真的生效。
## 为什么需要它
`sys_customer_assignment` 是**全平台**的"谁负责哪个客户"口径,不是记忆私有:
- 记忆召回(`app/core/memory_scope.py`):员工只能读**归属给自己**的客户记忆;
- `memory:read:customer` / `investment-goal:*:customer`(`data_scope='own_customers'`):
授权判定要求目标客户落在 `context.customer_ids` 内;
- 投顾"已发布方案"的可见客户 = `{自己} ∪ 归属客户`(`product_recommendation_service`)。
所以"给某人补归属"是一句**授权动作**,不该用一次性 SQL 随手写进库 —— 这里把它做成
幂等、可复跑、默认只预览的脚本,并**当场用真实身份链路回验**。
## 用法
python tools/assign_customer_scope.py # 只预览:9002(风控) ← 9001
python tools/assign_customer_scope.py --apply # 真正写入
python tools/assign_customer_scope.py --show # 只读:现状 + 每个员工的可见范围
python tools/assign_customer_scope.py --employee 9020 --customer 9001 --apply
python tools/assign_customer_scope.py --employee 9002 --customers 12001 12002 --apply
## 两个已踩过的坑(本脚本已内置)
1. **`assigned_at` 不能写"当前时间"**:MySQL `DATETIME(0)` 会四舍五入到秒、可能落进未来,
于是 `assigned_at <= now` 判定"尚未生效",归属拿不到 —— **而且不报错**
(`identity_repository` 的过滤条件就是这么写的)。统一往前留 5 秒。
2. **`employee_role` 是 NOT NULL 且参与唯一键 `(customer_id, employee_role)`**:
同一个客户、同一个角色只能有一行生效归属。本脚本从 RBAC 反查员工真实角色码,
不靠手填,避免"角色名写错 → 撞唯一键或语义错位"。
`id` 由本脚本在 9901+ 段分配(避开业务号段);该列**已恢复 AUTO_INCREMENT**
(`alembic/versions/20260914_baseline_auto_increment.py`),显式给 id 只是为了可复现。
## 注意:这不是"只影响记忆"
给**投顾/管理员**补归属会**同时**解锁 `investment-goal:*:customer`(客户投资目标的
读/写/确认),并改变投顾"已发布方案"的可见客户 —— 那是设计用途,但确实是权限扩大。
给**风控**补归属则只影响记忆召回(它的其余权限都是 `data_scope='all'`,
`scope_from_context()` 见 `all` 直接放行、不看 `customer_ids`)。
"""
from __future__ import annotations
import argparse
import asyncio
import sys
from datetime import UTC, datetime, timedelta
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from sqlalchemy import text # noqa: E402
from app.core.contracts import RequestContext # noqa: E402
from app.core.memory_scope import ( # noqa: E402
REQUIRED_EMPLOYEE_PERMISSION,
customer_memory_scope,
)
from app.infrastructure.db import SessionFactory # noqa: E402
from app.repository.identity_repository import IdentityRepository # noqa: E402
#: id 号段起点:避开业务号段(与 `seed_advisor_demo.py` 的 9901 同段)。
ASSIGNMENT_ID_BASE = 9901
#: `assigned_at` 往前留的秒数,见模块 docstring 第 1 条坑。
BACKDATE_SECONDS = 5
async def _employee_role(session, employee_id: int) -> str | None:
"""从 RBAC 反查员工的角色码(`employee_role` 列要求 NOT NULL 且参与唯一键)。"""
return await session.scalar(
text(
"SELECT r.role_code FROM sys_user_role ur JOIN sys_role r ON r.id = ur.role_id "
"WHERE ur.user_id = :uid AND r.status = 'active' ORDER BY r.role_code LIMIT 1"
),
{"uid": employee_id},
)
async def _next_assignment_id(session) -> int:
"""在 9901+ 段分配下一个 id;被占用就顺延,避免与既有行撞主键。"""
current = await session.scalar(
text("SELECT MAX(id) FROM sys_customer_assignment WHERE id >= :base"),
{"base": ASSIGNMENT_ID_BASE},
)
return max(ASSIGNMENT_ID_BASE, int(current or 0) + 1)
async def show() -> None:
async with SessionFactory() as session:
rows = (await session.execute(text(
"SELECT id, customer_id, employee_id, employee_role, assigned_at, unassigned_at "
"FROM sys_customer_assignment ORDER BY employee_id, customer_id"
))).all()
print(f"sys_customer_assignment 现有 {len(rows)} 行:")
for row in rows:
state = "生效" if row[5] is None else f"已解除({row[5]})"
print(f" id={row[0]} 员工 {row[2]} → 客户 {row[1]} 角色={row[3]} "
f"{state} assigned_at={row[4]}")
employees = sorted({int(row[2]) for row in rows})
print("\n每个员工的**实际可见范围**(走真实身份链路):")
for employee_id in employees:
context = await IdentityRepository(session).load_context(
RequestContext(user_id=str(employee_id), trace_id="assign-scope-show")
)
scope = customer_memory_scope(context)
has_capability = REQUIRED_EMPLOYEE_PERMISSION in context.permissions
print(f" 员工 {employee_id} roles={context.roles} "
f"customer_ids={context.customer_ids}")
print(f" 记忆可读范围={scope} 持有 {REQUIRED_EMPLOYEE_PERMISSION}="
f"{has_capability}")
async def apply(employee_id: int, customer_ids: list[int]) -> int:
inserted = 0
async with SessionFactory() as session:
role = await _employee_role(session, employee_id)
if role is None:
print(f"[失败] 员工 {employee_id} 没有生效角色,先绑角色再补归属")
return 0
for customer_id in customer_ids:
exists = await session.scalar(
text(
"SELECT COUNT(*) FROM sys_customer_assignment "
"WHERE customer_id = :cid AND employee_role = :role AND unassigned_at IS NULL"
),
{"cid": customer_id, "role": role},
)
if exists:
print(f" 客户 {customer_id} ← 角色 {role} 的归属已存在,跳过(幂等)")
continue
next_id = await _next_assignment_id(session)
assigned_at = datetime.now(UTC).replace(tzinfo=None) - timedelta(
seconds=BACKDATE_SECONDS
)
# 不能再用 `async with session.begin()`:上面的校验查询已经开了隐式事务,
# 再 begin 会抛 `InvalidRequestError: A transaction is already begun`(实测踩到)。
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, :role, :assigned_at, NULL)
"""
),
{"id": next_id, "customer_id": customer_id, "employee_id": employee_id,
"role": role, "assigned_at": assigned_at},
)
await session.commit()
inserted += 1
print(f" 已写入 id={next_id}:员工 {employee_id}({role})→ 客户 {customer_id}")
return inserted
async def verify(employee_id: int) -> None:
"""当场用真实身份链路回验:身份解析 → 记忆可读范围。"""
async with SessionFactory() as session:
context = await IdentityRepository(session).load_context(
RequestContext(user_id=str(employee_id), trace_id="assign-scope-verify")
)
print("\n回验(`IdentityRepository.load_context` → `customer_memory_scope`):")
print(f" roles = {context.roles}")
print(f" customer_ids = {context.customer_ids}")
print(f" 记忆可读范围 = {customer_memory_scope(context)}")
if REQUIRED_EMPLOYEE_PERMISSION not in context.permissions:
print(f" ⚠️ 该身份**没有** {REQUIRED_EMPLOYEE_PERMISSION} 能力码 ⇒ 记忆仍然读不到"
f"(归属关系不等于授权,见 app/core/memory_scope.py)")
def _parse_ids(raw: list[str]) -> list[int]:
return [int(item) for item in raw]
def main() -> int:
parser = argparse.ArgumentParser(description="给员工补一条客户归属(幂等,默认只预览)")
parser.add_argument("--employee", type=int, default=9002, help="员工 id(默认 9002 风控)")
parser.add_argument("--customer", type=int, default=9001, help="客户 id(默认 9001)")
parser.add_argument("--customers", nargs="*", default=None,
help="一次补多个客户(给了它就忽略 --customer)")
parser.add_argument("--apply", action="store_true", help="真正写入;不加只预览")
parser.add_argument("--show", action="store_true", help="只读:打印现状与每个员工的可见范围")
args = parser.parse_args()
if args.show:
asyncio.run(show())
return 0
customers = _parse_ids(args.customers) if args.customers else [args.customer]
if not args.apply:
print("[预览] 将执行:")
print(f" 员工 {args.employee} ← 客户 {customers}"
f"(角色码从 RBAC 反查;assigned_at 往前留 {BACKDATE_SECONDS} 秒)")
print(" 加 --apply 才真正写入。")
return 0
inserted = asyncio.run(apply(args.employee, customers))
print(f"\n写入 {inserted} 行。")
asyncio.run(verify(args.employee))
return 0
if __name__ == "__main__":
raise SystemExit(main())