From a7f9e182a4d6cb009164f777f4c5b4ab86b3a74c Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?=E6=AC=A7=E9=98=B3=E6=B4=8B?= <2443479321@qq.com>
Date: Mon, 14 Sep 2026 00:19:31 +0800
Subject: [PATCH] =?UTF-8?q?feat:=E7=94=A8=E6=88=B7=E7=94=BB=E5=83=8F?=
=?UTF-8?q?=E5=8A=9F=E8=83=BD?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
api/routers/profile.py | 42 +++
nl2sql/semantic_catalog.json | 10 +-
nl2sql/semantics.py | 6 +-
rag/intent.py | 10 +-
service/advisor/suitability.py | 6 +-
service/memory/composed_profile.py | 436 +++++++++++++++++++++++++++++
service/memory/context_builder.py | 2 +
service/memory/facade.py | 39 +++
service/memory/holdings.py | 112 ++++++++
service/memory/schemas.py | 1 +
service/purchase.py | 10 +-
service/questionnaire.py | 16 +-
service/redeem.py | 3 +
service/risk/engine.py | 7 +-
sql/risk_level_r1r5_migration.sql | 38 +++
sql/schema.sql | 268 +++++++++++++++++-
开发计划.md | 2 +-
17 files changed, 971 insertions(+), 37 deletions(-)
create mode 100644 api/routers/profile.py
create mode 100644 service/memory/composed_profile.py
create mode 100644 service/memory/holdings.py
create mode 100644 sql/risk_level_r1r5_migration.sql
diff --git a/api/routers/profile.py b/api/routers/profile.py
new file mode 100644
index 0000000..5d5e3a0
--- /dev/null
+++ b/api/routers/profile.py
@@ -0,0 +1,42 @@
+"""最终用户画像路由:仅员工/管理员可查询三源整合画像。
+
+权限约定:
+- 客户(CUSTOMER):**无权访问本接口**,整合画像属于内部经营数据;
+- 员工/管理员(EMPLOYEE/ADMIN):必须显式指定 customer_id,可查任意客户。
+
+整合逻辑在 ComposedProfileService(LLM 整合 + Redis 缓存 + 降级),
+本路由只做鉴权与编排;画像仅供参考,不用于适当性校验。
+"""
+from fastapi import APIRouter, Depends, Query
+from sqlalchemy.ext.asyncio import AsyncSession
+
+from api.deps import get_current_user
+from config.deps import get_db
+from model.sys_user import SysUser
+from service.memory.composed_profile import ComposedProfileService
+from utils.exceptions import ForbiddenError
+from utils.response import success
+
+router = APIRouter()
+
+
+async def require_profile_viewer(user: SysUser = Depends(get_current_user)) -> SysUser:
+ """仅员工或管理员可查询最终用户画像,客户一律拒绝。"""
+ if user.user_type == "ADMIN":
+ return user
+ if user.user_type != "EMPLOYEE":
+ raise ForbiddenError("仅员工可查询用户画像")
+ return user
+
+
+@router.get("/profile/composed", summary="查询最终用户画像(画像+记忆+持仓三源整合,仅员工)")
+async def get_composed_profile(
+ customer_id: int = Query(..., description="目标客户 ID"),
+ user: SysUser = Depends(require_profile_viewer),
+ db: AsyncSession = Depends(get_db),
+):
+ """返回整合后的最终用户画像 JSON 及降级 warnings。"""
+ service = ComposedProfileService()
+ # 完整响应体先查 Redis(20 分钟过期),重复请求直接命中缓存
+ data, _warnings = await service.compose_response(db, customer_id=customer_id)
+ return success(data)
diff --git a/nl2sql/semantic_catalog.json b/nl2sql/semantic_catalog.json
index b5cec37..9895b80 100644
--- a/nl2sql/semantic_catalog.json
+++ b/nl2sql/semantic_catalog.json
@@ -6,22 +6,24 @@
"enabled": true,
"aliases": ["基金产品", "产品", "基金", "基金名称", "产品名称"],
"tables": ["fin_product"],
- "fields": ["product_code", "product_name", "product_type", "risk_level", "status"]
+ "fields": ["product_code", "product_name", "product_type", "risk_level", "status"],
+ "value_hint": "product_type 实际枚举值为:货币型/债券型/混合型/股票型/指数型/QDII(注意是'XX型',不是 Schema 注释里的'股票基金');status 实际枚举值为:在售/暂停/到期/清盘,查在售产品用 status = '在售'"
},
{
"term": "股票型基金",
"enabled": true,
- "aliases": ["股票型", "股票基金", "权益类"],
+ "aliases": ["股票型", "股票基金", "股票型基金", "权益类"],
"tables": ["fin_product"],
"fields": ["product_type"],
- "value_hint": "product_type 通常使用业务库中的股票型标准值,不能自行创造枚举值"
+ "value_hint": "fin_product.product_type 实际枚举值为:货币型/债券型/混合型/股票型/指数型/QDII;股票型基金必须用 product_type = '股票型',不要使用 Schema 注释里的'股票基金'"
},
{
"term": "客户",
"enabled": true,
"aliases": ["客户", "投资人", "持有人"],
"tables": ["fin_customer_profile", "fin_holdings", "fin_risk_assessment"],
- "fields": ["customer_id", "risk_level", "customer_level"]
+ "fields": ["customer_id", "risk_level", "customer_level"],
+ "value_hint": "fin_customer_profile.risk_level / fin_risk_assessment.risk_level 实际枚举值为 R1-R5(R1 最保守,R5 最激进,与产品风险等级同一口径);查'保守型/稳健型客户'对应 R1/R2,'激进型客户'对应 R5"
},
{
"term": "持仓",
diff --git a/nl2sql/semantics.py b/nl2sql/semantics.py
index f4b0cd6..9ca7193 100644
--- a/nl2sql/semantics.py
+++ b/nl2sql/semantics.py
@@ -107,10 +107,12 @@ def resolve_semantics(question: str, *, catalog: dict[str, Any] | None = None) -
{
"term": item["term"],
"field": item["fields"][0],
- "hint": item.get("metric_hint", item.get("value_hint", "")),
+ "hint": item.get("metric_hint") or item.get("value_hint", ""),
}
+ # metric_hint(指标口径)与 value_hint(字段取值口径)都要进入
+ # SQL 生成上下文;只过滤 metric_hint 会让枚举值提示永远丢失。
for item in matched
- if "metric_hint" in item
+ if "metric_hint" in item or "value_hint" in item
]
return {
"terms": [item["term"] for item in matched],
diff --git a/rag/intent.py b/rag/intent.py
index 788077b..99b75c0 100644
--- a/rag/intent.py
+++ b/rag/intent.py
@@ -50,12 +50,14 @@ INTENT_SYSTEM_PROMPT = (
"- guide_purchase: 用户询问如何购买基金、开户、注册等引导类问题\n"
"- want_advisor: 用户希望获得个性化基金推荐或投资顾问服务\n"
"- knowledge_qa: 用户询问基金相关的知识性问题,如净值、费率、风险、申赎规则等\n"
- "- company_info: 用户询问华夏科技公司本身的信息,如公司全称、成立时间、牌照、总部地址、"
- "客服电话、服务时间、官网、投诉渠道等\n"
+ "- company_info: 用户询问华夏科技公司本身的静态档案信息,如公司全称、成立时间、牌照、总部地址、"
+ "客服电话、服务时间、官网、投诉渠道等;"
+ "注意:公司的产品清单、产品数量、产品数据不属于公司信息,应归为 nl2sql_request\n"
"- nl2sql_request: 用户要求查询具体数据或账户信息,"
"如“我的持仓有哪些”“我买了多少XX基金”“我的交易记录”"
- "“最近一周XX基金的净值数据”“XX基金最新的规模/费率数据”等"
- "要求数据本身而非知识解释的问题\n"
+ "“最近一周XX基金的净值数据”“XX基金最新的规模/费率数据”"
+ "“你们公司有哪些股票型基金产品”“公司在售的债券基金有哪些”等"
+ "询问公司产品清单/产品数据或个人账户数据、要求数据本身而非知识解释的问题\n"
"- chitchat: 普通寒暄,如问候、致谢、告别、询问你是谁/你能做什么、在吗等一两句话的闲聊\n"
"- off_topic: 用户要求你实质性地处理与金融、基金、公司业务无关的事情,"
"如写代码、讲笑话、写作文、问天气、聊政治、情感咨询、做数学题等\n"
diff --git a/service/advisor/suitability.py b/service/advisor/suitability.py
index eaf24c8..8275e4b 100644
--- a/service/advisor/suitability.py
+++ b/service/advisor/suitability.py
@@ -7,13 +7,11 @@ Agent 复用同一实现、禁止各自复制。在 Agent 侧共享包发布前
"""
from __future__ import annotations
-# 风险等级 → 序号(C_n / R_n 的 n)。兼容两套口径:标准 C1-C5/R1-R5 与历史中文等级,
-# 避免问卷侧尚未完成 C1-C5 归一化时发送终审误判。
-# 假设:中文等级与 C 级一一对应(保守=C1/稳健=C2/平衡=C3/进取=C4/激进=C5)。
+# 风险等级 → 序号(C_n / R_n 的 n)。客户风险能力口径 C1-C5,客户画像与产品口径
+# 统一为 R1-R5(历史中文等级已在数据迁移中转换,故不再收录中文键)。
_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,
}
diff --git a/service/memory/composed_profile.py b/service/memory/composed_profile.py
new file mode 100644
index 0000000..d3263fa
--- /dev/null
+++ b/service/memory/composed_profile.py
@@ -0,0 +1,436 @@
+"""最终用户画像整合服务:fin_customer_profile + memory_unit + 持仓 三源加权。
+
+- 权重语义是 LLM 整合时的注意力权重(配置于 sys_config:profile.weight.*),
+ 不做算术加权——三源是异构数据。
+- 只做派生视图:结果缓存在 Redis(成功 30 分钟 / 降级 5 分钟),
+ 不落库、不回写 fin_customer_profile,适当性校验仍以问卷风评原值为准。
+- LLM 失败时降级为规则拼接摘要,沿项目 warnings 机制上报,绝不阻塞调用方。
+"""
+from __future__ import annotations
+
+import json
+import logging
+import re
+from datetime import datetime
+from typing import Any
+
+from sqlalchemy.ext.asyncio import AsyncSession
+
+from repositories.sys_config import SysConfigRepo
+from tool.confidence_rank import FinalConfidenceRankTool
+
+from .holdings import CustomerHoldingsMemory
+from .long_term import LongTermMemoryService
+from .profile import CustomerProfileMemory
+
+logger = logging.getLogger(__name__)
+
+
+class ComposedProfileService:
+ """拉取三路数据,调用统一 LLM 整合为最终用户画像。"""
+
+ CACHE_TTL = 20 * 60
+ DEGRADED_CACHE_TTL = 20 * 60
+ RESPONSE_CACHE_TTL = 20 * 60
+ RESPONSE_CACHE_PREFIX = "composed_profile_response"
+ DEFAULT_WEIGHTS = {"profile": 0.5, "memory": 0.2, "holdings": 0.3}
+ WEIGHT_KEYS = {
+ "profile": "profile.weight.profile",
+ "memory": "profile.weight.memory",
+ "holdings": "profile.weight.holdings",
+ }
+ ENABLED_KEY = "profile.composed.enabled"
+ MEMORY_TOP_K_KEY = "profile.composed.memory_top_k"
+ DEFAULT_MEMORY_TOP_K = 10
+
+ SYSTEM_PROMPT = """
+你是用户画像整合器。根据三路客户数据生成最终用户画像 JSON,只输出 JSON,不要 Markdown:
+{{"summary": "一段话画像结论", "risk_tendency": "风险偏好结论",
+ "preference_tags": ["标签"], "behavior_traits": ["行为特征"], "advice_focus": "服务建议"}}
+权重指引:结构化画像信息可信度最高(权重 {w_profile});持仓反映真实风险行为(权重 {w_holdings});
+对话记忆是口头表达补充(权重 {w_memory})。三源冲突时按权重取舍,且必须在 summary 中显式说明冲突。
+status 为 candidate 的记忆只能作参考,不能作为确定性结论。
+只能基于输入事实归纳,禁止编造;数据缺失的字段留空数组或空字符串,不要推测。
+禁止在结果中出现内部字段名、ID、SQL 或证据引用。
+""".strip()
+
+ def __init__(
+ self,
+ *,
+ profile: CustomerProfileMemory | None = None,
+ long_term: LongTermMemoryService | None = None,
+ holdings: CustomerHoldingsMemory | None = None,
+ redis=None,
+ llm_client=None,
+ rank_tool: FinalConfidenceRankTool | None = None,
+ config_repo_cls=SysConfigRepo,
+ ):
+ from config.database.redis import client as redis_client
+ from tool.llm import llm as default_llm
+
+ self.profile = profile or CustomerProfileMemory()
+ self.long_term = long_term or LongTermMemoryService()
+ self.holdings = holdings or CustomerHoldingsMemory()
+ self.redis = redis or redis_client()
+ self.llm_client = llm_client or default_llm
+ self.rank_tool = rank_tool or FinalConfidenceRankTool()
+ self.config_repo_cls = config_repo_cls
+
+ @staticmethod
+ def cache_key(customer_id: int) -> str:
+ """生成画像负载缓存 Key。"""
+ return f"composed_profile:{customer_id}"
+
+ @staticmethod
+ def response_cache_key(customer_id: int) -> str:
+ """生成接口完整响应缓存 Key。"""
+ return f"{ComposedProfileService.RESPONSE_CACHE_PREFIX}:{customer_id}"
+
+ async def compose(
+ self,
+ db: AsyncSession,
+ *,
+ customer_id: int,
+ query: str | None = None,
+ profile: dict[str, Any] | None = None,
+ memories: list[dict[str, Any]] | None = None,
+ holdings_summary: dict[str, Any] | None = None,
+ ) -> tuple[dict[str, Any] | None, list[str]]:
+ """整合最终画像;可传入 facade 已召回的数据避免重复取数。"""
+ warnings: list[str] = []
+ enabled = await self._enabled(db, warnings)
+ if not enabled:
+ return None, warnings
+
+ cached = await self._cache_get(customer_id)
+ if cached is not None:
+ return cached, warnings
+
+ weights, weight_warnings = await self._load_weights(db)
+ warnings.extend(weight_warnings)
+ top_k = await self._memory_top_k(db, warnings)
+
+ if profile is None:
+ profile, profile_warnings = await self.profile.get(db, customer_id)
+ warnings.extend(profile_warnings)
+ if memories is None:
+ memories = await self._recall_memories(db, customer_id, query, top_k, warnings)
+ if holdings_summary is None:
+ holdings_summary, holdings_warnings = await self.holdings.summary(db, customer_id)
+ warnings.extend(holdings_warnings)
+
+ slim_memories = self._slim_memories(memories, top_k)
+ if profile is None and not slim_memories and holdings_summary is None:
+ return None, warnings
+
+ try:
+ result = await self._integrate(weights, profile, slim_memories, holdings_summary)
+ cache_ttl = self.CACHE_TTL
+ except Exception as exc: # noqa: BLE001 降级不阻塞
+ warnings.append(f"composed_profile_llm_failed:{type(exc).__name__}")
+ result = self._degraded(profile, slim_memories, holdings_summary)
+ cache_ttl = self.DEGRADED_CACHE_TTL
+
+ payload = {
+ "profile": result,
+ "weights": weights,
+ "generated_at": datetime.now().isoformat(timespec="seconds"),
+ }
+ await self._cache_set(customer_id, payload, cache_ttl)
+ return payload, warnings
+
+ async def invalidate(self, customer_id: int) -> list[str]:
+ """删除画像负载与接口响应两级缓存,供画像/持仓变更后调用。"""
+ warnings: list[str] = []
+ for key in (self.cache_key(customer_id), self.response_cache_key(customer_id)):
+ try:
+ await self.redis.delete(key)
+ except Exception as exc: # noqa: BLE001
+ warnings.append(f"composed_profile_cache_invalidate_failed:{type(exc).__name__}")
+ return warnings
+
+ async def compose_response(
+ self, db: AsyncSession, *, customer_id: int
+ ) -> tuple[dict[str, Any] | None, list[str]]:
+ """接口层入口:完整响应体先查 Redis(20 分钟过期),未命中才整合。
+
+ 返回结构即接口 data:
+ {"customer_id": ..., "profile": <画像负载或 None>, "warnings": [...]}
+ profile 为 None(功能关闭或无任何输入)时不缓存,便于开关立即生效。
+ """
+ key = self.response_cache_key(customer_id)
+ try:
+ raw = await self.redis.get(key)
+ if raw:
+ return json.loads(raw), []
+ except Exception: # noqa: BLE001 缓存读取失败视为未命中
+ pass
+
+ payload, warnings = await self.compose(db, customer_id=customer_id)
+ data = {
+ "customer_id": customer_id,
+ "profile": payload,
+ "warnings": warnings,
+ }
+ if payload is not None:
+ try:
+ await self.redis.set(
+ key,
+ json.dumps(data, ensure_ascii=False, default=str),
+ ex=self.RESPONSE_CACHE_TTL,
+ )
+ except Exception: # noqa: BLE001 缓存写入失败不影响返回
+ pass
+ return data, warnings
+
+ # ---- 配置 -----------------------------------------------------------
+
+ async def _enabled(self, db: AsyncSession, warnings: list[str]) -> bool:
+ """读取功能开关;配置读取失败时默认开启,保证功能可用。"""
+ try:
+ value = await self.config_repo_cls(db).get_value(self.ENABLED_KEY, "true")
+ except Exception as exc: # noqa: BLE001
+ warnings.append(f"composed_config_failed:{type(exc).__name__}")
+ return True
+ return str(value).strip().lower() not in {"false", "0", "off"}
+
+ async def _load_weights(self, db: AsyncSession) -> tuple[dict[str, float], list[str]]:
+ """读取三源权重;缺失用默认值,非法或权重和不为 1 时整体回退默认。"""
+ warnings: list[str] = []
+ try:
+ repo = self.config_repo_cls(db)
+ raw = {
+ name: await repo.get_value(key, str(default))
+ for name, key, default in (
+ (name, key, self.DEFAULT_WEIGHTS[name])
+ for name, key in self.WEIGHT_KEYS.items()
+ )
+ }
+ except Exception as exc: # noqa: BLE001
+ return dict(self.DEFAULT_WEIGHTS), [f"composed_config_failed:{type(exc).__name__}"]
+
+ weights: dict[str, float] = {}
+ valid = True
+ for name, value in raw.items():
+ try:
+ parsed = float(value)
+ except (TypeError, ValueError):
+ valid = False
+ break
+ if not 0.0 <= parsed <= 1.0:
+ valid = False
+ break
+ weights[name] = parsed
+ if not valid or abs(sum(weights.values()) - 1.0) > 0.01:
+ warnings.append("composed_weight_invalid")
+ return dict(self.DEFAULT_WEIGHTS), warnings
+ return weights, warnings
+
+ async def _memory_top_k(self, db: AsyncSession, warnings: list[str]) -> int:
+ """读取进入整合的记忆条数上限,非法回退默认。"""
+ try:
+ raw = await self.config_repo_cls(db).get_value(
+ self.MEMORY_TOP_K_KEY, str(self.DEFAULT_MEMORY_TOP_K)
+ )
+ top_k = int(str(raw))
+ except Exception: # noqa: BLE001
+ return self.DEFAULT_MEMORY_TOP_K
+ if top_k <= 0:
+ warnings.append("composed_top_k_invalid")
+ return self.DEFAULT_MEMORY_TOP_K
+ return top_k
+
+ # ---- 数据与整合 ------------------------------------------------------
+
+ async def _recall_memories(
+ self,
+ db: AsyncSession,
+ customer_id: int,
+ query: str | None,
+ top_k: int,
+ warnings: list[str],
+ ) -> list[dict[str, Any]]:
+ """召回长期记忆并按综合置信分重排;失败降级为空列表。"""
+ try:
+ dtos, memory_warnings = await self.long_term.recall(
+ db, customer_id, limit=max(top_k, 10), query=query
+ )
+ warnings.extend(memory_warnings)
+ ranked = self.rank_tool.rank(
+ [dto.model_dump(mode="json") for dto in dtos], top_k=top_k
+ )
+ return ranked
+ except Exception as exc: # noqa: BLE001
+ warnings.append(f"composed_memory_recall_failed:{type(exc).__name__}")
+ return []
+
+ @staticmethod
+ def _slim_memories(
+ memories: list[dict[str, Any]] | None, top_k: int
+ ) -> list[dict[str, Any]]:
+ """裁剪记忆字段:只保留整合所需的最小集合,不带证据引用。"""
+ slim = []
+ for item in (memories or [])[:top_k]:
+ slim.append(
+ {
+ "tag": item.get("tag"),
+ "memory_type": item.get("memory_type"),
+ "content": item.get("content"),
+ "status": item.get("status"),
+ "confidence": item.get("confidence"),
+ }
+ )
+ return slim
+
+ async def _integrate(
+ self,
+ weights: dict[str, float],
+ profile: dict[str, Any] | None,
+ memories: list[dict[str, Any]],
+ holdings: dict[str, Any] | None,
+ ) -> dict[str, Any]:
+ """调用统一 LLM 整合三源数据,解析失败时重试一次。"""
+ prompt = [
+ {
+ "role": "system",
+ "content": self.SYSTEM_PROMPT.format(
+ w_profile=weights["profile"],
+ w_memory=weights["memory"],
+ w_holdings=weights["holdings"],
+ ),
+ },
+ {
+ "role": "user",
+ "content": json.dumps(
+ {
+ "weights": weights,
+ "profile": profile,
+ "memories": memories,
+ "holdings": holdings,
+ },
+ ensure_ascii=False,
+ default=str,
+ ),
+ },
+ ]
+ last_error: Exception | None = None
+ for _ in range(2):
+ try:
+ # max_tokens 传 None:沿用全局 LLM_MAX_TOKENS。思考型模型(如
+ # qwen3)会先输出推理段,限太小会截断 JSON 导致解析失败。
+ text = await self.llm_client.chat(
+ prompt, temperature=0, max_tokens=None
+ )
+ result = self._parse_json(text)
+ except Exception as exc: # noqa: BLE001 网络或解析失败都重试一次
+ last_error = exc
+ logger.warning(
+ "composed profile integrate attempt failed: %s: %s; raw=%r",
+ type(exc).__name__,
+ exc,
+ (text if "text" in locals() else "")[:300],
+ )
+ continue
+ if isinstance(result.get("summary"), str) and result["summary"].strip():
+ return result
+ last_error = ValueError("整合结果缺少 summary 字段")
+ raise last_error or ValueError("LLM 整合结果非法")
+
+ _THINK_PATTERN = re.compile(r".*?", re.S | re.I)
+
+ @classmethod
+ def _parse_json(cls, text: str) -> dict[str, Any]:
+ """解析 LLM JSON 输出:剥离思考标签和围栏,容忍正文夹杂说明文字。"""
+ payload = text.strip()
+ payload = cls._THINK_PATTERN.sub("", payload).strip()
+ fenced = re.search(r"```(?:json)?\s*(.*?)\s*```", payload, re.S | re.I)
+ if fenced:
+ payload = fenced.group(1).strip()
+ try:
+ data = json.loads(payload)
+ except json.JSONDecodeError:
+ # 模型可能在 JSON 前后加了说明文字:提取首个 { 到最后一个 } 的片段重试
+ start = payload.find("{")
+ end = payload.rfind("}")
+ if start < 0 or end <= start:
+ raise
+ data = json.loads(payload[start : end + 1])
+ if not isinstance(data, dict):
+ raise ValueError("整合结果必须是 JSON 对象")
+ return data
+
+ @staticmethod
+ def _degraded(
+ profile: dict[str, Any] | None,
+ memories: list[dict[str, Any]],
+ holdings: dict[str, Any] | None,
+ ) -> dict[str, Any]:
+ """LLM 失败时的规则拼接摘要,保证接口始终有可用输出。"""
+ parts: list[str] = []
+ if profile:
+ parts.append(
+ f"画像:风险等级 {profile.get('risk_level') or '未知'}"
+ f"(评分 {profile.get('risk_score') or '未知'})"
+ )
+ if memories:
+ confirmed = [m["content"] for m in memories if m.get("status") == "confirmed"]
+ chosen = confirmed or [m["content"] for m in memories]
+ parts.append("记忆:" + ";".join(chosen[:3]))
+ if holdings:
+ parts.append(
+ f"持仓:{holdings.get('holding_count', 0)} 只在持"
+ f",总市值 {holdings.get('total_market_value', 0)}"
+ )
+ return {
+ "summary": "。".join(parts) if parts else "暂无可用客户信息",
+ "risk_tendency": (profile or {}).get("risk_level") or "",
+ "preference_tags": [m.get("tag") for m in memories if m.get("tag")][:5],
+ "behavior_traits": [],
+ "advice_focus": "",
+ "degraded": True,
+ "sources": {
+ "profile": profile is not None,
+ "memory": bool(memories),
+ "holdings": holdings is not None,
+ },
+ }
+
+ # ---- 缓存 -------------------------------------------------------------
+
+ async def _cache_get(self, customer_id: int) -> dict[str, Any] | None:
+ """读取缓存;读取或解析失败视为未命中。"""
+ try:
+ raw = await self.redis.get(self.cache_key(customer_id))
+ if raw:
+ return json.loads(raw)
+ except Exception: # noqa: BLE001
+ return None
+ return None
+
+ async def _cache_set(
+ self, customer_id: int, payload: dict[str, Any], ttl: int
+ ) -> None:
+ """写入缓存;失败静默(缓存只是加速手段)。"""
+ try:
+ await self.redis.set(
+ self.cache_key(customer_id),
+ json.dumps(payload, ensure_ascii=False, default=str),
+ ex=ttl,
+ )
+ except Exception: # noqa: BLE001
+ return
+
+
+async def invalidate_composed_profile(customer_id: int) -> None:
+ """轻量失效入口:画像回写、申购/赎回落账后调用,静默失败不阻塞业务。"""
+ from config.database.redis import client as redis_client
+
+ try:
+ redis = redis_client()
+ await redis.delete(ComposedProfileService.cache_key(customer_id))
+ await redis.delete(ComposedProfileService.response_cache_key(customer_id))
+ except Exception: # noqa: BLE001 缓存失效失败只影响新鲜度,不阻塞业务事务
+ return
+
+
+__all__ = ["ComposedProfileService", "invalidate_composed_profile"]
diff --git a/service/memory/context_builder.py b/service/memory/context_builder.py
index 6a009e3..0da4052 100644
--- a/service/memory/context_builder.py
+++ b/service/memory/context_builder.py
@@ -15,6 +15,7 @@ def build_customer_memory_context(
long_term_memories: list[MemoryUnitDTO],
customer_relations: list[dict] | None = None,
customer_products: list[dict] | None = None,
+ composed_profile: dict | None = None,
warnings: list[str] | None = None,
) -> CustomerMemoryContext:
"""构造稳定的客服记忆上下文结构。"""
@@ -23,6 +24,7 @@ def build_customer_memory_context(
session_id=session_id,
short_term_messages=short_term_messages,
customer_profile=customer_profile,
+ composed_profile=composed_profile,
work_orders=work_orders,
long_term_memories=long_term_memories,
customer_relations=customer_relations or [],
diff --git a/service/memory/facade.py b/service/memory/facade.py
index 98f2112..88a5dc9 100644
--- a/service/memory/facade.py
+++ b/service/memory/facade.py
@@ -9,9 +9,11 @@ from config.database.mysql import get_session_factory
from tool.confidence_rank import FinalConfidenceRankTool
from .archive import ConversationArchiver
+from .composed_profile import ComposedProfileService
from .customer_relation import CustomerRelationMemory
from .customer_product import CustomerProductMemory
from .context_builder import build_customer_memory_context
+from .holdings import CustomerHoldingsMemory
from .long_term import LongTermMemoryService
from .interest_topic import InterestTopicTracker
from .profile import CustomerProfileMemory
@@ -32,6 +34,8 @@ class MemoryService:
work_orders: WorkOrderMemory | None = None,
relations: CustomerRelationMemory | None = None,
products: CustomerProductMemory | None = None,
+ holdings: CustomerHoldingsMemory | None = None,
+ composed: ComposedProfileService | None = None,
long_term: LongTermMemoryService | None = None,
archiver: ConversationArchiver | None = None,
rank_tool: FinalConfidenceRankTool | None = None,
@@ -43,6 +47,13 @@ class MemoryService:
self.work_orders = work_orders or WorkOrderMemory()
self.relations = relations or CustomerRelationMemory()
self.products = products or CustomerProductMemory()
+ self.holdings = holdings or CustomerHoldingsMemory()
+ self.composed = composed or ComposedProfileService(
+ profile=self.profile,
+ long_term=self.long_term,
+ holdings=self.holdings,
+ redis=self.short_term.redis,
+ )
self.long_term = long_term or LongTermMemoryService()
self.archiver = archiver or ConversationArchiver(short_term=self.short_term)
self.rank_tool = rank_tool or FinalConfidenceRankTool()
@@ -109,11 +120,35 @@ class MemoryService:
top_k=limit,
)
ranked_memories = [MemoryUnitDTO.model_validate(item) for item in ranked]
+
+ holdings_summary = None
+ try:
+ holdings_summary, holdings_warnings = await self.holdings.summary(
+ db, customer_id
+ )
+ warnings.extend(holdings_warnings)
+ except Exception as exc:
+ warnings.append(f"holdings_recall_failed:{type(exc).__name__}")
+ composed_payload: dict[str, Any] | None = None
+ try:
+ composed_payload, composed_warnings = await self.composed.compose(
+ db,
+ customer_id=customer_id,
+ query=query,
+ profile=profile,
+ memories=ranked,
+ holdings_summary=holdings_summary,
+ )
+ warnings.extend(composed_warnings)
+ except Exception as exc:
+ warnings.append(f"composed_profile_failed:{type(exc).__name__}")
+
return build_customer_memory_context(
customer_id=customer_id,
session_id=session_id,
short_term_messages=short_term_messages,
customer_profile=profile,
+ composed_profile=composed_payload,
work_orders=work_orders,
customer_relations=customer_relations,
customer_products=customer_products,
@@ -155,6 +190,10 @@ class MemoryService:
"""记录一次重复兴趣主题信号,达到阈值后允许写入长期记忆。"""
return await self.interest_tracker.record(customer_id, tag)
+ async def invalidate_composed(self, customer_id: int) -> list[str]:
+ """失效最终画像缓存,供画像回写/持仓变更后调用。"""
+ return await self.composed.invalidate(customer_id)
+
async def close_session(
self,
*,
diff --git a/service/memory/holdings.py b/service/memory/holdings.py
new file mode 100644
index 0000000..948fc3b
--- /dev/null
+++ b/service/memory/holdings.py
@@ -0,0 +1,112 @@
+"""客户持仓记忆:fin_holdings 聚合摘要,供最终画像整合使用。
+
+只输出聚合结果(总市值、产品类型分布、盈亏概览、前三大持仓),
+不向 LLM 或接口消费方暴露原始持仓行,控制 token 并避免明细泄露。
+"""
+from __future__ import annotations
+
+from decimal import Decimal
+from typing import Any
+
+from sqlalchemy.ext.asyncio import AsyncSession
+
+from repositories.fin_holdings import FinHoldingsRepo
+
+
+class CustomerHoldingsMemory:
+ """读取并聚合客户当前持仓,输出画像整合所需的持仓摘要。"""
+
+ TOP_HOLDINGS_LIMIT = 3
+
+ def __init__(self, *, repository_factory=FinHoldingsRepo):
+ self.repository_factory = repository_factory
+
+ async def summary(
+ self, db: AsyncSession, customer_id: int
+ ) -> tuple[dict[str, Any] | None, list[str]]:
+ """聚合持有中持仓;无持仓返回 None,异常返回降级 warnings。"""
+ try:
+ rows = await self.repository_factory(db).list_with_products(
+ customer_id, include_closed=False
+ )
+ except Exception as exc: # noqa: BLE001 沿用记忆模块降级约定
+ return None, [f"holdings_recall_failed:{type(exc).__name__}"]
+ if not rows:
+ return None, []
+
+ total_value = Decimal("0")
+ total_cost = Decimal("0")
+ total_pnl = Decimal("0")
+ win_count = 0
+ lose_count = 0
+ mix: dict[str, Decimal] = {}
+ ranked: list[tuple[Decimal, str]] = []
+ for holding, product in rows:
+ value = self._to_decimal(holding.current_value)
+ total_value += value
+ total_cost += self._to_decimal(holding.cost_amount)
+ total_pnl += self._to_decimal(holding.profit_loss)
+ ratio = holding.profit_ratio
+ if ratio is not None:
+ if ratio > 0:
+ win_count += 1
+ elif ratio < 0:
+ lose_count += 1
+ product_type = (product.product_type if product else None) or "未知"
+ mix[product_type] = mix.get(product_type, Decimal("0")) + value
+ name = product.product_name if product else f"产品{holding.product_id}"
+ ranked.append((value, name))
+
+ payload: dict[str, Any] = {
+ "holding_count": len(rows),
+ "total_market_value": float(total_value),
+ "total_cost": float(total_cost),
+ "product_type_mix": self._ratio_map(mix, total_value),
+ "profit_summary": {
+ "total_pnl": float(total_pnl),
+ "total_pnl_ratio": (
+ round(float(total_pnl / total_cost), 4)
+ if total_cost > 0
+ else None
+ ),
+ "win_count": win_count,
+ "lose_count": lose_count,
+ },
+ "top_holdings": self._top_holdings(ranked, total_value),
+ }
+ return payload, []
+
+ @staticmethod
+ def _to_decimal(value: Any) -> Decimal:
+ """把 ORM Decimal / float / None 统一转换为 Decimal,防御脏数据类型。"""
+ if value is None:
+ return Decimal("0")
+ try:
+ return value if isinstance(value, Decimal) else Decimal(str(value))
+ except Exception: # noqa: BLE001
+ return Decimal("0")
+
+ @staticmethod
+ def _ratio_map(values: dict[str, Decimal], total: Decimal) -> dict[str, float]:
+ """把各维度市值转换为相对总市值的占比,总量为 0 时返回空表。"""
+ if total <= 0:
+ return {}
+ return {
+ key: round(float(value / total), 4)
+ for key, value in sorted(values.items(), key=lambda item: -item[1])
+ }
+
+ def _top_holdings(
+ self, ranked: list[tuple[Decimal, str]], total: Decimal
+ ) -> list[dict[str, Any]]:
+ """按市值取前三大持仓,返回名称与占比。"""
+ if total <= 0:
+ return []
+ ranked.sort(key=lambda item: -item[0])
+ return [
+ {"name": name, "weight": round(float(value / total), 4)}
+ for value, name in ranked[: self.TOP_HOLDINGS_LIMIT]
+ ]
+
+
+__all__ = ["CustomerHoldingsMemory"]
diff --git a/service/memory/schemas.py b/service/memory/schemas.py
index d3f3be2..e8072b1 100644
--- a/service/memory/schemas.py
+++ b/service/memory/schemas.py
@@ -99,6 +99,7 @@ class CustomerMemoryContext(BaseModel):
session_id: str
short_term_messages: list[ShortTermMessage] = Field(default_factory=list)
customer_profile: dict[str, Any] | None = None
+ composed_profile: dict[str, Any] | None = None
work_orders: list[dict[str, Any]] = Field(default_factory=list)
long_term_memories: list[MemoryUnitDTO] = Field(default_factory=list)
customer_relations: list[dict[str, Any]] = Field(default_factory=list)
diff --git a/service/purchase.py b/service/purchase.py
index 0933779..b2919f7 100644
--- a/service/purchase.py
+++ b/service/purchase.py
@@ -17,6 +17,7 @@ from repositories.fin_product import FinProductRepo
from repositories.trade_order import TradeOrderRepo
from schemas.holdings import HoldingResp
from schemas.purchase import PurchaseResp
+from service.memory.composed_profile import invalidate_composed_profile
from service.risk.engine import RiskEngine, summarize
from service.risk.settle import settle
from utils.exceptions import (
@@ -30,11 +31,8 @@ from utils.order_no import gen_order_no
_MONEY = Decimal("0.01")
_SHARES = Decimal("0.0001")
-# 风险等级 → 序号。兼容两套口径:R1~R5 与 保守~激进(同一映射)。
-_RISK_RANK = {
- "R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5,
- "保守": 1, "稳健": 2, "平衡": 3, "进取": 4, "激进": 5,
-}
+# 风险等级 → 序号。客户画像与产品统一 R1~R5 口径(历史中文等级已在数据迁移中转换)。
+_RISK_RANK = {"R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5}
def _risk_rank(level: str | None) -> int | None:
@@ -137,6 +135,8 @@ async def purchase(
order.id, status="已确认", confirm_time=now
)
await db.commit()
+ # 持仓变更后失效最终画像缓存(失败静默,不影响交易结果)
+ await invalidate_composed_profile(user.id)
account = await account_repo.get_by_customer_id(user.id)
holding = await holdings_repo.get_by_customer_product(user.id, product_id)
return PurchaseResp(
diff --git a/service/questionnaire.py b/service/questionnaire.py
index e1ec9e5..cb81bd6 100644
--- a/service/questionnaire.py
+++ b/service/questionnaire.py
@@ -8,15 +8,17 @@ from sqlalchemy.ext.asyncio import AsyncSession
from model.fin_risk_assessment import FinRiskAssessment
from repositories.questionnaire import QuestionRepo, QuestionnaireRepo
from repositories.risk_assessment import CustomerProfileRepo, RiskAssessmentRepo
+from service.memory.composed_profile import invalidate_composed_profile
from utils.exceptions import NotFoundError, ParamError
# 分数 → 风险等级映射(总分 0-100,20 分一档;后续可迁移到 sys_config 运营化)
+# 等级口径统一为 R1-R5:R1 最保守、R5 最激进,与产品 fin_product.risk_level 对齐
_SCORE_LEVELS = (
- (80, "激进"),
- (60, "进取"),
- (40, "平衡"),
- (20, "稳健"),
- (0, "保守"),
+ (80, "R5"),
+ (60, "R4"),
+ (40, "R3"),
+ (20, "R2"),
+ (0, "R1"),
)
@@ -24,7 +26,7 @@ def _score_to_level(score: int) -> str:
for threshold, level in _SCORE_LEVELS:
if score >= threshold:
return level
- return "保守"
+ return "R1"
async def submit_assessment(
@@ -76,6 +78,8 @@ async def submit_assessment(
# 回写画像:无画像初始化,有画像更新风险等级/评分并递增版本号
await CustomerProfileRepo(db).upsert_risk(customer_id, risk_level, total_score)
+ # 画像变更后失效最终画像缓存(失败静默,不影响问卷提交)
+ await invalidate_composed_profile(customer_id)
return {
"assessment_id": record.id,
diff --git a/service/redeem.py b/service/redeem.py
index 22f0c5e..8ac543c 100644
--- a/service/redeem.py
+++ b/service/redeem.py
@@ -16,6 +16,7 @@ from repositories.fin_product import FinProductRepo
from repositories.trade_order import TradeOrderRepo
from schemas.holdings import HoldingResp
from schemas.redeem import RedeemResp
+from service.memory.composed_profile import invalidate_composed_profile
from service.risk.engine import RiskEngine, summarize
from service.risk.settle import settle
from utils.exceptions import ForbiddenError, NotFoundError, ParamError
@@ -118,6 +119,8 @@ async def redeem(
order.id, status="已确认", confirm_time=now
)
await db.commit()
+ # 持仓变更后失效最终画像缓存(失败静默,不影响交易结果)
+ await invalidate_composed_profile(user.id)
account = await account_repo.get_by_customer_id(user.id)
current_holding = await holdings_repo.get_by_customer_product(
user.id, product_id
diff --git a/service/risk/engine.py b/service/risk/engine.py
index 57b83d3..a834d57 100644
--- a/service/risk/engine.py
+++ b/service/risk/engine.py
@@ -31,11 +31,8 @@ from repositories.fin_product import FinProductRepo
from repositories.fin_transaction import FinTransactionRepo
from repositories.risk_rule import RiskRuleRepo
-# 风险等级 → 序号(兼容 R1~R5 与 保守~激进 两套口径)。
-_RISK_RANK = {
- "R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5,
- "保守": 1, "稳健": 2, "平衡": 3, "进取": 4, "激进": 5,
-}
+# 风险等级 → 序号(客户画像与产品统一 R1~R5 口径,历史中文等级已迁移转换)。
+_RISK_RANK = {"R1": 1, "R2": 2, "R3": 3, "R4": 4, "R5": 5}
_LEVEL_RANK = {"低": 1, "中": 2, "高": 3}
diff --git a/sql/risk_level_r1r5_migration.sql b/sql/risk_level_r1r5_migration.sql
new file mode 100644
index 0000000..fe088b2
--- /dev/null
+++ b/sql/risk_level_r1r5_migration.sql
@@ -0,0 +1,38 @@
+-- ---------------------------------------------------------------------
+-- 风险等级口径统一迁移:保守/稳健/平衡/进取/激进 → R1-R5
+-- 映射:保守=R1 / 稳健=R2 / 平衡=R3 / 进取=R4 / 激进=R5
+-- 幂等:UPDATE 仅命中中文值,重复执行无副作用
+-- 覆盖表:fin_customer_profile / fin_risk_assessment / portfolio_benchmark
+-- ---------------------------------------------------------------------
+
+-- 01 数据迁移(中文值 → R 级)
+UPDATE fin_customer_profile SET risk_level = 'R1' WHERE risk_level = '保守';
+UPDATE fin_customer_profile SET risk_level = 'R2' WHERE risk_level = '稳健';
+UPDATE fin_customer_profile SET risk_level = 'R3' WHERE risk_level = '平衡';
+UPDATE fin_customer_profile SET risk_level = 'R4' WHERE risk_level = '进取';
+UPDATE fin_customer_profile SET risk_level = 'R5' WHERE risk_level = '激进';
+
+UPDATE fin_risk_assessment SET risk_level = 'R1' WHERE risk_level = '保守';
+UPDATE fin_risk_assessment SET risk_level = 'R2' WHERE risk_level = '稳健';
+UPDATE fin_risk_assessment SET risk_level = 'R3' WHERE risk_level = '平衡';
+UPDATE fin_risk_assessment SET risk_level = 'R4' WHERE risk_level = '进取';
+UPDATE fin_risk_assessment SET risk_level = 'R5' WHERE risk_level = '激进';
+
+UPDATE portfolio_benchmark SET risk_level = 'R1' WHERE risk_level = '保守';
+UPDATE portfolio_benchmark SET risk_level = 'R2' WHERE risk_level = '稳健';
+UPDATE portfolio_benchmark SET risk_level = 'R3' WHERE risk_level = '平衡';
+UPDATE portfolio_benchmark SET risk_level = 'R4' WHERE risk_level = '进取';
+UPDATE portfolio_benchmark SET risk_level = 'R5' WHERE risk_level = '激进';
+
+-- 02 列注释与数据口径对齐(MODIFY 需带完整列定义,不改动类型/约束)
+ALTER TABLE fin_customer_profile
+ MODIFY COLUMN risk_level VARCHAR(16) NULL
+ COMMENT '风险等级:R1-R5(R1 最保守,R5 最激进,与产品 risk_level 对齐)';
+
+ALTER TABLE fin_risk_assessment
+ MODIFY COLUMN risk_level VARCHAR(16) NOT NULL
+ COMMENT '评定:R1-R5(R1 最保守,R5 最激进)';
+
+ALTER TABLE portfolio_benchmark
+ MODIFY COLUMN risk_level VARCHAR(16) NOT NULL
+ COMMENT '客户风险等级:R1-R5(R1 最保守,R5 最激进)';
diff --git a/sql/schema.sql b/sql/schema.sql
index 2b6b9c3..8f2ced7 100644
--- a/sql/schema.sql
+++ b/sql/schema.sql
@@ -1,8 +1,9 @@
-- =====================================================================
--- 智能公募基金系统 · MySQL 建表脚本(24 张,见开发计划 §3.1)
+-- 智能公募基金系统 · MySQL 建表脚本(36 张,见开发计划 §3.1,与 model/ 目录 ORM 对齐)
-- 执行:mysql --default-character-set=utf8mb4 -u -p < sql/schema.sql
-- 约定:InnoDB / utf8mb4;采用"逻辑外键"(不建物理 FK,便于 Mock 数据灌入),
-- 注释即文档,NL2SQL 以 COMMENT 注入 Schema。
+-- 维护:新增/调整 ORM 模型时同步更新本脚本(01-24 基础业务域,25-36 扩展域)。
-- =====================================================================
SET NAMES utf8mb4;
@@ -36,6 +37,18 @@ DROP TABLE IF EXISTS fund_nav_history;
DROP TABLE IF EXISTS fin_product;
DROP TABLE IF EXISTS fin_customer_profile;
DROP TABLE IF EXISTS sys_user;
+DROP TABLE IF EXISTS nl2sql_query_history;
+DROP TABLE IF EXISTS nl2sql_sensitive_field;
+DROP TABLE IF EXISTS nl2sql_role_column_permission;
+DROP TABLE IF EXISTS nl2sql_role_table_permission;
+DROP TABLE IF EXISTS nl2sql_query_role;
+DROP TABLE IF EXISTS event_log;
+DROP TABLE IF EXISTS advisor_visit_record;
+DROP TABLE IF EXISTS advisor_report;
+DROP TABLE IF EXISTS advisor_draft;
+DROP TABLE IF EXISTS advisor_todo;
+DROP TABLE IF EXISTS sensitive_word;
+DROP TABLE IF EXISTS fin_account;
SET FOREIGN_KEY_CHECKS = 1;
-- ---------------------------------------------------------------------
@@ -64,7 +77,7 @@ CREATE TABLE IF NOT EXISTS sys_user (
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS fin_customer_profile (
customer_id BIGINT UNSIGNED NOT NULL COMMENT '客户ID,一对一关联 sys_user.id',
- risk_level VARCHAR(16) NULL COMMENT '风险等级:保守/稳健/平衡/进取/激进',
+ risk_level VARCHAR(16) NULL COMMENT '风险等级:R1-R5(R1 最保守,R5 最激进,与产品 risk_level 对齐)',
risk_score INT NULL COMMENT '风险评分 0-100,越高越激进',
investment_experience VARCHAR(16) NULL COMMENT '投资经验:0-1/1-3/3-5/5-10/10年以上',
annual_income_range VARCHAR(32) NULL COMMENT '年收入区间',
@@ -87,7 +100,7 @@ CREATE TABLE IF NOT EXISTS fin_product (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
product_code VARCHAR(32) NOT NULL COMMENT '基金代码,唯一,如 F000001',
product_name VARCHAR(128) NOT NULL COMMENT '基金名称',
- product_type VARCHAR(32) NOT NULL COMMENT '类型:货币/债券/混合/股票基金',
+ product_type VARCHAR(32) NOT NULL COMMENT '类型:货币型/债券型/混合型/股票型/指数型/QDII',
risk_level VARCHAR(8) NOT NULL COMMENT '风险等级 R1-R5',
expected_return DECIMAL(7,4) NULL COMMENT '预期年化收益率(%)',
nav DECIMAL(12,6) NULL COMMENT '最新单位净值',
@@ -191,7 +204,7 @@ CREATE TABLE IF NOT EXISTS fin_risk_assessment (
assessment_date DATE NOT NULL COMMENT '评估日期',
question_version VARCHAR(16) NULL COMMENT '问卷版本(题目变更后按版本追溯)',
total_score INT NOT NULL COMMENT '总分 0-100',
- risk_level VARCHAR(16) NOT NULL COMMENT '评定:保守/稳健/平衡/进取/激进',
+ risk_level VARCHAR(16) NOT NULL COMMENT '评定:R1-R5(R1 最保守,R5 最激进)',
answers JSON NULL COMMENT '答题详情 [{q:1,a:B,score:10}]',
assessor_type VARCHAR(16) NOT NULL DEFAULT 'AI评估' COMMENT 'AI评估/人工评估',
valid_until DATE NOT NULL COMMENT '有效期至(默认 +1年,监管一年一评)',
@@ -308,7 +321,7 @@ CREATE TABLE IF NOT EXISTS conversation_archive (
KEY idx_agent (agent_type),
UNIQUE KEY uk_session_message (session_id, message_id),
KEY idx_agent_run (agent_run_id)
-) COMMENT='会话归档表(审计回溯 + Agent 持续学习素材,归档前脱敏)';
+) COMMENT='会话归档表(审计回溯 + Agent 持续学习素材)';
-- ---------------------------------------------------------------------
-- 15 知识元数据表(原文存 MySQL,替代 MinIO)
@@ -492,7 +505,7 @@ CREATE TABLE IF NOT EXISTS fund_performance (
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS portfolio_benchmark (
id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
- risk_level VARCHAR(16) NOT NULL COMMENT '客户风险等级:保守/稳健/平衡/进取/激进',
+ risk_level VARCHAR(16) NOT NULL COMMENT '客户风险等级:R1-R5(R1 最保守,R5 最激进)',
target_allocation JSON NOT NULL COMMENT '目标配置 {货币:50,债券:40,股票:10}',
drift_threshold DECIMAL(5,2) NOT NULL DEFAULT 5.00 COMMENT '偏离阈值(%),超阈值触发调仓建议',
status VARCHAR(8) NOT NULL DEFAULT '启用' COMMENT '启用/停用',
@@ -500,3 +513,246 @@ CREATE TABLE IF NOT EXISTS portfolio_benchmark (
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_risk_level (risk_level)
) COMMENT='组合基准配置表(投顾Agent 再平衡参照,运营可调整)';
+
+-- =====================================================================
+-- 以下为扩展域建表(25-36),与 model/ 目录 ORM 一一对齐
+-- =====================================================================
+
+-- ---------------------------------------------------------------------
+-- 25 客户资金账户表(现金余额,一人一户)
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS fin_account (
+ customer_id BIGINT UNSIGNED NOT NULL COMMENT '客户ID,一人一户,关联 sys_user.id',
+ balance DECIMAL(18,2) NOT NULL DEFAULT 0 COMMENT '账户总余额(含冻结部分)',
+ frozen_amount DECIMAL(18,2) NOT NULL DEFAULT 0 COMMENT '冻结金额,可用余额=balance-frozen_amount(不落库)',
+ currency VARCHAR(8) NOT NULL DEFAULT 'CNY' COMMENT '币种',
+ status VARCHAR(16) NOT NULL DEFAULT '正常' COMMENT '账户状态,默认正常',
+ version INT NOT NULL DEFAULT 0 COMMENT '乐观锁版本号(并发扣款校验)',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+ PRIMARY KEY (customer_id)
+) COMMENT='客户资金账户表(现金余额,申购扣款/赎回入账的账务载体)';
+
+-- ---------------------------------------------------------------------
+-- 26 敏感词库表(合规拦截唯一来源)
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS sensitive_word (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ word VARCHAR(128) NOT NULL COMMENT '敏感词',
+ category VARCHAR(32) NOT NULL COMMENT '敏感词分类(工作台可配置)',
+ level VARCHAR(8) NOT NULL DEFAULT '高' COMMENT '级别,默认高',
+ status VARCHAR(8) NOT NULL DEFAULT '启用' COMMENT '启用/停用',
+ update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+ UNIQUE KEY uk_word (word)
+) COMMENT='敏感词库表(工作台与Agent同源,合规拦截唯一来源)';
+
+-- ---------------------------------------------------------------------
+-- 27 投顾待办表(事件/定时/计算三类来源统一承载)
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS advisor_todo (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ todo_type VARCHAR(32) NOT NULL COMMENT '待办类型(事件/定时/计算来源的细分类型)',
+ customer_id BIGINT UNSIGNED NULL COMMENT '关联客户(可为空)',
+ advisor_id BIGINT UNSIGNED NOT NULL COMMENT '所属投顾',
+ source VARCHAR(32) NOT NULL COMMENT '来源:事件/定时/计算',
+ biz_id VARCHAR(64) NULL COMMENT '关联业务ID(与 todo_type 组合定位业务)',
+ priority VARCHAR(8) NOT NULL DEFAULT '普通' COMMENT '普通/紧急/特急',
+ status VARCHAR(16) NOT NULL DEFAULT '待处理' COMMENT '待办状态,默认待处理',
+ due_at DATETIME NULL COMMENT '截止时间',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ handle_time DATETIME NULL COMMENT '处理时间',
+ UNIQUE KEY uk_todo (todo_type, customer_id, biz_id),
+ KEY idx_advisor_status (advisor_id, status)
+) COMMENT='投顾待办表(事件/定时/计算三类来源统一承载,uk_todo 幂等防重复建待办)';
+
+-- ---------------------------------------------------------------------
+-- 28 投顾 Agent 草稿表(不存 sent,不生成交易指令)
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS advisor_draft (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ draft_id VARCHAR(64) NOT NULL COMMENT '草稿ID,全局唯一',
+ customer_id BIGINT UNSIGNED NOT NULL COMMENT '目标客户',
+ advisor_id BIGINT UNSIGNED NOT NULL COMMENT '所属投顾',
+ intent VARCHAR(32) NOT NULL COMMENT '报告意图,如 调仓建议',
+ title VARCHAR(128) NOT NULL COMMENT '草稿标题',
+ content TEXT NOT NULL COMMENT '报告正文',
+ structured_data JSON NULL COMMENT '结构化数据(调仓明细等)',
+ status VARCHAR(16) NOT NULL DEFAULT 'draft' COMMENT '草稿状态,默认draft',
+ deviation DECIMAL(10,4) NULL COMMENT '组合偏离度(%)',
+ disclaimer_ok TINYINT(1) NOT NULL DEFAULT 0 COMMENT '免责声明确认',
+ warning VARCHAR(512) NULL COMMENT '合规提示/警告信息',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+ UNIQUE KEY uk_draft_id (draft_id),
+ CONSTRAINT ck_advisor_draft_status CHECK (status IN ('draft', 'discarded')),
+ KEY idx_customer (customer_id),
+ KEY idx_advisor (advisor_id),
+ KEY idx_intent (intent),
+ KEY idx_status (status),
+ KEY idx_create_time (create_time)
+) COMMENT='投顾 Agent 草稿(不存 sent,不生成交易指令)';
+
+-- ---------------------------------------------------------------------
+-- 29 建议报告表(工作台本地镜像:sent 状态 + 编辑/发送留痕)
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS advisor_report (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ report_id VARCHAR(64) NOT NULL COMMENT '报告ID,全局唯一',
+ draft_id VARCHAR(64) NOT NULL COMMENT '来源草稿ID(advisor_draft.draft_id)',
+ customer_id BIGINT UNSIGNED NOT NULL COMMENT '目标客户',
+ advisor_id BIGINT UNSIGNED NOT NULL COMMENT '发送投顾',
+ intent VARCHAR(32) NOT NULL COMMENT '报告意图(继承草稿)',
+ title VARCHAR(128) NULL COMMENT '报告标题',
+ content TEXT NULL COMMENT '报告正文',
+ edit_history JSON NULL COMMENT '编辑留痕(JSON数组)',
+ send_status VARCHAR(16) NOT NULL DEFAULT 'draft' COMMENT 'draft/sent/discarded,sent 为工作台本地状态不回写投顾Agent',
+ send_time DATETIME NULL COMMENT '发送时间',
+ send_by BIGINT UNSIGNED NULL COMMENT '发送人',
+ msg_id BIGINT UNSIGNED NULL COMMENT '发送后落 sys_message 的站内信ID',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+ UNIQUE KEY uk_report_id (report_id),
+ KEY idx_customer (customer_id),
+ KEY idx_advisor (advisor_id),
+ KEY idx_send_status (send_status)
+) COMMENT='建议报告表(工作台本地镜像,承载 sent 状态与发送/编辑留痕)';
+
+-- ---------------------------------------------------------------------
+-- 30 投顾回访记录表(人工留痕归档)
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS advisor_visit_record (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ customer_id BIGINT UNSIGNED NOT NULL COMMENT '回访客户',
+ advisor_id BIGINT UNSIGNED NOT NULL COMMENT '回访投顾',
+ visit_type VARCHAR(32) NOT NULL COMMENT '回访类型',
+ visit_time DATETIME NOT NULL COMMENT '回访时间',
+ summary TEXT NULL COMMENT '回访纪要',
+ audio_url VARCHAR(512) NULL COMMENT '录音文件地址',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ KEY idx_customer (customer_id),
+ KEY idx_advisor (advisor_id)
+) COMMENT='投顾回访记录表(人工留痕归档,不触发Agent记忆抽取)';
+
+-- ---------------------------------------------------------------------
+-- 31 事件日志表(事件双写落库,消费以 event_id 幂等)
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS event_log (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ event_id VARCHAR(64) NOT NULL COMMENT '事件ID,全局唯一,消费幂等键',
+ event_name VARCHAR(64) NOT NULL COMMENT '事件名',
+ payload JSON NULL COMMENT '事件负载',
+ trace_id VARCHAR(64) NULL COMMENT '请求链路追踪ID',
+ status VARCHAR(16) NOT NULL DEFAULT 'pending' COMMENT '消费状态,默认pending',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ consume_time DATETIME NULL COMMENT '消费完成时间',
+ UNIQUE KEY uk_event_id (event_id),
+ KEY idx_event_name (event_name),
+ KEY idx_status (status)
+) COMMENT='事件日志表(事件双写落库,Pub/Sub 仅作实时通知,消费幂等)';
+
+-- ---------------------------------------------------------------------
+-- 32 NL2SQL 查询角色(按 sys_user.employee_role 映射)
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS nl2sql_query_role (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ role_code VARCHAR(64) NOT NULL COMMENT '角色编码,全局唯一',
+ role_name VARCHAR(128) NOT NULL COMMENT '角色名称',
+ employee_role VARCHAR(32) NOT NULL COMMENT '对应 sys_user.employee_role,一对一映射',
+ can_query TINYINT(1) NOT NULL DEFAULT 0 COMMENT '是否允许 NL2SQL 查询:0/1',
+ max_rows BIGINT NOT NULL DEFAULT 1000 COMMENT '单次查询最大返回行数',
+ daily_quota BIGINT NOT NULL DEFAULT 0 COMMENT '每日查询配额',
+ status VARCHAR(16) NOT NULL DEFAULT 'active' COMMENT 'active/inactive',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+ UNIQUE KEY uk_role_code (role_code),
+ UNIQUE KEY uk_employee_role (employee_role),
+ CONSTRAINT ck_nl2sql_role_can_query CHECK (can_query IN (0, 1)),
+ CONSTRAINT ck_nl2sql_role_status CHECK (status IN ('active', 'inactive')),
+ KEY idx_status (status)
+) COMMENT='NL2SQL 查询角色,按 sys_user.employee_role 映射';
+
+-- ---------------------------------------------------------------------
+-- 33 NL2SQL 角色的表级 + 行级权限
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS nl2sql_role_table_permission (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ role_id BIGINT UNSIGNED NOT NULL COMMENT '角色ID(nl2sql_query_role.id)',
+ table_name VARCHAR(128) NOT NULL COMMENT '授权表名',
+ permission VARCHAR(16) NOT NULL DEFAULT 'SELECT' COMMENT '授权动作,仅 SELECT',
+ row_scope_type VARCHAR(32) NOT NULL DEFAULT 'none' COMMENT '行级范围:none/customer_ids/product_ids',
+ row_scope_column VARCHAR(128) NULL COMMENT '行级过滤列(如 customer_id)',
+ status VARCHAR(16) NOT NULL DEFAULT 'active' COMMENT 'active/inactive',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+ CONSTRAINT ck_nl2sql_table_permission CHECK (permission = 'SELECT'),
+ CONSTRAINT ck_nl2sql_row_scope_type CHECK (row_scope_type IN ('none', 'customer_ids', 'product_ids')),
+ CONSTRAINT ck_nl2sql_table_status CHECK (status IN ('active', 'inactive')),
+ KEY idx_role (role_id),
+ KEY idx_status (status)
+) COMMENT='NL2SQL 查询角色的表和行级权限';
+
+-- ---------------------------------------------------------------------
+-- 34 NL2SQL 角色的字段权限
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS nl2sql_role_column_permission (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ role_id BIGINT UNSIGNED NOT NULL COMMENT '角色ID(nl2sql_query_role.id)',
+ table_name VARCHAR(128) NOT NULL COMMENT '表名',
+ column_name VARCHAR(128) NOT NULL COMMENT '列名',
+ access_mode VARCHAR(16) NOT NULL DEFAULT 'allow' COMMENT 'allow/deny/mask',
+ mask_type VARCHAR(32) NULL COMMENT '脱敏类型(access_mode=mask 时生效)',
+ status VARCHAR(16) NOT NULL DEFAULT 'active' COMMENT 'active/inactive',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+ CONSTRAINT ck_nl2sql_column_access_mode CHECK (access_mode IN ('allow', 'deny', 'mask')),
+ CONSTRAINT ck_nl2sql_column_status CHECK (status IN ('active', 'inactive')),
+ KEY idx_role (role_id),
+ KEY idx_status (status)
+) COMMENT='NL2SQL 查询角色的字段权限';
+
+-- ---------------------------------------------------------------------
+-- 35 NL2SQL 全局敏感字段规则
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS nl2sql_sensitive_field (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ table_name VARCHAR(128) NOT NULL COMMENT '表名',
+ column_name VARCHAR(128) NOT NULL COMMENT '列名',
+ mask_type VARCHAR(32) NOT NULL DEFAULT 'partial' COMMENT '脱敏方式,如 partial',
+ status VARCHAR(16) NOT NULL DEFAULT 'active' COMMENT 'active/inactive',
+ description VARCHAR(255) NULL COMMENT '规则说明',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+ CONSTRAINT ck_nl2sql_sensitive_status CHECK (status IN ('active', 'inactive')),
+ KEY idx_table_col (table_name, column_name),
+ KEY idx_status (status)
+) COMMENT='NL2SQL 全局敏感字段规则';
+
+-- ---------------------------------------------------------------------
+-- 36 NL2SQL 查询归档(不保存完整结果行)
+-- ---------------------------------------------------------------------
+CREATE TABLE IF NOT EXISTS nl2sql_query_history (
+ id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
+ query_id VARCHAR(64) NOT NULL COMMENT '查询ID,全局唯一',
+ user_id BIGINT UNSIGNED NOT NULL COMMENT '发起用户',
+ session_id VARCHAR(64) NULL COMMENT '会话ID',
+ caller_agent VARCHAR(64) NULL COMMENT '调用方Agent标识',
+ question TEXT NOT NULL COMMENT '自然语言问题',
+ generated_sql TEXT NULL COMMENT '生成的SQL(未执行则为空)',
+ access_tables JSON NULL COMMENT '实际访问的表列表',
+ status VARCHAR(16) NOT NULL COMMENT 'success/failed/blocked/timeout',
+ error_code VARCHAR(64) NULL COMMENT '错误码',
+ error_message VARCHAR(512) NULL COMMENT '错误信息',
+ row_count BIGINT NOT NULL DEFAULT 0 COMMENT '返回行数',
+ truncated TINYINT(1) NOT NULL DEFAULT 0 COMMENT '结果是否被截断',
+ elapsed_ms FLOAT NULL COMMENT '耗时(毫秒)',
+ trace_id VARCHAR(64) NULL COMMENT '请求链路追踪ID',
+ create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ UNIQUE KEY uk_query_id (query_id),
+ CONSTRAINT ck_nl2sql_history_status CHECK (status IN ('success', 'failed', 'blocked', 'timeout')),
+ KEY idx_user (user_id),
+ KEY idx_session (session_id),
+ KEY idx_caller_agent (caller_agent),
+ KEY idx_status (status),
+ KEY idx_trace (trace_id),
+ KEY idx_create_time (create_time)
+) COMMENT='NL2SQL 查询归档,不保存完整结果行';
diff --git a/开发计划.md b/开发计划.md
index 3d8eea9..c497cf5 100644
--- a/开发计划.md
+++ b/开发计划.md
@@ -148,7 +148,7 @@ Agent 编排层(统一执行骨架: 记忆召回→意图路由→核心逻辑
| `sys_message` | 站内信/消息中心(user_id, type, is_read) |
| `memory_unit` | **记忆单元**(tag, content, info_type=FACT/OPINION, source, confidence, evidence/conflict/recall_count, status, valid_until)——记忆架构主体,见 §4 |
| `fund_performance` | 基金业绩指标缓存(近1月/3月/6月/1年/成立以来收益、年化波动率、最大回撤、夏普比率、计算日期)——定时任务从净值计算 |
-| `portfolio_benchmark` | **组合基准配置**(risk_level → 目标资产类别权重,如 保守=货币50/债券40/股票10,可运营调整)——再平衡的参照系 |
+| `portfolio_benchmark` | **组合基准配置**(risk_level → 目标资产类别权重,如 R1=货币50/债券40/股票10,可运营调整)——再平衡的参照系 |
| `sys_config` | 运营参数 KV(**高净值门槛**、偏离度阈值、风控开关等,免发版调整) |
### 3.2 Milvus 集合(3 个业务集 + 1 个记忆集)