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