From 1ddd44a6cb10f7e2fb497d37ea83ced74ada780f Mon Sep 17 00:00:00 2001 From: Andrew Date: Sat, 5 Sep 2026 17:39:16 +0800 Subject: [PATCH] Add core database support and enhance documentation - Introduced `MYSQL_CORE_DATABASE` in `.env.example` and `settings.py` for core database configuration. - Added `CoreReadOnlyRepository` for read-only access to the `jinrong_core` database. - Updated `AGENTS.md`, `README.md`, and various documentation files to reflect new agent onboarding processes and project structure. - Revised requirements in `requirements.txt` to include `langgraph` and `langchain-core`. - Enhanced `FLOW.md` with local bootstrap instructions for setting up the core simulation environment. - Added new scripts for database creation and seeding for the core simulation library. - Improved overall documentation for clarity on project architecture and memory management. - Updated `TODO.md` to reflect current development priorities and tasks. --- .env.example | 1 + AGENTS.md | 35 +- README.md | 4 +- app/config/settings.py | 1 + app/repository/__init__.py | 5 + app/repository/core_ro.py | 125 +++++++ app/service/agent_service.py | 2 +- docs/memory/ENVIRONMENT.md | 32 +- docs/memory/FLOW.md | 61 +++- docs/memory/FRAMEWORK.md | 70 ++-- docs/memory/ITERATION.md | 4 +- docs/memory/MEMORY.md | 91 ++++-- docs/memory/REQUIREMENTS.md | 7 +- docs/memory/TODO.md | 24 +- docs/业务记忆管理/README.md | 9 + docs/业务记忆管理/业务记忆管理手册.md | 309 ++++++++++++++++++ docs/项目框架设计/Core模拟底座/00-方案总览.md | 264 +++++++++++++++ .../Core模拟底座/01-表结构与种子说明.md | 66 ++++ .../技术选型和版本/01-技术栈与版本.md | 4 +- docs/项目框架设计/表设计/00-架构总览.md | 1 + .../表设计/05-多Agent共用底座清单.md | 1 + requirements.txt | 5 +- scripts/core/00-create-database.sql | 4 + scripts/core/01-ddl.sql | 129 ++++++++ scripts/core/02-seed-base.sql | 70 ++++ scripts/core/03-seed-customers.sql | 107 ++++++ scripts/core/04-seed-holdings.sql | 79 +++++ scripts/core/05-seed-trades.sql | 41 +++ scripts/core/06-seed-nav.sql | 16 + scripts/core/README.md | 43 +++ scripts/core/reset.ps1 | 43 +++ scripts/dev/rbac-seed-reference.md | 59 ++++ scripts/sync/sync_advisor_rel.py | 62 ++++ scripts/sync/sync_neo4j.py | 171 ++++++++++ 34 files changed, 1848 insertions(+), 97 deletions(-) create mode 100644 app/repository/__init__.py create mode 100644 app/repository/core_ro.py create mode 100644 docs/业务记忆管理/README.md create mode 100644 docs/业务记忆管理/业务记忆管理手册.md create mode 100644 docs/项目框架设计/Core模拟底座/00-方案总览.md create mode 100644 docs/项目框架设计/Core模拟底座/01-表结构与种子说明.md create mode 100644 scripts/core/00-create-database.sql create mode 100644 scripts/core/01-ddl.sql create mode 100644 scripts/core/02-seed-base.sql create mode 100644 scripts/core/03-seed-customers.sql create mode 100644 scripts/core/04-seed-holdings.sql create mode 100644 scripts/core/05-seed-trades.sql create mode 100644 scripts/core/06-seed-nav.sql create mode 100644 scripts/core/README.md create mode 100644 scripts/core/reset.ps1 create mode 100644 scripts/dev/rbac-seed-reference.md create mode 100644 scripts/sync/sync_advisor_rel.py create mode 100644 scripts/sync/sync_neo4j.py diff --git a/.env.example b/.env.example index 5848ddb..be8cba7 100644 --- a/.env.example +++ b/.env.example @@ -5,6 +5,7 @@ APP_ENV=development MYSQL_HOST=127.0.0.1 MYSQL_PORT=3306 MYSQL_DATABASE=jinrong_agent +MYSQL_CORE_DATABASE=jinrong_core MYSQL_USER=root MYSQL_PASSWORD= diff --git a/AGENTS.md b/AGENTS.md index 311e1b4..cfa7547 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1,12 +1,33 @@ # Agent -日常只读 [docs/memory/MEMORY.md](docs/memory/MEMORY.md)。 +**无上下文时先读** [docs/memory/MEMORY.md](docs/memory/MEMORY.md) **§0 交接清单**(5 分钟了解仓库地图、bootstrap、禁止项、下一步)。 -- 对照开发:`docs/memory/REQUIREMENTS.md`(场景 ID:F/C/A/D/R) -- 框架/选型:`docs/memory/FRAMEWORK.md` → `docs/项目框架设计/技术选型和版本/` -- 表结构:`docs/项目框架设计/表设计/` -- 业务需求原文:`docs/需求拆解/` +## 阅读顺序 + +| 顺序 | 文件 | 何时 | +| --- | --- | --- | +| 1 | `docs/memory/MEMORY.md` | 每次任务 | +| 2 | `docs/memory/TODO.md` | 取当前待办 | +| 3 | `docs/memory/REQUIREMENTS.md` | 对照场景 ID(F/C/A/D/R) | +| 4 | `docs/memory/FRAMEWORK.md` | 分层、选型、**实现状态表** | +| 5 | `docs/memory/FLOW.md` | 端到端链路 + §0 bootstrap | +| 6 | `docs/memory/ENVIRONMENT.md` | As-Is;Core 为模拟库 | + +## 原文与深度文档 + +- 业务需求:`docs/需求拆解/` +- 表结构 SQL:`docs/项目框架设计/表设计/` +- 技术选型 / JWT:`docs/项目框架设计/技术选型和版本/` +- Core 模拟底座:`docs/项目框架设计/Core模拟底座/` +- 业务记忆分层:`docs/业务记忆管理/业务记忆管理手册.md` + +## 代码入口 + +```text +app/main.py # FastAPI;当前仅 /health +app/config/settings.py # 双库 jinrong_agent + jinrong_core +app/repository/core_ro.py # Core 只读(已实现) +scripts/core/reset.ps1 # 本地灌 Core 模拟库 +``` 技术选型硬阀门见 MEMORY 第 3、7 节。Cursor 以 `.cursor/rules/project-memory.mdc` 为准。 - -后端入口:`app/main.py` · 分层见 `docs/memory/FRAMEWORK.md` §3。 diff --git a/README.md b/README.md index 0ca4c1a..ee0f6f5 100644 --- a/README.md +++ b/README.md @@ -6,7 +6,7 @@ | 目录 | 内容 | | --- | --- | -| [docs/memory/MEMORY.md](docs/memory/MEMORY.md) | 项目记忆入口(Agent 先读) | +| [docs/memory/MEMORY.md](docs/memory/MEMORY.md) | 项目记忆入口(**新 Agent 读 §0 交接清单**) | | [docs/需求拆解/](docs/需求拆解/) | 业务场景、数据矩阵、合规 | | [docs/项目框架设计/](docs/项目框架设计/) | 表设计、技术选型、JWT 手册 | @@ -18,6 +18,7 @@ app/ ├── service/ # Agent、RAG、记忆 ├── tool/ # 解析、Embedding、Milvus ├── model/ # Pydantic + ORM +├── repository/ # Core 只读 core_ro(已实现) ├── config/ # settings、database ├── utils/ └── main.py @@ -33,6 +34,7 @@ app/ ```bash cp .env.example .env pip install -r requirements.txt +# 首次:灌 Core 模拟库与同步 — 见 docs/memory/FLOW.md §0 uvicorn app.main:app --reload ``` diff --git a/app/config/settings.py b/app/config/settings.py index a9753ae..b96c208 100644 --- a/app/config/settings.py +++ b/app/config/settings.py @@ -11,6 +11,7 @@ class Settings(BaseSettings): mysql_host: str = "127.0.0.1" mysql_port: int = 3306 mysql_database: str = "jinrong_agent" + mysql_core_database: str = "jinrong_core" mysql_user: str = "root" mysql_password: str = "" diff --git a/app/repository/__init__.py b/app/repository/__init__.py new file mode 100644 index 0000000..17168ad --- /dev/null +++ b/app/repository/__init__.py @@ -0,0 +1,5 @@ +"""数据访问层。""" + +from app.repository.core_ro import CoreReadOnlyRepository + +__all__ = ["CoreReadOnlyRepository"] diff --git a/app/repository/core_ro.py b/app/repository/core_ro.py new file mode 100644 index 0000000..ec8594f --- /dev/null +++ b/app/repository/core_ro.py @@ -0,0 +1,125 @@ +"""Core 模拟库只读访问(jinrong_core · 无 HTTP API)。""" + +from __future__ import annotations + +from datetime import date, datetime +from typing import Any + +from sqlalchemy import create_engine, text +from sqlalchemy.engine import Engine + +from app.config.settings import settings + + +class CoreReadOnlyRepository: + """仅 SELECT jinrong_core;禁止写操作。""" + + def __init__(self, engine: Engine | None = None) -> None: + self._engine = engine or self._default_engine() + + @staticmethod + def _default_engine() -> Engine: + pwd = settings.mysql_password + auth = f"{settings.mysql_user}:{pwd}" if pwd else settings.mysql_user + url = ( + f"mysql+pymysql://{auth}@{settings.mysql_host}:{settings.mysql_port}" + f"/{settings.mysql_core_database}?charset=utf8mb4" + ) + return create_engine(url, pool_pre_ping=True) + + def get_customer_l0(self, customer_id: str) -> dict[str, Any] | None: + sql = text( + """ + SELECT c.customer_id, c.display_name, c.age, c.occupation, c.open_date, + r.risk_code, r.evaluated_at AS risk_evaluated_at + FROM core_customer c + LEFT JOIN core_customer_risk r ON r.customer_id = c.customer_id + WHERE c.customer_id = :cid AND c.is_active = 1 + """ + ) + with self._engine.connect() as conn: + row = conn.execute(sql, {"cid": customer_id}).mappings().first() + return dict(row) if row else None + + def list_holdings(self, customer_id: str) -> list[dict[str, Any]]: + sql = text( + """ + SELECT h.*, p.product_name, p.min_risk_code, p.product_type + FROM core_holding h + JOIN core_product p ON p.product_id = h.product_id + WHERE h.customer_id = :cid + ORDER BY h.market_value DESC + """ + ) + with self._engine.connect() as conn: + return [dict(r) for r in conn.execute(sql, {"cid": customer_id}).mappings()] + + def list_trades( + self, customer_id: str, since: date | None = None, limit: int = 50 + ) -> list[dict[str, Any]]: + sql = text( + """ + SELECT t.*, p.product_name + FROM core_trade t + JOIN core_product p ON p.product_id = t.product_id + WHERE t.customer_id = :cid + AND (:since IS NULL OR t.traded_at >= :since) + ORDER BY t.traded_at DESC + LIMIT :lim + """ + ) + with self._engine.connect() as conn: + return [ + dict(r) + for r in conn.execute( + sql, {"cid": customer_id, "since": since, "lim": limit} + ).mappings() + ] + + def get_product(self, product_id: str) -> dict[str, Any] | None: + sql = text("SELECT * FROM core_product WHERE product_id = :pid") + with self._engine.connect() as conn: + row = conn.execute(sql, {"pid": product_id}).mappings().first() + return dict(row) if row else None + + def get_latest_nav(self, product_id: str) -> dict[str, Any] | None: + sql = text( + """ + SELECT * FROM core_product_nav + WHERE product_id = :pid + ORDER BY nav_date DESC LIMIT 1 + """ + ) + with self._engine.connect() as conn: + row = conn.execute(sql, {"pid": product_id}).mappings().first() + return dict(row) if row else None + + def list_customers_by_advisor(self, advisor_id: str) -> list[str]: + sql = text( + """ + SELECT customer_id FROM core_customer_advisor + WHERE advisor_id = :aid AND rel_status = 'active' + """ + ) + with self._engine.connect() as conn: + return [r[0] for r in conn.execute(sql, {"aid": advisor_id})] + + def get_staff(self, staff_id: str) -> dict[str, Any] | None: + """RBAC 联调:查员工角色种子。""" + sql = text( + "SELECT staff_id, display_name, staff_type, roles FROM core_staff WHERE staff_id = :sid AND is_active = 1" + ) + with self._engine.connect() as conn: + row = conn.execute(sql, {"sid": staff_id}).mappings().first() + return dict(row) if row else None + + def is_advisor_assigned(self, advisor_id: str, customer_id: str) -> bool: + sql = text( + """ + SELECT 1 FROM core_customer_advisor + WHERE advisor_id = :aid AND customer_id = :cid AND rel_status = 'active' + LIMIT 1 + """ + ) + with self._engine.connect() as conn: + return conn.execute(sql, {"aid": advisor_id, "cid": customer_id}).first() is not None diff --git a/app/service/agent_service.py b/app/service/agent_service.py index 4d671b9..0e607f8 100644 --- a/app/service/agent_service.py +++ b/app/service/agent_service.py @@ -1 +1 @@ -"""Agent 编排:LangChain + DeepSeek;Tool 调用;四 Agent 能力边界。""" +"""Agent 编排:LangGraph StateGraph + DeepSeek;Tool 节点;四 Agent 能力边界。""" diff --git a/docs/memory/ENVIRONMENT.md b/docs/memory/ENVIRONMENT.md index c5e8b60..c163b23 100644 --- a/docs/memory/ENVIRONMENT.md +++ b/docs/memory/ENVIRONMENT.md @@ -1,6 +1,7 @@ # 服务对象大环境(现状 / As-Is) -> 来源:`docs/需求拆解/业务场景优先级清单.md`、`docs/需求拆解/用户故事/` +> 来源:`docs/需求拆解/业务场景优先级清单.md`、`docs/需求拆解/用户故事/` +> **Core 说明:** 生产环境应以真实 Core 为 L0 权威;**本项目无真实 Core**,开发用 `jinrong_core` 模拟库(见 `docs/项目框架设计/Core模拟底座/`)。 ------ @@ -13,7 +14,7 @@ | 内部员工 | 分析、运营、合规 | 问数、报表、抽检会话、协同风控 | | 风控 / 合规官 | 风控专员 | 监测大额/适当性/AML;人工审核,不自动冻户 | -**Core:** 持仓、流水、交易、正式风险等级、客户-代理人归属的权威来源。 +**L0 权威来源(本项目):** MySQL `jinrong_core` 模拟库 — 客户、C1~C5、持仓、流水、产品、代理人-客户归属。Agent 代码经 `CoreReadOnlyRepository` **只读 SELECT**,不提供 HTTP Core API。 ------ @@ -24,26 +25,28 @@ - **内部分析**:SQL 或 IT 导表。 - **风控**:事后/T+1 发现异常;适当性靠交易前规则;AML 批处理或人工。 -Agent **尚未上线**。 +Agent **尚未上线**;后端为脚手架 + 模拟数据,无生产接入。 ------ -## 3. 主路径(现状) +## 3. 主路径(现状 → 模拟后) -1. 开户/测评 → Core 写 L0(C1~C5、归属代理人) -2. 交易 → Core 流水;交易前适当性(若有) -3. 代理人服务 → 查 Core + 口头沟通(无统一 L2 画像库) -4. 风控 → 规则/人工台账;合规事后抽检 +1. 开户/测评 → Core(模拟:`core_customer` + `core_customer_risk`)写 L0 +2. 交易 → Core 流水(模拟:`core_trade`);交易前适当性(R-02 待实现) +3. 同步 → `sync_advisor_rel.py` 写入 agent 库 `customer_advisor_rel`;`sync_neo4j.py` 灌关系图 +4. 代理人服务 → 查 Core RO + 口头沟通(L2 画像表待业务实现) +5. 风控 → 规则/人工台账;合规事后抽检 ------ ## 4. 环境里的关键对象 -| 名称 | 含义 | -| --- | --- | -| L0 正式档案 | customer_id、C1~C5、年龄、职业、advisor | -| 持仓/流水 | 产品、份额、盈亏 | -| 代理人-客户归属 | customer_advisor_rel(Core 同步) | +| 名称 | 含义 | 本项目存储 | +| --- | --- | --- | +| L0 正式档案 | customer_id、C1~C5、年龄、职业、advisor | `jinrong_core.core_*` | +| 持仓/流水 | 产品、份额、盈亏 | `core_holding` / `core_trade` | +| 代理人-客户归属 | 服务关系 | Core 种子 → 同步 `jinrong_agent.customer_advisor_rel` | +| Agent 会话/画像/审计 | 平台业务数据 | `jinrong_agent`(11 共用表 + agent 专用表) | ------ @@ -51,5 +54,6 @@ Agent **尚未上线**。 ```text 痛点:客户看不清持仓与规则;代理人多系统慢、话术易违规;问数靠 SQL;风控感知滞后 -不能动:正式 C1~C5 仅测评/人工变更;事实账以 Core 为准;监管要求审计与适当性 +不能动:正式 C1~C5 仅测评/人工变更;事实账以 Core 为准(模拟库亦只读);监管要求审计与适当性 +模拟边界:不模拟 TA/清算/报盘;种子数据静态净值可接受;Neo4j 由同步脚本维护 ``` diff --git a/docs/memory/FLOW.md b/docs/memory/FLOW.md index a8ad714..e36e6e1 100644 --- a/docs/memory/FLOW.md +++ b/docs/memory/FLOW.md @@ -5,12 +5,49 @@ ------ -## 1. 链路总览 +## 0. 本地 bootstrap(新环境 / 新 Agent) ```text -Client → Gateway(JWT/RBAC) → api/chat → agent_service → Tools → 存储 → 响应 + audit_log +前置:MySQL 8、Redis、Neo4j Desktop、Ollama(bge-m3) 已安装并启动 + +① 配置 + cp .env.example .env + 必填:MYSQL_PASSWORD、NEO4J_PASSWORD、DEEPSEEK_API_KEY(对话阶段) + +② 依赖 + pip install -r requirements.txt + +③ Core 模拟库(L0) + .\scripts\core\reset.ps1 + → 创建 jinrong_core + 28 客户 / 12 产品 / 持仓交易种子 + +④ Agent 共用底座(若库未建) + mysql -u root -p < docs/项目框架设计/表设计/01-mysql-共用底座.sql + +⑤ 同步 + python scripts/sync/sync_advisor_rel.py → customer_advisor_rel + python scripts/sync/sync_neo4j.py → 关系图 + +⑥ 验证 + uvicorn app.main:app --reload + GET http://127.0.0.1:8000/health → {"status":"ok"} + (可选)Python 调 CoreReadOnlyRepository.get_customer_l0("CUST-10001") + +RBAC 联调账号:scripts/dev/rbac-seed-reference.md ``` +**尚未自动化:** agent 专用表 SQL、Milvus 建 Collection、JWT 中间件 — 见 `TODO.md`。 + +------ + +## 1. 链路总览(目标) + +```text +Client → Gateway(JWT/RBAC) → api/chat → agent_service(LangGraph) → Tools → 存储 → 响应 + audit_log +``` + +当前:`main.py` 仅 health;全链路待 Wave 0 起逐步挂载。 + ------ ## 2. 主链路 · 对话 @@ -18,19 +55,19 @@ Client → Gateway(JWT/RBAC) → api/chat → agent_service → Tools → 存储 ```text 输入:Authorization + X-Agent-Type + X-Trace-Id + 用户消息 ↓ -Gateway:验签、角色准入、注入 AuthContext +Gateway:验签、角色准入、注入 AuthContext 【待做 T-01】 ↓ -api/chat:SessionGuard;创建/续 agent_session +api/chat:SessionGuard;创建/续 agent_session 【待做】 ↓ -memory_service:Redis 读最近 N 轮;异步写 agent_message +memory_service:Redis 读最近 N 轮;异步写 agent_message 【待做】 ↓ -agent_service:LangChain + DeepSeek;按 Agent 类型选 Tool 集 +agent_service:LangGraph StateGraph + DeepSeek;按 Agent 类型选 Tool 节点 【待做】 ↓ Tool 示例: - - Core RO:持仓/流水(自动注入 customer_id) - - milvus_tool:产品规则 RAG + source_refs - - memory_service:读/写 L1/L2/L3(ProfileGuard) - - 风控 service 账号:R-02 适当性(无会话) + - CoreReadOnlyRepository:持仓/流水/L0(自动注入 customer_id) 【Repository 已有,Tool 未接】 + - milvus_tool:产品规则 RAG + source_refs 【待做】 + - memory_service:读/写 L1/L2/L3(ProfileGuard) 【待做】 + - 风控 service 账号:R-02 适当性(无会话) 【待做】 ↓ 输出:assistant 消息 + has_disclaimer(客户/对外) ↓ @@ -68,13 +105,15 @@ admin/knowledge 上传 → document_parser → embedding_tool(Ollama bge-m3) | 会话 | session_id | Redis ctx + MySQL session | actor_id 一致 | | RAG | question + product_id? | chunks + source_doc_id | 必须溯源 | | 画像读 | customer_id | L1/L2/L3 JSON | RBAC 矩阵 §8 | +| Core RO | customer_id | L0/持仓/流水 | 只读 jinrong_core | | 审计 | 任意判定 | audit_log INSERT | 不可删改 | ------ ## 6. 关键数据约定 -- MySQL 脚本:`docs/项目框架设计/表设计/01-mysql-共用底座.sql`、`02-mysql-agent专用.sql` +- Agent 库 SQL:`docs/项目框架设计/表设计/01-mysql-共用底座.sql`、`02-mysql-agent专用.sql` +- Core 模拟:`scripts/core/01-ddl.sql` · 说明 `docs/项目框架设计/Core模拟底座/01-表结构与种子说明.md` - Redis Key:`docs/项目框架设计/表设计/02-redis-keys.md` - Milvus:`kb_product_rules`、`kb_business_ops`;向量 **1024** 维 - 禁止:Agent 覆盖 L0 正式 C1~C5 diff --git a/docs/memory/FRAMEWORK.md b/docs/memory/FRAMEWORK.md index 1921a2e..6a9a79e 100644 --- a/docs/memory/FRAMEWORK.md +++ b/docs/memory/FRAMEWORK.md @@ -2,7 +2,8 @@ > 选型详情:`docs/项目框架设计/技术选型和版本/01-技术栈与版本.md` > 表设计:`docs/项目框架设计/表设计/` -> 鉴权:`docs/项目框架设计/技术选型和版本/02-JWT-RBAC鉴权手册.md` +> 鉴权:`docs/项目框架设计/技术选型和版本/02-JWT-RBAC鉴权手册.md` +> Core 模拟:`docs/项目框架设计/Core模拟底座/` ------ @@ -11,15 +12,15 @@ | 技术栈 | 已定(用户已同意) | 备注 | | --- | --- | --- | | 后端语言/框架 | Python 3.13.14 + FastAPI | 系统 Python | -| Agent 编排 | langchain 1.3.18 / langchain-openai 1.6.0 | DeepSeek API | -| 关系库 | MySQL 8.0.46 | jinrong_agent,:3306 | +| Agent 编排 | **LangGraph** 1.2.x + langchain-core / langchain-openai 1.6.0 | DeepSeek API;StateGraph + Tool 节点 | +| 关系库 | MySQL 8.0.46 | **双库**:`jinrong_agent` + `jinrong_core` | | 缓存 | Redis 8.10.1 | :6379 | -| 图库 | Neo4j 5.26.19 Enterprise | Desktop,关系查询 | -| 向量库 | Milvus Lite + pymilvus 3.0.1 | 本地文件;1024 维 | +| 图库 | Neo4j 5.26.19 Enterprise | Desktop;`sync_neo4j.py` 灌图 | +| 向量库 | Milvus Lite + pymilvus 3.0.1 | 本地 `./data/milvus.db`;1024 维 | | Embedding | Ollama bge-m3 | 本地,不出内网 | | LLM 生成 | DeepSeek API | 对话/推理 | -| 前端 | React 19 + Vite 7 + AntD 5 + TS strict | 不用 Streamlit | -| 对象存储 | 本地目录 data/kb/ | 不用 MinIO | +| 前端 | React 19 + Vite 7 + AntD 5 + TS strict | 不用 Streamlit;`web/` 待 init | +| 对象存储 | 本地目录 `data/kb/` | 不用 MinIO | | 部署 | Windows 原生 | Docker 仅答辩备选 | ```text @@ -31,30 +32,41 @@ ## 2. 能力模块 -| 模块 | 职责 | 依赖 | -| --- | --- | --- | -| Agent Gateway / Auth SDK | JWT、RBAC、归属校验 | Redis、MySQL customer_advisor_rel | -| 客户财富 Agent | L1 画像、事实查询、阈值提醒 | Core RO、Milvus 产品库 | -| 代理人助手 Agent | L2 画像、RAG、草稿 | L1 只读、Milvus | -| 数据分析 Agent | NL→SQL→解读 | Core RO、画像只读 | -| 风控监测 Agent | 预警、L3、R-02 适当性 | 交易事件、AML 名单 | -| 共用底座 | 会话、审计、输入防护 | MySQL 11 表 + Redis | +| 模块 | 职责 | 依赖 | 代码状态 | +| --- | --- | --- | --- | +| Agent Gateway / Auth SDK | JWT、RBAC、归属校验 | Redis、MySQL customer_advisor_rel | **未做**(见 JWT 手册) | +| 客户财富 Agent | L1 画像、事实查询、阈值提醒 | Core RO、Milvus 产品库 | 空壳 service | +| 代理人助手 Agent | L2 画像、RAG、草稿 | L1 只读、Milvus | 空壳 service | +| 数据分析 Agent | NL→SQL→解读 | Core RO、画像只读 | 空壳 service | +| 风控监测 Agent | 预警、L3、R-02 适当性 | 交易事件、AML 名单 | 空壳 service | +| Core 只读层 | L0 事实查询 | `jinrong_core` | **core_ro.py 已实现** | +| 共用底座 | 会话、审计、输入防护 | MySQL 11 表 + Redis | SQL 已定;代码未接 | +| 同步脚本 | 归属、Neo4j | Core → agent / 图库 | **sync_*.py 已实现** | ------ ## 3. 后端分层(app/) ```text -api/ → 路由(chat、knowledge、admin);薄,不含业务 -service/ → agent_service、rag_service、memory_service -tool/ → document_parser、embedding_tool、milvus_tool -model/ → schemas(Pydantic)、entities(ORM) -config/ → settings、database -utils/ → response、exceptions、logger -main.py → FastAPI 入口 +api/ → 路由(chat、knowledge、admin);薄,不含业务 【空壳】 +service/ → agent_service、rag_service、memory_service 【空壳】 +tool/ → document_parser、embedding_tool、milvus_tool 【空壳】 +repository/ → core_ro(Core 只读);后续 agent 库 Repository 【core_ro 已实现】 +model/ → schemas(Pydantic)、entities(ORM) 【占位】 +config/ → settings、database 【settings 已实现】 +utils/ → response、exceptions、logger 【占位】 +main.py → FastAPI 入口;当前仅 /health 【部分】 -允许:api → service → tool / model / config -禁止:api 直连 Milvus/MySQL 写复杂逻辑;tool 写业务流程 +允许:api → service → tool / repository / model / config +禁止:api 直连 Milvus/MySQL 写复杂逻辑;tool 写业务流程;repository 写 Core +``` + +**脚本(非 app 包):** + +```text +scripts/core/ → jinrong_core DDL + 种子 + reset.ps1 +scripts/sync/ → sync_advisor_rel.py、sync_neo4j.py +scripts/dev/ → rbac-seed-reference.md(联调账号) ``` ------ @@ -63,14 +75,16 @@ main.py → FastAPI 入口 | 层 | 存储 | 说明 | | --- | --- | --- | -| L0 | Core 只读 | 正式 C1~C5、持仓 | -| L1/L2/L3 | MySQL 画像表 | 客户/代理人/风控 enrich | +| L0 | `jinrong_core` 只读 | 正式 C1~C5、持仓、流水(模拟) | +| L1/L2/L3 | MySQL 画像表(agent 库) | 客户/代理人/风控 enrich | | 会话 | Redis + MySQL | 窗口 + 永久审计 | -| RAG | Milvus + 本地文件 | kb_product_rules 等 | +| RAG | Milvus + `data/kb/` | kb_product_rules 等 | | 关系 | Neo4j | 客户-产品-代理人(Core 同步) | ------ ## 5. 骨架一句话 -Gateway 鉴权 → FastAPI api → service 编排 Agent → tool 访问 MySQL/Redis/Milvus/Neo4j/Core → 审计落库。 +Gateway 鉴权 → FastAPI api → service(LangGraph)→ tool/repository → MySQL/Redis/Milvus/Neo4j/Core → 审计落库。 + +**记忆分层详解:** [业务记忆管理手册.md](../业务记忆管理/业务记忆管理手册.md) diff --git a/docs/memory/ITERATION.md b/docs/memory/ITERATION.md index 0a25f3f..c80ae1a 100644 --- a/docs/memory/ITERATION.md +++ b/docs/memory/ITERATION.md @@ -5,4 +5,6 @@ | 2026-09-05 | init 项目记忆 + app 脚手架 | 用户要求 memory kit + 目录结构 | MEMORY / ENVIRONMENT / REQUIREMENTS / FRAMEWORK / FLOW / TODO | | 2026-09-05 | 需求拆解四 Agent P0 场景定稿 | 用户故事归纳 | REQUIREMENTS(对照 docs/需求拆解/) | | 2026-09-05 | 技术栈 Windows 原生 + Milvus Lite | 本机内存限制 | FRAMEWORK | -| 2026-09-05 | MySQL 16 表 + L0~L3 画像 + JWT 统一鉴权 | 多 Agent 共用底座 | FRAMEWORK / FLOW | +| 2026-09-05 | Core 模拟库 jinrong_core + 扩大种子 + Neo4j 同步 | 无真实 Core | FRAMEWORK / FLOW / ENVIRONMENT | +| 2026-09-05 | Agent 编排依赖改为 LangGraph 为主 | 状态图 + Tool 节点 | FRAMEWORK | +| 2026-09-05 | memory 全量刷新:§0 新 Agent 交接清单 + 实现状态表 + bootstrap | 脚手架已落地,便于无上下文交接 | MEMORY / FRAMEWORK / FLOW / TODO / REQUIREMENTS / ENVIRONMENT | diff --git a/docs/memory/MEMORY.md b/docs/memory/MEMORY.md index 10b2296..ebbc586 100644 --- a/docs/memory/MEMORY.md +++ b/docs/memory/MEMORY.md @@ -1,15 +1,55 @@ # 项目记忆(人 + Agent 共用) -> 各阶段的**简化快照**(门闸索引)。只留基础信息,不写源文件正文。 -> 详细需求见 `docs/需求拆解/`;表设计与技术选型见 `docs/项目框架设计/`。 +> 各阶段的**简化快照**(门闸索引)。只留基础信息,不写源文件正文。 +> **新 Agent 无上下文:先读本节 → 第 0 节交接清单 → 按需打开源文件。** + +------ + +## 0. 新 Agent 交接(5 分钟) + +**项目是什么:** 金融四 Agent(客户财富 / 代理人 / 数据分析 / 风控)共用数据层与合规底座;**不**互调 LLM,跨 Agent 走 L1/L2/L3 画像与预警表。 + +**当前进度:** 需求与表设计已定 · 后端 **脚手架 + Core 模拟库脚本** 已落地 · **业务 API / JWT / LangGraph 图尚未实现**(多为空壳模块)。 + +**仓库地图:** + +| 路径 | 状态 | 说明 | +| --- | --- | --- | +| `app/main.py` | 可跑 | 仅 `/health`;路由未挂载 | +| `app/api/*.py` | 空壳 | chat / knowledge / admin 待实现 | +| `app/service/*.py` | 空壳 | agent / rag / memory 待实现 | +| `app/repository/core_ro.py` | **已实现** | Core 只读 SELECT(jinrong_core) | +| `app/config/settings.py` | **已实现** | 双库 `jinrong_agent` + `jinrong_core` | +| `scripts/core/*.sql` + `reset.ps1` | **已实现** | Core 模拟库 DDL + 种子 | +| `scripts/sync/*.py` | **已实现** | 归属同步 + Neo4j 全图 | +| `docs/需求拆解/` | 已定 | 场景 P0、矩阵、合规原文 | +| `docs/项目框架设计/表设计/` | 已定 | Agent 共用 11 表 + agent 专用 SQL | +| `docs/项目框架设计/Core模拟底座/` | 已定 | 无真实 Core 时的 L0 方案 | +| `web/` | **不存在** | 前端 React 待 init | + +**本地 bootstrap(首次):** + +```text +1. cp .env.example .env → 填 MYSQL_PASSWORD、NEO4J、DEEPSEEK_API_KEY +2. pip install -r requirements.txt +3. .\scripts\core\reset.ps1 # 建 jinrong_core + 种子 +4. mysql … < 表设计/01-mysql-共用底座.sql # jinrong_agent(若未建) +5. python scripts/sync/sync_advisor_rel.py +6. python scripts/sync/sync_neo4j.py +7. uvicorn app.main:app --reload # GET /health +``` + +**下一步开发(见 TODO):** T-01 JWT 中间件 → 挂载 api → LangGraph agent_service → Wave 0 验收。 + +**禁止(改代码前必记):** Core 正式 C1~C5 不可被画像覆盖 · 审计表只 INSERT · 代理人草稿不外发 · 仅 R-02 可阻断交易 · 四 Agent 不互调 LLM。 ------ ## 1. 项目简介 - **名称:** JinRong 金融四 Agent 智能管家 -- **当前阶段:** MVP(需求与底座设计已定,后端脚手架已 init) -- **当前优先级:** 交付速度(先 Wave 0 共用底座 + Wave 1 内部 Agent 闭环) +- **当前阶段:** MVP 脚手架期(设计已定,Wave 0 开发中) +- **当前优先级:** Wave 0 共用底座 → Wave 1 代理人/分析闭环 ------ @@ -17,10 +57,10 @@ | 阶段 | 当前一句话 | 源文件 | | --- | --- | --- | -| 服务对象大环境 | 公募场景下四类角色各自痛点明确,数据以 Core 为准、Agent 只读事实 | `ENVIRONMENT.md` | -| 需求目的 | 四个 Agent 服务四类人群,统一数据层交换画像与预警,合规可审计 | `REQUIREMENTS.md` | -| 项目框架 | FastAPI 分层 app + MySQL/Redis/Milvus Lite/Neo4j + React;Windows 原生 | `FRAMEWORK.md` | -| 实现流程 | Gateway JWT/RBAC → Agent 编排 → Tool/画像库 → 审计落库 | `FLOW.md` | +| 服务对象大环境 | 公募四类角色;**无真实 Core**,L0 由 `jinrong_core` 模拟库供给 | `ENVIRONMENT.md` | +| 需求目的 | 四 Agent + 统一数据层 + 合规可审计;Wave 0~3 P0 见 REQUIREMENTS | `REQUIREMENTS.md` | +| 项目框架 | FastAPI 分层 + LangGraph + MySQL 双库/Redis/Milvus Lite/Neo4j + React | `FRAMEWORK.md` | +| 实现流程 | Gateway JWT → chat → LangGraph → Tool → 审计;bootstrap 见 FLOW §0 | `FLOW.md` | ------ @@ -33,7 +73,7 @@ **禁止触碰:** ```text -Core 正式 C1~C5、持仓/交易真账(Agent 只读) +Core 正式 C1~C5、持仓/交易真账(Agent 只读;当前为模拟库) audit_log 等审计表(只 INSERT) 代理人草稿自动外发客户 风控自动冻户 / 改正式风险等级 @@ -66,22 +106,32 @@ audit_log 等审计表(只 INSERT) ## 4. 入口与运行 ```text -关键文件:app/main.py · app/config/settings.py · docs/项目框架设计/表设计/*.sql -如何启动:MySQL/Redis/Neo4j/Ollama 就绪后 → uvicorn app.main:app --reload -配置位置:.env(见 .env.example) +后端入口:app/main.py · 配置 app/config/settings.py · Core 只读 app/repository/core_ro.py +Agent 库 SQL:docs/项目框架设计/表设计/01-mysql-共用底座.sql +Core 模拟:scripts/core/reset.ps1 · 文档 docs/项目框架设计/Core模拟底座/ +依赖:requirements.txt(LangGraph + langchain-core/openai + FastAPI + SQLAlchemy) +启动:uvicorn app.main:app --reload → GET /health +配置:.env(见 .env.example) +RBAC 联调账号:scripts/dev/rbac-seed-reference.md ``` ------ ## 5. 文档资产约定 -| 文件 | 何时更新 | +| 文件 | 何时读 | | --- | --- | -| `docs/memory/*.md` | 见 project-memory-kit 约定 | -| `docs/需求拆解/` | 业务场景与合规(参考) | -| `docs/项目框架设计/` | 表结构、鉴权、技术版本 | +| `docs/memory/MEMORY.md` | **每次任务先读**(本文件) | +| `docs/memory/REQUIREMENTS.md` | 对照场景 ID(F/C/A/D/R)与 Wave | +| `docs/memory/FRAMEWORK.md` | 分层、选型、模块、**实现状态表** | +| `docs/memory/FLOW.md` | 端到端链路 + 本地 bootstrap | +| `docs/memory/ENVIRONMENT.md` | As-Is 业务环境与 Core 模拟说明 | +| `docs/memory/TODO.md` | 当前待办与已完成 | +| `docs/需求拆解/` | 业务原文(场景、矩阵、合规) | +| `docs/项目框架设计/` | 表结构、JWT 手册、Core 模拟、技术版本 | +| `docs/业务记忆管理/` | Redis 短期 vs MySQL/Milvus/Neo4j 权威记忆 | -日常改代码只读本文件;验收对照 `REQUIREMENTS.md`。 +缺 `docs/memory/*` 文件:按 project-memory-kit 同名补回,**禁止空模板盖进度**。 ------ @@ -94,8 +144,9 @@ audit_log 等审计表(只 INSERT) ## 7. 阅读完成门闸(Agent) 1. 四个 Agent 服务谁、禁止什么? -2. 改动属于 api / service / tool 哪一层? +2. 改动属于 api / service / tool / repository 哪一层? 3. 是否需 customer_id 归属与 JWT RBAC? -4. 如何验证? +4. Core 是模拟库只读还是 agent 库读写? +5. 如何验证?(health / SQL / sync 脚本 / 对照 REQUIREMENTS 验收列) -大任务:FRAMEWORK/FLOW 为空时先补再编码(用户确认跳过除外)。 +大任务:FRAMEWORK/FLOW 与实现状态不符时先更新 memory 再编码(用户确认跳过除外)。 diff --git a/docs/memory/REQUIREMENTS.md b/docs/memory/REQUIREMENTS.md index 4d23cb3..953f44b 100644 --- a/docs/memory/REQUIREMENTS.md +++ b/docs/memory/REQUIREMENTS.md @@ -20,8 +20,11 @@ | F-01 | JWT + RBAC + 数据归属 | 越权 403 + audit;JWT 手册 §13 | 未做 | T-01 | | F-02 | 全量审计留痕 | trace_id 可还原 | 未做 | T-02 | | F-03 | 输入防护 | input_guard_log | 未做 | T-03 | -| F-04 | Core 只读层 | 不改 Core 账 | 未做 | T-04 | -| R0-DB | MySQL 共用 11 表 + Redis | 01-mysql-共用底座.sql | 未做 | T-05 | +| F-04 | Core 只读层 | 不改 Core 账;Repository 只 SELECT | **部分** | T-04 | +| R0-DB | MySQL 共用 11 表 + Redis | 01-mysql-共用底座.sql | SQL 已定;灌库待验 | T-05 | +| R0-CORE | Core 模拟 + 同步 | reset.ps1 + sync 脚本 | **脚本已落地** | T-05 | + +**F-04 部分完成说明:** `app/repository/core_ro.py` + `scripts/core/*` + `settings.mysql_core_database` 已有;尚未接入 api/service Tool 与归属校验。 ## Wave 1 · 内部 Agent P0 diff --git a/docs/memory/TODO.md b/docs/memory/TODO.md index bed9135..67b997b 100644 --- a/docs/memory/TODO.md +++ b/docs/memory/TODO.md @@ -1,21 +1,29 @@ # TODO -> 大任务开始时更新;完成一步打钩并验证后再往下。 +> 大任务开始时更新;完成一步打钩并验证后再往下。 +> **新 Agent:** 先看 `MEMORY.md` §0 交接清单,再从此处取下一项。 ## 进行中 -- [ ] Wave 0 共用底座 — 验收:F-01~F-04 + MySQL 11 表可连 +- [ ] Wave 0 共用底座 — 验收:F-01~F-03 + agent 库灌库 + JWT 可联调 -## 待办 +## 待办(推荐顺序) -- [ ] T-01 Auth SDK / JWT 中间件(对照 JWT 手册) -- [ ] T-05 执行 01-mysql-共用底座.sql + database.py 连接 +- [ ] T-01 Auth SDK / JWT 中间件(对照 `02-JWT-RBAC鉴权手册.md`) +- [ ] T-02 audit_log 中间件 + trace_id 贯通 +- [ ] T-05 本机执行 `reset.ps1` + `01-mysql-共用底座.sql` + sync 脚本验收 +- [ ] T-06 挂载 `app.api` 路由到 `main.py`;chat 最小闭环 +- [ ] T-07 LangGraph `agent_service` StateGraph 骨架 + DeepSeek +- [ ] T-04 Core RO 封装为 Tool 节点;A-01 归属校验 - [ ] T-21 Milvus Lite + kb_product_rules 首批入库 -- [ ] T-20 代理人 A-01 持仓 Tool + 归属校验 -- [ ] 前端 React 多 Agent 入口(HashRouter) +- [ ] 前端 React 多 Agent 入口(HashRouter,`web/` init) ## 已完成 -- [x] 2026-09-05 init 项目记忆目录 + app/ 脚手架 +- [x] 2026-09-05 init 项目记忆目录 + `app/` 脚手架(api/service/tool/model/config/utils) - [x] 2026-09-05 需求拆解文档(业务场景/矩阵/合规) - [x] 2026-09-05 表设计 + 技术选型 + JWT 手册 +- [x] 2026-09-05 Core 模拟库 `scripts/core/*` + `reset.ps1` + 文档 +- [x] 2026-09-05 `core_ro.py` + `settings.mysql_core_database` + sync 脚本 +- [x] 2026-09-05 Agent 编排依赖改为 LangGraph(requirements.txt) +- [x] 2026-09-05 memory 文件夹更新(新 Agent 交接清单) diff --git a/docs/业务记忆管理/README.md b/docs/业务记忆管理/README.md new file mode 100644 index 0000000..f88315b --- /dev/null +++ b/docs/业务记忆管理/README.md @@ -0,0 +1,9 @@ +# 业务记忆管理 + +> 四 Agent 记忆分层:短期(Redis)vs 权威(MySQL)vs 知识(Milvus)vs 关系(Neo4j)vs 官方事实(Core) + +| 文档 | 说明 | +| --- | --- | +| [业务记忆管理手册.md](./业务记忆管理手册.md) | 主文档:什么放哪儿、决策流程、验收清单 | + +表结构 / Redis Key 细节见 [../项目框架设计/表设计/](../项目框架设计/表设计/)。 diff --git a/docs/业务记忆管理/业务记忆管理手册.md b/docs/业务记忆管理/业务记忆管理手册.md new file mode 100644 index 0000000..e9b9b85 --- /dev/null +++ b/docs/业务记忆管理/业务记忆管理手册.md @@ -0,0 +1,309 @@ +# 业务记忆管理手册 + +> 写给:产品、开发、测试、合规 +> 目的:说清楚 **什么叫短期记忆、什么必须长期保存**,以及 **Redis / MySQL / Neo4j / Milvus / Core / 本地文件** 各自放什么 +> 关联:[00-架构总览.md](../项目框架设计/表设计/00-架构总览.md) · [02-redis-keys.md](../项目框架设计/表设计/02-redis-keys.md) · [03-milvus-collections.md](../项目框架设计/表设计/03-milvus-collections.md) · [04-neo4j-model.md](../项目框架设计/表设计/04-neo4j-model.md) · [数据交互矩阵.md](../需求拆解/数据交互矩阵.md) + +--- + +## 1. 一句话总览 + +```text +短期记忆 = 为了「这一轮对话跑得顺」的临时数据,丢了能重建,不做法务证据。 +权威记忆 = 必须长期保存、能审计、跨 Agent 共享的业务结论,以 MySQL 为准。 +知识记忆 = 产品/制度文档的语义检索(Milvus),不是聊天记录。 +关系记忆 = 谁持有啥、谁管谁、产品要什么风险等级(Neo4j),金额事实仍以 Core 为准。 +官方事实 = 持仓、流水、正式 C1~C5(Core 只读,Agent 库不复制账表)。 +``` + +**铁律:** Redis 里的内容 **永远不是最终真相**;合规纠纷、监管检查、跨 Agent 交换,一律以 **MySQL + Core** 为准。 + +--- + +## 2. 记忆分层模型(业务语言) + +可以把整个系统的「记忆」分成五层,从「聊完就忘」到「必须留档」: + +| 层级 | 业务名称 | 技术载体 | 生命周期 | 丢了怎么办 | +| --- | --- | --- | --- | --- | +| **M0** | 对话草稿 | **Redis** | 分钟~小时(TTL) | 从 MySQL 最近消息重建上下文,略慢 | +| **M1** | 权威业务记忆 | **MySQL** | 永久(只增不改的审计表) | **不可接受丢失** | +| **M2** | 知识库记忆 | **Milvus** + 本地 `data/kb/` | 随文档版本更新 | 重新切片、embedding 导入 | +| **M3** | 关系记忆 | **Neo4j** | 随 Core 同步刷新 | 从 Core 重跑同步 Job | +| **M4** | 官方事实 | **Core(只读)** | 业务系统权威 | Agent 不建第二套账 | + +另外还有 **L0~L3 用户画像**(见 §4):L0 在 Core,L1/L2/L3 在 MySQL(Redis 只做热缓存)。 + +--- + +## 3. 什么叫「短期记忆」? + +### 3.1 定义 + +**短期记忆** = 当前会话进行中、为了少查库、少重复推理而放在 **Redis** 里的数据。 + +特征: + +- 有 **TTL**(过期自动删) +- **可重建**(源数据在 MySQL 或 Core) +- **不参与合规最终认定**(例如不能以 Redis 里的草稿代替 audit_log) +- **不跨 Agent 长期共享**(跨 Agent 共享走 MySQL 画像表) + +### 3.2 典型短期记忆(P0 全部在 Redis) + +| 业务场景 | Redis Key(示例) | 存什么 | TTL | 为何是短期 | +| --- | --- | --- | --- | --- | +| 多轮对话上下文 | `sess:{agent}:{session_id}:ctx` | 当前意图、槽位、上一轮 Tool 摘要 | 2h | 关页/超时后不需要;完整消息在 MySQL | +| 最近聊天窗口 | `sess:{agent}:{session_id}:msgs` | 最近 ≤20 轮 JSON | 2h | 滑动窗口;全量在 `agent_message` | +| 会话并发锁 | `sess:{agent}:{session_id}:lock` | 防双写 | 30s | 纯技术锁 | +| 画像热缓存 | `profile:l1/l2/l3:{...}` | MySQL 画像 JSON 副本 | 5~10m | 加速读;权威在 MySQL | +| 代理人资产快照缓存 | `cache:advisor:snapshot:...` | A-01 查 Core 后的摘要 | 15m | 权威副本在 L2 `asset_snapshot` | +| 风控预警推送 | `risk:pub:alert` | Pub/Sub 通知 | — | 实时通道;单据在 `risk_alert` | +| 预警去重 | `risk:dedup:...` | 同日同规则是否已报 | 24h | 防风暴;不是业务台账 | +| 限流/封禁 | `guard:rate` / `guard:block` | 计数、临时封禁 | 1m~15m | 安全控制 | +| Token 吊销 | `auth:revoked:{jti}` | JWT 黑名单 | 至 exp | 鉴权辅助 | + +### 3.3 短期记忆的读写规则 + +```text +写:对话每条消息 → 先/并行写 MySQL agent_message → 再更新 Redis 窗口 +读:优先 Redis 窗口拼上下文 → 不够再读 MySQL 最近 N 条 +删:TTL 到期自动删;画像 MySQL UPDATE 后主动 DEL 对应 profile:* Key +``` + +**禁止放进 Redis 的(必须 MySQL):** + +- 审计总账 `audit_log` +- 预警单最终状态 `risk_alert.status` +- 草稿审核结果 `advisor_draft.review_status` +- 分析 SQL 留痕 `analytics_query_log` +- 适当性阻断记录 `risk_suitability_log` + +--- + +## 4. 用户画像:L0~L3 存哪儿? + +| 层级 | 业务含义 | 权威存储 | Redis 缓存 | 谁写 | +| --- | --- | --- | --- | --- | +| **L0** | 正式 C1~C5、年龄、职业、资产规模 | **Core 只读** | 一般不缓存(或极短 TTL) | 非 Agent | +| **L1** | 客户偏好:风格、规划、阈值摘要、行为标签 | **MySQL** `customer_profile_l1` | `profile:l1:{customer_id}` | 客户 Agent | +| **L2** | 服务侧:诉求、待办、资产概况快照、服务标签 | **MySQL** `customer_profile_l2` | `profile:l2:{customer_id}:{advisor_id}` | 代理人 Agent | +| **L3** | 监测侧:正常/关注/高风险、评分维度 | **MySQL** `customer_profile_l3` | `profile:l3:{customer_id}` | 风控 Agent | + +**记忆管理要点:** + +- L1/L2/L3 **以 MySQL 为权威**;Redis 只是「刚查过」的副本。 +- **禁止**用 L1 客户口头偏好 **覆盖** L0 正式风险等级。 +- 客户 **不可见** L2/L3(API 层 404,不是 Redis 里藏一下就行)。 + +--- + +## 5. MySQL:权威业务记忆放什么? + +MySQL = **正式档案柜**。凡是 **要审计、要跨 Agent 读、要留痕** 的,都落 MySQL。 + +### 5.1 按业务类型对照表 + +| 业务类型 | 典型内容 | MySQL 表(P0) | 是否短期 | +| --- | --- | --- | --- | +| 会话归档 | 每句 user/assistant/tool | `agent_session`, `agent_message`, `agent_tool_call` | 长期 | +| 合规审计 | 谁问了什么、判定结果 | `audit_log`, `input_guard_log` | 长期,**只 INSERT** | +| 数据归属 | 客户归哪个代理人 | `customer_advisor_rel` | 长期(Core 同步) | +| 客户画像 L1 | 风格、规划、阈值偏好 | `customer_profile_l1` | 长期 | +| 服务画像 L2 | 诉求、待办、快照 | `customer_profile_l2` | 长期 | +| 监测画像 L3 | 分层、评分、标签 | `customer_profile_l3` | 长期 | +| 风控预警 | 预警单、人工处置 | `risk_alert` | 长期 | +| 适当性 | 匹配/阻断记录 | `risk_suitability_log` | 长期 | +| 客户阈值 | C-04 自设亏损线 | `customer_threshold_config` | 长期 | +| 客户提醒 | 阈值/波动提醒发送记录 | `customer_notify_log` | 长期 | +| 代理人草稿 | 话术/跟进草稿 | `advisor_draft` | 长期 | +| 违规命中 | A-07 合规词库 | `compliance_hit_log` | 长期 | +| 分析留痕 | NL→SQL→解读 | `analytics_query_log` | 长期 | + +### 5.2 MySQL 不存什么? + +| 不存 | 改存 | 原因 | +| --- | --- | --- | +| 产品手册全文块(大段 PDF 文本) | Milvus chunk + 本地 `data/kb/` | 体积大、要语义检索 | +| 客户-产品-持仓 **关系遍历** | Neo4j | 图遍历更高效 | +| 持仓/流水 **权威金额** | Core 只读 API | 禁止 Agent 各算一套 | +| 当前对话「最近 5 轮」窗口 | Redis | 临时、TTL | + +--- + +## 6. Milvus:知识库记忆放什么? + +Milvus = **按意思找文档**,不是存对话、不是存客户画像。 + +| 放什么 | Collection | 典型场景 | +| --- | --- | --- | +| 基金产品手册、费率、申赎规则、风险说明 | `kb_product_rules` | C-02, C-03, A-02 | +| 内部办事流程、办理条件 | `kb_business_ops` | A-04 | +| (P1)合规话术规范模板 | `kb_compliance_scripts` | A-03 | + +**每条向量记录应包含:** `chunk_text`、`embedding`(1024 维 bge-m3)、`source_doc_id`、`source_version`、`product_id` 等(见 [03-milvus-collections.md](../项目框架设计/表设计/03-milvus-collections.md))。 + +**本地文件:** 原始 PDF/Word 在 `data/kb/`(不用 MinIO);Milvus 存切块与向量,MySQL 可选存文档版本索引(P1)。 + +**不放 Milvus:** + +- 聊天记录、审计日志 +- 客户 L1/L2/L3 画像 JSON +- 预警单、适当性结果 + +--- + +## 7. Neo4j:关系记忆放什么? + +Neo4j = **关系网**,回答「谁持有啥、产品要什么等级、归哪个代理人管」。 + +| 放什么 | 节点/关系 | 典型场景 | +| --- | --- | --- | +| 客户、代理人、产品、风险等级、行业 | `Customer`, `Advisor`, `Product`, `RiskGrade` | 全 Agent 关系查询 | +| 客户归属代理人 | `Customer -[:ASSIGNED_TO]-> Advisor` | F-01 归属(与 MySQL rel 表一致) | +| 客户正式风险等级 L0 | `Customer -[:HAS_RISK_LEVEL]-> RiskGrade` | R-02 | +| 持仓关系(同步快照) | `Customer -[:HOLDS]-> Product` | C-01, A-01, R-01 | +| 产品最低适配等级 | `Product -[:REQUIRES_MIN_RISK]-> RiskGrade` | R-02 | + +**同步原则:** 金额、`as_of` 来自 Core 同步 Job;**Agent 写 L1/L2/L3 不进 Neo4j**,避免双写。 + +**不放 Neo4j:** + +- 会话、消息、审计 +- 向量、RAG 文档 +- 画像 JSON 字段 + +--- + +## 8. Core 与本地文件 + +| 存储 | 放什么 | Agent 权限 | +| --- | --- | --- | +| **Core** | L0 正式档案、持仓、流水、交易、行情/净值 | **只读** | +| **本地 `data/kb/`** | 原始制度/产品 PDF、Word | 平台入库用;Agent 通过 Milvus 检索 | +| **本地 `data/milvus.db`** | Milvus Lite 文件 | 向量索引文件,非业务表 | + +--- + +## 9. 决策流程:新数据放哪儿? + +```text + 新产生一条数据 + │ + ┌───────────────┼───────────────┐ + ▼ ▼ ▼ + 是否官方账事实? 是否多轮对话临时态? 是否文档/RAG? + │ │ │ + 是 是 是 + │ │ │ + ▼ ▼ ▼ + Core 只读 Redis(TTL) Milvus + data/kb/ + 不在 Agent 建表 + MySQL 消息归档 + source 溯源字段 + │ │ + 否 否 + │ │ + └───────┬───────┘ + ▼ + 是否需要审计 / 跨 Agent 共享? + │ + ┌────────┴────────┐ + 是 否 + ▼ ▼ + MySQL 是否关系遍历? + (选对 L1/L2/L3 │ + 或业务表) ┌─────┴─────┐ + 是 否 + ▼ ▼ + Neo4j 重新评估 + (Core 同步) 是否其实该进 MySQL +``` + +### 9.1 快速对照(常见误区) + +| 误区 | 正确做法 | +| --- | --- | +| 聊天记录只放 Redis | Redis 窗口 + **MySQL 全量归档** | +| 画像只放 Redis | **MySQL 权威** + Redis 热缓存 | +| 产品规则塞 MySQL TEXT | **Milvus** 向量检索 + 溯源 ID | +| 持仓金额以 Neo4j 为准 | **Core 为准**;Neo4j 是同步快照 + 关系 | +| 预警状态放 Redis | **MySQL `risk_alert`**;Redis 只做 Pub/Sub 通知 | +| 四 Agent 共享记忆靠互相调 LLM | **统一 MySQL 画像/预警表** + RBAC | + +--- + +## 10. 四 Agent × 记忆读写(摘要) + +| Agent | 短期(Redis 写) | 权威(MySQL 写) | 只读 | +| --- | --- | --- | --- | +| 客户财富 | 会话窗口 | L1、阈值、提醒日志、会话消息 | Core、Milvus 产品库、L3 不可见 | +| 代理人助手 | 会话、快照缓存 | L2、草稿、合规命中 | L1/L3、Core、Milvus | +| 数据分析 | 会话窗口 | 分析 SQL 留痕 | 画像 L1/L2/L3、预警、Core | +| 风控监测 | dedup、Pub/Sub | L3、预警、适当性 | L0/L1/L2、Core、Neo4j | + +详细矩阵见 [05-多Agent共用底座清单.md](../项目框架设计/表设计/05-多Agent共用底座清单.md) §九。 + +--- + +## 11. 一致性与失效策略 + +| 场景 | 策略 | +| --- | --- | +| Redis 挂了 | 从 MySQL 拉最近消息续聊;画像直接读 MySQL;服务降级 | +| MySQL 与 Redis 画像不一致 | **以 MySQL 为准**;修复后 DEL Redis Key | +| Milvus 与本地文件不一致 | 以 `source_version` 为准;A-09 知识库校验 | +| Neo4j 与 Core 不一致 | **以 Core 为准**;重跑同步 Job | +| 文档更新 | 重新 parse → embed → 写 Milvus;旧 version 过滤不返回 | + +--- + +## 12. 与代码模块的对应(app/) + +| 模块 | 记忆职责 | +| --- | --- | +| `memory_service.py` | Redis 会话/画像缓存;MySQL 会话与画像落盘 | +| `rag_service.py` | Milvus 检索 + 溯源 | +| `agent_service.py` | 编排上下文(读 Redis/MySQL);不写错层 | +| `milvus_tool.py` | Collection CRUD | +| `config/database.py` | MySQL / Redis / Neo4j 连接 | + +--- + +## 13. 验收检查(记忆管理) + +- [ ] 任意一条对话可在 MySQL 按 `trace_id` 还原,不仅靠 Redis +- [ ] Redis TTL 过期后,续聊仍可从 MySQL 恢复最近上下文 +- [ ] L1 更新后,`profile:l1:*` 被失效 +- [ ] 产品回答带 `source_doc_id`,来自 Milvus 而非模型幻觉 +- [ ] 持仓金额展示与 Core 一致;Neo4j 仅作关系辅助 +- [ ] 审计表无 UPDATE/DELETE 业务路径 + +--- + +## 14. 关联文档 + +| 文档 | 内容 | +| --- | --- | +| [02-redis-keys.md](../项目框架设计/表设计/02-redis-keys.md) | Redis Key 明细 | +| [01-mysql-共用底座.sql](../项目框架设计/表设计/01-mysql-共用底座.sql) | 共用表 DDL | +| [02-mysql-agent专用.sql](../项目框架设计/表设计/02-mysql-agent专用.sql) | 专用表 DDL | +| [03-milvus-collections.md](../项目框架设计/表设计/03-milvus-collections.md) | 向量 Collection | +| [04-neo4j-model.md](../项目框架设计/表设计/04-neo4j-model.md) | 图模型 | +| [数据交互矩阵.md](../需求拆解/数据交互矩阵.md) | 跨 Agent 数据对象 | +| [FRAMEWORK.md](../memory/FRAMEWORK.md) | 项目框架快照 | +| [FLOW.md](../memory/FLOW.md) | 端到端链路 | + +--- + +## 附录:记忆类型 × 存储 × 示例一句 + +| 记忆类型 | 存储 | 示例 | +| --- | --- | --- | +| 短期对话态 | Redis | 「这一轮用户在问 C-02 基金费率」 | +| 长期对话档案 | MySQL | 「2026-09-05 14:03 客户问了 XX 基金申赎规则」 | +| 客户偏好 | MySQL L1 + Redis 缓存 | 「风格偏价值;亏损 10% 要提醒」 | +| 服务记录 | MySQL L2 | 「客户近期关注赎回;待跟进电话」 | +| 监测标签 | MySQL L3 | 「监测分层:关注;最近有大额预警」 | +| 产品知识 | Milvus | 「某基金封闭期 18 个月(chunk + 出处)」 | +| 持有关系 | Neo4j | 「客户 A 持有产品 P1、P2」 | +| 真实持仓金额 | Core | 「截至 T 日市值 100 万」 | +| 合规审计 | MySQL | 「trace-xxx:代理人查客户 B 被拒绝,AUTH_403」 | diff --git a/docs/项目框架设计/Core模拟底座/00-方案总览.md b/docs/项目框架设计/Core模拟底座/00-方案总览.md new file mode 100644 index 0000000..8f36cf4 --- /dev/null +++ b/docs/项目框架设计/Core模拟底座/00-方案总览.md @@ -0,0 +1,264 @@ +# Core 模拟底座方案(无真实 Core、无 HTTP API) + +> 状态:**已落地**(2026-09-05) +> 详见 [01-表结构与种子说明.md](./01-表结构与种子说明.md) + +--- + +## 1. 要解决什么问题 + +| 需求 | 没有 Core 时的问题 | 本方案 | +| --- | --- | --- | +| C-01/A-01 查持仓 | 无数据 | 种子持仓表 | +| R-02 适当性 | 无 C1~C5 / 产品 R 等级 | 种子测评 + 产品风险等级 | +| F-01 代理人只看名下客户 | 无归属 | 种子归属 + 同步到 `customer_advisor_rel` | +| D-01 分析 Agent 查数 | 无宽表 | Core 库可被只读 SQL 查(同 MySQL 实例) | +| R-01 大额监测 | 无交易流 | 种子交易 + 可选「注入脚本」模拟新交易 | +| Neo4j 关系图 | 无源 | 从 Core 库 **同步脚本** 灌图 | + +**不做什么:** + +- 不模拟真实 TA/清算/报盘 +- 不提供 HTTP Core API(以后有真 Core 再换 Adapter,Repository 接口不变) +- 不让 Agent 写 Core 库(无 UPDATE/INSERT 权限或代码层禁止) + +--- + +## 2. 推荐架构(方案 A · 首选) + +```text +┌─────────────────────────────────────────────────────────┐ +│ MySQL 8.0(本机 :3306,同一实例两个库) │ +├─────────────────────────┬───────────────────────────────┤ +│ jinrong_core │ jinrong_agent │ +│ 【模拟 Core · L0 权威】 │ 【Agent 平台 · 已有 16 表】 │ +│ 客户/产品/持仓/流水/测评 │ 会话/画像/预警/审计 │ +│ ★ Agent 只读 SELECT │ Agent 读写 │ +└───────────┬─────────────┴───────────────┬───────────────┘ + │ │ + │ scripts/sync/*.py │ + └──────────► customer_advisor_rel (agent 库) + └──────────► Neo4j (可选 sync_neo4j.py) + +┌─────────────────────────────────────────────────────────┐ +│ app/repository/core_ro.py │ +│ 只读 SQL · 禁止写 · 供 Tool / 分析 Agent 白名单 schema │ +└─────────────────────────────────────────────────────────┘ +``` + +### 为何不用 API、不用 JSON 文件? + +| 方案 | 优点 | 缺点 | 结论 | +| --- | --- | --- | --- | +| **A. 独立库 jinrong_core + SQL 种子** | 贴近真 Core;D-01 可直接 SELECT;易 reset | 多一套 DDL | **推荐** | +| B. JSON/YAML fixtures | 简单 | 分析 Agent 难 SQL;大数据量慢 | 仅适合单元测试 | +| C. 全塞进 agent 库 | 一张库 | 违反「Core / Agent 分治」;合规演示不清 | 不推荐 | + +--- + +## 3. 库与目录规划 + +```text +scripts/ +└── core/ + ├── 00-create-database.sql # CREATE DATABASE jinrong_core + ├── 01-ddl.sql # Core 表结构 + ├── 02-seed-base.sql # 产品、代理人、风险等级字典 + ├── 03-seed-customers.sql # 客户 L0 + 测评 + ├── 04-seed-holdings.sql # 持仓 + ├── 05-seed-trades.sql # 流水/交易(含 R-01 大额样例) + ├── 06-seed-nav.sql # 净值/行情(C-05) + ├── reset.ps1 # Windows:drop 数据 + 重跑 seed + └── README.md # 执行顺序 + +scripts/sync/ + ├── sync_advisor_rel.py # core → jinrong_agent.customer_advisor_rel + └── sync_neo4j.py # core → Neo4j(P0 可选第二批) + +app/ +└── repository/ + ├── __init__.py + └── core_ro.py # 只读查询封装(无 HTTP) + +docs/项目框架设计/Core模拟底座/ + ├── 00-方案总览.md # 本文件 + └── 01-表结构与种子说明.md # 确认后补充字段级文档 +``` + +**执行顺序(开发/答辩前):** + +```text +1. mysql < scripts/core/00-create-database.sql +2. mysql < scripts/core/01-ddl.sql +3. mysql < scripts/core/02~06-seed-*.sql +4. mysql < docs/.../01-mysql-共用底座.sql # agent 库(若未建) +5. python scripts/sync/sync_advisor_rel.py +6. python scripts/sync/sync_neo4j.py # 可选 +``` + +--- + +## 4. Core 库表设计(P0 最小集) + +> 命名前缀 `core_`,与 agent 库表 **绝不重名**。 + +| 表名 | 对应 L0 能力 | P0 场景 | +| --- | --- | --- | +| `core_advisor` | 理财师/员工 | F-01、A-01 | +| `core_customer` | 客户主档(脱敏字段) | 全 Agent | +| `core_customer_risk` | 正式 C1~C5 + 测评时间 | R-02、C-07 | +| `core_customer_advisor` | 客户-代理人归属 | F-01 → 同步 agent 库 | +| `core_product` | 基金/产品 + R1~R5 | C-02、R-02 | +| `core_holding` | 持仓份额、成本、市值、盈亏% | C-01、A-01、C-04 | +| `core_trade` | 申购/赎回/转换流水 | R-01、D-01 | +| `core_cash_flow` | 资金进出 | C-01 | +| `core_product_nav` | 最新净值/日涨跌 | C-05 | + +**P1 可加:** `core_industry`、`core_dividend`、`core_aml_list`(R-03 名单) + +### 4.1 关键字段约定 + +- 主键统一字符串:`CUST-xxxx`、`STAFF-xxxx`、`PROD-xxxx`、`TRD-xxxx` +- 金额 `DECIMAL(18,2)`;份额 `DECIMAL(18,4)`;盈亏比例 `DECIMAL(8,4)` +- 所有事实表带 `as_of DATE` 或 `updated_at`,展示时必须带「数据截至」 +- `core_customer_risk.is_authoritative = 1` 表示 L0 正式等级(Agent 不可写此表) + +--- + +## 5. 种子数据设计(演示用 personas) + +建议 **少而全**:6 个客户 + 2 代理人 + 8~10 产品,覆盖 P0 验收。 + +| customer_id | 正式等级 | 代理人 | 用途 | +| --- | --- | --- | --- | +| `CUST-9527` | C3 | STAFF-10086 | 主 demo:正常持仓、问持仓/净值 | +| `CUST-1001` | C1 | STAFF-10086 | R-02:尝试买 R5 应阻断 | +| `CUST-1002` | C2 | STAFF-10086 | 亏损接近阈值 C-04 | +| `CUST-2001` | C4 | STAFF-10087 | 跨代理人:10086 不能查 | +| `CUST-3001` | C3 | STAFF-10086 | R-01:单笔 50 万+ 交易 seed | +| `CUST-4001` | C5 | STAFF-10087 | 高龄 + 高风险产品冲突(R-02) | + +| advisor_id | 名下客户 | +| --- | --- | +| `STAFF-10086` | 9527, 1001, 1002, 3001 | +| `STAFF-10087` | 2001, 4001 | + +| product_id | 风险 | 说明 | +| --- | --- | --- | +| `PROD-110022` | R2 | 债基,C1 可买 | +| `PROD-005827` | R3 | 混合 | +| `PROD-161725` | R4 | 行业主题 | +| `PROD-XYZ999` | R5 | 用于适当性拒单 demo | + +**Milvus 产品文档**仍走 `data/kb/` + 向量库;Core 的 `core_product` 只放 **结构化字段**(代码、名称、R 等级、费率摘要),与 RAG 文档通过 `product_id` 关联。 + +--- + +## 6. 应用层:只读 Repository(非 API) + +```python +# app/repository/core_ro.py 职责边界 + +class CoreReadOnlyRepository: + """仅 SELECT jinrong_core.*;连接只读账号或同一连接但方法内禁止写。""" + + def get_customer_l0(customer_id: str) -> CustomerL0 + def list_holdings(customer_id: str) -> list[Holding] + def list_trades(customer_id: str, since: date) -> list[Trade] + def get_product(product_id: str) -> Product + def get_latest_nav(product_id: str) -> NavQuote + def list_customers_by_advisor(advisor_id: str) -> list[str] # 归属 +``` + +- **Tool 层**(agent_service)只调 Repository,不拼裸 SQL。 +- **数据分析 Agent**:SQL 网关白名单仅允许 `jinrong_core` + `jinrong_agent` 指定表;默认禁止 JOIN 跨库写。 +- **以后有真 Core**:新增 `CoreHttpAdapter` 实现同一接口,种子库保留作集成测试。 + +--- + +## 7. 同步脚本(模拟「Core → Agent/Neo4j」) + +### 7.1 `sync_advisor_rel.py` + +```text +读 core_customer_advisor (status=active) +→ UPSERT jinrong_agent.customer_advisor_rel +→ 供 JWT 归属校验 / F-01 +``` + +### 7.2 `sync_neo4j.py`(第二批) + +```text +core_customer / core_advisor / core_product / core_holding +→ Neo4j Customer/Advisor/Product/HOLDS/ASSIGNED_TO/HAS_RISK_LEVEL/REQUIRES_MIN_RISK +``` + +同步频率:**开发期手动跑**;README 写「改 seed 后必跑 sync」。 + +### 7.3 `inject_trade.py`(可选) + +```text +命令行:模拟一笔新交易写入 core_trade +→ 触发风控 Agent 测试 R-01(不必真 API) +``` + +--- + +## 8. 权限与安全(开发期) + +| 项 | 做法 | +| --- | --- | +| MySQL 用户 | `agent_app`:对 `jinrong_agent` 读写;对 `jinrong_core` **仅 SELECT** | +| 代码 | Repository 内仅 `session.execute(text("SELECT ..."))` | +| 测试 | 断言无 `INSERT/UPDATE/DELETE` 指向 `jinrong_core` | + +答辩环境可共用 root;生产模拟仍建议分权限。 + +--- + +## 9. 与现有文档的用语对齐 + +| 原说法 | 落地后 | +| --- | --- | +| Core 只读 API | **Core 只读 Repository(查 jinrong_core)** | +| Core 同步 | **scripts/sync/*.py** | +| L0 | **jinrong_core 表内正式字段** | +| F-04 | Agent 不写 core 库,只 SELECT | + +建议在 `业务记忆管理手册.md` §8 加一句:「开发期 Core = `jinrong_core` 模拟库」。 + +--- + +## 10. 实施分期 + +| 阶段 | 交付 | 工期估 | +| --- | --- | --- | +| **P0-1** | DDL + base/customer/holding seed + sync_advisor_rel | 1~2 天 | +| **P0-2** | core_ro.py + 2 个 Tool(持仓/产品)+ C-01/A-01 可 demo | 1~2 天 | +| **P0-3** | trade seed + inject_trade + R-01/R-02 联调 | 1 天 | +| **P1** | sync_neo4j + nav seed + D-01 跨表 SQL | 1~2 天 | + +--- + +## 11. 验收清单 + +- [ ] `jinrong_core` 与 `jinrong_agent` 同实例不同库,Agent 代码无 core 写操作 +- [ ] 客户 CUST-9527 可查持仓;CUST-2001 代理人 10086 查不到(403) +- [ ] CUST-1001(C1)买 PROD-XYZ999(R5)→ R-02 阻断记录 +- [ ] `reset.ps1` 一键清空并重灌 seed,同步 rel 后代理人权限仍正确 +- [ ] 分析 Agent 可对 `jinrong_core.core_holding` 做 COUNT/GROUP BY + +--- + +## 12. 待你确认的点 + +~~已确认:28 客户、RBAC 多角色、Neo4j P0、静态净值。~~ + +--- + +## 13. 关联文档 + +- [05-多Agent共用底座清单.md](../表设计/05-多Agent共用底座清单.md) §七 Core 只读接口 → 本方案替代实现 +- [04-neo4j-model.md](../表设计/04-neo4j-model.md) +- [业务记忆管理手册.md](../../业务记忆管理/业务记忆管理手册.md) §8 M4 官方事实 +- [业务场景优先级清单.md](../../需求拆解/业务场景优先级清单.md) P0 场景 diff --git a/docs/项目框架设计/Core模拟底座/01-表结构与种子说明.md b/docs/项目框架设计/Core模拟底座/01-表结构与种子说明.md new file mode 100644 index 0000000..2ed901c --- /dev/null +++ b/docs/项目框架设计/Core模拟底座/01-表结构与种子说明.md @@ -0,0 +1,66 @@ +# Core 模拟底座 · 表结构与种子说明 + +> 状态:**已落地** · 库名 `jinrong_core` +> 方案:[00-方案总览.md](./00-方案总览.md) + +--- + +## 1. 已确认决策 + +| 项 | 决定 | +| --- | --- | +| 库名 | `jinrong_core` | +| 访问方式 | `CoreReadOnlyRepository` 只读 SQL,无 HTTP API | +| 种子规模 | 28 客户、13 员工(含 RBAC 多角色)、12 产品 | +| Neo4j | **P0 必做** · `scripts/sync/sync_neo4j.py` | +| 行情 | 静态 `core_product_nav` seed | + +--- + +## 2. 表清单 + +见 [scripts/core/01-ddl.sql](../../../scripts/core/01-ddl.sql) + +| 表 | 说明 | +| --- | --- | +| `core_risk_grade` | C1~C5 / R1~R5 字典 | +| `core_staff` | 内部员工 + **roles JSON**(RBAC 种子) | +| `core_customer` | 客户主档 | +| `core_customer_risk` | L0 正式测评 | +| `core_customer_advisor` | 归属 → sync 到 agent 库 | +| `core_product` | 产品 + 最低风险等级 | +| `core_holding` | 持仓快照 | +| `core_trade` / `core_cash_flow` | 流水 | +| `core_product_nav` | 净值 C-05 | +| `core_industry` | 行业(Neo4j BELONGS_TO) | + +--- + +## 3. 脚本与代码 + +| 路径 | 作用 | +| --- | --- | +| `scripts/core/00~06-*.sql` | 建库 + 种子 | +| `scripts/core/reset.ps1` | 一键重置 | +| `scripts/sync/sync_advisor_rel.py` | → `customer_advisor_rel` | +| `scripts/sync/sync_neo4j.py` | → Neo4j 图 | +| `app/repository/core_ro.py` | 只读 Repository | +| `scripts/dev/rbac-seed-reference.md` | RBAC 验收对照 | + +--- + +## 4. 环境变量 + +```env +MYSQL_CORE_DATABASE=jinrong_core +NEO4J_URI=bolt://localhost:7687 +NEO4J_PASSWORD=... +``` + +--- + +## 5. 关联文档 + +- [RBAC 手册](../技术选型和版本/02-JWT-RBAC鉴权手册.md) +- [Neo4j 模型](../表设计/04-neo4j-model.md) +- [业务记忆管理手册](../../业务记忆管理/业务记忆管理手册.md) diff --git a/docs/项目框架设计/技术选型和版本/01-技术栈与版本.md b/docs/项目框架设计/技术选型和版本/01-技术栈与版本.md index 79bff6c..af0628b 100644 --- a/docs/项目框架设计/技术选型和版本/01-技术栈与版本.md +++ b/docs/项目框架设计/技术选型和版本/01-技术栈与版本.md @@ -10,7 +10,7 @@ | 组件 | 版本 / 状态 | 备注 | | --- | --- | --- | -| **后端** | Python 3.13.14 + FastAPI | 系统 Python;已装 `langchain` 1.3.18 / `langchain-openai` 1.6.0 / `pymilvus` 3.0.1 | +| **后端** | Python 3.13.14 + FastAPI | 系统 Python;Agent 编排 **LangGraph** 1.2.x + `langchain-core` / `langchain-openai`(DeepSeek)/ `pymilvus` 3.0.1 | | **关系库** | MySQL 8.0.46 | 原生安装,端口 **3306** ✓ | | **缓存** | Redis 8.10.1 | 原生安装,路径 `F:\Redis\...`,端口 **6379** ✓ | | **图库** | Neo4j 5.26.19(Enterprise) | Neo4j Desktop 2,已建库 ✓ | @@ -106,7 +106,7 @@ Agent 业务库脚本:[01-mysql-共用底座.sql](../项目框架设计/表设 | 能力 | 选型 | 数据出境 | | --- | --- | --- | -| 对话 / 推理 / Tool 编排 | DeepSeek API | 按 API 协议;敏感字段需脱敏 | +| 对话 / 推理 / Tool 编排 | DeepSeek API + **LangGraph** StateGraph | 按 API 协议;敏感字段需脱敏 | | 文档 Embedding | Ollama bge-m3 本地 | **不出内网** | --- diff --git a/docs/项目框架设计/表设计/00-架构总览.md b/docs/项目框架设计/表设计/00-架构总览.md index 079752e..d6cd3e3 100644 --- a/docs/项目框架设计/表设计/00-架构总览.md +++ b/docs/项目框架设计/表设计/00-架构总览.md @@ -16,6 +16,7 @@ | **共用底座(先建)** | MySQL **11 张** + Redis 会话/画像 + Milvus 产品库 + Neo4j | [01-mysql-共用底座.sql](./01-mysql-共用底座.sql) | | **各 Agent 专用(后建)** | MySQL **5 张** | [02-mysql-agent专用.sql](./02-mysql-agent专用.sql) | | **技术栈与版本** | 组件版本、Windows 原生部署 | [01-技术栈与版本.md](../../技术选型和版本/01-技术栈与版本.md) | +| **业务记忆管理** | 短期/长期记忆、Redis vs SQL vs 图 vs 向量 | [业务记忆管理手册.md](../../业务记忆管理/业务记忆管理手册.md) | --- diff --git a/docs/项目框架设计/表设计/05-多Agent共用底座清单.md b/docs/项目框架设计/表设计/05-多Agent共用底座清单.md index f264879..ef4507f 100644 --- a/docs/项目框架设计/表设计/05-多Agent共用底座清单.md +++ b/docs/项目框架设计/表设计/05-多Agent共用底座清单.md @@ -212,4 +212,5 @@ - Redis / Milvus / Neo4j 字段细节 → `02`~`04` 号文档 - JWT + RBAC 统一鉴权 → [02-JWT-RBAC鉴权手册.md](../../技术选型和版本/02-JWT-RBAC鉴权手册.md) - 技术栈与版本 → [01-技术栈与版本.md](../../技术选型和版本/01-技术栈与版本.md) +- 业务记忆分层(Redis/MySQL/Milvus/Neo4j)→ [业务记忆管理手册.md](../../业务记忆管理/业务记忆管理手册.md) - 业务需求来源 → `docs/需求拆解/` diff --git a/requirements.txt b/requirements.txt index 477ffd6..83d473f 100644 --- a/requirements.txt +++ b/requirements.txt @@ -9,9 +9,10 @@ pymysql>=1.1.1 redis>=5.2.0 neo4j>=5.26.0 -# Vector & LLM +# Vector & LLM / Agent 编排(LangGraph) pymilvus>=3.0.1 -langchain>=1.3.18 +langgraph>=1.2.11 +langchain-core>=1.6.0 langchain-openai>=1.6.0 httpx>=0.28.0 diff --git a/scripts/core/00-create-database.sql b/scripts/core/00-create-database.sql new file mode 100644 index 0000000..bfe1bc2 --- /dev/null +++ b/scripts/core/00-create-database.sql @@ -0,0 +1,4 @@ +-- 模拟 Core 库(L0 权威事实,Agent 只读) +CREATE DATABASE IF NOT EXISTS jinrong_core + DEFAULT CHARACTER SET utf8mb4 + COLLATE utf8mb4_unicode_ci; diff --git a/scripts/core/01-ddl.sql b/scripts/core/01-ddl.sql new file mode 100644 index 0000000..6e47471 --- /dev/null +++ b/scripts/core/01-ddl.sql @@ -0,0 +1,129 @@ +-- jinrong_core 表结构 · P0 +USE jinrong_core; + +-- 风险等级字典(客户 C1~C5 · 产品 R1~R5) +CREATE TABLE core_risk_grade ( + code VARCHAR(4) NOT NULL PRIMARY KEY COMMENT 'C1~C5 或 R1~R5', + grade_type ENUM('customer','product') NOT NULL, + display_name VARCHAR(32) NOT NULL, + sort_order TINYINT UNSIGNED NOT NULL +) ENGINE=InnoDB COMMENT='风险等级字典'; + +CREATE TABLE core_industry ( + industry_code VARCHAR(16) NOT NULL PRIMARY KEY, + industry_name VARCHAR(64) NOT NULL +) ENGINE=InnoDB COMMENT='行业分类'; + +-- 内部员工(含 RBAC 角色种子,供 JWT/鉴权联调对照) +CREATE TABLE core_staff ( + staff_id VARCHAR(64) NOT NULL PRIMARY KEY, + display_name VARCHAR(64) NOT NULL, + staff_type ENUM('advisor','analyst','risk_officer','compliance','ops') NOT NULL, + roles JSON NOT NULL COMMENT 'JWT roles 数组,如 ["advisor"]', + tenant_id VARCHAR(32) NOT NULL DEFAULT 'TENANT-001', + is_active TINYINT(1) NOT NULL DEFAULT 1, + created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) +) ENGINE=InnoDB COMMENT='内部员工主档(模拟 IdP 账号源)'; + +CREATE TABLE core_customer ( + customer_id VARCHAR(64) NOT NULL PRIMARY KEY, + display_name VARCHAR(64) NOT NULL COMMENT '脱敏展示名', + age TINYINT UNSIGNED NULL, + occupation VARCHAR(64) NULL, + phone_mask VARCHAR(16) NULL COMMENT '138****9527', + tenant_id VARCHAR(32) NOT NULL DEFAULT 'TENANT-001', + open_date DATE NOT NULL, + is_active TINYINT(1) NOT NULL DEFAULT 1, + created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) +) ENGINE=InnoDB COMMENT='客户主档 L0'; + +CREATE TABLE core_customer_risk ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + customer_id VARCHAR(64) NOT NULL, + risk_code CHAR(2) NOT NULL COMMENT 'C1~C5', + is_authoritative TINYINT(1) NOT NULL DEFAULT 1, + evaluated_at DATE NOT NULL, + source VARCHAR(32) NOT NULL DEFAULT 'risk_questionnaire', + UNIQUE KEY uk_customer_current (customer_id), + KEY idx_risk (risk_code), + CONSTRAINT fk_cust_risk_customer FOREIGN KEY (customer_id) REFERENCES core_customer(customer_id), + CONSTRAINT fk_cust_risk_code FOREIGN KEY (risk_code) REFERENCES core_risk_grade(code) +) ENGINE=InnoDB COMMENT='客户正式风险测评 L0'; + +CREATE TABLE core_customer_advisor ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + customer_id VARCHAR(64) NOT NULL, + advisor_id VARCHAR(64) NOT NULL COMMENT '对应 core_staff.staff_id', + rel_status ENUM('active','transferred','closed') NOT NULL DEFAULT 'active', + effective_from DATE NOT NULL, + effective_to DATE NULL, + UNIQUE KEY uk_cust_advisor_from (customer_id, advisor_id, effective_from), + KEY idx_advisor (advisor_id, rel_status), + CONSTRAINT fk_ca_customer FOREIGN KEY (customer_id) REFERENCES core_customer(customer_id), + CONSTRAINT fk_ca_advisor FOREIGN KEY (advisor_id) REFERENCES core_staff(staff_id) +) ENGINE=InnoDB COMMENT='客户-代理人归属'; + +CREATE TABLE core_product ( + product_id VARCHAR(64) NOT NULL PRIMARY KEY, + product_name VARCHAR(128) NOT NULL, + product_type ENUM('bond','mixed','stock','index','money') NOT NULL, + min_risk_code CHAR(2) NOT NULL COMMENT 'R1~R5 最低适配', + industry_code VARCHAR(16) NULL, + fee_rate DECIMAL(6,4) NULL, + is_open TINYINT(1) NOT NULL DEFAULT 1, + created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3), + KEY idx_min_risk (min_risk_code), + CONSTRAINT fk_product_risk FOREIGN KEY (min_risk_code) REFERENCES core_risk_grade(code), + CONSTRAINT fk_product_industry FOREIGN KEY (industry_code) REFERENCES core_industry(industry_code) +) ENGINE=InnoDB COMMENT='产品主档'; + +CREATE TABLE core_holding ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + customer_id VARCHAR(64) NOT NULL, + product_id VARCHAR(64) NOT NULL, + qty DECIMAL(18,4) NOT NULL, + cost_amount DECIMAL(18,2) NOT NULL, + market_value DECIMAL(18,2) NOT NULL, + pnl_pct DECIMAL(8,4) NOT NULL COMMENT '盈亏比例', + as_of DATE NOT NULL, + UNIQUE KEY uk_cust_product (customer_id, product_id), + KEY idx_customer (customer_id), + CONSTRAINT fk_hold_customer FOREIGN KEY (customer_id) REFERENCES core_customer(customer_id), + CONSTRAINT fk_hold_product FOREIGN KEY (product_id) REFERENCES core_product(product_id) +) ENGINE=InnoDB COMMENT='持仓快照'; + +CREATE TABLE core_trade ( + trade_id VARCHAR(64) NOT NULL PRIMARY KEY, + customer_id VARCHAR(64) NOT NULL, + product_id VARCHAR(64) NOT NULL, + trade_type ENUM('subscribe','redeem','convert') NOT NULL, + amount DECIMAL(18,2) NOT NULL, + qty DECIMAL(18,4) NULL, + trade_status ENUM('confirmed','pending','cancelled') NOT NULL DEFAULT 'confirmed', + traded_at DATETIME(3) NOT NULL, + KEY idx_customer_time (customer_id, traded_at), + KEY idx_amount (amount, traded_at), + CONSTRAINT fk_trade_customer FOREIGN KEY (customer_id) REFERENCES core_customer(customer_id), + CONSTRAINT fk_trade_product FOREIGN KEY (product_id) REFERENCES core_product(product_id) +) ENGINE=InnoDB COMMENT='交易流水'; + +CREATE TABLE core_cash_flow ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + customer_id VARCHAR(64) NOT NULL, + flow_type ENUM('in','out') NOT NULL, + amount DECIMAL(18,2) NOT NULL, + remark VARCHAR(128) NULL, + occurred_at DATETIME(3) NOT NULL, + KEY idx_customer (customer_id, occurred_at), + CONSTRAINT fk_cf_customer FOREIGN KEY (customer_id) REFERENCES core_customer(customer_id) +) ENGINE=InnoDB COMMENT='资金进出'; + +CREATE TABLE core_product_nav ( + id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT PRIMARY KEY, + product_id VARCHAR(64) NOT NULL, + nav DECIMAL(10,4) NOT NULL, + daily_chg_pct DECIMAL(8,4) NOT NULL, + nav_date DATE NOT NULL, + UNIQUE KEY uk_product_date (product_id, nav_date), + CONSTRAINT fk_nav_product FOREIGN KEY (product_id) REFERENCES core_product(product_id) +) ENGINE=InnoDB COMMENT='产品净值'; diff --git a/scripts/core/02-seed-base.sql b/scripts/core/02-seed-base.sql new file mode 100644 index 0000000..7dbbea0 --- /dev/null +++ b/scripts/core/02-seed-base.sql @@ -0,0 +1,70 @@ +USE jinrong_core; + +-- 客户风险等级 C1~C5 +INSERT INTO core_risk_grade (code, grade_type, display_name, sort_order) VALUES +('C1', 'customer', '保守型 C1', 1), +('C2', 'customer', '稳健型 C2', 2), +('C3', 'customer', '平衡型 C3', 3), +('C4', 'customer', '成长型 C4', 4), +('C5', 'customer', '进取型 C5', 5); + +-- 产品风险等级 R1~R5 +INSERT INTO core_risk_grade (code, grade_type, display_name, sort_order) VALUES +('R1', 'product', '低风险 R1', 1), +('R2', 'product', '中低风险 R2', 2), +('R3', 'product', '中风险 R3', 3), +('R4', 'product', '中高风险 R4', 4), +('R5', 'product', '高风险 R5', 5); + +INSERT INTO core_industry (industry_code, industry_name) VALUES +('IND-BOND', '债券固收'), +('IND-MIX', '混合平衡'), +('IND-TECH', '科技成长'), +('IND-CONS', '消费'), +('IND-MED', '医药健康'), +('IND-INDEX', '宽基指数'), +('IND-MONEY', '货币现金'); + +-- ========== 内部员工 · RBAC 种子 ========== +-- 5 理财代理人 +INSERT INTO core_staff (staff_id, display_name, staff_type, roles) VALUES +('STAFF-10086', '张理财', 'advisor', '["advisor"]'), +('STAFF-10087', '李顾问', 'advisor', '["advisor"]'), +('STAFF-10088', '王经理', 'advisor', '["advisor"]'), +('STAFF-10089', '赵专员', 'advisor', '["advisor"]'), +('STAFF-10090', '刘助理', 'advisor', '["advisor"]'); + +-- 2 数据分析员 +INSERT INTO core_staff (staff_id, display_name, staff_type, roles) VALUES +('STAFF-20001', '陈分析', 'analyst', '["analyst"]'), +('STAFF-20002', '周数据', 'analyst', '["analyst"]'); + +-- 2 风控专员 +INSERT INTO core_staff (staff_id, display_name, staff_type, roles) VALUES +('STAFF-30001', '吴风控', 'risk_officer', '["risk_officer"]'), +('STAFF-30002', '郑监测', 'risk_officer', '["risk_officer"]'); + +-- 2 合规专员(兼 advisor 审计台场景用 compliance 单角色) +INSERT INTO core_staff (staff_id, display_name, staff_type, roles) VALUES +('STAFF-40001', '孙合规', 'compliance', '["compliance"]'), +('STAFF-40002', '钱监察', 'compliance', '["compliance"]'); + +-- 1 运营 + 1 双角色(代理人+合规,测权限并集与归属最窄) +INSERT INTO core_staff (staff_id, display_name, staff_type, roles) VALUES +('STAFF-50001', '冯运营', 'ops', '["ops"]'), +('STAFF-10091', '双角色顾问', 'advisor', '["advisor", "compliance"]'); + +-- ========== 产品 12 只 ========== +INSERT INTO core_product (product_id, product_name, product_type, min_risk_code, industry_code, fee_rate) VALUES +('PROD-110022', '稳健债基 A', 'bond', 'R1', 'IND-BOND', 0.0030), +('PROD-110023', '信用债精选', 'bond', 'R2', 'IND-BOND', 0.0040), +('PROD-005827', '平衡混合一号', 'mixed', 'R3', 'IND-MIX', 0.0150), +('PROD-005828', '稳健增利混合', 'mixed', 'R2', 'IND-MIX', 0.0120), +('PROD-161725', '科技成长主题', 'stock', 'R4', 'IND-TECH', 0.0150), +('PROD-161726', '消费升级主题', 'stock', 'R4', 'IND-CONS', 0.0150), +('PROD-003095', '医药健康精选', 'stock', 'R4', 'IND-MED', 0.0120), +('PROD-510300', '沪深300指数', 'index', 'R3', 'IND-INDEX', 0.0050), +('PROD-510500', '中证500指数', 'index', 'R4', 'IND-INDEX', 0.0060), +('PROD-XYZ999', '进取成长五号', 'stock', 'R5', 'IND-TECH', 0.0180), +('PROD-000001', '现金宝货币', 'money', 'R1', 'IND-MONEY', 0.0010), +('PROD-000002', '同业存单基金', 'bond', 'R1', 'IND-BOND', 0.0020); diff --git a/scripts/core/03-seed-customers.sql b/scripts/core/03-seed-customers.sql new file mode 100644 index 0000000..94319c3 --- /dev/null +++ b/scripts/core/03-seed-customers.sql @@ -0,0 +1,107 @@ +USE jinrong_core; + +-- 28 位客户(脱敏展示名) +INSERT INTO core_customer (customer_id, display_name, age, occupation, phone_mask, open_date) VALUES +('CUST-9527', '客户·林**', 35, '工程师', '138****9527', '2020-03-15'), +('CUST-1001', '客户·王**', 28, '职员', '139****1001', '2021-06-01'), +('CUST-1002', '客户·赵**', 42, '教师', '137****1002', '2019-11-20'), +('CUST-1003', '客户·孙**', 72, '退休', '136****1003', '2018-05-10'), +('CUST-1004', '客户·周**', 55, '企业主', '135****1004', '2017-08-22'), +('CUST-1005', '客户·吴**', 31, '设计师', '133****1005', '2022-01-08'), +('CUST-1006', '客户·郑**', 48, '会计', '132****1006', '2016-12-30'), +('CUST-1007', '客户·冯**', 39, '销售', '131****1007', '2020-09-14'), +('CUST-1008', '客户·陈**', 26, '程序员', '130****1008', '2023-04-02'), +('CUST-1009', '客户·褚**', 61, '退休干部', '189****1009', '2015-07-18'), +('CUST-1010', '客户·卫**', 44, '医生', '188****1010', '2019-02-28'), +('CUST-1011', '客户·蒋**', 33, '律师', '187****1011', '2021-10-11'), +('CUST-1012', '客户·沈**', 50, '经理', '186****1012', '2018-03-03'), +('CUST-1013', '客户·韩**', 29, '职员', '185****1013', '2022-08-19'), +('CUST-1014', '客户·杨**', 37, '研究员', '184****1014', '2020-12-01'), +('CUST-1015', '客户·朱**', 58, '个体户', '183****1015', '2016-06-06'), +('CUST-1016', '客户·秦**', 45, '工程师', '182****1016', '2019-09-09'), +('CUST-1017', '客户·尤**', 32, '产品经理', '181****1017', '2021-03-21'), +('CUST-1018', '客户·许**', 67, '退休', '180****1018', '2014-11-11'), +('CUST-1019', '客户·何**', 41, '公务员', '177****1019', '2018-08-08'), +('CUST-1020', '客户·吕**', 36, '金融从业', '176****1020', '2020-05-05'), +('CUST-1021', '客户·施**', 27, '研究生', '175****1021', '2023-01-15'), +('CUST-1022', '客户·张**', 53, '主管', '174****1022', '2017-04-04'), +('CUST-1023', '客户·孔**', 38, '运营', '173****1023', '2019-07-07'), +('CUST-1024', '客户·曹**', 46, '顾问', '172****1024', '2018-10-10'), +('CUST-3001', '客户·大额**', 49, '企业主', '171****3001', '2016-01-01'), +('CUST-4001', '客户·高龄进取**', 70, '退休', '170****4001', '2015-05-05'), +('CUST-4002', '客户·冲突测**', 71, '退休', '169****4002', '2015-06-06'); + +-- 正式风险测评 L0 +INSERT INTO core_customer_risk (customer_id, risk_code, evaluated_at) VALUES +('CUST-9527', 'C3', '2025-06-01'), +('CUST-1001', 'C1', '2025-03-15'), +('CUST-1002', 'C2', '2025-04-20'), +('CUST-1003', 'C2', '2024-12-01'), +('CUST-1004', 'C4', '2025-01-10'), +('CUST-1005', 'C3', '2025-05-05'), +('CUST-1006', 'C3', '2024-08-18'), +('CUST-1007', 'C2', '2025-02-14'), +('CUST-1008', 'C3', '2025-07-01'), +('CUST-1009', 'C1', '2024-06-30'), +('CUST-1010', 'C3', '2025-03-03'), +('CUST-1011', 'C4', '2025-04-04'), +('CUST-1012', 'C3', '2024-11-11'), +('CUST-1013', 'C2', '2025-08-08'), +('CUST-1014', 'C4', '2025-01-20'), +('CUST-1015', 'C2', '2024-09-09'), +('CUST-1016', 'C3', '2025-02-02'), +('CUST-1017', 'C3', '2025-06-06'), +('CUST-1018', 'C1', '2024-03-03'), +('CUST-1019', 'C2', '2025-05-15'), +('CUST-1020', 'C4', '2025-03-30'), +('CUST-1021', 'C3', '2025-09-01'), +('CUST-1022', 'C3', '2024-07-07'), +('CUST-1023', 'C2', '2025-04-18'), +('CUST-1024', 'C4', '2025-02-28'), +('CUST-3001', 'C3', '2025-01-01'), +('CUST-4001', 'C5', '2025-06-15'), +('CUST-4002', 'C1', '2025-06-20'); + +-- 客户-代理人归属(RBAC:代理人只能查名下) +-- STAFF-10086: 10 人 +INSERT INTO core_customer_advisor (customer_id, advisor_id, effective_from) VALUES +('CUST-9527', 'STAFF-10086', '2020-03-15'), +('CUST-1001', 'STAFF-10086', '2021-06-01'), +('CUST-1002', 'STAFF-10086', '2019-11-20'), +('CUST-1003', 'STAFF-10086', '2018-05-10'), +('CUST-1005', 'STAFF-10086', '2022-01-08'), +('CUST-1007', 'STAFF-10086', '2020-09-14'), +('CUST-1008', 'STAFF-10086', '2023-04-02'), +('CUST-1013', 'STAFF-10086', '2022-08-19'), +('CUST-3001', 'STAFF-10086', '2016-01-01'), +('CUST-4002', 'STAFF-10086', '2015-06-06'); + +-- STAFF-10087: 6 人 +INSERT INTO core_customer_advisor (customer_id, advisor_id, effective_from) VALUES +('CUST-1004', 'STAFF-10087', '2017-08-22'), +('CUST-1010', 'STAFF-10087', '2019-02-28'), +('CUST-1011', 'STAFF-10087', '2021-10-11'), +('CUST-1014', 'STAFF-10087', '2025-01-20'), +('CUST-1020', 'STAFF-10087', '2020-05-05'), +('CUST-4001', 'STAFF-10087', '2015-05-05'); + +-- STAFF-10088: 5 人 +INSERT INTO core_customer_advisor (customer_id, advisor_id, effective_from) VALUES +('CUST-1006', 'STAFF-10088', '2016-12-30'), +('CUST-1009', 'STAFF-10088', '2015-07-18'), +('CUST-1012', 'STAFF-10088', '2018-03-03'), +('CUST-1015', 'STAFF-10088', '2016-06-06'), +('CUST-1018', 'STAFF-10088', '2014-11-11'); + +-- STAFF-10089: 4 人 +INSERT INTO core_customer_advisor (customer_id, advisor_id, effective_from) VALUES +('CUST-1016', 'STAFF-10089', '2019-09-09'), +('CUST-1017', 'STAFF-10089', '2021-03-21'), +('CUST-1019', 'STAFF-10089', '2018-08-08'), +('CUST-1023', 'STAFF-10089', '2019-07-07'); + +-- STAFF-10090: 3 人 +INSERT INTO core_customer_advisor (customer_id, advisor_id, effective_from) VALUES +('CUST-1021', 'STAFF-10090', '2023-01-15'), +('CUST-1022', 'STAFF-10090', '2017-04-04'), +('CUST-1024', 'STAFF-10090', '2018-10-10'); diff --git a/scripts/core/04-seed-holdings.sql b/scripts/core/04-seed-holdings.sql new file mode 100644 index 0000000..70c41d6 --- /dev/null +++ b/scripts/core/04-seed-holdings.sql @@ -0,0 +1,79 @@ +USE jinrong_core; + +-- 持仓 as_of = 2026-09-04 +INSERT INTO core_holding (customer_id, product_id, qty, cost_amount, market_value, pnl_pct, as_of) VALUES +-- CUST-9527 主 demo +('CUST-9527', 'PROD-005827', 50000.0000, 50000.00, 47500.00, -5.0000, '2026-09-04'), +('CUST-9527', 'PROD-110022', 80000.0000, 80000.00, 82400.00, 3.0000, '2026-09-04'), +('CUST-9527', 'PROD-000001', 20000.0000, 20000.00, 20040.00, 0.2000, '2026-09-04'), +-- CUST-1001 C1 保守 +('CUST-1001', 'PROD-110022', 30000.0000, 30000.00, 30900.00, 3.0000, '2026-09-04'), +('CUST-1001', 'PROD-000001', 50000.0000, 50000.00, 50100.00, 0.2000, '2026-09-04'), +-- CUST-1002 接近阈值 -12% +('CUST-1002', 'PROD-005828', 40000.0000, 40000.00, 35200.00, -12.0000, '2026-09-04'), +('CUST-1002', 'PROD-110023', 25000.0000, 25000.00, 24500.00, -2.0000, '2026-09-04'), +-- CUST-1003 高龄 C2 +('CUST-1003', 'PROD-110022', 120000.0000, 120000.00, 123600.00, 3.0000, '2026-09-04'), +('CUST-1003', 'PROD-005828', 30000.0000, 30000.00, 29100.00, -3.0000, '2026-09-04'), +-- CUST-1004 ~10087 +('CUST-1004', 'PROD-161725', 60000.0000, 60000.00, 55800.00, -7.0000, '2026-09-04'), +('CUST-1004', 'PROD-510300', 40000.0000, 40000.00, 41200.00, 3.0000, '2026-09-04'), +-- CUST-1005 +('CUST-1005', 'PROD-005827', 35000.0000, 35000.00, 36750.00, 5.0000, '2026-09-04'), +('CUST-1005', 'PROD-510300', 15000.0000, 15000.00, 15300.00, 2.0000, '2026-09-04'), +-- CUST-1006 +('CUST-1006', 'PROD-110023', 90000.0000, 90000.00, 91800.00, 2.0000, '2026-09-04'), +('CUST-1006', 'PROD-000002', 40000.0000, 40000.00, 40200.00, 0.5000, '2026-09-04'), +-- CUST-1007 +('CUST-1007', 'PROD-005828', 45000.0000, 45000.00, 43650.00, -3.0000, '2026-09-04'), +-- CUST-1008 +('CUST-1008', 'PROD-161726', 20000.0000, 20000.00, 19000.00, -5.0000, '2026-09-04'), +('CUST-1008', 'PROD-510500', 10000.0000, 10000.00, 9800.00, -2.0000, '2026-09-04'), +-- CUST-1009 C1 +('CUST-1009', 'PROD-000001', 150000.0000, 150000.00, 150300.00, 0.2000, '2026-09-04'), +('CUST-1009', 'PROD-110022', 80000.0000, 80000.00, 82400.00, 3.0000, '2026-09-04'), +-- CUST-1010 ~10087 跨权限测 +('CUST-1010', 'PROD-003095', 55000.0000, 55000.00, 52250.00, -5.0000, '2026-09-04'), +('CUST-1010', 'PROD-005827', 30000.0000, 30000.00, 31500.00, 5.0000, '2026-09-04'), +-- CUST-1011 +('CUST-1011', 'PROD-161725', 70000.0000, 70000.00, 73500.00, 5.0000, '2026-09-04'), +-- CUST-1012 +('CUST-1012', 'PROD-510300', 65000.0000, 65000.00, 66950.00, 3.0000, '2026-09-04'), +('CUST-1012', 'PROD-110023', 35000.0000, 35000.00, 35700.00, 2.0000, '2026-09-04'), +-- CUST-1013 +('CUST-1013', 'PROD-005828', 28000.0000, 28000.00, 27440.00, -2.0000, '2026-09-04'), +-- CUST-1014 +('CUST-1014', 'PROD-510500', 90000.0000, 90000.00, 87300.00, -3.0000, '2026-09-04'), +('CUST-1014', 'PROD-161726', 40000.0000, 40000.00, 42000.00, 5.0000, '2026-09-04'), +-- CUST-1015 +('CUST-1015', 'PROD-110022', 100000.0000, 100000.00, 103000.00, 3.0000, '2026-09-04'), +-- CUST-1016 +('CUST-1016', 'PROD-005827', 42000.0000, 42000.00, 43260.00, 3.0000, '2026-09-04'), +('CUST-1016', 'PROD-003095', 18000.0000, 18000.00, 17100.00, -5.0000, '2026-09-04'), +-- CUST-1017 +('CUST-1017', 'PROD-510300', 32000.0000, 32000.00, 32960.00, 3.0000, '2026-09-04'), +-- CUST-1018 C1 +('CUST-1018', 'PROD-000001', 200000.0000, 200000.00, 200400.00, 0.2000, '2026-09-04'), +-- CUST-1019 +('CUST-1019', 'PROD-110023', 48000.0000, 48000.00, 47040.00, -2.0000, '2026-09-04'), +('CUST-1019', 'PROD-005828', 22000.0000, 22000.00, 22440.00, 2.0000, '2026-09-04'), +-- CUST-1020 +('CUST-1020', 'PROD-161725', 85000.0000, 85000.00, 89250.00, 5.0000, '2026-09-04'), +('CUST-1020', 'PROD-XYZ999', 15000.0000, 15000.00, 13500.00, -10.0000, '2026-09-04'), +-- CUST-1021 +('CUST-1021', 'PROD-005827', 15000.0000, 15000.00, 15750.00, 5.0000, '2026-09-04'), +-- CUST-1022 +('CUST-1022', 'PROD-510300', 72000.0000, 72000.00, 74160.00, 3.0000, '2026-09-04'), +('CUST-1022', 'PROD-110022', 48000.0000, 48000.00, 49440.00, 3.0000, '2026-09-04'), +-- CUST-1023 +('CUST-1023', 'PROD-005828', 38000.0000, 38000.00, 37160.00, -2.2000, '2026-09-04'), +-- CUST-1024 +('CUST-1024', 'PROD-161726', 95000.0000, 95000.00, 99750.00, 5.0000, '2026-09-04'), +-- CUST-3001 大额客户 +('CUST-3001', 'PROD-510300', 200000.0000, 200000.00, 206000.00, 3.0000, '2026-09-04'), +('CUST-3001', 'PROD-005827', 150000.0000, 150000.00, 157500.00, 5.0000, '2026-09-04'), +-- CUST-4001 高龄+C5 +('CUST-4001', 'PROD-XYZ999', 80000.0000, 80000.00, 72000.00, -10.0000, '2026-09-04'), +('CUST-4001', 'PROD-161725', 50000.0000, 50000.00, 52500.00, 5.0000, '2026-09-04'), +-- CUST-4002 C1 却持有 R4(适当性冲突测) +('CUST-4002', 'PROD-161725', 10000.0000, 10000.00, 9500.00, -5.0000, '2026-09-04'); diff --git a/scripts/core/05-seed-trades.sql b/scripts/core/05-seed-trades.sql new file mode 100644 index 0000000..1274307 --- /dev/null +++ b/scripts/core/05-seed-trades.sql @@ -0,0 +1,41 @@ +USE jinrong_core; + +-- 交易流水(含 R-01 大额、R-04 频繁样例) +INSERT INTO core_trade (trade_id, customer_id, product_id, trade_type, amount, qty, traded_at) VALUES +('TRD-20260801-001', 'CUST-9527', 'PROD-005827', 'subscribe', 10000.00, 10000.0000, '2026-08-01 10:15:00.000'), +('TRD-20260815-001', 'CUST-9527', 'PROD-110022', 'subscribe', 20000.00, 20000.0000, '2026-08-15 14:20:00.000'), +('TRD-20260901-001', 'CUST-1001', 'PROD-110022', 'subscribe', 5000.00, 5000.0000, '2026-09-01 09:30:00.000'), +('TRD-20260902-001', 'CUST-1002', 'PROD-005828', 'redeem', 8000.00, 8000.0000, '2026-09-02 11:00:00.000'), +-- CUST-3001 大额 R-01 +('TRD-20260903-001', 'CUST-3001', 'PROD-510300', 'subscribe', 520000.00, 520000.0000, '2026-09-03 10:00:00.000'), +('TRD-20260903-002', 'CUST-3001', 'PROD-005827', 'subscribe', 80000.00, 80000.0000, '2026-09-03 15:30:00.000'), +-- 同日多笔接近阈值(R-04 P1 样例) +('TRD-20260904-001', 'CUST-3001', 'PROD-510300', 'subscribe', 450000.00, 450000.0000, '2026-09-04 09:05:00.000'), +('TRD-20260904-002', 'CUST-3001', 'PROD-510300', 'subscribe', 420000.00, 420000.0000, '2026-09-04 09:12:00.000'), +('TRD-20260904-003', 'CUST-3001', 'PROD-510300', 'subscribe', 410000.00, 410000.0000, '2026-09-04 09:18:00.000'), +-- 分散到其他客户 +('TRD-20260901-010', 'CUST-1010', 'PROD-003095', 'subscribe', 30000.00, 30000.0000, '2026-09-01 13:00:00.000'), +('TRD-20260901-011', 'CUST-1014', 'PROD-510500', 'subscribe', 55000.00, 55000.0000, '2026-09-01 14:00:00.000'), +('TRD-20260902-010', 'CUST-1020', 'PROD-161725', 'subscribe', 120000.00, 120000.0000, '2026-09-02 10:30:00.000'), +('TRD-20260902-011', 'CUST-1004', 'PROD-161725', 'redeem', 25000.00, 25000.0000, '2026-09-02 16:00:00.000'), +('TRD-20260903-010', 'CUST-1018', 'PROD-000001', 'subscribe', 100000.00, 100000.0000, '2026-09-03 11:00:00.000'), +('TRD-20260904-010', 'CUST-4001', 'PROD-XYZ999', 'subscribe', 200000.00, 200000.0000, '2026-09-04 10:00:00.000'), +('TRD-20260904-011', 'CUST-4002', 'PROD-161725', 'subscribe', 50000.00, 50000.0000, '2026-09-04 11:30:00.000'), +('TRD-20260820-001', 'CUST-1006', 'PROD-110023', 'subscribe', 18000.00, 18000.0000, '2026-08-20 09:00:00.000'), +('TRD-20260825-001', 'CUST-1012', 'PROD-510300', 'subscribe', 22000.00, 22000.0000, '2026-08-25 10:00:00.000'), +('TRD-20260828-001', 'CUST-1022', 'PROD-110022', 'subscribe', 15000.00, 15000.0000, '2026-08-28 14:00:00.000'), +('TRD-20260901-020', 'CUST-1021', 'PROD-005827', 'subscribe', 8000.00, 8000.0000, '2026-09-01 15:00:00.000'), +('TRD-20260902-020', 'CUST-1016', 'PROD-005827', 'subscribe', 12000.00, 12000.0000, '2026-09-02 09:45:00.000'), +('TRD-20260903-020', 'CUST-1017', 'PROD-510300', 'subscribe', 16000.00, 16000.0000, '2026-09-03 13:20:00.000'), +('TRD-20260904-020', 'CUST-1019', 'PROD-110023', 'redeem', 6000.00, 6000.0000, '2026-09-04 09:00:00.000'), +('TRD-20260904-021', 'CUST-1023', 'PROD-005828', 'subscribe', 9000.00, 9000.0000, '2026-09-04 14:30:00.000'), +('TRD-20260904-022', 'CUST-1024', 'PROD-161726', 'subscribe', 35000.00, 35000.0000, '2026-09-04 15:00:00.000'); + +INSERT INTO core_cash_flow (customer_id, flow_type, amount, remark, occurred_at) VALUES +('CUST-9527', 'in', 50000.00, '银行转入', '2026-08-01 09:00:00.000'), +('CUST-9527', 'out', 10000.00, '申购扣款', '2026-08-01 10:15:00.000'), +('CUST-3001', 'in', 600000.00, '大额转入', '2026-09-03 09:30:00.000'), +('CUST-1010', 'in', 30000.00, '工资结余', '2026-09-01 08:00:00.000'), +('CUST-1018', 'in', 100000.00, '退休年金', '2026-09-03 10:30:00.000'), +('CUST-1002', 'out', 8000.00, '赎回到账', '2026-09-02 11:30:00.000'), +('CUST-1020', 'in', 120000.00, '经营收入', '2026-09-02 10:00:00.000'); diff --git a/scripts/core/06-seed-nav.sql b/scripts/core/06-seed-nav.sql new file mode 100644 index 0000000..439bd5c --- /dev/null +++ b/scripts/core/06-seed-nav.sql @@ -0,0 +1,16 @@ +USE jinrong_core; + +-- 最新净值 nav_date = 2026-09-04 +INSERT INTO core_product_nav (product_id, nav, daily_chg_pct, nav_date) VALUES +('PROD-110022', 1.0300, 0.1500, '2026-09-04'), +('PROD-110023', 1.0200, 0.1000, '2026-09-04'), +('PROD-005827', 0.9500, -0.8000, '2026-09-04'), +('PROD-005828', 1.0200, 0.2000, '2026-09-04'), +('PROD-161725', 1.0500, 1.2000, '2026-09-04'), +('PROD-161726', 1.0500, 0.9000, '2026-09-04'), +('PROD-003095', 0.9500, -1.1000, '2026-09-04'), +('PROD-510300', 1.0300, 0.3500, '2026-09-04'), +('PROD-510500', 0.9800, -0.5000, '2026-09-04'), +('PROD-XYZ999', 0.9000, -2.0000, '2026-09-04'), +('PROD-000001', 1.0020, 0.0100, '2026-09-04'), +('PROD-000002', 1.0050, 0.0200, '2026-09-04'); diff --git a/scripts/core/README.md b/scripts/core/README.md new file mode 100644 index 0000000..28a2192 --- /dev/null +++ b/scripts/core/README.md @@ -0,0 +1,43 @@ +# Core 模拟库 · 初始化 + +> 库名 `jinrong_core` · 与 `jinrong_agent` 同 MySQL 实例 · Agent **只读** + +## 一键重置(Windows) + +```powershell +.\scripts\core\reset.ps1 +``` + +## 手动顺序 + +```powershell +mysql -u root -p < scripts/core/00-create-database.sql +mysql -u root -p < scripts/core/01-ddl.sql +mysql -u root -p < scripts/core/02-seed-base.sql +mysql -u root -p < scripts/core/03-seed-customers.sql +mysql -u root -p < scripts/core/04-seed-holdings.sql +mysql -u root -p < scripts/core/05-seed-trades.sql +mysql -u root -p < scripts/core/06-seed-nav.sql + +# Agent 库(若未建) +mysql -u root -p < docs/项目框架设计/表设计/01-mysql-共用底座.sql + +python scripts/sync/sync_advisor_rel.py +python scripts/sync/sync_neo4j.py +``` + +## 种子规模(当前) + +| 类型 | 数量 | +| --- | --- | +| 客户 | 28 | +| 理财代理人 | 5 | +| 分析/风控/合规/运营 | 2+2+2+1 | +| 双角色 staff | 1(STAFF-10091 advisor+compliance) | +| 产品 | 12 | +| 持仓记录 | ~45 | +| 交易 | 25+ | + +## RBAC 联调账号见 + +[rbac-seed-reference.md](../dev/rbac-seed-reference.md) diff --git a/scripts/core/reset.ps1 b/scripts/core/reset.ps1 new file mode 100644 index 0000000..993dfb8 --- /dev/null +++ b/scripts/core/reset.ps1 @@ -0,0 +1,43 @@ +#Requires -Version 5.1 +<# +.SYNOPSIS + 重置 jinrong_core 模拟库并同步到 agent / Neo4j +#> +param( + [string]$MysqlUser = "root", + [string]$MysqlHost = "127.0.0.1", + [switch]$SkipNeo4j +) + +$ErrorActionPreference = "Stop" +$Root = Split-Path (Split-Path $PSScriptRoot -Parent) -Parent +Set-Location $Root + +function Invoke-SqlFile($path) { + Write-Host ">> $path" + Get-Content -LiteralPath $path -Encoding UTF8 | mysql -h $MysqlHost -u $MysqlUser -p +} + +Write-Host "=== Drop & recreate jinrong_core ===" +mysql -h $MysqlHost -u $MysqlUser -p -e "DROP DATABASE IF EXISTS jinrong_core;" + +$files = @( + "scripts/core/00-create-database.sql", + "scripts/core/01-ddl.sql", + "scripts/core/02-seed-base.sql", + "scripts/core/03-seed-customers.sql", + "scripts/core/04-seed-holdings.sql", + "scripts/core/05-seed-trades.sql", + "scripts/core/06-seed-nav.sql" +) +foreach ($f in $files) { Invoke-SqlFile (Join-Path $Root $f) } + +Write-Host "=== Sync advisor rel ===" +python scripts/sync/sync_advisor_rel.py + +if (-not $SkipNeo4j) { + Write-Host "=== Sync Neo4j ===" + python scripts/sync/sync_neo4j.py +} + +Write-Host "Done." diff --git a/scripts/dev/rbac-seed-reference.md b/scripts/dev/rbac-seed-reference.md new file mode 100644 index 0000000..c8a69c7 --- /dev/null +++ b/scripts/dev/rbac-seed-reference.md @@ -0,0 +1,59 @@ +# RBAC 联调种子对照 + +> 开发期 JWT 可手工签发;`sub` / `roles` 与下表对齐。 +> 客户 Token:`sub` = `customer_id`。员工 Token:`sub` = `staff_id`。 + +## 内部员工(core_staff) + +| staff_id | 角色 roles | 用途 | +| --- | --- | --- | +| STAFF-10086 | advisor | 主代理人,名下 10 客户 | +| STAFF-10087 | advisor | 跨权限测:1010 仅其名下 | +| STAFF-10088 | advisor | 5 客户 | +| STAFF-10089 | advisor | 4 客户 | +| STAFF-10090 | advisor | 3 客户 | +| STAFF-10091 | advisor, compliance | 双角色权限并集 | +| STAFF-20001 | analyst | 分析 Agent | +| STAFF-20002 | analyst | 分析 Agent | +| STAFF-30001 | risk_officer | 风控 Agent | +| STAFF-30002 | risk_officer | 风控 Agent | +| STAFF-40001 | compliance | 合规审计台 | +| STAFF-40002 | compliance | 合规审计台 | +| STAFF-50001 | ops | 运营统计 A-08 | + +## 推荐验收用例 + +| 操作者 | 目标 customer_id | 预期 | +| --- | --- | --- | +| STAFF-10086 | CUST-9527 | 允许 | +| STAFF-10086 | CUST-1010 | **403**(归属 10087) | +| CUST-9527 | CUST-1001 数据 | **403**(非本人) | +| STAFF-20001 | 聚合 SQL | 允许(脱敏/聚合) | +| STAFF-30001 | 任意客户读 | 允许 | +| STAFF-40001 | 审计 API | 允许;写 L2 **403** | + +## 客户(登录 C 端) + +| customer_id | 正式等级 | 代理人 | 特殊场景 | +| --- | --- | --- | --- | +| CUST-9527 | C3 | 10086 | 主 demo | +| CUST-1001 | C1 | 10086 | R-02 买 PROD-XYZ999 | +| CUST-1002 | C2 | 10086 | C-04 盈亏 -12% | +| CUST-1010 | C3 | **10087** | 10086 越权测 | +| CUST-3001 | C3 | 10086 | R-01 大额交易 | +| CUST-4001 | C5 | 10087 | 高龄 70 + R5 持仓 | +| CUST-4002 | C1 | 10086 | 持有 R4 产品冲突 | + +## 产品适当性 + +| product_id | min_risk | 说明 | +| --- | --- | --- | +| PROD-110022 | R1 | C1 可买 | +| PROD-XYZ999 | R5 | C1/C2 应拒 | + +查询员工角色: + +```python +from app.repository.core_ro import CoreReadOnlyRepository +CoreReadOnlyRepository().get_staff("STAFF-10086") +``` diff --git a/scripts/sync/sync_advisor_rel.py b/scripts/sync/sync_advisor_rel.py new file mode 100644 index 0000000..47aba29 --- /dev/null +++ b/scripts/sync/sync_advisor_rel.py @@ -0,0 +1,62 @@ +#!/usr/bin/env python3 +"""Core → jinrong_agent.customer_advisor_rel 同步。""" + +from __future__ import annotations + +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from sqlalchemy import create_engine, text + +from app.config.settings import settings + + +def db_url(database: str) -> str: + pwd = settings.mysql_password + auth = f"{settings.mysql_user}:{pwd}" if pwd else settings.mysql_user + return ( + f"mysql+pymysql://{auth}@{settings.mysql_host}:{settings.mysql_port}" + f"/{database}?charset=utf8mb4" + ) + + +def main() -> None: + core = create_engine(db_url(settings.mysql_core_database), pool_pre_ping=True) + agent = create_engine(db_url(settings.mysql_database), pool_pre_ping=True) + + with core.connect() as cconn: + rows = cconn.execute( + text( + """ + SELECT customer_id, advisor_id, rel_status, effective_from, effective_to + FROM core_customer_advisor + WHERE rel_status = 'active' + """ + ) + ).mappings().all() + + upsert_sql = text( + """ + INSERT INTO customer_advisor_rel + (customer_id, advisor_id, rel_status, effective_from, effective_to, synced_at) + VALUES + (:customer_id, :advisor_id, :rel_status, :effective_from, :effective_to, NOW(3)) + ON DUPLICATE KEY UPDATE + rel_status = VALUES(rel_status), + effective_to = VALUES(effective_to), + synced_at = NOW(3) + """ + ) + + with agent.begin() as aconn: + for row in rows: + aconn.execute(upsert_sql, dict(row)) + + print(f"sync_advisor_rel: upserted {len(rows)} rows → {settings.mysql_database}.customer_advisor_rel") + + +if __name__ == "__main__": + main() diff --git a/scripts/sync/sync_neo4j.py b/scripts/sync/sync_neo4j.py new file mode 100644 index 0000000..dad3124 --- /dev/null +++ b/scripts/sync/sync_neo4j.py @@ -0,0 +1,171 @@ +#!/usr/bin/env python3 +"""Core → Neo4j 关系图同步(P0 节点与关系)。""" + +from __future__ import annotations + +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from neo4j import GraphDatabase +from sqlalchemy import create_engine, text + +from app.config.settings import settings + + +def db_url(database: str) -> str: + pwd = settings.mysql_password + auth = f"{settings.mysql_user}:{pwd}" if pwd else settings.mysql_user + return ( + f"mysql+pymysql://{auth}@{settings.mysql_host}:{settings.mysql_port}" + f"/{database}?charset=utf8mb4" + ) + + +def main() -> None: + if not settings.neo4j_password: + print("WARN: NEO4J_PASSWORD empty, skip neo4j sync") + return + + engine = create_engine(db_url(settings.mysql_core_database), pool_pre_ping=True) + driver = GraphDatabase.driver( + settings.neo4j_uri, + auth=(settings.neo4j_user, settings.neo4j_password), + ) + + with engine.connect() as conn: + customers = conn.execute( + text("SELECT customer_id, display_name FROM core_customer WHERE is_active=1") + ).mappings().all() + advisors = conn.execute( + text( + "SELECT staff_id, display_name FROM core_staff WHERE staff_type='advisor' AND is_active=1" + ) + ).mappings().all() + products = conn.execute( + text("SELECT product_id, product_name, min_risk_code, industry_code FROM core_product") + ).mappings().all() + grades = conn.execute( + text("SELECT code, display_name, grade_type FROM core_risk_grade") + ).mappings().all() + industries = conn.execute( + text("SELECT industry_code, industry_name FROM core_industry") + ).mappings().all() + assignments = conn.execute( + text( + """ + SELECT customer_id, advisor_id, effective_from, rel_status + FROM core_customer_advisor WHERE rel_status='active' + """ + ) + ).mappings().all() + cust_risks = conn.execute( + text("SELECT customer_id, risk_code, evaluated_at FROM core_customer_risk") + ).mappings().all() + holdings = conn.execute( + text( + """ + SELECT customer_id, product_id, qty, market_value, cost_amount, pnl_pct, as_of + FROM core_holding + """ + ) + ).mappings().all() + + with driver.session() as session: + session.run("CREATE CONSTRAINT customer_id IF NOT EXISTS FOR (c:Customer) REQUIRE c.customer_id IS UNIQUE") + session.run("CREATE CONSTRAINT advisor_id IF NOT EXISTS FOR (a:Advisor) REQUIRE a.advisor_id IS UNIQUE") + session.run("CREATE CONSTRAINT product_id IF NOT EXISTS FOR (p:Product) REQUIRE p.product_id IS UNIQUE") + session.run("CREATE CONSTRAINT risk_code IF NOT EXISTS FOR (r:RiskGrade) REQUIRE r.code IS UNIQUE") + session.run("CREATE CONSTRAINT industry_code IF NOT EXISTS FOR (i:Industry) REQUIRE i.industry_code IS UNIQUE") + + for g in grades: + session.run( + "MERGE (r:RiskGrade {code: $code}) SET r.display_name=$display_name, r.grade_type=$grade_type", + **g, + ) + for ind in industries: + session.run( + "MERGE (i:Industry {industry_code: $industry_code}) SET i.industry_name=$industry_name", + **ind, + ) + for c in customers: + session.run( + "MERGE (c:Customer {customer_id: $customer_id}) SET c.display_name=$display_name", + **c, + ) + for a in advisors: + session.run( + "MERGE (a:Advisor {advisor_id: $staff_id}) SET a.display_name=$display_name", + staff_id=a["staff_id"], + display_name=a["display_name"], + ) + for p in products: + session.run( + """ + MERGE (p:Product {product_id: $product_id}) + SET p.product_name=$product_name, p.min_risk_code=$min_risk_code + WITH p + OPTIONAL MATCH (i:Industry {industry_code: $industry_code}) + FOREACH (_ IN CASE WHEN i IS NOT NULL THEN [1] ELSE [] END | + MERGE (p)-[:BELONGS_TO]->(i) + ) + """, + **p, + ) + session.run( + """ + MATCH (p:Product {product_id: $product_id}), (r:RiskGrade {code: $min_risk_code}) + MERGE (p)-[:REQUIRES_MIN_RISK {rule_id: 'R02-PROD-MIN'}]->(r) + """, + product_id=p["product_id"], + min_risk_code=p["min_risk_code"], + ) + for row in assignments: + session.run( + """ + MATCH (c:Customer {customer_id: $customer_id}), (a:Advisor {advisor_id: $advisor_id}) + MERGE (c)-[r:ASSIGNED_TO]->(a) + SET r.since=$effective_from, r.status=$rel_status + """, + **row, + ) + for row in cust_risks: + session.run( + """ + MATCH (c:Customer {customer_id: $customer_id}), (r:RiskGrade {code: $risk_code}) + MERGE (c)-[hr:HAS_RISK_LEVEL]->(r) + SET hr.source='l0', hr.evaluated_at=$evaluated_at + """, + customer_id=row["customer_id"], + risk_code=row["risk_code"], + evaluated_at=str(row["evaluated_at"]), + ) + session.run("MATCH ()-[h:HOLDS]->() DELETE h") + for h in holdings: + session.run( + """ + MATCH (c:Customer {customer_id: $customer_id}), (p:Product {product_id: $product_id}) + MERGE (c)-[h:HOLDS]->(p) + SET h.qty=$qty, h.market_value=$market_value, h.cost=$cost_amount, + h.pnl_pct=$pnl_pct, h.as_of=$as_of + """, + customer_id=h["customer_id"], + product_id=h["product_id"], + qty=float(h["qty"]), + market_value=float(h["market_value"]), + cost_amount=float(h["cost_amount"]), + pnl_pct=float(h["pnl_pct"]), + as_of=str(h["as_of"]), + ) + + driver.close() + print( + f"sync_neo4j: customers={len(customers)} advisors={len(advisors)} " + f"products={len(products)} holdings={len(holdings)}" + ) + + +if __name__ == "__main__": + main()