Files
group_fqcd_jr/tools/seed_knowledge_demo.py
张胜宇 e239eb778b docs: 品牌全量口径统一为「南方基金」+ 作废文档清理
1) 客服 Agent 四份交付文档 + 构建脚手架:品牌由包装占位 XX科技 / 旧名 南方财富
   统一为南方基金(热线 400-889-8899 / 官网 nffund.com),系统名改为「智能服务系统」;
   同步追加 §0.4 修订记录行,工程记录行保留原占位字面以支撑硬编码扫描验收。
2) 开发文档:清理 28 份已作废/残留文档(14 份移出归档 + 14 份仓库副本),
   新增《文档规整方案与开发前待决事项-2026-09-17》。
3) 客服agent 四份交付文档首次纳入本分支。
2026-09-17 15:15:22 +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()))