Files
group_fqcd_jr/tools/configure_embedding_endpoint.py
lzf_0626 d2aff7c129 feat: 建客服知识库(Milvus 三集合)并修复模型端点筛选缺陷
一、知识库建设
- 新增 tools/build_knowledge_chunks.py:把 knowledge/ 下文档切成可检索知识块。
  采用「叶子标题」策略(其后没有更深标题的标题即切分点),同时覆盖三种真实结构:
  带子条款的按子条款切、无子条款的条款单独成块、无小节的章整章成块。
  第一版按固定标题级别切是失败的——适当性指南的条款是 ### 而没有 ####,产品手册的
  ### 1.1 又不匹配「第X条」,两条规则互相打架,导致 4 个文件一块都没切出来。
- 新增 tools/load_knowledge_milvus.py:向量化并写入 Milvus,用 upsert 保证幂等。
  schema 按方案 §4.2 统一字段,另加 chapter/section/source_file/doc_no/visibility 五个
  检索与合规必需字段;索引 IVF_FLAT + COSINE + nlist=128;向量输入取「标题+正文」,
  标题含条款号与章节名,是比正文更干净的检索信号。
- 知识内容按业务范围裁剪:反洗钱合规操作手册不入客服知识库(业务只做公募基金、
  不涉及资金划付,且该手册标注内部机密、禁止向客户透露可疑交易信息),留给后续风控;
  高净值客户服务规范只保留「客户分层标准」与「各层级专属权益」两章,
  家族信托、资产配置流程、客户经理考核、隐私应急预案等内部管理章节不入库。
- 入库现状:fin_faq_collection 61 块、fin_product_collection 26 块、
  fin_policy_collection 73 块,合计 160 块。检索自检 5/6——未命中的一条分数 0.660
  落在中置信区间,按三档兜底策略本应提示信息可能不完整,属于预期行为。

二、embedding 端点
- 新增 tools/configure_embedding_endpoint.py:走管理 API(draft→approved→active)
  配置并激活 qwen-embedding 端点,而不是直接写库。理由是状态机与审计都要留痕,
  且 DatabaseModelGateway 只认 status='active',手工写错状态会报成与病因无关的
  「模型端点未注册或未激活」。脚本先查 endpoint_code 是否已存在,幂等可重跑。

三、修复模型端点筛选缺陷(app/service/model_gateway.py)
- 原 DatabaseModelEndpointResolver 忽略 agent_type 与 task_type、直接返回全部 active
  端点,而 ModelDispatchService 只按顺序尝试前 max_attempts(默认 2)个。两者叠加使
  「能否选到支持该任务的端点」取决于端点表顺序:实测每次 embedding 都先拿文本生成
  端点失败一次再落到向量端点(0.61s,修复后 0.42s)。
- 新增 TASK_CAPABILITY 显式映射后按能力筛选。用映射而不是同名筛选是必需的:
  memory_extraction 并不是任何端点的能力名(deepseek 声明的是 text_generation 等),
  按同名筛会得到空集、把记忆抽取打成失败关闭——这是本次修复最容易引入的回归。
- 保守兜底:未映射的 task_type、以及没有任何端点声明该能力时,都退回全部端点,
  让配置缺口表现为调用失败,而不是让上层收到「解析为空」这种与病因无关的报错。
- 验证结果:embedding→[qwen-embedding]、intent_classification→[deepseek-flash]、
  memory_extraction→[deepseek-flash]、未映射 task_type→全部;ruff 通过、
  mypy 103 文件无错、unit+contract 447 passed。

四、需求文档提取物
- 新增 _flows/:三份流程文档(智能客服 Agent 专项设计方案、投资顾问流程、基金运营流程)
  的纯文本提取,供开发期对照。原始 .docx/.html 保留在业务方目录侧。

说明:本次仅本地提交,未推送远程仓库。knowledge/ 内含公司内部制度与产品资料,
是否入远程库待确认。
2026-09-10 20:15:09 +08:00

138 lines
5.2 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.
"""配置并激活 embedding 模型端点(走管理 API,不直接写库)。
为什么走 API 而不是 INSERT:端点配置要经过 draft → approved → active 状态机并留下
`interaction_audit`。直接写库会绕过审核与审计,而且 `DatabaseModelGateway` 只认
`status='active'`,手工写错状态会表现为「模型端点未注册或未激活」这种与病因无关的报错。
幂等:脚本先查该 endpoint_code 是否已存在,已存在则跳过创建,只做后续状态推进。
用法:python tools/configure_embedding_endpoint.py
"""
import asyncio
import datetime as dt
import sys
import uuid
from pathlib import Path
import asyncmy
import httpx
import jwt
from app.core.config import get_settings
from app.main import create_app
ADMIN = "9003"
ENDPOINT_CODE = "qwen-embedding"
MYSQL_DSN_HOST = "127.0.0.1"
MYSQL_USER, MYSQL_PASSWORD, MYSQL_DB = "root", "123456", "jr"
def token(subject: str) -> str:
settings = get_settings()
private_key = Path(settings.jwt_private_key_path).read_text(encoding="utf-8")
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 endpoint_id_if_exists() -> int | None:
connection = await asyncmy.connect(
host=MYSQL_DSN_HOST, port=3306, user=MYSQL_USER, password=MYSQL_PASSWORD, db=MYSQL_DB
)
try:
cursor = connection.cursor()
await cursor.execute(
"SELECT id FROM model_endpoint_config WHERE endpoint_code=%s", (ENDPOINT_CODE,)
)
row = await cursor.fetchone()
return int(row[0]) if row else None
finally:
connection.close()
async def post(
client: httpx.AsyncClient, path: str, *, auth: dict[str, str],
payload: dict[str, object] | None = None, if_match: str | None = None,
) -> httpx.Response:
headers = {**auth, "Idempotency-Key": uuid.uuid4().hex}
if if_match:
headers["If-Match"] = if_match
return await client.post(path, json=payload, headers=headers)
async def etag_of(client: httpx.AsyncClient, path: str, auth: dict[str, str]) -> str | None:
response = await client.get(path, headers=auth)
return response.headers.get("ETag")
async def main() -> int:
app = create_app()
auth = {"Authorization": f"Bearer {token(ADMIN)}"}
async with httpx.AsyncClient(
transport=httpx.ASGITransport(app=app), base_url="http://test", timeout=60
) as client:
endpoint_id = await endpoint_id_if_exists()
if endpoint_id is None:
created = await post(
client, "/api/v1/admin/model-endpoints", auth=auth,
payload={
"endpoint_code": ENDPOINT_CODE,
"provider": "dashscope",
"model_name": "qwen3.7-text-embedding-flash",
"base_url": "https://dashscope.aliyuncs.com/compatible-mode/v1",
"secret_ref": "env:QWEN_EMBEDDING_API_KEY",
"capabilities": ["embedding"],
"allowed_data_levels": ["public", "internal"],
"context_window": 8192,
"timeout_ms": 30000,
},
)
print(f"创建端点:{created.status_code} {created.text[:160]}")
if created.status_code != 201:
return 1
endpoint_id = int(created.json()["data"]["id"])
else:
print(f"端点已存在,复用 id={endpoint_id}")
detail_path = f"/api/v1/admin/model-endpoints/{endpoint_id}"
current = (await client.get(detail_path, headers=auth)).json()["data"]
print(f"当前状态:{current['status']}")
if current["status"] == "draft":
reviewed = await post(
client, f"{detail_path}/reviews", auth=auth,
payload={"decision": "approved", "comment": "embedding 端点配置"},
if_match=await etag_of(client, detail_path, auth),
)
print(f"审核:{reviewed.status_code} {reviewed.text[:160]}")
if reviewed.status_code not in (200, 201):
return 1
current = (await client.get(detail_path, headers=auth)).json()["data"]
if current["status"] == "approved":
activated = await post(
client, f"{detail_path}/activations", auth=auth,
# activations 的 body 是必填的 EmptyPayload(extra=forbid),
# 不传 body 会得到 422 "Field required",因此显式给空对象。
payload={},
if_match=await etag_of(client, detail_path, auth),
)
print(f"激活:{activated.status_code} {activated.text[:200]}")
if activated.status_code not in (200, 201):
return 1
final = (await client.get(detail_path, headers=auth)).json()["data"]
print(f"最终状态:{final['status']} 模型={final['model_name']}")
return 0 if final["status"] == "active" else 1
sys.exit(asyncio.run(main()))