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 个记忆集)