diff --git a/_conflict_backup/model/customer_relation.py b/_conflict_backup/model/customer_relation.py new file mode 100644 index 0000000..48f6111 --- /dev/null +++ b/_conflict_backup/model/customer_relation.py @@ -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)) diff --git a/_conflict_backup/repositories/customer_relation.py b/_conflict_backup/repositories/customer_relation.py new file mode 100644 index 0000000..562d8cf --- /dev/null +++ b/_conflict_backup/repositories/customer_relation.py @@ -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) diff --git a/api/advisor/__init__.py b/api/advisor/__init__.py new file mode 100644 index 0000000..bff65c2 --- /dev/null +++ b/api/advisor/__init__.py @@ -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) diff --git a/api/advisor/_auth.py b/api/advisor/_auth.py new file mode 100644 index 0000000..8b5adb2 --- /dev/null +++ b/api/advisor/_auth.py @@ -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 diff --git a/api/advisor/audit.py b/api/advisor/audit.py new file mode 100644 index 0000000..113628f --- /dev/null +++ b/api/advisor/audit.py @@ -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"'}, + ) diff --git a/api/advisor/customers.py b/api/advisor/customers.py new file mode 100644 index 0000000..d18e9e2 --- /dev/null +++ b/api/advisor/customers.py @@ -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)) diff --git a/api/advisor/dashboard.py b/api/advisor/dashboard.py new file mode 100644 index 0000000..98adcc0 --- /dev/null +++ b/api/advisor/dashboard.py @@ -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)) diff --git a/api/advisor/diagnosis.py b/api/advisor/diagnosis.py new file mode 100644 index 0000000..ef86a2d --- /dev/null +++ b/api/advisor/diagnosis.py @@ -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, + ) + ) diff --git a/api/advisor/drafts.py b/api/advisor/drafts.py new file mode 100644 index 0000000..ff9791a --- /dev/null +++ b/api/advisor/drafts.py @@ -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)) diff --git a/api/advisor/report.py b/api/advisor/report.py new file mode 100644 index 0000000..2669ab3 --- /dev/null +++ b/api/advisor/report.py @@ -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)) diff --git a/api/advisor/todos.py b/api/advisor/todos.py new file mode 100644 index 0000000..1fba54d --- /dev/null +++ b/api/advisor/todos.py @@ -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)) diff --git a/api/advisor/visits.py b/api/advisor/visits.py new file mode 100644 index 0000000..e52b156 --- /dev/null +++ b/api/advisor/visits.py @@ -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 + ) + ) diff --git a/common_const.py b/common_const.py new file mode 100644 index 0000000..541ba3e --- /dev/null +++ b/common_const.py @@ -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" # 结束服务 diff --git a/model/advisor_report.py b/model/advisor_report.py new file mode 100644 index 0000000..9088bef --- /dev/null +++ b/model/advisor_report.py @@ -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() + ) diff --git a/model/advisor_todo.py b/model/advisor_todo.py new file mode 100644 index 0000000..d02f0d3 --- /dev/null +++ b/model/advisor_todo.py @@ -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) diff --git a/model/advisor_visit_record.py b/model/advisor_visit_record.py new file mode 100644 index 0000000..94615dd --- /dev/null +++ b/model/advisor_visit_record.py @@ -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()) diff --git a/model/audit_log.py b/model/audit_log.py new file mode 100644 index 0000000..cf10582 --- /dev/null +++ b/model/audit_log.py @@ -0,0 +1,26 @@ +"""audit_log 操作审计日志表 ORM(企业级合规留痕,工作台台账只读)。""" +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 AuditLog(Base): + __tablename__ = "audit_log" + __table_args__ = {"comment": "操作审计日志表(企业级合规留痕)"} + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + user_id: Mapped[int | None] = mapped_column(BigInteger) + username: Mapped[str | None] = mapped_column(String(64)) + module: Mapped[str] = mapped_column(String(32)) + action: Mapped[str] = mapped_column(String(64)) + target: Mapped[str | None] = mapped_column(String(128)) + detail: Mapped[str | None] = mapped_column(Text) + ip: Mapped[str | None] = mapped_column(String(64)) + trace_id: Mapped[str | None] = mapped_column(String(32)) + status: Mapped[str] = mapped_column(String(8), server_default="成功") + create_time: Mapped[datetime] = mapped_column(DateTime, server_default=func.now()) diff --git a/model/customer_relation.py b/model/customer_relation.py index 7ab50cd..24ddf2d 100644 --- a/model/customer_relation.py +++ b/model/customer_relation.py @@ -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)) - diff --git a/model/event_log.py b/model/event_log.py new file mode 100644 index 0000000..4a6c387 --- /dev/null +++ b/model/event_log.py @@ -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) diff --git a/model/portfolio_benchmark.py b/model/portfolio_benchmark.py new file mode 100644 index 0000000..ffd8973 --- /dev/null +++ b/model/portfolio_benchmark.py @@ -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() + ) diff --git a/model/sensitive_word.py b/model/sensitive_word.py new file mode 100644 index 0000000..08d4763 --- /dev/null +++ b/model/sensitive_word.py @@ -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() + ) diff --git a/model/sys_message.py b/model/sys_message.py new file mode 100644 index 0000000..b6a1c42 --- /dev/null +++ b/model/sys_message.py @@ -0,0 +1,23 @@ +"""sys_message 站内信/消息中心表 ORM(投顾触达客户通道,发送报告时写入)。""" +from __future__ import annotations + +from datetime import datetime + +from sqlalchemy import BigInteger, DateTime, Integer, String, func +from sqlalchemy.orm import Mapped, mapped_column + +from model.base import Base + + +class SysMessage(Base): + __tablename__ = "sys_message" + __table_args__ = {"comment": "站内信/消息中心表(投顾触达客户通道)"} + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True) + user_id: Mapped[int] = mapped_column(BigInteger) + msg_type: Mapped[str] = mapped_column(String(32)) + title: Mapped[str] = mapped_column(String(128)) + content: Mapped[str | None] = mapped_column(String(512)) + biz_id: Mapped[str | None] = mapped_column(String(64)) + is_read: Mapped[int] = mapped_column(Integer, server_default="0") + create_time: Mapped[datetime] = mapped_column(DateTime, server_default=func.now()) diff --git a/repositories/advisor_report.py b/repositories/advisor_report.py new file mode 100644 index 0000000..e896f7a --- /dev/null +++ b/repositories/advisor_report.py @@ -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()) diff --git a/repositories/advisor_todo.py b/repositories/advisor_todo.py new file mode 100644 index 0000000..e0b41c2 --- /dev/null +++ b/repositories/advisor_todo.py @@ -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 diff --git a/repositories/advisor_visit_record.py b/repositories/advisor_visit_record.py new file mode 100644 index 0000000..8293e61 --- /dev/null +++ b/repositories/advisor_visit_record.py @@ -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 diff --git a/repositories/audit_log.py b/repositories/audit_log.py new file mode 100644 index 0000000..f728274 --- /dev/null +++ b/repositories/audit_log.py @@ -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 diff --git a/repositories/customer_relation.py b/repositories/customer_relation.py index 17ef156..dad3410 100644 --- a/repositories/customer_relation.py +++ b/repositories/customer_relation.py @@ -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"] - diff --git a/repositories/event_log.py b/repositories/event_log.py new file mode 100644 index 0000000..6d16107 --- /dev/null +++ b/repositories/event_log.py @@ -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() diff --git a/repositories/portfolio_benchmark.py b/repositories/portfolio_benchmark.py new file mode 100644 index 0000000..2535759 --- /dev/null +++ b/repositories/portfolio_benchmark.py @@ -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() + ) diff --git a/repositories/sensitive_word.py b/repositories/sensitive_word.py new file mode 100644 index 0000000..5b07778 --- /dev/null +++ b/repositories/sensitive_word.py @@ -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() + ) diff --git a/repositories/sys_message.py b/repositories/sys_message.py new file mode 100644 index 0000000..65cd11b --- /dev/null +++ b/repositories/sys_message.py @@ -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()) diff --git a/schemas/advisor.py b/schemas/advisor.py new file mode 100644 index 0000000..6d9e810 --- /dev/null +++ b/schemas/advisor.py @@ -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)$") diff --git a/service/advisor/__init__.py b/service/advisor/__init__.py new file mode 100644 index 0000000..1f18fd5 --- /dev/null +++ b/service/advisor/__init__.py @@ -0,0 +1,6 @@ +"""投顾工作台业务层(B 端)。 + +模块职责:草稿代理与发送终审、客户 360、待办、回访、驾驶舱、审计台账、个人报表, +以及事件消费与定时调度。对投顾Agent 的调用统一走 agent_client;本地合规校验统一走 +send_review;敏感信息脱敏统一走 masking。 +""" diff --git a/service/advisor/agent_client.py b/service/advisor/agent_client.py new file mode 100644 index 0000000..2a91693 --- /dev/null +++ b/service/advisor/agent_client.py @@ -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 diff --git a/service/advisor/audit.py b/service/advisor/audit.py new file mode 100644 index 0000000..28c0c61 --- /dev/null +++ b/service/advisor/audit.py @@ -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]) diff --git a/service/advisor/audit_writer.py b/service/advisor/audit_writer.py new file mode 100644 index 0000000..473fef3 --- /dev/null +++ b/service/advisor/audit_writer.py @@ -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() diff --git a/service/advisor/customers.py b/service/advisor/customers.py new file mode 100644 index 0000000..bed7e9f --- /dev/null +++ b/service/advisor/customers.py @@ -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} diff --git a/service/advisor/dashboard.py b/service/advisor/dashboard.py new file mode 100644 index 0000000..ce69dbf --- /dev/null +++ b/service/advisor/dashboard.py @@ -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}, + } diff --git a/service/advisor/diagnosis.py b/service/advisor/diagnosis.py new file mode 100644 index 0000000..38e7e80 --- /dev/null +++ b/service/advisor/diagnosis.py @@ -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", + ) diff --git a/service/advisor/drafts.py b/service/advisor/drafts.py new file mode 100644 index 0000000..53e3fdb --- /dev/null +++ b/service/advisor/drafts.py @@ -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 diff --git a/service/advisor/event_consumer.py b/service/advisor/event_consumer.py new file mode 100644 index 0000000..25518e5 --- /dev/null +++ b/service/advisor/event_consumer.py @@ -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 diff --git a/service/advisor/masking.py b/service/advisor/masking.py new file mode 100644 index 0000000..8e5642c --- /dev/null +++ b/service/advisor/masking.py @@ -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 diff --git a/service/advisor/permissions.py b/service/advisor/permissions.py new file mode 100644 index 0000000..1a31c5f --- /dev/null +++ b/service/advisor/permissions.py @@ -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) diff --git a/service/advisor/report.py b/service/advisor/report.py new file mode 100644 index 0000000..d10e26e --- /dev/null +++ b/service/advisor/report.py @@ -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, + } diff --git a/service/advisor/scheduler.py b/service/advisor/scheduler.py new file mode 100644 index 0000000..576d87d --- /dev/null +++ b/service/advisor/scheduler.py @@ -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 + ) diff --git a/service/advisor/send_review.py b/service/advisor/send_review.py new file mode 100644 index 0000000..2f30210 --- /dev/null +++ b/service/advisor/send_review.py @@ -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("报告缺少免责声明,禁止发送") diff --git a/service/advisor/suitability.py b/service/advisor/suitability.py new file mode 100644 index 0000000..eaf24c8 --- /dev/null +++ b/service/advisor/suitability.py @@ -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": ""} diff --git a/service/advisor/todos.py b/service/advisor/todos.py new file mode 100644 index 0000000..9753b80 --- /dev/null +++ b/service/advisor/todos.py @@ -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} diff --git a/service/advisor/visits.py b/service/advisor/visits.py new file mode 100644 index 0000000..8c5f143 --- /dev/null +++ b/service/advisor/visits.py @@ -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], + } diff --git a/投顾工作台开发文档/common_const.md b/投顾工作台开发文档/common_const.md new file mode 100644 index 0000000..bdbd0f3 --- /dev/null +++ b/投顾工作台开发文档/common_const.md @@ -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. 业务规则(如适当性匹配逻辑)写在业务代码中,不在此文档写业务逻辑,本文件只存放**常量字符串、枚举、模板文本**。 \ No newline at end of file