Files

158 lines
6.3 KiB
Python
Raw Permalink Normal View History

"""把演示用的场内基金知识灌进知识库(幂等:先清同源旧数据再上传)。
## 为什么需要它
知识库是客服能答对问题的前提,但它是**演示数据里唯一没有脚本化的一环** ——
之前是手工调 `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()))