Files
group_fqcd_jr/tools/seed_knowledge_demo.py
lzf_0626 bfd3964ef8 feat(demo): 演示数据一键准备、启动脚本与演示流程文档
交付"别人能自己把系统跑起来做演示"所需的四件东西:

- 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
2026-09-13 21:56:27 +08:00

158 lines
6.3 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""把演示用的场内基金知识灌进知识库(幂等:先清同源旧数据再上传)。
## 为什么需要它
知识库是客服能答对问题的前提,但它是**演示数据里唯一没有脚本化的一环** ——
之前是手工调 `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()))