交付"别人能自己把系统跑起来做演示"所需的四件东西: - tools/seed_demo_data.py 演示数据一键准备:10 步按依赖排序 (此前散在 10 个脚本里,没人知道该跑哪些、按什么顺序跑,且知识库那步 根本没有脚本、靠手工调接口) - tools/seed_knowledge_demo.py 知识库演示素材灌入,幂等(先删同 source_file 再灌) - start.ps1 启动 API + Worker,含依赖检查与**自动刷新行情** - docs/44-演示流程.md 8 个主线场景的照读流程 + 排障表 + 账号/命令速查 两条硬约束同时写进了脚本和文档: 1. **Worker 必须常驻**:没有它客服对话一直停在 queued(前端只显示"超时")、 新知识不进 Milvus 且**没有任何报错**。 2. **行情有效期仅 15 分钟**(`trade_service.MAX_QUOTE_AGE`):超时后所有委托 直接 503「行情已过期」,而系统**没有自动刷新机制**。故 start.ps1 启动时刷一次, 文档另给"演示中途 503 时补刷、无需重启服务"的处置方法(已实测)。 实测证据: - start.ps1 在 8099 完整启动,`/docs` 与 `/portal/` 均 HTTP 200(测试进程已清理) - `tools/e2e_smoke_test.py` → 40/40 通过 - `tools/seed_demo_data.py --dry-run` → 10 步全部正常列出 - 文档中的下单命令实跑:510300 成交价 4.579、金额 457.90、手续费 0.05
158 lines
6.3 KiB
Python
158 lines
6.3 KiB
Python
"""把演示用的场内基金知识灌进知识库(幂等:先清同源旧数据再上传)。
|
||
|
||
## 为什么需要它
|
||
|
||
知识库是客服能答对问题的前提,但它是**演示数据里唯一没有脚本化的一环** ——
|
||
之前是手工调 `POST /api/v1/knowledge/upload` 灌的,换台机器没人知道该灌什么、
|
||
灌到哪个集合。本脚本把 `docs/43-场内基金产品手册(知识库入库版).md` 定为唯一素材源。
|
||
|
||
## 三个必须讲清的点
|
||
|
||
1. **上传后不会立刻可检索**:入库只写 MySQL 元数据 + 投 outbox 事件,向量由
|
||
**Agent Worker** 消费事件后写 Milvus。所以**必须先起 Worker**
|
||
(`python -m app.worker`),否则知识永远检索不到 —— 而且现象是"客服照旧答不上",
|
||
没有任何报错。
|
||
2. **素材必须是"客户可见版"**:用 `docs/43`,**不要**用 `docs/42` —— 后者含内部决策
|
||
备注("待你确认""库里 vs 真实"),灌进去可能被客户问题检索出来。
|
||
3. **走进程内调用**(ASGITransport),与 `tools/publish_*.py` 一致,因此**不需要先起 API**。
|
||
|
||
## 用法
|
||
|
||
python tools/seed_knowledge_demo.py # 幂等重灌
|
||
python tools/seed_knowledge_demo.py --dry-run # 只看会做什么,不调写接口
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import asyncio
|
||
import base64
|
||
import sys
|
||
import uuid
|
||
from pathlib import Path
|
||
from typing import Any
|
||
|
||
import httpx
|
||
import jwt
|
||
|
||
PROJECT_ROOT = Path(__file__).resolve().parents[1]
|
||
if str(PROJECT_ROOT) not in sys.path:
|
||
sys.path.insert(0, str(PROJECT_ROOT))
|
||
|
||
from app.core.config import get_settings # noqa: E402
|
||
from app.main import create_app # noqa: E402
|
||
|
||
if hasattr(sys.stdout, "reconfigure"):
|
||
sys.stdout.reconfigure(errors="replace") # type: ignore[union-attr]
|
||
|
||
ADMIN = "9003"
|
||
SOURCE = PROJECT_ROOT / "docs" / "43-场内基金产品手册(知识库入库版).md"
|
||
#: 上传后的文件名,也是**幂等清理的判据**:同名的旧行会被先删掉。
|
||
UPLOAD_FILENAME = "场内基金产品手册.md"
|
||
KNOWLEDGE_TYPE = "product"
|
||
|
||
|
||
def token(subject: str) -> str:
|
||
settings = get_settings()
|
||
private_key = Path(settings.jwt_private_key_path).read_text(encoding="utf-8")
|
||
import datetime as dt
|
||
|
||
now = dt.datetime.now(dt.UTC)
|
||
return jwt.encode(
|
||
{
|
||
"sub": subject, "iss": settings.jwt_issuer, "aud": settings.jwt_audience,
|
||
"exp": now + dt.timedelta(minutes=30), "nbf": now - dt.timedelta(seconds=5),
|
||
"jti": str(uuid.uuid4()),
|
||
},
|
||
private_key,
|
||
algorithm="RS256",
|
||
)
|
||
|
||
|
||
async def list_existing(client: httpx.AsyncClient, auth: dict[str, str]) -> list[dict[str, Any]]:
|
||
"""列出当前知识。
|
||
|
||
⚠️ 这个端点**不套 `data` 信封**(直接返回 `{"items": [...], "count": N}`),
|
||
按 `data.items` 解包会得到空列表、看着像"库里没数据" —— 实测踩过。
|
||
"""
|
||
response = await client.get("/api/v1/knowledge/list?limit=100", headers=auth)
|
||
if response.status_code != 200:
|
||
print(f"[警告] 列表接口返回 {response.status_code},按『无现存知识』继续")
|
||
return []
|
||
body = response.json()
|
||
items = body.get("items")
|
||
if items is None and isinstance(body.get("data"), dict):
|
||
items = body["data"].get("items")
|
||
return items or []
|
||
|
||
|
||
async def main() -> int:
|
||
parser = argparse.ArgumentParser(description="灌入演示用的场内基金知识")
|
||
parser.add_argument("--dry-run", action="store_true", help="只看会做什么,不调写接口")
|
||
args = parser.parse_args()
|
||
|
||
if not SOURCE.exists():
|
||
print(f"[失败] 素材不存在:{SOURCE}")
|
||
return 1
|
||
text = SOURCE.read_text(encoding="utf-8")
|
||
print(f"素材:{SOURCE.name}({len(text)} 字符)")
|
||
print(f"目标:knowledge_type={KNOWLEDGE_TYPE} → fin_product_collection")
|
||
print(f"文件名:{UPLOAD_FILENAME}(同名旧行会被先删除,以保证幂等)\n")
|
||
|
||
app = create_app()
|
||
auth = {"Authorization": f"Bearer {token(ADMIN)}"}
|
||
async with httpx.AsyncClient(
|
||
transport=httpx.ASGITransport(app=app), base_url="http://test", timeout=120
|
||
) as client:
|
||
existing = await list_existing(client, auth)
|
||
stale = [row for row in existing if row.get("source_file") == UPLOAD_FILENAME]
|
||
print(f"现存知识 {len(existing)} 块,其中本素材的旧版 {len(stale)} 块")
|
||
|
||
if args.dry_run:
|
||
print("\n[dry-run] 将会:")
|
||
print(f" 1. 删除 {len(stale)} 块旧版(id: {[r.get('knowledge_id') for r in stale]})")
|
||
print(f" 2. 上传 {SOURCE.name} 并切块入库")
|
||
print(" 未调用任何写接口。")
|
||
return 0
|
||
|
||
deleted = 0
|
||
for row in stale:
|
||
knowledge_id = row.get("knowledge_id")
|
||
# 写接口必须带幂等键:平台对缺失键的写请求按失败关闭处理。
|
||
response = await client.delete(
|
||
f"/api/v1/knowledge/{knowledge_id}",
|
||
headers={**auth, "Idempotency-Key": uuid.uuid4().hex},
|
||
)
|
||
if response.status_code in (200, 204):
|
||
deleted += 1
|
||
else:
|
||
print(f" [警告] 删除 {knowledge_id} 返回 {response.status_code}")
|
||
if stale:
|
||
print(f"已清理旧版 {deleted}/{len(stale)} 块")
|
||
|
||
response = await client.post(
|
||
"/api/v1/knowledge/upload",
|
||
headers={**auth, "Idempotency-Key": uuid.uuid4().hex},
|
||
json={
|
||
"filename": UPLOAD_FILENAME,
|
||
"knowledge_type": KNOWLEDGE_TYPE,
|
||
"content_base64": base64.b64encode(text.encode("utf-8")).decode("ascii"),
|
||
},
|
||
)
|
||
if response.status_code not in (200, 201):
|
||
print(f"[失败] 上传返回 {response.status_code}:{response.text[:300]}")
|
||
return 1
|
||
print(f"上传成功(HTTP {response.status_code})")
|
||
|
||
print(
|
||
"\n完成。⚠️ 接下来必须:\n"
|
||
" 1. 确认 Agent Worker 在跑(python -m app.worker)—— 向量由它写进 Milvus;\n"
|
||
" 2. 等几秒让向量同步完成;\n"
|
||
" 3. 用 python tools/e2e_smoke_test.py 验证 A 线(访客问答)不再转人工。"
|
||
)
|
||
return 0
|
||
|
||
|
||
if __name__ == "__main__":
|
||
sys.exit(asyncio.run(main()))
|