Merge branch 'develop' of http://47.106.207.27:3000/AI260626/Mutual_Fund into develop_feature_risk

# Conflicts:
#	model/audit_log.py
#	model/sys_message.py
This commit is contained in:
zhangyongcai
2026-09-13 15:23:38 +08:00
49 changed files with 3386 additions and 14 deletions
@@ -0,0 +1,23 @@
"""customer_relation 客户-投顾关系表 ORM(签约状态 unsigned/signed/closed,工作台唯一写入方)。"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import BigInteger, DateTime, String, func
from sqlalchemy.orm import Mapped, mapped_column
from model.base import Base
class CustomerRelation(Base):
__tablename__ = "customer_relation"
__table_args__ = {"comment": "客户-投顾关系表(状态驱动:签约后投顾Agent方案正式触达客户)"}
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
customer_id: Mapped[int] = mapped_column(BigInteger)
advisor_id: Mapped[int] = mapped_column(BigInteger)
assign_time: Mapped[datetime] = mapped_column(DateTime, server_default=func.now())
signed_time: Mapped[datetime | None] = mapped_column(DateTime)
end_time: Mapped[datetime | None] = mapped_column(DateTime)
status: Mapped[str] = mapped_column(String(16), server_default="unsigned")
reason: Mapped[str | None] = mapped_column(String(128))
@@ -0,0 +1,142 @@
"""customer_relation 客户-投顾关系仓储:数据权限根(仅本人名下客户)+ 客户列表 + AUM 聚合。"""
from __future__ import annotations
from decimal import Decimal
from sqlalchemy import func, or_, select
from model.customer_relation import CustomerRelation
from model.fin_customer_profile import FinCustomerProfile
from model.fin_holdings import FinHoldings
from model.sys_user import SysUser
from repositories.base import BaseRepository
# 持仓中状态:仅统计「持有中」市值,与 service/holdings.py 口径一致(DDL 注释,未入 common_const)
_HOLDING_STATUS = "持有中"
class CustomerRelationRepo(BaseRepository):
model = CustomerRelation
async def get_by_customer_advisor(
self, customer_id: int, advisor_id: int
) -> CustomerRelation | None:
"""取指定客户与投顾的关系(数据权限判断用)。"""
return await self.db.scalar(
select(CustomerRelation).where(
CustomerRelation.customer_id == customer_id,
CustomerRelation.advisor_id == advisor_id,
)
)
async def get_current_by_customer(self, customer_id: int) -> CustomerRelation | None:
"""取客户当前有效关系(历史重分配可能多行,取最近一条,按 id 倒序)。"""
return await self.db.scalar(
select(CustomerRelation)
.where(CustomerRelation.customer_id == customer_id)
.order_by(CustomerRelation.id.desc())
.limit(1)
)
async def list_by_status(self, status: str) -> list[CustomerRelation]:
"""按状态取全部关系(定时调度遍历 signed 客户用,跨投顾)。"""
return list(
(
await self.db.scalars(
select(CustomerRelation).where(CustomerRelation.status == status)
)
).all()
)
async def list_by_advisor(
self, advisor_id: int, status: str | None = None
) -> list[CustomerRelation]:
"""名下客户关系列表(status=None 不过滤)。"""
stmt = select(CustomerRelation).where(CustomerRelation.advisor_id == advisor_id)
if status:
stmt = stmt.where(CustomerRelation.status == status)
return list((await self.db.scalars(stmt)).all())
async def list_customer_rows(
self,
*,
advisor_id: int,
status: str | None = None,
keyword: str | None = None,
limit: int = 100,
offset: int = 0,
) -> list[tuple[CustomerRelation, SysUser, FinCustomerProfile | None]]:
"""客户列表:关系 + 客户账号 + 画像(左连),返回三元组供 service 拼装响应。"""
conds = [CustomerRelation.advisor_id == advisor_id]
if status:
conds.append(CustomerRelation.status == status)
if keyword:
like = f"%{keyword}%"
conds.append(or_(SysUser.real_name.like(like), SysUser.phone.like(like)))
stmt = (
select(CustomerRelation, SysUser, FinCustomerProfile)
.join(SysUser, SysUser.id == CustomerRelation.customer_id)
.outerjoin(
FinCustomerProfile,
FinCustomerProfile.customer_id == CustomerRelation.customer_id,
)
.where(*conds)
.order_by(CustomerRelation.id.desc())
.limit(limit)
.offset(offset)
)
return list((await self.db.execute(stmt)).all())
async def count_customer_rows(
self,
*,
advisor_id: int,
status: str | None = None,
keyword: str | None = None,
) -> int:
conds = [CustomerRelation.advisor_id == advisor_id]
if status:
conds.append(CustomerRelation.status == status)
if keyword:
like = f"%{keyword}%"
conds.append(or_(SysUser.real_name.like(like), SysUser.phone.like(like)))
stmt = (
select(func.count())
.select_from(CustomerRelation)
.join(SysUser, SysUser.id == CustomerRelation.customer_id)
.where(*conds)
)
return (await self.db.scalar(stmt)) or 0
async def risk_level_distribution(
self, advisor_id: int
) -> list[tuple[str | None, int]]:
"""名下客户风险等级分布(C1-C5 各档人数),供驾驶舱客户分层。"""
stmt = (
select(FinCustomerProfile.risk_level, func.count())
.select_from(CustomerRelation)
.join(
FinCustomerProfile,
FinCustomerProfile.customer_id == CustomerRelation.customer_id,
)
.where(CustomerRelation.advisor_id == advisor_id)
.group_by(FinCustomerProfile.risk_level)
)
return list((await self.db.execute(stmt)).all())
async def sum_holdings_value(
self, advisor_id: int, relation_status: str | None = None
) -> Decimal:
"""AUM:名下客户「持有中」持仓当前市值之和(relation_status 过滤,如 signed)。"""
stmt = (
select(func.coalesce(func.sum(FinHoldings.current_value), 0))
.select_from(CustomerRelation)
.join(FinHoldings, FinHoldings.customer_id == CustomerRelation.customer_id)
.where(
CustomerRelation.advisor_id == advisor_id,
FinHoldings.status == _HOLDING_STATUS,
)
)
if relation_status:
stmt = stmt.where(CustomerRelation.status == relation_status)
return await self.db.scalar(stmt)
+14
View File
@@ -0,0 +1,14 @@
"""投顾工作台路由聚合(前缀 /api/advisor,在 api/router.py 以 /api 挂载)。"""
from fastapi import APIRouter
from api.advisor import audit, customers, dashboard, diagnosis, drafts, report, todos, visits
router = APIRouter(prefix="/advisor", tags=["投顾工作台"])
router.include_router(dashboard.router)
router.include_router(customers.router)
router.include_router(diagnosis.router)
router.include_router(drafts.router)
router.include_router(todos.router)
router.include_router(visits.router)
router.include_router(audit.router)
router.include_router(report.router)
+13
View File
@@ -0,0 +1,13 @@
"""路由层透传信息提取:投顾 JWT 与 trace_id(供 Agent 客户端透传)。"""
from __future__ import annotations
from fastapi import Request
from utils.request_id import get_request_id, new_request_id
def extract_auth(request: Request) -> tuple[str, str]:
"""返回 (Authorization 头原值, trace_id);trace_id 优先透传,无则生成。"""
auth = request.headers.get("Authorization", "")
trace_id = request.headers.get("X-Trace-Id") or get_request_id() or new_request_id()
return auth, trace_id
+52
View File
@@ -0,0 +1,52 @@
"""审计台账路由(业务在 service/advisor/audit.py)。导出为 CSV 下载。"""
from datetime import datetime
from fastapi import APIRouter, Depends, Query
from fastapi.responses import Response
from sqlalchemy.ext.asyncio import AsyncSession
from api.deps import require_advisor
from config.deps import get_db
from model.sys_user import SysUser
from service.advisor import audit as audit_service
from utils.response import success
router = APIRouter()
@router.get("/audit/ledger", summary="本人审计台账筛选")
async def ledger(
action: str | None = Query(None, max_length=64),
keyword: str | None = Query(None, max_length=128),
start: datetime | None = Query(None),
end: datetime | None = Query(None),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(
await audit_service.list_ledger(
db, user, action=action, keyword=keyword, start=start, end=end,
page=page, page_size=page_size,
)
)
@router.get("/audit/export", summary="导出本人审计台账(CSV)")
async def export(
action: str | None = Query(None, max_length=64),
keyword: str | None = Query(None, max_length=128),
start: datetime | None = Query(None),
end: datetime | None = Query(None),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
csv_text = await audit_service.export_ledger(
db, user, action=action, keyword=keyword, start=start, end=end
)
return Response(
content=csv_text,
media_type="text/csv; charset=utf-8",
headers={"Content-Disposition": 'attachment; filename="advisor_audit.csv"'},
)
+72
View File
@@ -0,0 +1,72 @@
"""客户 360 路由(业务在 service/advisor/customers.py,路由只做编排)。"""
from fastapi import APIRouter, Depends, Query
from sqlalchemy.ext.asyncio import AsyncSession
from api.deps import require_advisor
from config.deps import get_db
from model.sys_user import SysUser
from schemas.advisor import RelationReq
from service.advisor import customers as customers_service
from utils.response import success
router = APIRouter()
@router.get("/customers", summary="名下客户列表(只读 + 脱敏)")
async def list_customers(
status: str | None = Query(None, max_length=16),
keyword: str | None = Query(None, max_length=64),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(
await customers_service.list_customers(
db, user, status=status, keyword=keyword, page=page, page_size=page_size
)
)
@router.get("/customers/{customer_id}", summary="客户详情(unmask=true 查看完整手机号并留痕)")
async def get_customer(
customer_id: int,
unmask: bool = Query(False),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(await customers_service.get_customer(db, user, customer_id, unmask=unmask))
@router.get("/customers/{customer_id}/holdings", summary="客户持仓(只读)")
async def get_holdings(
customer_id: int,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(await customers_service.get_customer_holdings(db, user, customer_id))
@router.get("/customers/{customer_id}/reports", summary="客户历史建议报告")
async def get_reports(
customer_id: int,
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(
await customers_service.get_customer_reports(
db, user, customer_id, page=page, page_size=page_size
)
)
@router.post("/customers/{customer_id}/relation", summary="签约/结束服务(工作台唯一写入方)")
async def update_relation(
customer_id: int,
req: RelationReq,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(await customers_service.update_relation(db, user, customer_id, req))
+19
View File
@@ -0,0 +1,19 @@
"""驾驶舱路由(业务在 service/advisor/dashboard.py,路由只做编排)。"""
from fastapi import APIRouter, Depends
from sqlalchemy.ext.asyncio import AsyncSession
from api.deps import require_advisor
from config.deps import get_db
from model.sys_user import SysUser
from service.advisor import dashboard as dashboard_service
from utils.response import success
router = APIRouter()
@router.get("/dashboard", summary="首页数据与待办聚合")
async def dashboard(
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(await dashboard_service.get_dashboard(db, user))
+46
View File
@@ -0,0 +1,46 @@
"""资产配置与组合诊断路由(业务在 service/advisor/diagnosis.py,路由只做编排)。"""
from fastapi import APIRouter, Depends, Query
from sqlalchemy.ext.asyncio import AsyncSession
from api.deps import require_advisor
from config.deps import get_db
from model.sys_user import SysUser
from service.advisor import diagnosis as diagnosis_service
from utils.response import success
router = APIRouter()
@router.get("/diagnosis/{customer_id}", summary="持仓诊断(分布 + 集中度,只读)")
async def diagnose(
customer_id: int,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(await diagnosis_service.diagnose(db, user, customer_id))
@router.get("/strategies", summary="标准策略库(只读)")
async def strategies(
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(await diagnosis_service.list_strategies(db))
@router.get("/funds", summary="白名单产品筛选")
async def funds(
page: int = Query(1, ge=1),
page_size: int = Query(10, ge=1, le=100),
keyword: str | None = Query(None, max_length=64),
product_type: str | None = Query(None, max_length=32),
risk_level: str | None = Query(None, max_length=8),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(
await diagnosis_service.list_funds(
db, page=page, page_size=page_size, keyword=keyword,
product_type=product_type, risk_level=risk_level,
)
)
+99
View File
@@ -0,0 +1,99 @@
"""草稿/调仓/发送路由(业务在 service/advisor/drafts.py,路由只做编排 + 透传 JWT/trace_id)。"""
from fastapi import APIRouter, Depends, Query, Request
from sqlalchemy.ext.asyncio import AsyncSession
from api.advisor._auth import extract_auth
from api.deps import require_advisor
from config.deps import get_db
from model.sys_user import SysUser
from schemas.advisor import DraftSaveReq, RebalanceRunReq, TalkScriptReq
from service.advisor import drafts as drafts_service
from utils.response import success
router = APIRouter()
@router.get("/drafts", summary="草稿列表(代理 Agent draft/list)")
async def list_drafts(
request: Request,
customer_id: int | None = Query(None),
status: str | None = Query(None, max_length=16),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
auth, trace_id = extract_auth(request)
return success(
await drafts_service.list_drafts(
db, user, auth_header=auth, trace_id=trace_id,
customer_id=customer_id, status=status, page=page, page_size=page_size,
)
)
@router.get("/drafts/{draft_id}", summary="草稿详情(代理 Agent draft/{id})")
async def get_draft(
draft_id: str,
request: Request,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
auth, trace_id = extract_auth(request)
return success(await drafts_service.get_draft(db, user, auth_header=auth, trace_id=trace_id, draft_id=draft_id))
@router.put("/drafts/{draft_id}/save", summary="编辑保存(代理 Agent + 本地镜像快照)")
async def save_draft(
draft_id: str,
req: DraftSaveReq,
request: Request,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
auth, trace_id = extract_auth(request)
return success(await drafts_service.save_draft(db, user, auth_header=auth, trace_id=trace_id, draft_id=draft_id, req=req))
@router.post("/drafts/{draft_id}/discard", summary="废弃草稿(回写 Agent)")
async def discard_draft(
draft_id: str,
request: Request,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
auth, trace_id = extract_auth(request)
return success(await drafts_service.discard_draft(db, user, auth_header=auth, trace_id=trace_id, draft_id=draft_id))
@router.post("/drafts/{draft_id}/send", summary="发送终审 + 手动发送(工作台本地)")
async def send_draft(
draft_id: str,
request: Request,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
auth, trace_id = extract_auth(request)
return success(await drafts_service.send_draft(db, user, auth_header=auth, trace_id=trace_id, draft_id=draft_id))
@router.post("/rebalance/run", summary="手动触发调仓再平衡(异步受理)")
async def run_rebalance(
req: RebalanceRunReq,
request: Request,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
auth, trace_id = extract_auth(request)
return success(await drafts_service.run_rebalance(db, user, auth_header=auth, trace_id=trace_id, req=req))
@router.post("/talk-script", summary="生成沟通话术草稿(同步,超时降级)")
async def talk_script(
req: TalkScriptReq,
request: Request,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
auth, trace_id = extract_auth(request)
return success(await drafts_service.generate_talk_script(db, user, auth_header=auth, trace_id=trace_id, req=req))
+19
View File
@@ -0,0 +1,19 @@
"""个人绩效报表路由(业务在 service/advisor/report.py)。"""
from fastapi import APIRouter, Depends
from sqlalchemy.ext.asyncio import AsyncSession
from api.deps import require_advisor
from config.deps import get_db
from model.sys_user import SysUser
from service.advisor import report as report_service
from utils.response import success
router = APIRouter()
@router.get("/report/personal", summary="个人绩效统计")
async def personal_report(
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(await report_service.personal_report(db, user))
+38
View File
@@ -0,0 +1,38 @@
"""投顾待办路由(业务在 service/advisor/todos.py,路由只做编排)。"""
from fastapi import APIRouter, Depends, Query
from sqlalchemy.ext.asyncio import AsyncSession
from api.deps import require_advisor
from config.deps import get_db
from model.sys_user import SysUser
from schemas.advisor import TodoHandleReq
from service.advisor import todos as todos_service
from utils.response import success
router = APIRouter()
@router.get("/todos", summary="本人待办列表")
async def list_todos(
status: str | None = Query(None, max_length=16),
todo_type: str | None = Query(None, max_length=32),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(
await todos_service.list_todos(
db, user, status=status, todo_type=todo_type, page=page, page_size=page_size
)
)
@router.post("/todos/{todo_id}/handle", summary="处理待办(process/done)")
async def handle_todo(
todo_id: int,
req: TodoHandleReq,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(await todos_service.handle_todo(db, user, todo_id, req))
+56
View File
@@ -0,0 +1,56 @@
"""投后服务路由(回访/话术库/触达日志,业务在 service/advisor/visits.py)。"""
from fastapi import APIRouter, Depends, Query
from sqlalchemy.ext.asyncio import AsyncSession
from api.deps import require_advisor
from config.deps import get_db
from model.sys_user import SysUser
from schemas.advisor import VisitCreateReq
from service.advisor import visits as visits_service
from utils.response import success
router = APIRouter()
@router.get("/visits", summary="回访记录列表")
async def list_visits(
customer_id: int | None = Query(None),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(
await visits_service.list_visits(
db, user, customer_id=customer_id, page=page, page_size=page_size
)
)
@router.post("/visits", summary="新增回访留痕(不触发 Agent 记忆抽取)")
async def create_visit(
req: VisitCreateReq,
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(await visits_service.create_visit(db, user, req))
@router.get("/talk-templates", summary="合规话术库(内置,投顾参考)")
async def talk_templates(user: SysUser = Depends(require_advisor)):
return success(visits_service.list_talk_templates())
@router.get("/touch-logs", summary="触达留痕(站内信 + 回访)")
async def touch_logs(
customer_id: int = Query(...),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
user: SysUser = Depends(require_advisor),
db: AsyncSession = Depends(get_db),
):
return success(
await visits_service.list_touch_logs(
db, user, customer_id, page=page, page_size=page_size
)
)
+183
View File
@@ -0,0 +1,183 @@
"""公共常量定义(单一事实来源,与《投顾工作台开发文档/common_const.md》v1.2 对齐)。
用途:投顾工作台(B 端)与投顾Agent(AI 组件)统一引用,避免枚举/字符串硬编码不一致。
规范:业务代码禁止写死字符串字面量,一律 `from common_const import ...`。
说明:本文件只存放「常量字符串 / 枚举 / 模板文本」;业务规则(如适当性匹配逻辑)
写在业务代码中(见 service/advisor/suitability.py),不在此处写业务逻辑。
"""
from __future__ import annotations
# ---------------------------------------------------------------------------
# §1 客户-投顾关系状态(customer_relation.status)
# ---------------------------------------------------------------------------
CUSTOMER_REL_STATUS_UNSIGNED = "unsigned" # 未签约,仅分配,未开通投顾正式服务
CUSTOMER_REL_STATUS_SIGNED = "signed" # 已签约,可下发推荐/调仓方案给客户
CUSTOMER_REL_STATUS_CLOSED = "closed" # 服务已终止
# ---------------------------------------------------------------------------
# §2 草稿状态(advisor_draft.status,Agent 侧;不可物理删除,仅状态流转)
# ---------------------------------------------------------------------------
DRAFT_STATUS_DRAFT = "draft"
DRAFT_STATUS_DISCARDED = "discarded"
# 建议报告本地状态(advisor_report.send_status,工作台本地;sent 为工作台独有,Agent 不感知)
REPORT_SEND_STATUS_DRAFT = "draft"
REPORT_SEND_STATUS_SENT = "sent"
REPORT_SEND_STATUS_DISCARDED = "discarded"
# ---------------------------------------------------------------------------
# §3 记忆单元类型(memory_unit.info_type;工作台只读,不直接读写该表)
# ---------------------------------------------------------------------------
MEMORY_INFO_TYPE_FACT = "FACT" # 客观事实,作为业务硬规则依据
MEMORY_INFO_TYPE_OPINION = "OPINION" # 主观观点,仅做排序参考,不能绕过适当性
# ---------------------------------------------------------------------------
# §4 Agent 意图枚举(intent)
# ---------------------------------------------------------------------------
AGENT_INTENT_RECOMMEND = "recommend" # 基金推荐
AGENT_INTENT_REBALANCE = "rebalance" # 持仓诊断 & 调仓再平衡
AGENT_INTENT_FUND_ANALYSIS = "fund_analysis" # 基金深度分析
AGENT_INTENT_DIALOGUE_SCRIPT = "dialogue-script" # 生成沟通话术
# ---------------------------------------------------------------------------
# §5 沟通话术场景(generate-talk-script 入参 scene_type)
# ---------------------------------------------------------------------------
TALK_SCENE_RISK_BLOCK_ORDER = "risk_block_order" # 订单被风控拦截
TALK_SCENE_MARKET_FLUCTUATION = "market_fluctuation" # 市场波动安抚
TALK_SCENE_PORTFOLIO_DIVERGENCE = "portfolio_divergence" # 组合大幅偏离基准
TALK_SCENE_CUSTOMER_COMPLAINT = "customer_complaint" # 客户投诉
# ---------------------------------------------------------------------------
# §6 Agent 业务错误码(Agent 自有码域,成功码=0;工作台 service 层须映射,不得透传)
# ---------------------------------------------------------------------------
ERR_CODE_OK = 0
ERR_CODE_FORBIDDEN_CUSTOMER = 40001 # customer_id 不属于当前投顾,越权访问
ERR_CODE_SUITABILITY_INVALID = 40020 # 适当性校验不通过
ERR_CODE_NOT_SIGNED_REBALANCE = 40030 # 客户未签约,禁止生成 rebalance 草稿
ERR_CODE_DRAFT_NOT_FOUND = 40401 # draft_id 不存在或已废弃
ERR_CODE_LLM_ERROR = 50001 # Agent 内部 LLM 调用异常
ERR_CODE_GRAPH_ERROR = 50002 # GraphRAG 查询异常,触发降级
# Agent 错误码 → 默认提示文案(业务判断必须用 code,禁止依赖 message 文案)
AGENT_ERR_MESSAGE = {
ERR_CODE_FORBIDDEN_CUSTOMER: "无权操作该客户数据",
ERR_CODE_SUITABILITY_INVALID: "方案适当性校验不通过,包含超出客户风险等级的产品",
ERR_CODE_NOT_SIGNED_REBALANCE: "客户尚未签约,禁止生成调仓草稿",
ERR_CODE_DRAFT_NOT_FOUND: "草稿不存在或者已废弃",
ERR_CODE_LLM_ERROR: "AI服务调用异常,请稍后重试",
ERR_CODE_GRAPH_ERROR: "图谱查询异常,已降级返回部分结果",
}
# ---------------------------------------------------------------------------
# §7 风险等级(存储编码统一 C1-C5 / R1-R5;中文仅展示层,不入库)
# ---------------------------------------------------------------------------
C_RISK_C1, C_RISK_C2, C_RISK_C3, C_RISK_C4, C_RISK_C5 = "C1", "C2", "C3", "C4", "C5"
PROD_RISK_R1, PROD_RISK_R2, PROD_RISK_R3, PROD_RISK_R4, PROD_RISK_R5 = "R1", "R2", "R3", "R4", "R5"
# ---------------------------------------------------------------------------
# §8 系统配置 sys_config key(调仓/客户经营参数,代码不写死 key)
# ---------------------------------------------------------------------------
SYS_KEY_REBALANCE_DEVIATION_THRESHOLD = "rebalance.deviation.threshold" # 组合再平衡偏离阈值(全局兜底)
SYS_KEY_HIGH_NET_ASSET_THRESHOLD = "customer.high_net.asset.threshold" # 高净值客户资产门槛
SYS_KEY_RISK_QUESTIONNAIRE_EXPIRE_DAY = "risk.questionnaire.expire.day" # 风险测评过期天数
SYS_KEY_LARGE_FLOW_THRESHOLD = "customer.large_flow.threshold" # 大额申赎阈值
# ---------------------------------------------------------------------------
# §9 事件总线 Pub/Sub 事件名(双写 event_log + Redis Pub/Sub,消费以 event_id 幂等)
# ---------------------------------------------------------------------------
EVENT_REBALANCE_DRAFT_CREATED = "event:rebalance_draft_created"
EVENT_PROFILE_UPDATE = "event:profile_update"
# 工作台消费的投顾域事件(仅前两者;work_order_change/risk_alert 属其它域)
ADVISOR_EVENTS = (EVENT_REBALANCE_DRAFT_CREATED, EVENT_PROFILE_UPDATE)
# 事件消费状态(event_log.status)
EVENT_STATUS_PENDING = "待消费"
EVENT_STATUS_CONSUMED = "已消费"
# ---------------------------------------------------------------------------
# §10 固定文本模板
# ---------------------------------------------------------------------------
# 10.1 报告强制免责声明(完整性判定:正文必须「逐字包含」完整原文,精确匹配)
DISCLAIMER_TEXT = (
"【免责声明】本报告由AI辅助生成,仅供持牌投顾内部参考,不构成任何投资建议。"
"基金有风险,投资需谨慎。所有投资决策请结合自身风险承受能力审慎判断。"
)
# 10.2 SSE 事件类型(Agent 流式;工作台同步接口不涉及,保留枚举以对齐契约)
SSE_EVENT_TYPE_TEXT = "text"
SSE_EVENT_TYPE_META = "meta"
SSE_EVENT_TYPE_DONE = "done"
SSE_EVENT_TYPE_ERROR = "error"
# ---------------------------------------------------------------------------
# §11 定时任务 cron
# ---------------------------------------------------------------------------
CRON_MEMORY_MAINTENANCE = "0 2 * * 1" # Agent 侧每周记忆维护(工作台不调度)
CRON_PORTFOLIO_REBALANCE = "0 1 * * *" # 工作台每日组合再平衡(遍历 signed 客户)
CRON_FUND_NAV_UPDATE = "30 1 * * *" # 基金净值更新
# ---------------------------------------------------------------------------
# §12 审计日志 action(audit_log.action)
# ---------------------------------------------------------------------------
AUDIT_AGENT_CHAT_CALL = "agent_chat_call" # 调用投顾Agent对话
AUDIT_DRAFT_SAVE = "draft_save" # 草稿保存
AUDIT_DRAFT_DISCARD = "draft_discard" # 草稿废弃
AUDIT_CUSTOMER_SIGN = "customer_sign" # 客户签约
AUDIT_VIEW_SENSITIVE = "view_sensitive" # 投顾查看脱敏字段完整值(留痕)
# ---------------------------------------------------------------------------
# §13 投顾角色(sys_user.employee_role)
# ---------------------------------------------------------------------------
EMPLOYEE_ROLE_ADVISOR = "投顾" # 工作台唯一业务角色;「理财顾问」为历史表述,不参与鉴权
# ---------------------------------------------------------------------------
# §14 站内信消息类型(sys_message.msg_type)
# ---------------------------------------------------------------------------
MSG_TYPE_RECOMMEND = "推荐"
MSG_TYPE_REBALANCE = "调仓"
MSG_TYPE_RISK = "风控"
MSG_TYPE_SIGN = "签约"
MSG_TYPE_SYSTEM = "系统"
# 草稿 intent → 站内信 msg_type 映射(recommend→推荐 / rebalance→调仓)
INTENT_TO_MSG_TYPE = {
AGENT_INTENT_RECOMMEND: MSG_TYPE_RECOMMEND,
AGENT_INTENT_REBALANCE: MSG_TYPE_REBALANCE,
}
# ---------------------------------------------------------------------------
# 以下为「工作台本地」常量:DDL COMMENT / PRD 定义,common_const.md 未集中列出。
# 假设:按 sql/schema.sql 与 PRD 语义落地;后续若统一进 common_const.md 再收敛。
# ---------------------------------------------------------------------------
# 待办类型(advisor_todo.todo_type,见 schema.sql 表 27)
TODO_TYPE_NEW_REBALANCE_DRAFT = "new_rebalance_draft" # 新增调仓建议
TODO_TYPE_VISIT_DUE = "visit_due" # 回访到期
TODO_TYPE_RISK_EXPIRE = "risk_expire" # 风险测评过期
TODO_TYPE_LARGE_FLOW = "large_flow" # 客户大额申赎
TODO_TYPE_DRAWDOWN = "drawdown" # 持仓大幅回撤
# 待办来源(advisor_todo.source)
TODO_SOURCE_AGENT_EVENT = "agent_event" # 事件消费生成
TODO_SOURCE_LOCAL_TIMER = "local_timer" # 本地定时任务生成
TODO_SOURCE_LOCAL_CALC = "local_calc" # 本地计算生成
# 待办优先级(advisor_todo.priority)
TODO_PRIORITY_HIGH = "高"
TODO_PRIORITY_NORMAL = "普通"
TODO_PRIORITY_LOW = "低"
# 待办状态(advisor_todo.status;单向流转 待处理→处理中→已完成)
TODO_STATUS_PENDING = "待处理"
TODO_STATUS_PROCESSING = "处理中"
TODO_STATUS_DONE = "已完成"
# 回访方式(advisor_visit_record.visit_type,见 PRD §4.4.2)
VISIT_TYPE_PHONE = "电话"
VISIT_TYPE_WECHAT = "企微"
VISIT_TYPE_SMS = "短信"
VISIT_TYPE_ONSITE = "现场"
# 关系流转动作(工作台写入 customer_relation.status 的动作标识)
RELATION_ACTION_SIGN = "sign" # 签约
RELATION_ACTION_CLOSE = "close" # 结束服务
+37
View File
@@ -0,0 +1,37 @@
"""advisor_report 建议报告表 ORM(工作台本地镜像:sent 状态 + 编辑/发送留痕)。
send_status 流转:draft(草稿快照,随 save 镜像)→ sent(投顾手动发送)/ discarded(废弃归档)。
sent 是工作台本地状态,不回写投顾Agent。
"""
from __future__ import annotations
from datetime import datetime
from typing import Any
from sqlalchemy import BigInteger, DateTime, JSON, String, Text, func
from sqlalchemy.orm import Mapped, mapped_column
from model.base import Base
class AdvisorReport(Base):
__tablename__ = "advisor_report"
__table_args__ = {"comment": "建议报告表(工作台本地镜像,承载 sent 状态与发送/编辑留痕)"}
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
report_id: Mapped[str] = mapped_column(String(64), unique=True)
draft_id: Mapped[str] = mapped_column(String(64))
customer_id: Mapped[int] = mapped_column(BigInteger)
advisor_id: Mapped[int] = mapped_column(BigInteger)
intent: Mapped[str] = mapped_column(String(32))
title: Mapped[str | None] = mapped_column(String(128))
content: Mapped[str | None] = mapped_column(Text)
edit_history: Mapped[list[Any] | None] = mapped_column(JSON)
send_status: Mapped[str] = mapped_column(String(16), server_default="draft")
send_time: Mapped[datetime | None] = mapped_column(DateTime)
send_by: Mapped[int | None] = mapped_column(BigInteger)
msg_id: Mapped[int | None] = mapped_column(BigInteger)
create_time: Mapped[datetime] = mapped_column(DateTime, server_default=func.now())
update_time: Mapped[datetime] = mapped_column(
DateTime, server_default=func.now(), onupdate=func.now()
)
+30
View File
@@ -0,0 +1,30 @@
"""advisor_todo 投顾待办表 ORM(事件/定时/计算三类来源统一承载)。
幂等唯一键 uk_todo(todo_type, customer_id, biz_id) 由 DDL 维护;应用层在写入前
以相同三元组查询去重,避免重复建待办。
"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import BigInteger, DateTime, String, func
from sqlalchemy.orm import Mapped, mapped_column
from model.base import Base
class AdvisorTodo(Base):
__tablename__ = "advisor_todo"
__table_args__ = {"comment": "投顾待办表(事件/定时/计算三类来源统一承载)"}
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
todo_type: Mapped[str] = mapped_column(String(32))
customer_id: Mapped[int | None] = mapped_column(BigInteger)
advisor_id: Mapped[int] = mapped_column(BigInteger)
source: Mapped[str] = mapped_column(String(32))
biz_id: Mapped[str | None] = mapped_column(String(64))
priority: Mapped[str] = mapped_column(String(8), server_default="普通")
status: Mapped[str] = mapped_column(String(16), server_default="待处理")
due_at: Mapped[datetime | None] = mapped_column(DateTime)
create_time: Mapped[datetime] = mapped_column(DateTime, server_default=func.now())
handle_time: Mapped[datetime | None] = mapped_column(DateTime)
+23
View File
@@ -0,0 +1,23 @@
"""advisor_visit_record 投顾回访记录表 ORM(人工留痕,不触发 Agent 记忆抽取)。"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import BigInteger, DateTime, String, Text, func
from sqlalchemy.orm import Mapped, mapped_column
from model.base import Base
class AdvisorVisitRecord(Base):
__tablename__ = "advisor_visit_record"
__table_args__ = {"comment": "投顾回访记录表(人工留痕归档,不触发Agent记忆抽取)"}
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
customer_id: Mapped[int] = mapped_column(BigInteger)
advisor_id: Mapped[int] = mapped_column(BigInteger)
visit_type: Mapped[str] = mapped_column(String(32))
visit_time: Mapped[datetime] = mapped_column(DateTime)
summary: Mapped[str | None] = mapped_column(Text)
audio_url: Mapped[str | None] = mapped_column(String(512))
create_time: Mapped[datetime] = mapped_column(DateTime, server_default=func.now())
+14 -8
View File
@@ -1,27 +1,33 @@
"""customer_relation 客户-投顾关系表 ORM 模型。"""
"""customer_relation 客户-投顾关系表 ORM 模型。
同时服务于:
- 投顾工作台(advisor):签约状态 unsigned/signed/closed,工作台为唯一写入方;
- 记忆/client_agent 模块:读取客户与投顾的当前或历史关系。
"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import BigInteger, DateTime, String
from sqlalchemy import BigInteger, DateTime, String, func
from sqlalchemy.orm import Mapped, mapped_column
from model.base import Base
class CustomerRelation(Base):
"""客户与投顾的当前或历史关系。"""
__tablename__ = "customer_relation"
__table_args__ = {"comment": "客户-投顾关系表"}
__table_args__ = {"comment": "客户-投顾关系表(状态驱动:签约后投顾Agent方案正式触达客户)"}
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
customer_id: Mapped[int] = mapped_column(BigInteger, nullable=False)
advisor_id: Mapped[int] = mapped_column(BigInteger, nullable=False)
assign_time: Mapped[datetime] = mapped_column(DateTime, nullable=False)
assign_time: Mapped[datetime] = mapped_column(
DateTime, nullable=False, server_default=func.now()
)
signed_time: Mapped[datetime | None] = mapped_column(DateTime)
end_time: Mapped[datetime | None] = mapped_column(DateTime)
status: Mapped[str] = mapped_column(String(16), nullable=False)
status: Mapped[str] = mapped_column(
String(16), nullable=False, server_default="unsigned"
)
reason: Mapped[str | None] = mapped_column(String(128))
+24
View File
@@ -0,0 +1,24 @@
"""event_log 事件日志表 ORM(事件双写落库,Pub/Sub 仅作实时通知,消费以 event_id 幂等)。"""
from __future__ import annotations
from datetime import datetime
from typing import Any
from sqlalchemy import BigInteger, DateTime, JSON, String, func
from sqlalchemy.orm import Mapped, mapped_column
from model.base import Base
class EventLog(Base):
__tablename__ = "event_log"
__table_args__ = {"comment": "事件日志表(事件双写落库,Pub/Sub 仅作实时通知,消费幂等)"}
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
event_id: Mapped[str] = mapped_column(String(64), unique=True)
event_name: Mapped[str] = mapped_column(String(64))
payload: Mapped[dict[str, Any] | None] = mapped_column(JSON)
trace_id: Mapped[str | None] = mapped_column(String(64))
status: Mapped[str] = mapped_column(String(16), server_default="待消费")
create_time: Mapped[datetime] = mapped_column(DateTime, server_default=func.now())
consume_time: Mapped[datetime | None] = mapped_column(DateTime)
+26
View File
@@ -0,0 +1,26 @@
"""portfolio_benchmark 组合基准配置表 ORM(投顾Agent 再平衡参照,运营可调整;工作台只读展示)。"""
from __future__ import annotations
from datetime import datetime
from decimal import Decimal
from typing import Any
from sqlalchemy import BigInteger, DateTime, JSON, Numeric, String, func
from sqlalchemy.orm import Mapped, mapped_column
from model.base import Base
class PortfolioBenchmark(Base):
__tablename__ = "portfolio_benchmark"
__table_args__ = {"comment": "组合基准配置表(投顾Agent 再平衡参照,运营可调整)"}
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
risk_level: Mapped[str] = mapped_column(String(16))
target_allocation: Mapped[dict[str, Any] | None] = mapped_column(JSON)
drift_threshold: Mapped[Decimal] = mapped_column(Numeric(5, 2), server_default="5.00")
status: Mapped[str] = mapped_column(String(8), server_default="启用")
create_time: Mapped[datetime] = mapped_column(DateTime, server_default=func.now())
update_time: Mapped[datetime] = mapped_column(
DateTime, server_default=func.now(), onupdate=func.now()
)
+23
View File
@@ -0,0 +1,23 @@
"""sensitive_word 敏感词库表 ORM(工作台与 Agent 同源,合规拦截唯一来源)。"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import BigInteger, DateTime, String, func
from sqlalchemy.orm import Mapped, mapped_column
from model.base import Base
class SensitiveWord(Base):
__tablename__ = "sensitive_word"
__table_args__ = {"comment": "敏感词库表(工作台与Agent同源,合规拦截唯一来源)"}
id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
word: Mapped[str] = mapped_column(String(128), unique=True)
category: Mapped[str] = mapped_column(String(32))
level: Mapped[str] = mapped_column(String(8), server_default="高")
status: Mapped[str] = mapped_column(String(8), server_default="启用")
update_time: Mapped[datetime] = mapped_column(
DateTime, server_default=func.now(), onupdate=func.now()
)
+1
View File
@@ -1,4 +1,5 @@
"""sys_message 站内信 ORM 模型(风控拦截/冻结后触达客户的通道)。"""
"""sys_message 站内信/消息中心表 ORM(投顾触达客户通道,发送报告时写入)。"""
from __future__ import annotations
from datetime import datetime
+73
View File
@@ -0,0 +1,73 @@
"""advisor_report 建议报告仓储:本地镜像查询 + 列表(工作台本地,不直连 Agent 表)。"""
from __future__ import annotations
from sqlalchemy import func, select
from model.advisor_report import AdvisorReport
from repositories.base import BaseRepository
class AdvisorReportRepo(BaseRepository):
model = AdvisorReport
async def get_by_report_id(self, report_id: str) -> AdvisorReport | None:
return await self.db.scalar(
select(AdvisorReport).where(AdvisorReport.report_id == report_id)
)
async def get_by_draft_id(self, draft_id: str) -> AdvisorReport | None:
"""按草稿 draft_id 取本地镜像(草稿快照 / 已发送报告)。"""
return await self.db.scalar(
select(AdvisorReport).where(AdvisorReport.draft_id == draft_id)
)
async def list_by_advisor(
self,
*,
advisor_id: int,
customer_id: int | None = None,
send_status: str | None = None,
limit: int = 100,
offset: int = 0,
) -> list[AdvisorReport]:
conds = [AdvisorReport.advisor_id == advisor_id]
if customer_id is not None:
conds.append(AdvisorReport.customer_id == customer_id)
if send_status:
conds.append(AdvisorReport.send_status == send_status)
stmt = (
select(AdvisorReport)
.where(*conds)
.order_by(AdvisorReport.id.desc())
.limit(limit)
.offset(offset)
)
return list((await self.db.scalars(stmt)).all())
async def count_by_advisor(
self,
*,
advisor_id: int,
customer_id: int | None = None,
send_status: str | None = None,
) -> int:
conds = [AdvisorReport.advisor_id == advisor_id]
if customer_id is not None:
conds.append(AdvisorReport.customer_id == customer_id)
if send_status:
conds.append(AdvisorReport.send_status == send_status)
stmt = select(func.count()).select_from(AdvisorReport).where(*conds)
return (await self.db.scalar(stmt)) or 0
async def list_by_customer(
self, customer_id: int, limit: int = 100, offset: int = 0
) -> list[AdvisorReport]:
"""某客户的历史建议报告(360 视图用)。"""
stmt = (
select(AdvisorReport)
.where(AdvisorReport.customer_id == customer_id)
.order_by(AdvisorReport.id.desc())
.limit(limit)
.offset(offset)
)
return list((await self.db.scalars(stmt)).all())
+64
View File
@@ -0,0 +1,64 @@
"""advisor_todo 投顾待办仓储:按唯一键幂等建待办 + 列表查询。"""
from __future__ import annotations
from sqlalchemy import func, select
from model.advisor_todo import AdvisorTodo
from repositories.base import BaseRepository
class AdvisorTodoRepo(BaseRepository):
model = AdvisorTodo
async def get_by_unique(
self, todo_type: str, customer_id: int | None, biz_id: str | None
) -> AdvisorTodo | None:
"""按 DDL 唯一键 uk_todo(todo_type, customer_id, biz_id) 查重(幂等建待办)。
注:customer_id/biz_id 为 None 时 SQLAlchemy 生成 IS NULL 条件,语义与 DDL 一致。
"""
return await self.db.scalar(
select(AdvisorTodo).where(
AdvisorTodo.todo_type == todo_type,
AdvisorTodo.customer_id == customer_id,
AdvisorTodo.biz_id == biz_id,
)
)
async def list_by_advisor(
self,
*,
advisor_id: int,
status: str | None = None,
todo_type: str | None = None,
limit: int = 100,
offset: int = 0,
) -> list[AdvisorTodo]:
conds = [AdvisorTodo.advisor_id == advisor_id]
if status:
conds.append(AdvisorTodo.status == status)
if todo_type:
conds.append(AdvisorTodo.todo_type == todo_type)
stmt = (
select(AdvisorTodo)
.where(*conds)
.order_by(AdvisorTodo.id.desc())
.limit(limit)
.offset(offset)
)
return list((await self.db.scalars(stmt)).all())
async def count_by_advisor(
self,
*,
advisor_id: int,
status: str | None = None,
todo_type: str | None = None,
) -> int:
conds = [AdvisorTodo.advisor_id == advisor_id]
if status:
conds.append(AdvisorTodo.status == status)
if todo_type:
conds.append(AdvisorTodo.todo_type == todo_type)
stmt = select(func.count()).select_from(AdvisorTodo).where(*conds)
return (await self.db.scalar(stmt)) or 0
+40
View File
@@ -0,0 +1,40 @@
"""advisor_visit_record 回访记录仓储(人工留痕归档,不触发 Agent 记忆抽取)。"""
from __future__ import annotations
from sqlalchemy import func, select
from model.advisor_visit_record import AdvisorVisitRecord
from repositories.base import BaseRepository
class AdvisorVisitRecordRepo(BaseRepository):
model = AdvisorVisitRecord
async def list_by_advisor(
self,
*,
advisor_id: int,
customer_id: int | None = None,
limit: int = 100,
offset: int = 0,
) -> list[AdvisorVisitRecord]:
conds = [AdvisorVisitRecord.advisor_id == advisor_id]
if customer_id is not None:
conds.append(AdvisorVisitRecord.customer_id == customer_id)
stmt = (
select(AdvisorVisitRecord)
.where(*conds)
.order_by(AdvisorVisitRecord.visit_time.desc())
.limit(limit)
.offset(offset)
)
return list((await self.db.scalars(stmt)).all())
async def count_by_advisor(
self, *, advisor_id: int, customer_id: int | None = None
) -> int:
conds = [AdvisorVisitRecord.advisor_id == advisor_id]
if customer_id is not None:
conds.append(AdvisorVisitRecord.customer_id == customer_id)
stmt = select(func.count()).select_from(AdvisorVisitRecord).where(*conds)
return (await self.db.scalar(stmt)) or 0
+64
View File
@@ -0,0 +1,64 @@
"""audit_log 审计日志仓储:本人操作记录筛选(台账)。"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import func, or_, select
from model.audit_log import AuditLog
from repositories.base import BaseRepository
class AuditLogRepo(BaseRepository):
model = AuditLog
def _conds(self, *, user_id, module, action, keyword, start, end):
conds = [AuditLog.user_id == user_id, AuditLog.module == module]
if action:
conds.append(AuditLog.action == action)
if start:
conds.append(AuditLog.create_time >= start)
if end:
conds.append(AuditLog.create_time <= end)
if keyword:
like = f"%{keyword}%"
conds.append(or_(AuditLog.target.like(like), AuditLog.detail.like(like)))
return conds
async def list_by_advisor(
self,
*,
user_id: int,
module: str = "advisor",
action: str | None = None,
keyword: str | None = None,
start: datetime | None = None,
end: datetime | None = None,
limit: int = 100,
offset: int = 0,
) -> list[AuditLog]:
stmt = (
select(AuditLog)
.where(*self._conds(user_id=user_id, module=module, action=action, keyword=keyword, start=start, end=end))
.order_by(AuditLog.id.desc())
.limit(limit)
.offset(offset)
)
return list((await self.db.scalars(stmt)).all())
async def count_by_advisor(
self,
*,
user_id: int,
module: str = "advisor",
action: str | None = None,
keyword: str | None = None,
start: datetime | None = None,
end: datetime | None = None,
) -> int:
stmt = (
select(func.count())
.select_from(AuditLog)
.where(*self._conds(user_id=user_id, module=module, action=action, keyword=keyword, start=start, end=end))
)
return (await self.db.scalar(stmt)) or 0
+141 -6
View File
@@ -1,18 +1,51 @@
"""客户-投顾关系查询仓储。"""
"""customer_relation 客户-投顾关系仓储。
from sqlalchemy import select
同时服务于:
- 投顾工作台(advisor):数据权限根(仅本人名下客户)+ 客户列表 + AUM 聚合;
- 记忆/client_agent 模块:按客户 ID 读取有效关系。
"""
from __future__ import annotations
from decimal import Decimal
from sqlalchemy import func, or_, select
from model.customer_relation import CustomerRelation
from model.fin_customer_profile import FinCustomerProfile
from model.fin_holdings import FinHoldings
from model.sys_user import SysUser
from repositories.base import BaseRepository
# 持仓中状态:仅统计「持有中」市值,与 service/holdings.py 口径一致(DDL 注释,未入 common_const)
_HOLDING_STATUS = "持有中"
class CustomerRelationRepo(BaseRepository):
"""按客户 ID 隔离读取有效关系。"""
model = CustomerRelation
async def get_by_customer_advisor(
self, customer_id: int, advisor_id: int
) -> CustomerRelation | None:
"""取指定客户与投顾的关系(数据权限判断用)。"""
return await self.db.scalar(
select(CustomerRelation).where(
CustomerRelation.customer_id == customer_id,
CustomerRelation.advisor_id == advisor_id,
)
)
async def get_current_by_customer(self, customer_id: int) -> CustomerRelation | None:
"""取客户当前有效关系(历史重分配可能多行,取最近一条,按 id 倒序)。"""
return await self.db.scalar(
select(CustomerRelation)
.where(CustomerRelation.customer_id == customer_id)
.order_by(CustomerRelation.id.desc())
.limit(1)
)
async def list_by_customer(self, customer_id: int) -> list[CustomerRelation]:
"""返回指定客户尚未结束的关系。"""
"""返回指定客户尚未结束的关系(记忆/client_agent 模块用)。"""
statement = (
select(CustomerRelation)
.where(
@@ -23,6 +56,108 @@ class CustomerRelationRepo(BaseRepository):
)
return list((await self.db.scalars(statement)).all())
async def list_by_status(self, status: str) -> list[CustomerRelation]:
"""按状态取全部关系(定时调度遍历 signed 客户用,跨投顾)。"""
return list(
(
await self.db.scalars(
select(CustomerRelation).where(CustomerRelation.status == status)
)
).all()
)
async def list_by_advisor(
self, advisor_id: int, status: str | None = None
) -> list[CustomerRelation]:
"""名下客户关系列表(status=None 不过滤)。"""
stmt = select(CustomerRelation).where(CustomerRelation.advisor_id == advisor_id)
if status:
stmt = stmt.where(CustomerRelation.status == status)
return list((await self.db.scalars(stmt)).all())
async def list_customer_rows(
self,
*,
advisor_id: int,
status: str | None = None,
keyword: str | None = None,
limit: int = 100,
offset: int = 0,
) -> list[tuple[CustomerRelation, SysUser, FinCustomerProfile | None]]:
"""客户列表:关系 + 客户账号 + 画像(左连),返回三元组供 service 拼装响应。"""
conds = [CustomerRelation.advisor_id == advisor_id]
if status:
conds.append(CustomerRelation.status == status)
if keyword:
like = f"%{keyword}%"
conds.append(or_(SysUser.real_name.like(like), SysUser.phone.like(like)))
stmt = (
select(CustomerRelation, SysUser, FinCustomerProfile)
.join(SysUser, SysUser.id == CustomerRelation.customer_id)
.outerjoin(
FinCustomerProfile,
FinCustomerProfile.customer_id == CustomerRelation.customer_id,
)
.where(*conds)
.order_by(CustomerRelation.id.desc())
.limit(limit)
.offset(offset)
)
return list((await self.db.execute(stmt)).all())
async def count_customer_rows(
self,
*,
advisor_id: int,
status: str | None = None,
keyword: str | None = None,
) -> int:
conds = [CustomerRelation.advisor_id == advisor_id]
if status:
conds.append(CustomerRelation.status == status)
if keyword:
like = f"%{keyword}%"
conds.append(or_(SysUser.real_name.like(like), SysUser.phone.like(like)))
stmt = (
select(func.count())
.select_from(CustomerRelation)
.join(SysUser, SysUser.id == CustomerRelation.customer_id)
.where(*conds)
)
return (await self.db.scalar(stmt)) or 0
async def risk_level_distribution(
self, advisor_id: int
) -> list[tuple[str | None, int]]:
"""名下客户风险等级分布(C1-C5 各档人数),供驾驶舱客户分层。"""
stmt = (
select(FinCustomerProfile.risk_level, func.count())
.select_from(CustomerRelation)
.join(
FinCustomerProfile,
FinCustomerProfile.customer_id == CustomerRelation.customer_id,
)
.where(CustomerRelation.advisor_id == advisor_id)
.group_by(FinCustomerProfile.risk_level)
)
return list((await self.db.execute(stmt)).all())
async def sum_holdings_value(
self, advisor_id: int, relation_status: str | None = None
) -> Decimal:
"""AUM:名下客户「持有中」持仓当前市值之和(relation_status 过滤,如 signed)。"""
stmt = (
select(func.coalesce(func.sum(FinHoldings.current_value), 0))
.select_from(CustomerRelation)
.join(FinHoldings, FinHoldings.customer_id == CustomerRelation.customer_id)
.where(
CustomerRelation.advisor_id == advisor_id,
FinHoldings.status == _HOLDING_STATUS,
)
)
if relation_status:
stmt = stmt.where(CustomerRelation.status == relation_status)
return await self.db.scalar(stmt)
__all__ = ["CustomerRelationRepo"]
+43
View File
@@ -0,0 +1,43 @@
"""event_log 事件日志仓储:幂等消费(event_id 唯一)+ 补拉未消费事件。"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import select, update
from common_const import EVENT_STATUS_CONSUMED, EVENT_STATUS_PENDING
from model.event_log import EventLog
from repositories.base import BaseRepository
class EventLogRepo(BaseRepository):
model = EventLog
async def get_by_event_id(self, event_id: str) -> EventLog | None:
return await self.db.scalar(
select(EventLog).where(EventLog.event_id == event_id)
)
async def list_pending(
self, event_names: tuple[str, ...], limit: int = 100
) -> list[EventLog]:
"""补拉指定事件名中仍未消费的记录(Redis Pub/Sub 丢消息的兜底)。"""
stmt = (
select(EventLog)
.where(
EventLog.event_name.in_(event_names),
EventLog.status == EVENT_STATUS_PENDING,
)
.order_by(EventLog.id.asc())
.limit(limit)
)
return list((await self.db.scalars(stmt)).all())
async def mark_consumed(self, event_id: str) -> None:
"""标记事件已消费(含消费时间),并提交。"""
await self.db.execute(
update(EventLog)
.where(EventLog.event_id == event_id)
.values(status=EVENT_STATUS_CONSUMED, consume_time=datetime.now())
)
await self.db.commit()
+26
View File
@@ -0,0 +1,26 @@
"""portfolio_benchmark 组合基准仓储(工作台只读:标准策略库展示)。"""
from __future__ import annotations
from sqlalchemy import select
from model.portfolio_benchmark import PortfolioBenchmark
from repositories.base import BaseRepository
# 启用态(DDL:status 启用/停用)
_ACTIVE_STATUS = "启用"
class PortfolioBenchmarkRepo(BaseRepository):
model = PortfolioBenchmark
async def list_enabled(self) -> list[PortfolioBenchmark]:
"""取全部启用态基准(标准化策略库,投顾只读不可改)。"""
return list(
(
await self.db.scalars(
select(PortfolioBenchmark).where(
PortfolioBenchmark.status == _ACTIVE_STATUS
)
)
).all()
)
+24
View File
@@ -0,0 +1,24 @@
"""sensitive_word 敏感词仓储:取启用态词库(发送终审合规拦截唯一来源,与 Agent 同源)。"""
from __future__ import annotations
from sqlalchemy import select
from model.sensitive_word import SensitiveWord
from repositories.base import BaseRepository
# 词库启用态(DDL:status 启用/停用;未入 common_const,与 schema 注释一致)
_ACTIVE_STATUS = "启用"
class SensitiveWordRepo(BaseRepository):
model = SensitiveWord
async def list_active(self) -> list[SensitiveWord]:
"""取全部启用态敏感词,供发送终审扫描。"""
return list(
(
await self.db.scalars(
select(SensitiveWord).where(SensitiveWord.status == _ACTIVE_STATUS)
)
).all()
)
+24
View File
@@ -0,0 +1,24 @@
"""sys_message 站内信仓储(投顾触达客户,发送报告时写入)。"""
from __future__ import annotations
from sqlalchemy import select
from model.sys_message import SysMessage
from repositories.base import BaseRepository
class SysMessageRepo(BaseRepository):
model = SysMessage
async def list_by_user(
self, user_id: int, limit: int = 100, offset: int = 0
) -> list[SysMessage]:
"""某用户(客户)的站内信,倒序分页。"""
stmt = (
select(SysMessage)
.where(SysMessage.user_id == user_id)
.order_by(SysMessage.id.desc())
.limit(limit)
.offset(offset)
)
return list((await self.db.scalars(stmt)).all())
+57
View File
@@ -0,0 +1,57 @@
"""投顾工作台请求 DTO(入参校验;响应多以 dict 返回,随业务拼装)。"""
from __future__ import annotations
from datetime import datetime
from typing import Any
from pydantic import BaseModel, Field
class DraftSaveReq(BaseModel):
"""草稿编辑保存入参(代理到 Agent PUT /draft/{id}/save,Agent 重新适当性校验)。"""
title: str | None = Field(default=None, max_length=128)
content: str = Field(default="")
suggestions: Any = None # 结构化建议清单(超配赎回/低配申购),结构由 Agent 定义
class DraftOperateReq(BaseModel):
"""草稿操作入参(仅支持 discard 废弃)。"""
action: str = Field(pattern="^(discard)$")
class RebalanceRunReq(BaseModel):
"""触发调仓再平衡(异步受理,结果靠事件回执)。"""
customer_id: int = Field(gt=0)
class TalkScriptReq(BaseModel):
"""生成沟通话术草稿(同步,超时由工作台降级为「稍后重试」)。"""
customer_id: int = Field(gt=0)
scene_type: str = Field(min_length=1, max_length=64)
class RelationReq(BaseModel):
"""客户关系流转(sign=签约 / close=结束服务),工作台是签约状态唯一写入方。"""
action: str = Field(pattern="^(sign|close)$")
reason: str | None = Field(default=None, max_length=128)
class VisitCreateReq(BaseModel):
"""人工回访留痕入参(不触发 Agent 记忆抽取)。"""
customer_id: int = Field(gt=0)
visit_type: str = Field(min_length=1, max_length=32)
visit_time: datetime
summary: str | None = None
audio_url: str | None = Field(default=None, max_length=512)
class TodoHandleReq(BaseModel):
"""待办处理入参(process=开始处理 / done=完成)。"""
action: str = Field(pattern="^(process|done)$")
+6
View File
@@ -0,0 +1,6 @@
"""投顾工作台业务层(B 端)。
模块职责:草稿代理与发送终审、客户 360、待办、回访、驾驶舱、审计台账、个人报表,
以及事件消费与定时调度。对投顾Agent 的调用统一走 agent_client;本地合规校验统一走
send_review;敏感信息脱敏统一走 masking。
"""
+201
View File
@@ -0,0 +1,201 @@
"""投顾Agent HTTP 客户端 + 错误码映射层(工作台 → Agent 的唯一出口)。
职责:
1. 通过 httpx 调用投顾Agent 独立服务(统一前缀 /api/advisor-agent,Agent 文档 §5);
2. 透传上游投顾 JWT 与 X-Trace-Id;
3. 快接口超时/重试(仅幂等 GET 重试,写操作不盲目重试);
4. 将 Agent 自有错误码(0/40001/40020/40030/40401/50001/50002)映射为工作台异常,
绝不把 Agent 码透传给上层调用方(common_const §6 错误码域边界)。
"""
from __future__ import annotations
import httpx
from common_const import (
AGENT_ERR_MESSAGE,
ERR_CODE_DRAFT_NOT_FOUND,
ERR_CODE_FORBIDDEN_CUSTOMER,
ERR_CODE_GRAPH_ERROR,
ERR_CODE_LLM_ERROR,
ERR_CODE_NOT_SIGNED_REBALANCE,
ERR_CODE_OK,
ERR_CODE_SUITABILITY_INVALID,
)
from config.settings import settings
from utils.exceptions import (
ForbiddenError,
LLMFailError,
NotFoundError,
NotSuitableError,
ParamError,
)
# Agent 统一前缀(Agent 文档 §5)
_AGENT_PREFIX = "/api/advisor-agent"
def translate_agent_error(code: int, message: str | None = None) -> str | None:
"""把 Agent 业务码映射为工作台结果。
- code == 0:正常,返回 None;
- code == 50002:降级但成功,返回告警文案(不抛异常);
- 其余:抛出映射后的 ApiError(工作台码域,不透传 Agent 码)。
"""
if code == ERR_CODE_OK:
return None
if code == ERR_CODE_GRAPH_ERROR:
return message or AGENT_ERR_MESSAGE.get(code, "图谱查询异常,已降级返回部分结果")
default = message or AGENT_ERR_MESSAGE.get(code, "AI 服务调用异常,请稍后重试")
if code == ERR_CODE_FORBIDDEN_CUSTOMER:
raise ForbiddenError(default)
if code == ERR_CODE_SUITABILITY_INVALID:
raise NotSuitableError(default)
if code == ERR_CODE_NOT_SIGNED_REBALANCE:
raise ParamError(default)
if code == ERR_CODE_DRAFT_NOT_FOUND:
raise NotFoundError(default)
if code == ERR_CODE_LLM_ERROR:
raise LLMFailError(default)
# 未知 Agent 业务码:统一按 AI 服务异常兜底,不透传
raise LLMFailError(default)
class AdvisorAgentClient:
"""投顾Agent 客户端(模块级单例 get_agent_client() 获取)。"""
def __init__(self, base_url: str, timeout: float, retry: int):
self.base_url = (base_url or "").rstrip("/")
self.timeout = timeout
self.retry = max(0, retry)
@property
def configured(self) -> bool:
return bool(self.base_url)
def _url(self, path: str) -> str:
return f"{self.base_url}{_AGENT_PREFIX}{path}"
def _headers(self, auth_header: str, trace_id: str) -> dict:
return {
"Authorization": auth_header,
"X-Trace-Id": trace_id,
"Content-Type": "application/json",
}
async def _request(
self,
method: str,
path: str,
*,
auth_header: str,
trace_id: str,
params: dict | None = None,
json: dict | None = None,
allow_retry: bool = False,
) -> dict:
"""统一请求:校验配置 → 超时重试 → 解析返回体 → 映射业务码。
返回 {"data": ..., "warning": ...};业务失败抛出映射后的 ApiError。
"""
if not self.configured:
raise LLMFailError("投顾Agent 服务未配置,请联系管理员")
url = self._url(path)
headers = self._headers(auth_header, trace_id)
attempts = self.retry + 1 if allow_retry else 1
for attempt in range(attempts):
try:
async with httpx.AsyncClient(timeout=self.timeout) as client:
resp = await client.request(
method, url, headers=headers, params=params, json=json
)
except (httpx.TimeoutException, httpx.TransportError) as exc:
if attempt < attempts - 1:
continue
raise LLMFailError("AI 服务调用异常,请稍后重试") from exc
if resp.status_code != 200:
if attempt < attempts - 1:
continue
raise LLMFailError("AI 服务调用异常,请稍后重试")
try:
body = resp.json()
except ValueError as exc:
raise LLMFailError("AI 服务返回格式异常,请稍后重试") from exc
code = body.get("code", ERR_CODE_LLM_ERROR)
warning = translate_agent_error(code, body.get("message"))
return {"data": body.get("data"), "warning": warning}
raise LLMFailError("AI 服务调用异常,请稍后重试") # 理论不可达,防御
# ---- 各 Agent 接口(字段与 Agent 文档 §5 对齐) ----
async def draft_list(
self,
*,
auth_header: str,
trace_id: str,
advisor_id: int,
customer_id: int | None = None,
status: str | None = None,
page: int = 1,
page_size: int = 20,
) -> dict:
params = {"advisor_id": advisor_id, "page": page, "page_size": page_size}
if customer_id is not None:
params["customer_id"] = customer_id
if status:
params["status"] = status
return await self._request(
"GET", "/draft/list", auth_header=auth_header, trace_id=trace_id, params=params,
allow_retry=True,
)
async def draft_detail(self, draft_id: str, *, auth_header: str, trace_id: str) -> dict:
return await self._request(
"GET", f"/draft/{draft_id}", auth_header=auth_header, trace_id=trace_id,
allow_retry=True,
)
async def draft_save(
self, draft_id: str, payload: dict, *, auth_header: str, trace_id: str
) -> dict:
return await self._request(
"PUT", f"/draft/{draft_id}/save", auth_header=auth_header,
trace_id=trace_id, json=payload,
)
async def draft_operate(
self, draft_id: str, action: str, *, auth_header: str, trace_id: str
) -> dict:
return await self._request(
"POST", f"/draft/{draft_id}/operate", auth_header=auth_header,
trace_id=trace_id, json={"action": action},
)
async def rebalance_run(
self, customer_id: int, *, auth_header: str, trace_id: str
) -> dict:
return await self._request(
"POST", "/rebalance/run", auth_header=auth_header,
trace_id=trace_id, json={"customer_id": customer_id},
)
async def generate_talk_script(
self, customer_id: int, scene_type: str, *, auth_header: str, trace_id: str
) -> dict:
return await self._request(
"POST", "/generate-talk-script", auth_header=auth_header,
trace_id=trace_id, json={"customer_id": customer_id, "scene_type": scene_type},
)
_client: AdvisorAgentClient | None = None
def get_agent_client() -> AdvisorAgentClient:
"""取全局客户端单例(按 config.settings 组装)。"""
global _client
if _client is None:
cfg = settings.advisor_agent
_client = AdvisorAgentClient(cfg.base_url, cfg.timeout, cfg.retry)
return _client
+75
View File
@@ -0,0 +1,75 @@
"""审计台账(PRD §4.6):本人操作记录筛选 + CSV 导出。"""
from __future__ import annotations
import csv
import io
from datetime import datetime
from sqlalchemy.ext.asyncio import AsyncSession
from model.sys_user import SysUser
from repositories.audit_log import AuditLogRepo
def _audit_item(a) -> dict:
return {
"id": a.id,
"module": a.module,
"action": a.action,
"target": a.target,
"detail": a.detail,
"trace_id": a.trace_id,
"status": a.status,
"create_time": a.create_time.isoformat() if a.create_time else None,
}
async def list_ledger(
db: AsyncSession,
user: SysUser,
*,
action: str | None = None,
keyword: str | None = None,
start: datetime | None = None,
end: datetime | None = None,
page: int = 1,
page_size: int = 20,
) -> dict:
repo = AuditLogRepo(db)
items = await repo.list_by_advisor(
user_id=user.id, action=action, keyword=keyword, start=start, end=end,
limit=page_size, offset=(page - 1) * page_size,
)
total = await repo.count_by_advisor(
user_id=user.id, action=action, keyword=keyword, start=start, end=end
)
return {"total": total, "page": page, "page_size": page_size, "items": [_audit_item(a) for a in items]}
def _to_csv(rows: list[dict]) -> str:
"""行列表 → CSV 文本(空列表返回空串)。"""
if not rows:
return ""
buf = io.StringIO()
writer = csv.DictWriter(buf, fieldnames=list(rows[0].keys()))
writer.writeheader()
writer.writerows(rows)
return buf.getvalue()
async def export_ledger(
db: AsyncSession,
user: SysUser,
*,
action: str | None = None,
keyword: str | None = None,
start: datetime | None = None,
end: datetime | None = None,
) -> str:
"""导出本人审计台账为 CSV 文本(全量,不限分页大小,上限 10000 条)。"""
repo = AuditLogRepo(db)
items = await repo.list_by_advisor(
user_id=user.id, action=action, keyword=keyword, start=start, end=end,
limit=10000, offset=0,
)
return _to_csv([_audit_item(a) for a in items])
+43
View File
@@ -0,0 +1,43 @@
"""审计日志写入助手(复用 audit_log 表,参照 service/customer_agent/audit.py 范式)。
工作台审计留痕统一走本模块:敏感信息查看、签约/结束服务等合规动作全程可追溯。
"""
from __future__ import annotations
import json
from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession
from utils.request_id import get_request_id
async def write_audit(
db: AsyncSession,
*,
user_id: int | None,
username: str | None,
module: str,
action: str,
target: str | None = None,
detail: dict | None = None,
status: str = "成功",
) -> None:
"""插入一条审计日志并提交(trace_id 自动取当前请求链路)。"""
await db.execute(
text(
"INSERT INTO audit_log (user_id, username, module, action, target, detail, trace_id, status) "
"VALUES (:user_id, :username, :module, :action, :target, :detail, :trace_id, :status)"
),
{
"user_id": user_id,
"username": username,
"module": module,
"action": action,
"target": target,
"detail": json.dumps(detail, ensure_ascii=False) if detail is not None else None,
"trace_id": get_request_id(),
"status": status,
},
)
await db.commit()
+176
View File
@@ -0,0 +1,176 @@
"""客户 360 全景管理(只读 + 数据权限 + 脱敏 + 签约/结束服务写入)。
PRD §4.2:投顾仅可查看自身名下客户;敏感信息两层脱敏;画像只读。
数据权限以 customer_relation.advisor_id = current_user.id 硬过滤。
"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy.ext.asyncio import AsyncSession
from common_const import (
AUDIT_CUSTOMER_SIGN,
AUDIT_VIEW_SENSITIVE,
CUSTOMER_REL_STATUS_CLOSED,
CUSTOMER_REL_STATUS_SIGNED,
RELATION_ACTION_SIGN,
)
from model.sys_user import SysUser
from repositories.advisor_report import AdvisorReportRepo
from repositories.customer_relation import CustomerRelationRepo
from repositories.fin_holdings import FinHoldingsRepo
from repositories.fin_product import FinProductRepo
from repositories.risk_assessment import CustomerProfileRepo
from repositories.sys_user import SysUserRepo
from schemas.advisor import RelationReq
from service.advisor.audit_writer import write_audit
from service.advisor.masking import mask_name, mask_phone
from service.advisor.permissions import ensure_customer_owned, require_owned_relation
from utils.exceptions import ForbiddenError, NotFoundError
# 持仓中状态(与 service/holdings.py 口径一致)
_HOLDING_STATUS = "持有中"
_AUDIT_MODULE = "advisor"
def _fmt_holding(h, product) -> dict:
"""持仓 + 产品信息拼装(金额/份额按字符串输出,避免浮点精度)。"""
return {
"product_id": h.product_id,
"product_code": product.product_code if product else None,
"product_name": product.product_name if product else None,
"risk_level": product.risk_level if product else None,
"shares": f"{h.shares:.4f}",
"cost_amount": f"{h.cost_amount:.2f}",
"current_value": f"{h.current_value:.2f}",
"profit_loss": f"{h.profit_loss:.2f}",
"profit_ratio": f"{h.profit_ratio:.4f}",
"status": h.status,
}
async def list_customers(
db: AsyncSession,
user: SysUser,
*,
status: str | None = None,
keyword: str | None = None,
page: int = 1,
page_size: int = 20,
) -> dict:
repo = CustomerRelationRepo(db)
rows = await repo.list_customer_rows(
advisor_id=user.id, status=status, keyword=keyword,
limit=page_size, offset=(page - 1) * page_size,
)
total = await repo.count_customer_rows(advisor_id=user.id, status=status, keyword=keyword)
items = []
for rel, account, profile in rows:
items.append(
{
"customer_id": rel.customer_id,
"real_name": mask_name(account.real_name),
"phone": mask_phone(account.phone),
"risk_level": profile.risk_level if profile else None,
"customer_level": account.customer_level,
"relation_status": rel.status,
"total_assets": float(profile.total_assets)
if profile and profile.total_assets is not None
else None,
}
)
return {"total": total, "page": page, "page_size": page_size, "items": items}
async def get_customer(
db: AsyncSession, user: SysUser, customer_id: int, *, unmask: bool = False
) -> dict:
rel = await require_owned_relation(db, user.id, customer_id)
account = await SysUserRepo(db).get(customer_id)
if account is None:
raise NotFoundError("客户不存在")
profile = await CustomerProfileRepo(db).get_by_customer(customer_id)
# 常规敏感(手机号)默认掩码;unmask=true 时返回完整值并记审计留痕。
# 强敏感(身份证/银行卡):sys_user 当前无对应字段,V1.0 无数据可暴露,规则保留。
phone = mask_phone(account.phone)
if unmask:
phone = account.phone
await write_audit(
db, user_id=user.id, username=user.username, module=_AUDIT_MODULE,
action=AUDIT_VIEW_SENSITIVE, target=str(customer_id),
detail={"field": "phone"},
)
return {
"customer_id": customer_id,
"real_name": mask_name(account.real_name),
"phone": phone,
"risk_level": profile.risk_level if profile else None,
"risk_score": profile.risk_score if profile else None,
"customer_level": account.customer_level,
"total_assets": float(profile.total_assets)
if profile and profile.total_assets is not None
else None,
"relation_status": rel.status,
"signed_time": rel.signed_time.isoformat() if rel.signed_time else None,
"assigned_time": rel.assign_time.isoformat() if rel.assign_time else None,
}
async def get_customer_holdings(
db: AsyncSession, user: SysUser, customer_id: int
) -> list[dict]:
await ensure_customer_owned(db, user.id, customer_id)
holdings = await FinHoldingsRepo(db).list_by_customer(customer_id, _HOLDING_STATUS)
product_repo = FinProductRepo(db)
items = []
for h in holdings:
product = await product_repo.get(h.product_id)
items.append(_fmt_holding(h, product))
return items
async def get_customer_reports(
db: AsyncSession, user: SysUser, customer_id: int, *, page: int = 1, page_size: int = 20
) -> dict:
await ensure_customer_owned(db, user.id, customer_id)
repo = AdvisorReportRepo(db)
items = await repo.list_by_customer(
customer_id, limit=page_size, offset=(page - 1) * page_size
)
return {
"total": await repo.count_by_advisor(advisor_id=user.id, customer_id=customer_id),
"items": [
{
"report_id": r.report_id,
"intent": r.intent,
"title": r.title,
"send_status": r.send_status,
"send_time": r.send_time.isoformat() if r.send_time else None,
}
for r in items
],
}
async def update_relation(
db: AsyncSession, user: SysUser, customer_id: int, req: RelationReq
) -> dict:
"""签约 / 结束服务(customer_relation.status 唯一写入方为工作台)。"""
rel = await require_owned_relation(db, user.id, customer_id)
if req.action == RELATION_ACTION_SIGN:
rel.status = CUSTOMER_REL_STATUS_SIGNED
rel.signed_time = datetime.now()
rel.reason = None
else: # close
rel.status = CUSTOMER_REL_STATUS_CLOSED
rel.end_time = datetime.now()
rel.reason = req.reason
await db.commit()
await write_audit(
db, user_id=user.id, username=user.username, module=_AUDIT_MODULE,
action=AUDIT_CUSTOMER_SIGN, target=str(customer_id), detail={"action": req.action},
)
return {"customer_id": customer_id, "status": rel.status}
+69
View File
@@ -0,0 +1,69 @@
"""驾驶舱聚合(首页数据与待办聚合,PRD §4.1)。
口径说明(E1 最小口径):总 AUM=名下客户「持有中」持仓市值之和;有效客户=signed;
流失率=closed/总客户;报告产出/发送=advisor_report 计数。回访完成率等精细口径 P2 后置。
"""
from __future__ import annotations
from sqlalchemy.ext.asyncio import AsyncSession
from common_const import (
CUSTOMER_REL_STATUS_CLOSED,
CUSTOMER_REL_STATUS_SIGNED,
REPORT_SEND_STATUS_DRAFT,
REPORT_SEND_STATUS_SENT,
TODO_STATUS_PENDING,
)
from model.sys_user import SysUser
from repositories.advisor_report import AdvisorReportRepo
from repositories.advisor_todo import AdvisorTodoRepo
from repositories.advisor_visit_record import AdvisorVisitRecordRepo
from repositories.customer_relation import CustomerRelationRepo
async def get_dashboard(db: AsyncSession, user: SysUser) -> dict:
relation_repo = CustomerRelationRepo(db)
todo_repo = AdvisorTodoRepo(db)
report_repo = AdvisorReportRepo(db)
visit_repo = AdvisorVisitRecordRepo(db)
pending = await todo_repo.count_by_advisor(advisor_id=user.id, status=TODO_STATUS_PENDING)
total = await relation_repo.count_customer_rows(advisor_id=user.id)
signed = await relation_repo.count_customer_rows(
advisor_id=user.id, status=CUSTOMER_REL_STATUS_SIGNED
)
closed = await relation_repo.count_customer_rows(
advisor_id=user.id, status=CUSTOMER_REL_STATUS_CLOSED
)
aum = await relation_repo.sum_holdings_value(user.id)
sent = await report_repo.count_by_advisor(
advisor_id=user.id, send_status=REPORT_SEND_STATUS_SENT
)
draft = await report_repo.count_by_advisor(
advisor_id=user.id, send_status=REPORT_SEND_STATUS_DRAFT
)
visits = await visit_repo.count_by_advisor(advisor_id=user.id)
# 客户分层:C1-C5 各档人数(未知等级归入 unknown)
levels = {k: 0 for k in ("C1", "C2", "C3", "C4", "C5")}
unknown = 0
for level, cnt in await relation_repo.risk_level_distribution(user.id):
if level in levels:
levels[level] = cnt
else:
unknown += cnt
return {
"todos": {"pending": pending},
"overview": {
"total_customers": total,
"signed_customers": signed,
"closed_customers": closed,
"churn_rate": round(closed / total, 4) if total else 0.0,
"total_aum": f"{aum:.2f}",
"reports_sent": sent,
"reports_draft": draft,
"visits_count": visits,
},
"customer_levels": {**levels, "unknown": unknown},
}
+96
View File
@@ -0,0 +1,96 @@
"""资产配置与组合诊断(PRD §4.3,只读)。
边界:持仓诊断仅做「持仓分布 + 集中度」只读概览;偏离度/调仓建议以投顾Agent 草稿为准,
工作台不重复实现调仓逻辑。标准策略库=portfolio_benchmark 只读;白名单基金=fin_product 筛选。
"""
from __future__ import annotations
from decimal import Decimal
from sqlalchemy.ext.asyncio import AsyncSession
from model.sys_user import SysUser
from repositories.fin_holdings import FinHoldingsRepo
from repositories.fin_product import FinProductRepo
from repositories.portfolio_benchmark import PortfolioBenchmarkRepo
from service.advisor.permissions import ensure_customer_owned
from service.product import list_products
# 持仓中状态(与 service/holdings.py 口径一致)
_HOLDING_STATUS = "持有中"
async def diagnose(db: AsyncSession, user: SysUser, customer_id: int) -> dict:
await ensure_customer_owned(db, user.id, customer_id)
holdings = await FinHoldingsRepo(db).list_by_customer(customer_id, _HOLDING_STATUS)
product_repo = FinProductRepo(db)
by_type: dict[str, Decimal] = {}
total = Decimal("0")
items: list[dict] = []
for h in holdings:
product = await product_repo.get(h.product_id)
ptype = product.product_type if product else "未知"
total += h.current_value
by_type[ptype] = by_type.get(ptype, Decimal("0")) + h.current_value
items.append(
{
"product_id": h.product_id,
"product_code": product.product_code if product else None,
"product_name": product.product_name if product else None,
"product_type": ptype,
"current_value": f"{h.current_value:.2f}",
}
)
# 集中度:单一产品市值占比最高者
items.sort(key=lambda x: float(x["current_value"]), reverse=True)
top_ratio = 0.0
if total and items:
top_ratio = float(items[0]["current_value"]) / float(total)
return {
"total_value": f"{total:.2f}",
"allocation": {k: f"{v:.2f}" for k, v in by_type.items()},
"concentration": {
"top_product_ratio": round(top_ratio, 4),
"top_product": items[0] if items else None,
},
"holdings": items,
}
async def list_strategies(db: AsyncSession) -> list[dict]:
"""标准策略库(portfolio_benchmark 只读,投顾不可私自新建策略)。"""
rows = await PortfolioBenchmarkRepo(db).list_enabled()
return [
{
"risk_level": b.risk_level,
"target_allocation": b.target_allocation,
"drift_threshold": float(b.drift_threshold),
}
for b in rows
]
async def list_funds(
db: AsyncSession,
*,
page: int = 1,
page_size: int = 10,
keyword: str | None = None,
product_type: str | None = None,
risk_level: str | None = None,
) -> dict:
"""白名单产品筛选(复用 service/product.list_products,仅在售)。"""
return await list_products(
db,
page=page,
page_size=page_size,
keyword=keyword,
product_type=product_type,
risk_level=risk_level,
status="在售",
sort_by="create_time",
sort_order="desc",
)
+332
View File
@@ -0,0 +1,332 @@
"""草稿代理 + 本地镜像 + 发送终审(投顾工作台核心闭环)。
- 草稿 list/detail/save/discard 代理到投顾Agent(错误码映射见 agent_client);
- save 时本地镜像一份 advisor_report(send_status=draft),实现 A5「以最后保存快照为准」;
- send 为工作台本地实现:签约兜底 + 三检终审(敏感词/适当性/免责)+ 写 sys_message
与 advisor_report(sent) 同事务;Agent 不感知 sent。
"""
from __future__ import annotations
import uuid
from datetime import datetime
from typing import Any
from sqlalchemy.ext.asyncio import AsyncSession
from common_const import (
AGENT_INTENT_RECOMMEND,
CUSTOMER_REL_STATUS_SIGNED,
INTENT_TO_MSG_TYPE,
MSG_TYPE_RECOMMEND,
REPORT_SEND_STATUS_DISCARDED,
REPORT_SEND_STATUS_DRAFT,
REPORT_SEND_STATUS_SENT,
)
from model.advisor_report import AdvisorReport
from model.sys_message import SysMessage
from model.sys_user import SysUser
from repositories.advisor_report import AdvisorReportRepo
from repositories.product import ProductRepo
from repositories.risk_assessment import CustomerProfileRepo
from repositories.sensitive_word import SensitiveWordRepo
from schemas.advisor import DraftSaveReq, RebalanceRunReq, TalkScriptReq
from service.advisor.agent_client import get_agent_client
from service.advisor.permissions import ensure_customer_owned, require_owned_relation
from service.advisor.send_review import review_send
from utils.exceptions import ForbiddenError, ParamError
# 建议清单中「低配申购」侧键(Agent 侧 suggestions 结构,文档未给字段级定义,按合理约定)
# 假设:suggestions 为 dict,申购清单在 buy/purchase/低配申购/申购 任一键下,项含 product_code。
_BUY_KEYS = ("buy", "purchase", "低配申购", "申购")
def gen_report_id() -> str:
"""生成本地报告唯一 ID(advisor_report.report_id)。"""
return "rpt-" + uuid.uuid4().hex
def _build_report_from_detail(detail: dict, advisor_id: int) -> AdvisorReport:
"""由 Agent 草稿详情构建本地报告镜像(未 save 直接 send 的兜底,D6)。"""
return AdvisorReport(
report_id=gen_report_id(),
draft_id=str(detail.get("draft_id") or ""),
customer_id=int(detail.get("customer_id") or 0),
advisor_id=advisor_id,
intent=detail.get("intent") or AGENT_INTENT_RECOMMEND,
title=detail.get("title"),
content=detail.get("content"),
edit_history=[],
send_status=REPORT_SEND_STATUS_DRAFT,
)
def _extract_buy_codes(suggestions: Any) -> list[str]:
"""从建议清单提取「低配申购」产品代码(发送终审适当性校验对象)。
假设:结构为 dict,申购侧键见 _BUY_KEYS;项为 {product_code} 或 {code}/{fund_code}。
解析失败返回空列表(视为无结构化产品建议,适当性空过,不误拦)。
"""
if not isinstance(suggestions, dict):
return []
items: list | None = None
for key in _BUY_KEYS:
value = suggestions.get(key)
if isinstance(value, list):
items = value
break
codes: list[str] = []
for it in items or []:
if isinstance(it, dict):
code = it.get("product_code") or it.get("code") or it.get("fund_code")
if code:
codes.append(str(code))
return codes
async def _get_customer_risk(db: AsyncSession, customer_id: int) -> str | None:
"""客户当前风险等级(适当性硬依据)。
假设:取 fin_customer_profile.risk_level(问卷服务在风评后回写,与 fin_risk_assessment
保持一致),作为「当前有效 C 级」。
"""
profile = await CustomerProfileRepo(db).get_by_customer(customer_id)
return profile.risk_level if profile else None
async def _resolve_product_risks(
db: AsyncSession, draft_id: str, auth_header: str, trace_id: str
) -> list[str | None]:
"""适当性输入:取草稿建议清单的申购产品,查 fin_product 得 R 级。
本地镜像不存 suggestions,故发送时向 Agent 实时取一次结构化清单;Agent 不可用则
发送被 fail-closed 阻断(合规优先,符合 PRD「不可完全依赖 Agent」+ 双重防护)。
"""
result = await get_agent_client().draft_detail(
draft_id, auth_header=auth_header, trace_id=trace_id
)
codes = _extract_buy_codes((result["data"] or {}).get("suggestions"))
product_repo = ProductRepo(db)
risks: list[str | None] = []
for code in codes:
product = await product_repo.get_by_code(code)
risks.append(product.risk_level if product else None)
return risks
def _sent_payload(report: AdvisorReport) -> dict:
return {
"report_id": report.report_id,
"send_time": report.send_time.isoformat() if report.send_time else None,
"msg_id": report.msg_id,
}
# ---------------------------------------------------------------------------
# 草稿代理(读取/编辑/废弃)
# ---------------------------------------------------------------------------
async def list_drafts(
db: AsyncSession,
user: SysUser,
*,
auth_header: str,
trace_id: str,
customer_id: int | None = None,
status: str | None = None,
page: int = 1,
page_size: int = 20,
) -> dict:
result = await get_agent_client().draft_list(
auth_header=auth_header,
trace_id=trace_id,
advisor_id=user.id,
customer_id=customer_id,
status=status,
page=page,
page_size=page_size,
)
data = result["data"] or {}
if result["warning"]:
data["warning"] = result["warning"]
return data
async def get_draft(
db: AsyncSession, user: SysUser, *, auth_header: str, trace_id: str, draft_id: str
) -> dict:
result = await get_agent_client().draft_detail(
draft_id, auth_header=auth_header, trace_id=trace_id
)
data = result["data"] or {}
# 数据权限:草稿归属客户必须属于当前投顾;无法确认归属时拒绝(fail-closed)
customer_id = data.get("customer_id")
if customer_id is None:
raise ForbiddenError("无法校验草稿归属客户")
await ensure_customer_owned(db, user.id, int(customer_id))
if result["warning"]:
data["warning"] = result["warning"]
return data
async def _upsert_report_mirror(
db: AsyncSession, user: SysUser, draft_id: str, detail: dict, req: DraftSaveReq
) -> None:
"""save 后本地镜像快照(send_status=draft),记录编辑留痕。
已发送(sent)的快照不回退为 draft,避免覆盖已交付内容。
"""
repo = AdvisorReportRepo(db)
report = await repo.get_by_draft_id(draft_id)
edit_entry = {
"editor": user.id,
"time": datetime.now().isoformat(),
"title": req.title,
}
if report is None:
report = AdvisorReport(
report_id=gen_report_id(),
draft_id=draft_id,
customer_id=int(detail["customer_id"]),
advisor_id=user.id,
intent=detail.get("intent") or AGENT_INTENT_RECOMMEND,
title=req.title,
content=req.content,
edit_history=[edit_entry],
send_status=REPORT_SEND_STATUS_DRAFT,
)
db.add(report)
elif report.send_status != REPORT_SEND_STATUS_SENT:
report.title = req.title
report.content = req.content
report.edit_history = list(report.edit_history or []) + [edit_entry]
report.send_status = REPORT_SEND_STATUS_DRAFT
await db.commit()
async def save_draft(
db: AsyncSession,
user: SysUser,
*,
auth_header: str,
trace_id: str,
draft_id: str,
req: DraftSaveReq,
) -> dict:
# 1) 先取草稿确认归属(避免对无权限草稿执行写操作),并拿到 intent/customer_id
detail = await get_draft(db, user, auth_header=auth_header, trace_id=trace_id, draft_id=draft_id)
# 2) 调 Agent 保存(Agent 重新适当性校验,违规 40020;缺免责仅告警不阻断)
payload = {"title": req.title, "content": req.content, "suggestions": req.suggestions}
result = await get_agent_client().draft_save(
draft_id, payload, auth_header=auth_header, trace_id=trace_id
)
# 3) 本地镜像快照
await _upsert_report_mirror(db, user, draft_id, detail, req)
data = result["data"] or {}
if result["warning"]:
data["warning"] = result["warning"]
return data
async def discard_draft(
db: AsyncSession, user: SysUser, *, auth_header: str, trace_id: str, draft_id: str
) -> dict:
# 1) 确认归属(避免越权废弃)
await get_draft(db, user, auth_header=auth_header, trace_id=trace_id, draft_id=draft_id)
# 2) 回写 Agent 归档(Agent 不可用则废弃失败,本地不翻转,保持一致)
result = await get_agent_client().draft_operate(
draft_id, "discard", auth_header=auth_header, trace_id=trace_id
)
# 3) 本地镜像翻转 discarded(已 sent 的终态不回退)
report = await AdvisorReportRepo(db).get_by_draft_id(draft_id)
if report is not None and report.send_status != REPORT_SEND_STATUS_SENT:
report.send_status = REPORT_SEND_STATUS_DISCARDED
await db.commit()
data = result["data"] or {}
if result["warning"]:
data["warning"] = result["warning"]
return data
# ---------------------------------------------------------------------------
# 发送(工作台本地,Agent 不感知)
# ---------------------------------------------------------------------------
async def send_draft(
db: AsyncSession, user: SysUser, *, auth_header: str, trace_id: str, draft_id: str
) -> dict:
repo = AdvisorReportRepo(db)
report = await repo.get_by_draft_id(draft_id)
# 无快照:未 save 直接发送,先拉 Agent 草稿建快照(D6 fallback)
if report is None:
detail = await get_draft(db, user, auth_header=auth_header, trace_id=trace_id, draft_id=draft_id)
report = _build_report_from_detail(detail, user.id)
elif report.advisor_id != user.id:
raise ForbiddenError("无权操作该客户数据")
# 幂等:已发送直接返回(重复点击不重复发站内信)
if report.send_status == REPORT_SEND_STATUS_SENT:
return {**_sent_payload(report), "duplicated": True}
# 数据权限 + 签约状态实时兜底(发送前实时查询,PRD §4.5.2)
rel = await require_owned_relation(db, user.id, report.customer_id)
if rel.status != CUSTOMER_REL_STATUS_SIGNED:
raise ForbiddenError("客户尚未签约,禁止发送报告")
# 发送终审(三检全过才放行;适当性需 Agent 侧结构化建议清单)
customer_risk = await _get_customer_risk(db, report.customer_id)
product_risks = await _resolve_product_risks(db, draft_id, auth_header, trace_id)
sensitive_words = [w.word for w in await SensitiveWordRepo(db).list_active()]
review_send(
content=report.content,
title=report.title,
customer_risk=customer_risk,
product_risks=product_risks,
sensitive_words=sensitive_words,
)
# 写站内信 + 报告置 sent,同事务(失败可重试,避免假送达)
msg_type = INTENT_TO_MSG_TYPE.get(report.intent, MSG_TYPE_RECOMMEND)
message = SysMessage(
user_id=report.customer_id,
msg_type=msg_type,
title=report.title or msg_type,
content=(report.content or "")[:512],
biz_id=report.report_id,
)
report.send_status = REPORT_SEND_STATUS_SENT
report.send_time = datetime.now()
report.send_by = user.id
db.add(message)
db.add(report) # 已跟踪对象时无副作用,新对象时入 session
await db.flush() # 生成 message.id,供回填 msg_id
report.msg_id = message.id
await db.commit()
return _sent_payload(report)
# ---------------------------------------------------------------------------
# 调仓触发 / 话术(代理 Agent)
# ---------------------------------------------------------------------------
async def run_rebalance(
db: AsyncSession, user: SysUser, *, auth_header: str, trace_id: str, req: RebalanceRunReq
) -> dict:
# 本地先校验归属 + 签约,快速失败(Agent 侧 40030 为兜底);未签约不生成待办
rel = await require_owned_relation(db, user.id, req.customer_id)
if rel.status != CUSTOMER_REL_STATUS_SIGNED:
raise ParamError("该客户尚未签约,不支持生成调仓建议")
result = await get_agent_client().rebalance_run(
req.customer_id, auth_header=auth_header, trace_id=trace_id
)
return result["data"] or {"accepted": True}
async def generate_talk_script(
db: AsyncSession, user: SysUser, *, auth_header: str, trace_id: str, req: TalkScriptReq
) -> dict:
await ensure_customer_owned(db, user.id, req.customer_id)
result = await get_agent_client().generate_talk_script(
req.customer_id, req.scene_type, auth_header=auth_header, trace_id=trace_id
)
data = result["data"] or {}
if result["warning"]:
data["warning"] = result["warning"]
return data
+150
View File
@@ -0,0 +1,150 @@
"""Redis 事件消费者(投顾域事件:rebalance_draft_created / profile_update)。
双写原则(common_const §9):Agent 双写 event_log + Redis Pub/Sub;Pub/Sub 仅实时通知,
消费以 event_id 幂等,补拉靠 event_log 扫描兜底。process_event 是幂等入口,实时订阅与
补拉两条路径共用。
"""
from __future__ import annotations
import asyncio
import json
import logging
from common_const import (
ADVISOR_EVENTS,
EVENT_PROFILE_UPDATE,
EVENT_REBALANCE_DRAFT_CREATED,
EVENT_STATUS_CONSUMED,
TODO_SOURCE_AGENT_EVENT,
TODO_TYPE_NEW_REBALANCE_DRAFT,
)
from model.advisor_todo import AdvisorTodo
from repositories.advisor_todo import AdvisorTodoRepo
from repositories.event_log import EventLogRepo
logger = logging.getLogger("service.advisor.event_consumer")
async def handle_rebalance_created(db, payload: dict) -> None:
"""消费 rebalance_draft_created → 生成待办(uk_todo 三元组去重,幂等)。"""
advisor_id = payload.get("advisor_id")
customer_id = payload.get("customer_id")
draft_id = payload.get("draft_id")
if not advisor_id:
logger.warning("rebalance_draft_created 缺少 advisor_id,忽略: %s", payload)
return
repo = AdvisorTodoRepo(db)
existing = await repo.get_by_unique(TODO_TYPE_NEW_REBALANCE_DRAFT, customer_id, draft_id)
if existing is not None:
return # 已存在,幂等跳过
todo = AdvisorTodo(
todo_type=TODO_TYPE_NEW_REBALANCE_DRAFT,
customer_id=customer_id,
advisor_id=int(advisor_id),
source=TODO_SOURCE_AGENT_EVENT,
biz_id=str(draft_id) if draft_id else None,
)
await repo.add(todo) # add 内 commit + refresh
async def handle_profile_update(db, payload: dict) -> None:
"""消费 profile_update → 失效该客户 Redis 画像缓存(不落库,画像只读)。
注:V1.0 工作台尚未建设 profile:{customer_id} 画像热缓存读取,此处仅占位;后续
360 读画像接入缓存后生效。事件仍需标记已消费。
"""
return
async def process_event(db, *, event_id: str | None, event_name: str | None, payload: dict) -> None:
"""处理单条事件(幂等入口,实时订阅与补拉共用)。"""
# 幂等兜底:event_log 已消费则跳过(实时订阅与补拉并发时防重)
if event_id:
existing = await EventLogRepo(db).get_by_event_id(event_id)
if existing is not None and existing.status == EVENT_STATUS_CONSUMED:
return
if event_name == EVENT_REBALANCE_DRAFT_CREATED:
await handle_rebalance_created(db, payload)
elif event_name == EVENT_PROFILE_UPDATE:
await handle_profile_update(db, payload)
if event_id:
await EventLogRepo(db).mark_consumed(event_id)
def parse_message(raw: str) -> tuple[str | None, str | None, dict] | None:
"""解析 Redis Pub/Sub 消息为 (event_id, event_name, payload)。
容忍两种结构:带内层 payload 的标准结构;无内层 payload 时把外层整体当 payload。
"""
try:
msg = json.loads(raw)
except (ValueError, TypeError):
return None
if not isinstance(msg, dict):
return None
event_name = msg.get("event_name")
event_id = msg.get("event_id") or msg.get("id")
payload = msg.get("payload")
if not isinstance(payload, dict):
payload = msg
return event_id, event_name, payload
async def pull_pending_events(db) -> int:
"""补拉 event_log 中未消费的投顾域事件(Pub/Sub 丢消息兜底),返回处理条数。"""
events = await EventLogRepo(db).list_pending(ADVISOR_EVENTS, limit=100)
for ev in events:
await process_event(
db, event_id=ev.event_id, event_name=ev.event_name, payload=ev.payload or {}
)
return len(events)
class EventConsumer:
"""Redis 订阅循环(后台任务,dev 默认关闭,避免 --reload 重复订阅)。"""
def __init__(self, redis):
self.redis = redis
self._task: asyncio.Task | None = None
async def _run(self) -> None:
from config.database.mysql import get_session_factory
pubsub = self.redis.pubsub()
await pubsub.subscribe(*ADVISOR_EVENTS)
try:
async for message in pubsub.listen():
if message.get("type") != "message":
continue
parsed = parse_message(message.get("data"))
if parsed is None:
continue
event_id, event_name, payload = parsed
if event_name not in ADVISOR_EVENTS:
continue
# 每事件开独立 session,避免长事务占用连接
async with get_session_factory()() as session:
try:
await process_event(
session, event_id=event_id, event_name=event_name, payload=payload
)
except Exception:
logger.exception("process event failed: %s %s", event_name, event_id)
finally:
await pubsub.unsubscribe(*ADVISOR_EVENTS)
await pubsub.aclose()
def start(self) -> None:
if self._task is None:
self._task = asyncio.create_task(self._run())
async def stop(self) -> None:
if self._task is not None:
self._task.cancel()
try:
await self._task
except asyncio.CancelledError:
pass
self._task = None
+33
View File
@@ -0,0 +1,33 @@
"""敏感信息脱敏工具(工作台统一出口,避免各接口遗漏导致合规风险)。
两层脱敏(PRD §4.2.2):
- 常规敏感(手机号中 4 位、姓名/证件常规掩码):默认掩码展示;投顾查看完整值时,
调用方负责写审计 action=view_sensitive;
- 强敏感(完整身份证号、银行卡号):V1.0 一律只返回掩码,不提供完整值。
"""
from __future__ import annotations
def mask_phone(phone: str | None) -> str:
"""手机号:保留前 3 后 4,中间打码(138****1234)。"""
if not phone:
return ""
if len(phone) >= 7:
return phone[:3] + "****" + phone[-4:]
return "****"
def mask_idcard(idcard: str | None) -> str:
"""身份证:保留前 4 后 4,中间打码。"""
if not idcard:
return ""
if len(idcard) < 8:
return "****"
return idcard[:4] + "**********" + idcard[-4:]
def mask_name(name: str | None) -> str:
"""姓名:仅保留首字符(姓氏),其余打码。"""
if not name:
return ""
return name[0] + "*" * (len(name) - 1) if len(name) > 1 else name
+28
View File
@@ -0,0 +1,28 @@
"""数据权限校验(仅本人名下客户,PRD §五「最小权限原则」)。
customer_relation.advisor_id = 当前投顾 id 是数据权限唯一判据;越权一律 403。
"""
from __future__ import annotations
from sqlalchemy.ext.asyncio import AsyncSession
from model.customer_relation import CustomerRelation
from repositories.customer_relation import CustomerRelationRepo
from utils.exceptions import ForbiddenError
async def require_owned_relation(
db: AsyncSession, advisor_id: int, customer_id: int
) -> CustomerRelation:
"""校验客户归属并返回关系记录(调用方可继续判断签约状态);越权抛 403。"""
rel = await CustomerRelationRepo(db).get_by_customer_advisor(customer_id, advisor_id)
if rel is None:
raise ForbiddenError("无权操作该客户数据")
return rel
async def ensure_customer_owned(
db: AsyncSession, advisor_id: int, customer_id: int
) -> None:
"""仅校验归属,不返回关系(不需要签约状态时用)。"""
await require_owned_relation(db, advisor_id, customer_id)
+38
View File
@@ -0,0 +1,38 @@
"""个人绩效统计(PRD §4.7)。
口径说明(E1 最小口径):个人 AUM=名下客户持有中市值之和;报告发送数=advisor_report
sent 计数;回访次数=advisor_visit_record 计数。服务完成率/回访达标率等 P2 后置。
"""
from __future__ import annotations
from sqlalchemy.ext.asyncio import AsyncSession
from common_const import CUSTOMER_REL_STATUS_SIGNED, REPORT_SEND_STATUS_SENT
from model.sys_user import SysUser
from repositories.advisor_report import AdvisorReportRepo
from repositories.advisor_visit_record import AdvisorVisitRecordRepo
from repositories.customer_relation import CustomerRelationRepo
async def personal_report(db: AsyncSession, user: SysUser) -> dict:
relation_repo = CustomerRelationRepo(db)
report_repo = AdvisorReportRepo(db)
visit_repo = AdvisorVisitRecordRepo(db)
total = await relation_repo.count_customer_rows(advisor_id=user.id)
signed = await relation_repo.count_customer_rows(
advisor_id=user.id, status=CUSTOMER_REL_STATUS_SIGNED
)
aum = await relation_repo.sum_holdings_value(user.id)
sent = await report_repo.count_by_advisor(
advisor_id=user.id, send_status=REPORT_SEND_STATUS_SENT
)
visits = await visit_repo.count_by_advisor(advisor_id=user.id)
return {
"total_aum": f"{aum:.2f}",
"total_customers": total,
"signed_customers": signed,
"reports_sent": sent,
"visits_count": visits,
}
+83
View File
@@ -0,0 +1,83 @@
"""定时调度(工作台侧,dev 默认关):每日组合再平衡(遍历 signed 客户)。
架构 §7 已选型 APScheduler;仅当 config.advisor.scheduler_enabled=true 时启动。
dev 常以 uvicorn --reload 运行(重载会重复调度),故默认关闭,生产单进程开启。
边界:调度器无需前端/无需请求上下文,需自造投顾 JWT(service.auth.create_token)与
trace_id(new_request_id)调用 Agent;大额申赎扫描/回访到期/风评到期等本地定时任务
P1 后续补(依赖 fin_transaction 模型等)。
"""
from __future__ import annotations
import logging
from common_const import CRON_PORTFOLIO_REBALANCE, CUSTOMER_REL_STATUS_SIGNED
from config.settings import settings
from repositories.customer_relation import CustomerRelationRepo
logger = logging.getLogger("service.advisor.scheduler")
class AdvisorScheduler:
def __init__(self):
self._scheduler = None
def start(self) -> None:
if not settings.advisor.scheduler_enabled:
return
# 懒加载:未安装 apscheduler 时不阻塞应用启动(默认关闭本调度器)
from apscheduler.schedulers.asyncio import AsyncIOScheduler
from apscheduler.triggers.cron import CronTrigger
self._scheduler = AsyncIOScheduler()
self._scheduler.add_job(
self.daily_rebalance,
CronTrigger.from_crontab(CRON_PORTFOLIO_REBALANCE),
id="advisor_daily_rebalance",
name="每日组合再平衡(遍历 signed 客户)",
)
self._scheduler.start()
logger.info("advisor scheduler started (cron=%s)", CRON_PORTFOLIO_REBALANCE)
def shutdown(self) -> None:
if self._scheduler is not None:
self._scheduler.shutdown(wait=False)
self._scheduler = None
async def daily_rebalance(self) -> None:
"""每日遍历已签约客户,按投顾逐个调 Agent rebalance/run(异步受理)。
结果靠事件 event:rebalance_draft_created 回执 → 工作台消费生成待办。
"""
from config.database.mysql import get_session_factory
from service.advisor.agent_client import get_agent_client
from service.auth import create_token
from utils.request_id import new_request_id
client = get_agent_client()
if not client.configured:
logger.warning("投顾Agent 未配置,跳过每日 rebalance")
return
async with get_session_factory()() as session:
relations = await CustomerRelationRepo(session).list_by_status(
CUSTOMER_REL_STATUS_SIGNED
)
by_advisor: dict[int, list[int]] = {}
for rel in relations:
by_advisor.setdefault(rel.advisor_id, []).append(rel.customer_id)
for advisor_id, customer_ids in by_advisor.items():
# 后台任务无请求上下文,自造投顾 JWT + trace_id
token = create_token(int(advisor_id))
auth = f"Bearer {token}"
for customer_id in customer_ids:
try:
await client.rebalance_run(
customer_id, auth_header=auth, trace_id=new_request_id()
)
except Exception:
logger.exception(
"rebalance run failed advisor=%s customer=%s", advisor_id, customer_id
)
+51
View File
@@ -0,0 +1,51 @@
"""发送前合规终审(工作台本地独立校验,不依赖 Agent 返回值)。
PRD §4.5.2 / §4.8:发送环节必须再次独立执行全套校验(敏感词 + 适当性 + 免责声明),
不可完全依赖 Agent 保存阶段的结果。三检全过才放行,任一不过抛异常拦截(fail-closed)。
"""
from __future__ import annotations
from common_const import DISCLAIMER_TEXT
from service.advisor.suitability import check_suitability
from utils.exceptions import ForbiddenError, NotSuitableError
def check_disclaimer(content: str | None) -> bool:
"""免责声明完整性:正文必须「逐字包含」完整原文(精确子串匹配)。"""
return bool(content) and DISCLAIMER_TEXT in content
def check_sensitive_words(content: str, words: list[str]) -> list[str]:
"""敏感词扫描:返回命中的词列表(空列表=通过)。"""
text = content or ""
return [w for w in words if w and w in text]
def review_send(
*,
content: str | None,
title: str | None,
customer_risk: str | None,
product_risks: list[str | None],
sensitive_words: list[str],
) -> None:
"""发送终审主入口:三项校验,任一不过抛异常。
- 敏感词:命中即 ForbiddenError(403);
- 适当性:任一建议产品风险高于客户风险即 NotSuitableError(1005);
- 免责声明:正文未逐字包含完整原文即 ForbiddenError(403)。
"""
# 1) 敏感词(标题 + 正文合并扫描)
hit = check_sensitive_words(f"{title or ''}\n{content or ''}", sensitive_words)
if hit:
raise ForbiddenError(f"报告包含敏感词:{'、'.join(hit[:5])}")
# 2) 适当性:对建议清单中每个产品风险等级逐一校验
for pr in product_risks:
result = check_suitability(customer_risk, pr)
if not result["ok"]:
raise NotSuitableError(result["reason"])
# 3) 免责声明完整性(精确匹配;编辑改动导致不匹配即拦截)
if not check_disclaimer(content):
raise ForbiddenError("报告缺少免责声明,禁止发送")
+44
View File
@@ -0,0 +1,44 @@
"""适当性校验(工作台本地桩实现,契约与投顾Agent 共享包保持一致)。
common_const §7 决议:公共适当性校验函数沉淀为独立共享包 `suitability`,工作台与
Agent 复用同一实现、禁止各自复制。在 Agent 侧共享包发布前,本模块提供**契约一致的
本地桩**:`check_suitability(customer_risk, product_risk) -> {ok, reason}`;
共享包就绪后仅替换 import,不改变调用方。
"""
from __future__ import annotations
# 风险等级 → 序号(C_n / R_n 的 n)。兼容两套口径:标准 C1-C5/R1-R5 与历史中文等级,
# 避免问卷侧尚未完成 C1-C5 归一化时发送终审误判。
# 假设:中文等级与 C 级一一对应(保守=C1/稳健=C2/平衡=C3/进取=C4/激进=C5)。
_RISK_RANK = {
"C1": 1, "C2": 2, "C3": 3, "C4": 4, "C5": 5,
"R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5,
"保守": 1, "稳健": 2, "平衡": 3, "进取": 4, "激进": 5,
}
def _rank(level: str | None) -> int | None:
"""风险等级统一转序号;未知/空返回 None(表示无法判定)。"""
if not level:
return None
return _RISK_RANK.get(str(level).strip())
def check_suitability(customer_risk: str | None, product_risk: str | None) -> dict:
"""适当性硬规则:产品风险 R 不得高于客户风险 C(R_n ≤ C_n)。
返回 {"ok": bool, "reason": str}。无法判定(缺等级/未知编码)视为不通过并说明原因,
宁可拦截不放行(合规红线,发送环节 fail-closed)。
"""
customer_rank = _rank(customer_risk)
product_rank = _rank(product_risk)
if customer_rank is None:
return {"ok": False, "reason": "客户无有效风险等级,无法进行适当性校验"}
if product_rank is None:
return {"ok": False, "reason": f"产品风险等级无法识别: {product_risk}"}
if product_rank > customer_rank:
return {
"ok": False,
"reason": f"产品风险等级 {product_risk} 高于客户风险等级 {customer_risk}",
}
return {"ok": True, "reason": ""}
+71
View File
@@ -0,0 +1,71 @@
"""投顾待办服务(事件/定时/计算三类来源统一承载,advisor_todo)。"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy.ext.asyncio import AsyncSession
from common_const import TODO_STATUS_DONE, TODO_STATUS_PENDING, TODO_STATUS_PROCESSING
from model.advisor_todo import AdvisorTodo
from model.sys_user import SysUser
from repositories.advisor_todo import AdvisorTodoRepo
from schemas.advisor import TodoHandleReq
from utils.exceptions import ForbiddenError, NotFoundError, ParamError
def _todo_item(t: AdvisorTodo) -> dict:
return {
"id": t.id,
"todo_type": t.todo_type,
"customer_id": t.customer_id,
"advisor_id": t.advisor_id,
"source": t.source,
"biz_id": t.biz_id,
"priority": t.priority,
"status": t.status,
"due_at": t.due_at.isoformat() if t.due_at else None,
"create_time": t.create_time.isoformat() if t.create_time else None,
"handle_time": t.handle_time.isoformat() if t.handle_time else None,
}
async def list_todos(
db: AsyncSession,
user: SysUser,
*,
status: str | None = None,
todo_type: str | None = None,
page: int = 1,
page_size: int = 20,
) -> dict:
repo = AdvisorTodoRepo(db)
items = await repo.list_by_advisor(
advisor_id=user.id, status=status, todo_type=todo_type,
limit=page_size, offset=(page - 1) * page_size,
)
total = await repo.count_by_advisor(advisor_id=user.id, status=status, todo_type=todo_type)
return {"total": total, "page": page, "page_size": page_size, "items": [_todo_item(t) for t in items]}
async def handle_todo(db: AsyncSession, user: SysUser, todo_id: int, req: TodoHandleReq) -> dict:
"""待办处理:process=开始处理 / done=完成;单向流转,重复提交同目标态幂等成功。"""
todo = await AdvisorTodoRepo(db).get(todo_id)
if todo is None:
raise NotFoundError("待办不存在")
if todo.advisor_id != user.id:
raise ForbiddenError("无权操作该待办")
if req.action == "process":
if todo.status == TODO_STATUS_PROCESSING:
return {"todo_id": todo.id, "status": todo.status}
if todo.status != TODO_STATUS_PENDING:
raise ParamError("仅待处理状态的待办可开始处理")
todo.status = TODO_STATUS_PROCESSING
else: # done
if todo.status == TODO_STATUS_DONE:
return {"todo_id": todo.id, "status": todo.status}
todo.status = TODO_STATUS_DONE
todo.handle_time = datetime.now()
await db.commit()
return {"todo_id": todo.id, "status": todo.status}
+98
View File
@@ -0,0 +1,98 @@
"""投后服务与陪伴:回访留痕 + 合规话术库 + 触达日志 + AI 话术(代理 Agent)。
PRD §4.4:回访对话不经过投顾Agent,人工录入仅留痕,不更新记忆单元;
AI 话术仅作参考,正式发送走工作台消息模块。
"""
from __future__ import annotations
from sqlalchemy.ext.asyncio import AsyncSession
from model.advisor_visit_record import AdvisorVisitRecord
from model.sys_user import SysUser
from repositories.advisor_visit_record import AdvisorVisitRecordRepo
from repositories.sys_message import SysMessageRepo
from schemas.advisor import VisitCreateReq
from service.advisor.permissions import ensure_customer_owned
from utils.exceptions import ParamError
# 合规话术库(内置,标准化投教/市场解读/调仓沟通)。
# 假设:V1.0 内置常量,后续可迁移 sys_config 运营化;内容不含承诺收益等敏感词。
TALK_TEMPLATES = [
{"scene": "市场波动安抚", "content": "市场短期波动属正常现象,请结合自身风险承受能力理性看待,勿因短期涨跌追涨杀跌。"},
{"scene": "调仓沟通", "content": "本次调整旨在使组合回归目标配置,不构成任何收益承诺,请结合自身情况审慎判断。"},
{"scene": "风险测评到期提醒", "content": "您的风险测评即将到期,为保障适当性匹配,请及时完成最新测评。"},
]
def _visit_item(v: AdvisorVisitRecord) -> dict:
return {
"id": v.id,
"customer_id": v.customer_id,
"advisor_id": v.advisor_id,
"visit_type": v.visit_type,
"visit_time": v.visit_time.isoformat(),
"summary": v.summary,
"audio_url": v.audio_url,
"create_time": v.create_time.isoformat() if v.create_time else None,
}
async def list_visits(
db: AsyncSession, user: SysUser, *, customer_id: int | None = None,
page: int = 1, page_size: int = 20,
) -> dict:
if customer_id is not None:
await ensure_customer_owned(db, user.id, customer_id)
repo = AdvisorVisitRecordRepo(db)
items = await repo.list_by_advisor(
advisor_id=user.id, customer_id=customer_id,
limit=page_size, offset=(page - 1) * page_size,
)
total = await repo.count_by_advisor(advisor_id=user.id, customer_id=customer_id)
return {"total": total, "page": page, "page_size": page_size, "items": [_visit_item(v) for v in items]}
async def create_visit(db: AsyncSession, user: SysUser, req: VisitCreateReq) -> dict:
await ensure_customer_owned(db, user.id, req.customer_id)
record = AdvisorVisitRecord(
customer_id=req.customer_id,
advisor_id=user.id,
visit_type=req.visit_type,
visit_time=req.visit_time,
summary=req.summary,
audio_url=req.audio_url,
)
record = await AdvisorVisitRecordRepo(db).add(record)
return {"visit_id": record.id}
def list_talk_templates() -> list[dict]:
"""合规话术库(内置,投顾参考;不自动发送)。"""
return TALK_TEMPLATES
async def list_touch_logs(
db: AsyncSession, user: SysUser, customer_id: int, *, page: int = 1, page_size: int = 20
) -> dict:
"""触达留痕:站内信(发送触达)+ 回访记录(人工沟通)聚合。
假设:V1.0 触达通道仅站内信与回访;短信/企微/电话统一落在回访记录。
"""
await ensure_customer_owned(db, user.id, customer_id)
messages = await SysMessageRepo(db).list_by_user(
customer_id, limit=page_size, offset=(page - 1) * page_size
)
visits = await AdvisorVisitRecordRepo(db).list_by_advisor(
advisor_id=user.id, customer_id=customer_id,
limit=page_size, offset=(page - 1) * page_size,
)
return {
"messages": [
{
"id": m.id, "msg_type": m.msg_type, "title": m.title,
"content": m.content, "create_time": m.create_time.isoformat() if m.create_time else None,
}
for m in messages
],
"visits": [_visit_item(v) for v in visits],
}
+282
View File
@@ -0,0 +1,282 @@
# common_const.md
> 公共常量定义 用途:投顾工作台PRD(文档A)、投顾Agent需求文档(文档B)共同引用;后端、AI组件统一一份,避免枚举/字符串硬编码不一致。 版本:v1.2(架构评审决议修订) 维护方:产品 + 后端 + AI开发 注意:业务代码禁止直接写死字符串字面量,全部引用本文件定义常量。
------
## 1. customer_relation 客户‑投顾关系状态
表:`customer_relation`
```
# 客户投顾签约状态
CUSTOMER_REL_STATUS_UNSIGNED = "unsigned" # 未签约,仅分配,未开通投顾正式服务
CUSTOMER_REL_STATUS_SIGNED = "signed" # 已签约,可下发推荐/调仓方案给客户
CUSTOMER_REL_STATUS_CLOSED = "closed" # 服务已终止,投顾不再提供服务
```
说明:
1. 只有 `signed` 状态,允许工作台将草稿方案下发至客户Web端;
2. `unsigned`:Agent的rebalance接口直接拦截;recommend可以生成草稿供内部预览,但工作台禁止发送;
3. `closed`:禁止调用投顾Agent为该客户生成对外方案;
4. 签约状态**唯一写入方为工作台**(投顾认领/签约/结束服务),发送前工作台必须实时查询 `status` 兜底校验;
5. 存储值统一为英文 `unsigned/signed/closed`,数据库与代码禁止使用中文「已分配/已签约/已结束」。
------
## 2. Agent草稿 draft 状态枚举
> 存储:Agent内部draft草稿模型,持久化MySQL;草稿不能物理删除,仅做状态流转归档
```
DRAFT_STATUS_DRAFT = "draft" # 草稿初始状态,可编辑修改
DRAFT_STATUS_DISCARDED = "discarded"# 投顾废弃该草稿,归档保留,不可编辑、不可下发
```
> 说明:`sent`为工作台本地状态,Agent不维护。 状态流转:
- draft → discarded
- discarded **不允许回退到draft**。
------
## 3. memory_unit 记忆单元类型
表:`memory_unit`
```
MEMORY_INFO_TYPE_FACT = "FACT" # 客观事实(问卷、交易行为、基础身份信息,作为业务硬规则依据)
MEMORY_INFO_TYPE_OPINION = "OPINION"# 主观观点/情绪(客户自述感受,仅做排序参考,不能绕过适当性校验)
```
## 4. Agent意图枚举(intent)
> 投顾Agent内部意图路由识别结果,SSE‑meta、草稿元数据中使用
```
AGENT_INTENT_RECOMMEND = "recommend" # 基金推荐意图
AGENT_INTENT_REBALANCE = "rebalance" # 持仓诊断&调仓再平衡意图
AGENT_INTENT_FUND_ANALYSIS = "fund_analysis"# 基金深度分析意图
AGENT_INTENT_DIALOGUE_SCRIPT = "dialogue-script" # 生成沟通话术意图
```
## 5. 沟通话术场景 scene_type(generate‑talk‑script接口入参)
```
TALK_SCENE_RISK_BLOCK_ORDER = "risk_block_order" # 订单被风控拦截场景
TALK_SCENE_MARKET_FLUCTUATION = "market_fluctuation"# 市场波动安抚客户
TALK_SCENE_PORTFOLIO_DIVERGENCE = "portfolio_divergence" # 组合大幅偏离基准
TALK_SCENE_CUSTOMER_COMPLAINT = "customer_complaint"# 客户投诉场景
```
## 6. Agent业务错误码定义
> HTTP返回 `code` 字段,业务层错误码;http状态码统一200,靠code区分业务异常
```
ERR_CODE_OK = 0 # 成功
ERR_CODE_FORBIDDEN_CUSTOMER = 40001 # customer_id不属于当前投顾,越权访问
ERR_CODE_SUITABILITY_INVALID = 40020 # 适当性校验不通过:方案包含高于客户风险等级的产品
ERR_CODE_NOT_SIGNED_REBALANCE = 40030 # 客户未签约,禁止生成rebalance调仓草稿
ERR_CODE_DRAFT_NOT_FOUND = 40401 # draft_id不存在或已废弃
ERR_CODE_LLM_ERROR = 50001 # Agent内部LLM调用异常
ERR_CODE_GRAPH_ERROR = 50002 # GraphRAG查询异常,触发降级
```
| code | message(默认提示文案) |
| ----- | ------------------------------------------------ |
| 0 | success |
| 40001 | 无权操作该客户数据 |
| 40020 | 方案适当性校验不通过,包含超出客户风险等级的产品 |
| 40030 | 客户尚未签约,禁止生成调仓草稿 |
| 40401 | 草稿不存在或者已废弃 |
| 50001 | AI服务调用异常,请稍后重试 |
| 50002 | 图谱查询异常,已降级返回部分结果 |
> message 为后端默认文案;业务逻辑判断必须使用 `code`,禁止依赖 message 文案。
> **错误码域边界(重要)**:投顾Agent 返回的 `code`(0/40001/40020/40030/40401/50001/50002)为 Agent 自有错误码域;投顾工作台后端自有错误码域为 200/400/401/403/404/500/1001-1005(成功码为 200)。工作台调用 Agent 后**不得将 Agent 内部码原样透传给上层调用方**,须在工作台 service 层将 Agent 码映射为工作台业务分支 + 可读提示(例如 40030 → 提示"客户尚未签约,不支持生成调仓建议")。
------
## 7. 风险等级常量
> 风险测评问卷输出,客户风险等级C1‑C5;产品风险等级R1‑R5
```
# 客户风险等级(C‑Customer)
C_RISK_C1 = "C1"
C_RISK_C2 = "C2"
C_RISK_C3 = "C3"
C_RISK_C4 = "C4"
C_RISK_C5 = "C5"
# 基金产品风险等级(R‑Risk)
PROD_RISK_R1 = "R1"
PROD_RISK_R2 = "R2"
PROD_RISK_R3 = "R3"
PROD_RISK_R4 = "R4"
PROD_RISK_R5 = "R5"
```
> 适当性匹配规则:客户Cn,可购买产品R ≤ n;**硬规则在B端业务层、Agent层两处都要校验,双重防护**,且两处必须复用同一公共适当性校验函数、读取同一数据源(客户C级取自 `fin_risk_assessment.risk_level`,产品R级取自 `fin_product.risk_level`)。
> **共享包落地(评审决议)**:公共适当性校验函数沉淀为独立共享包 `suitability`(如 `common/suitability`),工作台与 Agent 均以内部依赖引入,**禁止各自复制实现**。函数契约:`check_suitability(customer_risk: str, product_risk: str) -> {ok: bool, reason: str}`;判定规则 `R_n ≤ C_n`。
> **存储编码统一**:客户风险等级统一存 `C1-C5`,产品风险等级统一存 `R1-R5`。历史中文等级「保守/稳健/平衡/进取/激进」仅允许存在于展示层,与 C 级一一对应(保守=C1、稳健=C2、平衡=C3、进取=C4、激进=C5)。数据库字段 `fin_customer_profile.risk_level`、`fin_risk_assessment.risk_level`、`portfolio_benchmark.risk_level` 必须存 `C1-C5`,`fin_product.risk_level` 存 `R1-R5`。
------
## 8. 系统配置 sys_config key常量
> sys_config表 KV配置,调仓、Agent相关参数key,代码不要写死字符串key
```
SYS_KEY_REBALANCE_DEVIATION_THRESHOLD = "rebalance.deviation.threshold" # 组合再平衡偏离阈值(全局兜底),浮点数,例0.1代表±10%
SYS_KEY_HIGH_NET_ASSET_THRESHOLD = "customer.high_net.asset.threshold" # 高净值客户资产门槛
SYS_KEY_RISK_QUESTIONNAIRE_EXPIRE_DAY = "risk.questionnaire.expire.day" # 风险测评过期天数
SYS_KEY_LARGE_FLOW_THRESHOLD = "customer.large_flow.threshold" # 大额申赎阈值(工作台本地定时任务扫描 fin_transaction)
```
------
## 9. 事件总线 Pub/Sub 事件名称
> Redis Pub/Sub事件频道名,业务同时双写落库,保证消息不丢失
```
EVENT_ADVISOR_REBALANCE_DRAFT_CREATED = "event:rebalance_draft_created" # rebalance草稿生成(工作台消费→待办)
EVENT_PROFILE_UPDATE = "event:profile_update" # 客户画像/记忆更新(工作台只读刷新)
EVENT_WORK_ORDER_CHANGE = "event:work_order_change" # 工单状态变更(工单域,非投顾工作台事件契约)
EVENT_RISK_ALERT = "event:risk_alert" # 风控预警事件(风控域,非投顾工作台事件契约)
> 投顾工作台与投顾Agent 的事件契约仅包含前两者:`event:rebalance_draft_created`、`event:profile_update`;`work_order_change`、`risk_alert` 属其他域事件。
```
### 事件消息公共字段约定(JSON payload)
> 事件结构 = 外层公共字段 + 内层 payload。以下为外层公共字段:
```
{
"event_name": "",
"trace_id": "",
"trigger_user_id": "",
"customer_id": "",
"payload": {}
}
```
> 内层 payload 按事件名定义。`event:rebalance_draft_created` 的内层 payload 结构如下(`advisor_id` 放内层):
```
{
"draft_id": "",
"customer_id": "",
"advisor_id": "",
"deviation": 0.0,
"created_at": ""
}
```
> 事件双写 `event_log` 表(event_name + event_id 幂等 + payload + trace_id + 消费状态),Pub/Sub 仅作实时通知;工作台消费以 event_id 幂等去重。
------
## 10. 固定文本模板
### 10.1 Agent输出报告强制免责声明
> ⚠️ Agent生成markdown报告必须拼接此文本;草稿保存接口校验,如果缺失该声明返回警告;**真正拦截发送发生在工作台发送接口**。
> **完整性判定规则(决议)**:报告正文必须**逐字包含**上述免责声明完整原文(精确匹配);编辑器中免责声明为**只读区**,投顾改动导致不匹配时,工作台发送环节直接拦截。
```
【免责声明】本报告由AI辅助生成,仅供持牌投顾内部参考,不构成任何投资建议。基金有风险,投资需谨慎。所有投资决策请结合自身风险承受能力审慎判断。
```
### 10.2 SSE事件类型常量
> Agent流式SSE推送事件type
```
SSE_EVENT_TYPE_TEXT = "text" # 增量文本片段
SSE_EVENT_TYPE_META = "meta" # 结构化元数据
SSE_EVENT_TYPE_DONE = "done" # 会话正常结束
SSE_EVENT_TYPE_ERROR = "error" # 会话异常中断
```
------
## 11. 定时任务相关常量
```
# 每周记忆维护任务:置信度重算、矛盾检测、记忆遗忘归档
CRON_MEMORY_MAINTENANCE = "0 2 * * 1"
# 每日客户组合再平衡任务【由工作台后端调度,Agent不调度】
CRON_PORTFOLIO_REBALANCE = "0 1 * * *"
# 基金净值更新定时任务
CRON_FUND_NAV_UPDATE = "30 1 * * *"
```
------
## 12. 数据库审计日志 action 动作常量 audit_log.action
```
AUDIT_AGENT_CHAT_CALL = "agent_chat_call" # 调用投顾Agent对话
AUDIT_DRAFT_SAVE = "draft_save" # 草稿保存
AUDIT_DRAFT_DISCARD = "draft_discard" # 草稿废弃
AUDIT_CUSTOMER_SIGN = "customer_sign" # 客户签约
AUDIT_VIEW_SENSITIVE = "view_sensitive" # 投顾查看脱敏字段完整值(留痕)
```
------
## 13. 投顾角色常量(sys_user.employee_role)
> 工作台 RBAC 与投顾Agent 鉴权统一使用本值,「理财顾问」为同岗位历史表述,不再参与鉴权。
```
EMPLOYEE_ROLE_ADVISOR = "投顾" # 投顾(工作台唯一业务角色)
```
## 14. 站内信消息类型常量(sys_message.msg_type)
```
MSG_TYPE_RECOMMEND = "推荐" # 基金推荐报告
MSG_TYPE_REBALANCE = "调仓" # 调仓建议报告
MSG_TYPE_RISK = "风控" # 风控预警通知
MSG_TYPE_SIGN = "签约" # 签约相关通知
MSG_TYPE_SYSTEM = "系统" # 系统通知
```
> 草稿 intent → msg_type 映射:`recommend`→`推荐`,`rebalance`→`调仓`。
## 15. 新增数据表(本次评审决议)
| 表名 | 归属 | 用途 |
| ---- | ---- | ---- |
| `advisor_draft` | 投顾Agent | draft 草稿持久化(draft/discarded,不存 sent) |
| `advisor_report` | 投顾工作台 | 建议报告本地镜像:sent 状态 + 编辑/发送留痕 |
| `advisor_todo` | 投顾工作台 | 待办任务(事件/定时/计算三类来源统一承载) |
| `advisor_visit_record` | 投顾工作台 | 人工回访纪要留痕归档 |
| `event_log` | 共享 | 事件双写落库(event_id 幂等,保证不丢) |
| `sensitive_word` | 共享 | 敏感词库(工作台与 Agent 同源,合规拦截唯一来源) |
> 字段结构以 `sql/schema.sql` 为准;`advisor_draft.status` 仅 `draft/discarded`,`advisor_report.send_status` 为 `draft/sent/discarded`。
------
## 16. 敏感词库数据源约定
> 敏感词唯一数据源为 `sensitive_word` 表;工作台发送终审与 Agent 生成校验均读取该表,禁止各自维护词库。V1.0 由后端灌种子数据 + 管理脚本维护,运营端维护界面随 Phase 5 补。
------
# 使用规范(重要)
1. 所有 Python/后端代码,import 本常量文件,禁止硬编码字符串;
2. 所有状态、错误码、事件名称修改,必须同步更新此文档,通知 AI、B 端后端;
3. 业务规则(如适当性匹配逻辑)写在业务代码中,不在此文档写业务逻辑,本文件只存放**常量字符串、枚举、模板文本**。