595 lines
27 KiB
Markdown
595 lines
27 KiB
Markdown
# 知识检索接入方案(客服 RAG)
|
||
|
||
> 版本:v1.0
|
||
> 编制日期:2026-09-09
|
||
> 范围:补齐 `05-接口文档` K001 知识引用链路的检索端,解锁客服域 RAG 能力
|
||
> **当前状态(2026-09-11 复核):本方案已落地闭环。**
|
||
> `app/service/knowledge_service.py` 已实现 HMAC-SHA256 引用令牌校验(`_signature` / `_verify_token`,
|
||
> 覆盖 token 前四段、`hmac.compare_digest` 定时安全比较,并按发布状态与有效期做二次校验);
|
||
> `app/infrastructure/milvus_adapter.py`、`app/infrastructure/milvus_knowledge_writer.py`、
|
||
> `app/service/knowledge_retrieval_service.py` 均已落地;`query_knowledge` 已注册为公共只读工具
|
||
> (`app/service/agent/bootstrap.py` 的 `ToolRegistry`)。Milvus 检索已在 Phase 1 验收中命中
|
||
> (score 0.7837)。**方案正文保留为设计参考,不再是待办。**
|
||
|
||
---
|
||
|
||
## 1. 背景与目标
|
||
|
||
客服域的 5 类意图中,`faq`、`product_inquiry`、`policy_explain` 三类依赖知识检索(RAG)。方案编制时
|
||
底座已具备意图分类、模型生成、工具执行器、审计与引用回填能力,唯独检索端缺失。**该缺口现已闭环**
|
||
(见文首状态行),本节保留当时的问题描述作为设计背景:
|
||
|
||
- K001 `GET /api/v1/knowledge-references/{reference_token}` 当时恒返回 404;现由
|
||
`KnowledgeReferenceService.resolve()` 真实解析,未发布/已过期/跨用户令牌仍一律 404;
|
||
- 客服 Agent 当时无法检索知识;现经 `query_knowledge` 工具检索,FAQ 直返、产品问答、
|
||
政策解读三条链路均已实现(`CustomerServiceAgent`);
|
||
- `conversation_message.source_references` 现由 `ToolExecutor` 自动生成并回填知识来源。
|
||
|
||
**目标**:按客服专项设计文档的三集合路由,落地「意图 → 集合 → embedding → Milvus 检索 → 重排 → MySQL 降级」完整链路,并作为**公共只读工具**供所有业务域复用。
|
||
|
||
---
|
||
|
||
## 2. 现状诊断
|
||
|
||
### 2.1 已有基础(可复用)
|
||
|
||
| 能力 | 位置 | 可复用点 |
|
||
|---|---|---|
|
||
| 受控外部边界范式 | `service/relationship_service.py` | 白名单 + 参数化 + 异常降级返回 degraded |
|
||
| 向量降级壳子 | `infrastructure/vector_memory.py` | `VectorClient` Protocol + 异常降级 |
|
||
| 业务工具注册 | `service/agent/bootstrap.py` | `ToolRegistry.register(ToolDefinition(...))` |
|
||
| 工具执行治理 | `service/tool_executor.py` | 白名单/权限/角色/参数/超时/审计/引用自动回填 |
|
||
| 配置读取范式 | `service/runtime_config_service.py` | `PlatformConfigItem` + namespace + config_key |
|
||
| 业务契约范式 | `core/fund_contracts.py` | frozen Pydantic + 字段校验 + `Literal` 来源标记 |
|
||
| 业务工具样板 | `service/fund_quote_service.py` | 运行时配置 + 降级标记 + 只读约束 + 延迟建连 |
|
||
| 知识元数据表 | `fin_knowledge_meta`(00 基线 + 02 §6.2) | `content_text` / `tags` / `milvus_collection` / `review_status` / 有效期 |
|
||
|
||
### 2.2 缺失清单(编制时;「现状」列由 2026-09-11 复核补注)
|
||
|
||
| # | 缺失项 | 影响 | 优先级 | 现状 |
|
||
|---|---|---|---|---|
|
||
| 1 | `ModelGateway.embed()` | 无向量,检索链路断头 | **P0 前置** | 已闭环(`ModelEmbeddingService`,`task_type=embedding` 端点解析) |
|
||
| 2 | Milvus 适配器 | 无真实检索通道 | P0 | 已闭环(`infrastructure/milvus_adapter.py`、`milvus_knowledge_writer.py`) |
|
||
| 3 | `KnowledgeRetrievalService` | 无检索服务 | P0 | 已闭环(`service/knowledge_retrieval_service.py`) |
|
||
| 4 | 集合路由(配置驱动) | 意图无法映射到集合 | P0 | 已闭环 |
|
||
| 5 | MySQL 降级检索 | Milvus 故障时无兜底 | P1 | 已闭环 |
|
||
| 6 | `query_knowledge` 工具注册 | 业务 Agent 无法调用 | P1 | 已闭环(`bootstrap.py` 的 `ToolRegistry`;`service/knowledge_tool.py`) |
|
||
| 7 | `KnowledgeReferenceService.resolve()` 占位 | 引用查询永远 404 | P1 | 已闭环(`service/knowledge_service.py`,HMAC-SHA256 令牌校验) |
|
||
|
||
---
|
||
|
||
## 3. 目标链路
|
||
|
||
```text
|
||
业务 Agent.handle()
|
||
│
|
||
├─ self.call_tool("query_knowledge", {query, intents}, intent=..., context=...)
|
||
│ │
|
||
│ ├─ ToolExecutor:白名单 → 权限 → 角色 → 参数校验 → 超时 → 审计
|
||
│ │
|
||
│ └─ KnowledgeRetrievalService.search()
|
||
│ ├─ 集合白名单校验(三集合)
|
||
│ ├─ 意图 → 集合 + TopK(配置驱动)
|
||
│ ├─ embedding(模型路由 task_type=embedding)
|
||
│ ├─ Milvus TopK 检索
|
||
│ │ └─ 失败 → MySQL content_text 降级检索
|
||
│ ├─ 重排去重 + 相似度阈值过滤
|
||
│ └─ fin_knowledge_meta 有效期/发布状态二次过滤
|
||
│
|
||
└─ 返回 hits,ToolExecutor 自动生成 SourceReference 并写入结果
|
||
└─ AgentPersistenceService.complete_run()
|
||
└─ conversation_message.source_references 回填
|
||
└─ K001 引用可解析
|
||
```
|
||
|
||
---
|
||
|
||
## 4. 数据契约设计
|
||
|
||
### 4.1 新增契约文件 `app/core/knowledge_contracts.py`
|
||
|
||
```python
|
||
"""知识检索公共契约。字段稳定,不暴露 Milvus/pymilvus 概念。"""
|
||
|
||
from typing import Literal
|
||
|
||
from pydantic import BaseModel, ConfigDict, Field, field_validator
|
||
|
||
ALLOWED_COLLECTIONS = frozenset({
|
||
"fin_faq_collection",
|
||
"fin_product_collection",
|
||
"fin_policy_collection",
|
||
})
|
||
|
||
|
||
class KnowledgeQuery(BaseModel):
|
||
"""工具入参。集合由意图映射,调用方不得直接指定集合名。"""
|
||
|
||
model_config = ConfigDict(extra="forbid", frozen=True)
|
||
|
||
query: str = Field(min_length=1, max_length=2000)
|
||
intents: tuple[str, ...] = Field(min_length=1, max_length=4)
|
||
top_k: int = Field(default=5, ge=1, le=20)
|
||
|
||
@field_validator("query")
|
||
@classmethod
|
||
def query_must_not_be_blank(cls, value: str) -> str:
|
||
if not value.strip():
|
||
raise ValueError("query must not be blank")
|
||
return value
|
||
|
||
|
||
class KnowledgeHit(BaseModel):
|
||
model_config = ConfigDict(extra="forbid", frozen=True)
|
||
|
||
knowledge_id: str
|
||
collection: str
|
||
title: str | None = None
|
||
snippet: str
|
||
score: float | None = Field(default=None, ge=0, le=1)
|
||
tags: tuple[str, ...] = ()
|
||
version: str | None = None
|
||
|
||
|
||
class KnowledgeSearchResult(BaseModel):
|
||
model_config = ConfigDict(extra="forbid", frozen=True)
|
||
|
||
hits: tuple[KnowledgeHit, ...]
|
||
degraded: bool = False
|
||
degradation_reason: str | None = None
|
||
searched_collections: tuple[str, ...] = ()
|
||
```
|
||
|
||
### 4.2 集合与维度约定(取自客服专项设计文档)
|
||
|
||
| 意图 | 集合 | TopK |
|
||
|---|---|---|
|
||
| `faq` | `fin_faq_collection` | 3 |
|
||
| `product_inquiry` | `fin_product_collection` | 5 |
|
||
| `policy_explain` | `fin_policy_collection` | 5 |
|
||
|
||
- **向量维度固定 1024**(`VECTOR_DIM`,与 Qwen Embedding 输出一致)。维度不符必须失败关闭,不得静默降级为错误结果。
|
||
- 集合内字段:`knowledge_id` / `title` / `snippet` / `tags` / `version` / `embedding`。
|
||
|
||
### 4.3 配置项设计
|
||
|
||
新增 `PlatformConfigItem`:
|
||
|
||
```text
|
||
namespace = "knowledge"
|
||
config_key = "default"
|
||
value_json = {
|
||
"collections": {
|
||
"faq": {"collection": "fin_faq_collection", "top_k": 3},
|
||
"product_inquiry": {"collection": "fin_product_collection", "top_k": 5},
|
||
"policy_explain": {"collection": "fin_policy_collection", "top_k": 5}
|
||
},
|
||
"embedding_endpoint": "embedding-primary",
|
||
"vector_dim": 1024,
|
||
"similarity_threshold": 0.60,
|
||
"fallback_enabled": true,
|
||
"milvus_timeout_seconds": 2
|
||
}
|
||
```
|
||
|
||
读取方式沿用 `RuntimeConfigService` 范式,新增 `knowledge(release_id)` 方法;配置缺失时使用与客服文档一致的安全默认值(上表)。
|
||
|
||
---
|
||
|
||
## 5. 分步实施方案
|
||
|
||
### 步骤 1:补齐 Embedding 出口(P0 前置)
|
||
|
||
**文件**:`app/service/model_gateway.py`
|
||
|
||
```python
|
||
class ModelGateway(Protocol):
|
||
async def generate(self, *, endpoint_code: str, prompt: str, timeout_ms: int) -> str: ...
|
||
async def embed(
|
||
self, *, endpoint_code: str, texts: list[str], timeout_ms: int
|
||
) -> list[list[float]]: ...
|
||
|
||
|
||
class OpenAICompatibleGateway:
|
||
async def embed(
|
||
self, *, endpoint_code: str, texts: list[str], timeout_ms: int
|
||
) -> list[list[float]]:
|
||
endpoint = self.endpoints.get(endpoint_code)
|
||
if endpoint is None:
|
||
raise RecoverableAgentError("嵌入端点未注册")
|
||
if not texts:
|
||
raise RecoverableAgentError("嵌入输入为空")
|
||
token = self.secret_resolver.resolve(endpoint.secret_ref)
|
||
client = self.client or httpx.AsyncClient()
|
||
try:
|
||
response = await client.post(
|
||
endpoint.base_url.rstrip("/") + "/embeddings",
|
||
headers={"Authorization": f"Bearer {token}",
|
||
"Content-Type": "application/json"},
|
||
json={"model": endpoint.model_name, "input": texts},
|
||
timeout=httpx.Timeout(timeout_ms / 1000),
|
||
)
|
||
response.raise_for_status()
|
||
body: Any = response.json()
|
||
vectors = [item.get("embedding") for item in body.get("data", [])]
|
||
if len(vectors) != len(texts) or not all(isinstance(v, list) for v in vectors):
|
||
raise RecoverableAgentError("嵌入响应与输入数量不一致")
|
||
return [list(map(float, vector)) for vector in vectors]
|
||
except httpx.HTTPError as exc:
|
||
raise RecoverableAgentError("嵌入端点调用失败") from exc
|
||
finally:
|
||
if self._owns_client:
|
||
await client.aclose()
|
||
|
||
|
||
class ModelGenerationService:
|
||
async def embed(self, endpoints: list[Any], texts: list[str]) -> list[list[float]]:
|
||
if not endpoints:
|
||
raise RecoverableAgentError("没有可用的已批准嵌入端点")
|
||
return await self.dispatch.embed(endpoints, texts)
|
||
```
|
||
|
||
**要点**:
|
||
- 复用已有模型路由(01 §9.1 已定义 `task_type=embedding`),端点通过 `ModelEndpointConfig.capabilities` 声明嵌入能力;
|
||
- 密钥仍只从 `secret_ref`(`env:` 前缀)解析;
|
||
- 不新增独立配置体系。
|
||
|
||
### 步骤 2:Milvus 适配器
|
||
|
||
**文件**:`app/infrastructure/milvus_adapter.py`(新增)
|
||
|
||
```python
|
||
"""Milvus 只读适配器:仅执行检索,不建集合、不写入。"""
|
||
|
||
from typing import Any
|
||
|
||
from app.core.errors import RecoverableAgentError
|
||
|
||
|
||
class MilvusKnowledgeClient:
|
||
def __init__(self, uri: str, token: str = "", *, timeout: float = 2.0) -> None:
|
||
self._uri, self._token, self._timeout = uri, token, timeout
|
||
self._client: Any = None
|
||
|
||
async def _ensure(self) -> Any:
|
||
if self._client is None:
|
||
from pymilvus import AsyncMilvusClient # 延迟导入,避免底座启动即建连
|
||
self._client = AsyncMilvusClient(uri=self._uri, token=self._token or None)
|
||
return self._client
|
||
|
||
async def search(
|
||
self, *, collection: str, vector: list[float], top_k: int,
|
||
output_fields: tuple[str, ...] = ("knowledge_id", "title", "snippet", "tags", "version"),
|
||
) -> list[dict[str, Any]]:
|
||
client = await self._ensure()
|
||
try:
|
||
results = await client.search(
|
||
collection_name=collection,
|
||
data=[vector],
|
||
limit=top_k,
|
||
output_fields=list(output_fields),
|
||
search_params={"metric_type": "COSINE"},
|
||
)
|
||
except Exception as exc:
|
||
raise RecoverableAgentError("知识检索不可用") from exc
|
||
return [dict(hit) for batch in results for hit in batch]
|
||
|
||
async def close(self) -> None:
|
||
if self._client is not None:
|
||
await self._client.close()
|
||
self._client = None
|
||
```
|
||
|
||
**要点**:
|
||
- **只读**:不提供 create/insert/delete;
|
||
- 延迟导入 + 延迟建连,遵循 `EastmoneyAdapterFactory` 的做法;
|
||
- 具体调用签名按 `pymilvus 2.6` 的 `AsyncMilvusClient` 校对(若该版本无异步客户端,则用 `MilvusClient` 配合 `anyio.to_thread.run_sync` 包装,避免阻塞事件循环)。
|
||
|
||
### 步骤 3:检索服务(受控边界)
|
||
|
||
**文件**:`app/service/knowledge_retrieval_service.py`(新增)
|
||
|
||
```python
|
||
"""知识检索受控边界:调用方不能指定任意集合、原始查询或未授权字段。"""
|
||
|
||
from collections.abc import Mapping
|
||
from typing import Any
|
||
|
||
from app.core.contracts import RequestContext
|
||
from app.core.errors import ForbiddenAgentError, RecoverableAgentError
|
||
from app.core.knowledge_contracts import (
|
||
ALLOWED_COLLECTIONS, KnowledgeHit, KnowledgeQuery, KnowledgeSearchResult,
|
||
)
|
||
from app.infrastructure.milvus_adapter import MilvusKnowledgeClient
|
||
from app.infrastructure.db import SessionFactory
|
||
from app.model.knowledge import FinKnowledgeMeta # 需按 00 基线补 ORM 映射
|
||
|
||
|
||
class KnowledgeRuntimeConfig:
|
||
"""与客服专项设计一致的安全默认值;配置缺失或非法时回退。"""
|
||
|
||
DEFAULT_ROUTES: Mapping[str, tuple[str, int]] = {
|
||
"faq": ("fin_faq_collection", 3),
|
||
"product_inquiry": ("fin_product_collection", 5),
|
||
"policy_explain": ("fin_policy_collection", 5),
|
||
}
|
||
|
||
def __init__(self, *, routes=None, vector_dim: int = 1024,
|
||
similarity_threshold: float = 0.60, fallback_enabled: bool = True) -> None:
|
||
self.routes = dict(routes or self.DEFAULT_ROUTES)
|
||
self.vector_dim = vector_dim
|
||
self.similarity_threshold = similarity_threshold
|
||
self.fallback_enabled = fallback_enabled
|
||
|
||
@classmethod
|
||
def from_mapping(cls, value: object) -> "KnowledgeRuntimeConfig":
|
||
... # 逐字段校验,非法值回退默认,与 FundQuoteRuntimeConfig.from_mapping 同风格
|
||
|
||
|
||
class KnowledgeRetrievalService:
|
||
def __init__(self, *, client: MilvusKnowledgeClient, embedder: Any,
|
||
config: KnowledgeRuntimeConfig) -> None:
|
||
self._client, self._embedder, self._config = client, embedder, config
|
||
|
||
async def search(
|
||
self, query: KnowledgeQuery, context: RequestContext, *,
|
||
embedding_endpoints: list[Any], embedding_timeout_ms: int = 5000,
|
||
) -> KnowledgeSearchResult:
|
||
targets = self._resolve_targets(query.intents) # 意图 → (集合, top_k)
|
||
self._assert_collections_allowed(targets) # 白名单
|
||
if not targets:
|
||
return KnowledgeSearchResult(hits=(), searched_collections=())
|
||
|
||
try:
|
||
vectors = await self._embedder.embed(embedding_endpoints, [query.query])
|
||
vector = vectors[0]
|
||
if len(vector) != self._config.vector_dim: # 维度失败关闭
|
||
raise RecoverableAgentError("嵌入维度与集合定义不一致")
|
||
except Exception as exc:
|
||
if not self._config.fallback_enabled:
|
||
raise
|
||
return await self._fallback_search(query, targets, reason=type(exc).__name__)
|
||
|
||
hits: list[KnowledgeHit] = []
|
||
try:
|
||
for collection, top_k in targets:
|
||
raw = await self._client.search(
|
||
collection=collection, vector=vector, top_k=top_k)
|
||
hits.extend(self._to_hits(raw, collection))
|
||
except Exception as exc:
|
||
if not self._config.fallback_enabled:
|
||
raise
|
||
return await self._fallback_search(query, targets, reason=type(exc).__name__)
|
||
|
||
ranked = self._rank_and_dedupe(hits)
|
||
verified = await self._filter_published(ranked) # MySQL 二次校验
|
||
return KnowledgeSearchResult(
|
||
hits=tuple(verified), searched_collections=tuple(c for c, _ in targets))
|
||
|
||
def _resolve_targets(self, intents: tuple[str, ...]) -> list[tuple[str, int]]: ...
|
||
def _assert_collections_allowed(self, targets) -> None:
|
||
if not {c for c, _ in targets} <= ALLOWED_COLLECTIONS:
|
||
raise ForbiddenAgentError("未授权的知识集合")
|
||
def _to_hits(self, raw: list[dict], collection: str) -> list[KnowledgeHit]: ...
|
||
def _rank_and_dedupe(self, hits: list[KnowledgeHit]) -> list[KnowledgeHit]: ...
|
||
async def _filter_published(self, hits: list[KnowledgeHit]) -> list[KnowledgeHit]: ...
|
||
async def _fallback_search(self, query, targets, *, reason: str) -> KnowledgeSearchResult: ...
|
||
```
|
||
|
||
**必须遵守的三条**(对齐 `RelationshipService` 与 02 §6.2):
|
||
|
||
1. **集合白名单**:`ALLOWED_COLLECTIONS` 之外的集合名一律拒绝;
|
||
2. **MySQL 二次校验**:Milvus 命中的条目必须回查 `fin_knowledge_meta`,确认 `review_status='published' AND status='active'` 且在有效期内——**防止已下线知识仍能通过向量检索命中**;
|
||
3. **降级必带有效期过滤**:`_fallback_search` 走 MySQL 时同样执行 02 §6.2 的四段过滤条件。
|
||
|
||
### 步骤 4:配置驱动集合路由
|
||
|
||
**文件**:`app/service/runtime_config_service.py`
|
||
|
||
```python
|
||
async def knowledge(self, release_id: int) -> dict[str, object]:
|
||
item = await self.session.scalar(
|
||
select(PlatformConfigItem).where(
|
||
PlatformConfigItem.release_id == release_id,
|
||
PlatformConfigItem.namespace == "knowledge",
|
||
PlatformConfigItem.config_key == "default",
|
||
)
|
||
)
|
||
return item.value_json if item is not None else {}
|
||
|
||
|
||
async def load_knowledge_config() -> dict[str, object]:
|
||
async with SessionFactory() as session:
|
||
release = await session.scalar(
|
||
select(ConfigRelease).where(ConfigRelease.status == "active"))
|
||
if release is None:
|
||
return {}
|
||
return await RuntimeConfigService(session).knowledge(release.id)
|
||
```
|
||
|
||
**要点**:集合路由走配置发布流程(双人审核 + 原子激活),不在代码里散落硬编码;`agent_intent_config.collection_routes` 字段(02 §7.6)保留为按意图覆盖的补充来源。
|
||
|
||
### 步骤 5:MySQL 降级检索
|
||
|
||
```sql
|
||
SELECT id, title, content_text, tags, version
|
||
FROM fin_knowledge_meta
|
||
WHERE milvus_collection IN (:collections)
|
||
AND review_status = 'published'
|
||
AND status = 'active'
|
||
AND (effective_date IS NULL OR effective_date <= UTC_DATE())
|
||
AND (expire_date IS NULL OR expire_date > UTC_DATE())
|
||
AND content_text LIKE CONCAT('%', :keyword, '%')
|
||
ORDER BY id DESC
|
||
LIMIT :limit
|
||
```
|
||
|
||
**要点**:
|
||
- 关键词从 `query.query` 提取(去除停用词、长度上限 64);
|
||
- `LIKE` 必须参数化,禁止拼接;
|
||
- 降级结果打 `degraded=True, degradation_reason="milvus_unavailable"`,供可观测性统计(01 §14 已要求"MySQL 降级次数"指标)。
|
||
|
||
### 步骤 6:注册为公共工具
|
||
|
||
**文件**:`app/service/agent/bootstrap.py`
|
||
|
||
```python
|
||
registry.register(ToolDefinition(
|
||
name="query_knowledge",
|
||
input_model=KnowledgeQuery,
|
||
handler=cast(Any, query_knowledge_tool),
|
||
required_permission="knowledge:query",
|
||
allowed_roles=("customer", "operator", "advisor", "risk_operator", "admin"),
|
||
timeout_seconds=8,
|
||
))
|
||
```
|
||
|
||
对应工具处理器(放在 `knowledge_retrieval_service.py` 末尾,与 `query_fund_quote_tool` 同风格):
|
||
|
||
```python
|
||
async def query_knowledge_tool(arguments: KnowledgeQuery, context: RequestContext) -> list[dict]:
|
||
"""ToolExecutor 使用的公共只读知识检索工具。"""
|
||
from app.service.model_router_service import resolve_embedding_endpoints
|
||
config = KnowledgeRuntimeConfig.from_mapping(await _safe_load_config())
|
||
endpoints = await resolve_embedding_endpoints()
|
||
result = await KnowledgeRetrievalService(
|
||
client=MilvusKnowledgeClient(get_settings().milvus_uri, get_settings().milvus_token),
|
||
embedder=ModelGenerationService(ModelDispatchService(DatabaseModelGateway())),
|
||
config=config,
|
||
).search(arguments, context, embedding_endpoints=endpoints)
|
||
return [hit.model_dump(mode="json") for hit in result.hits]
|
||
```
|
||
|
||
**无需**自行拼 `SourceReference`——`ToolExecutor` 已自动生成并回填(`tool_executor.py:94-96`)。
|
||
|
||
### 步骤 7:修复知识引用占位实现
|
||
|
||
**文件**:`app/service/knowledge_service.py`
|
||
|
||
**状态:已按下列 5 条改造完成**(2026-09-11 复核:`signing_secret()`、`_signature()` 用
|
||
`hmac.new(..., hashlib.sha256)`,`_verify_token()` 用 `hmac.compare_digest` 定时安全比较前四段签名,
|
||
`_assert_resolvable()` 校验发布状态与有效期,`_redact()` 只返回脱敏元数据)。
|
||
|
||
改造前(占位实现,**仅作对照,已不存在于代码中**):
|
||
|
||
```python
|
||
async def resolve(self, context, token):
|
||
await AuthorizationService.require(context, "knowledge:reference:read")
|
||
pieces = token.split(".")
|
||
if len(pieces) != 5 or pieces[0] != "kr1" or pieces[1] != context.user_id:
|
||
raise ResourceNotFoundError("引用不存在")
|
||
try:
|
||
if int(pieces[3]) < int(datetime.now(UTC).timestamp()):
|
||
raise ResourceNotFoundError("引用不存在")
|
||
except ValueError as exc:
|
||
raise ResourceNotFoundError("引用不存在") from exc
|
||
raise ResourceNotFoundError("引用不存在") # ← 无条件抛出(改造前)
|
||
```
|
||
|
||
**改造**:
|
||
1. 定义 token 格式:`kr1.{user_id}.{knowledge_id}.{expires_at_epoch}.{signature}`;
|
||
2. 校验签名(HMAC,密钥来自 `env:`),防止伪造与越权拼接;
|
||
3. 校验后按 `knowledge_id` 查 `fin_knowledge_meta`,**仍要求 `published + active + 有效期内`**;
|
||
4. 返回脱敏元数据(`title` / `version` / `tags` / `collection`),**不返回全文**(避免绕过工具审计);
|
||
5. 生成侧:在 `SourceReference` 落库时同步签发 token(`AgentPersistenceService` 内或引用组装处)。
|
||
|
||
### 步骤 8:测试
|
||
|
||
| 类型 | 文件 | 覆盖 |
|
||
|---|---|---|
|
||
| 单元 | `tests/unit/service/test_knowledge_retrieval.py` | 意图路由选中正确集合与 TopK;非法集合名拒绝;维度不符失败关闭;重排去重;相似度阈值过滤 |
|
||
| 单元 | `tests/unit/service/test_knowledge_fallback.py` | Milvus 异常 → MySQL 降级;降级结果带 `degraded`;降级仍执行有效期过滤 |
|
||
| 单元 | `tests/unit/service/test_knowledge_reference.py` | token 签名校验;过期/伪造/跨用户 token 返回 404;不返回全文 |
|
||
| 契约 | `tests/contract/test_knowledge_tool_contract.py` | 工具输入模型 `extra="forbid"`;返回结构稳定 |
|
||
| 集成 | `tests/integration/test_knowledge_mysql.py` | 真实 MySQL:只返回 `published + active + 有效期内`;未发布条目不可检索 |
|
||
|
||
**实际落地的测试文件(2026-09-11 复核,与上表计划名不完全一致)**:
|
||
`tests/unit/service/test_knowledge_retrieval_service.py`、`tests/unit/service/test_knowledge_reference.py`、
|
||
`tests/unit/service/test_knowledge_tool.py`、`tests/unit/service/test_knowledge_ingest_service.py`、
|
||
`tests/unit/service/test_knowledge_management_service.py`、`tests/unit/core/test_knowledge_contracts.py`、
|
||
`tests/unit/model/test_knowledge_orm.py`、`tests/unit/worker/test_knowledge_vector_worker.py`、
|
||
`tests/unit/worker/test_runtime_knowledge_wiring.py`。降级与契约覆盖并入上述文件,未单独建
|
||
`test_knowledge_fallback.py` 与 `tests/contract/test_knowledge_tool_contract.py`。
|
||
|
||
---
|
||
|
||
## 6. 风险与规避
|
||
|
||
| # | 风险 | 规避 |
|
||
|---|---|---|
|
||
| 1 | **照搬客服文档的路由代码** | 那份是独立 demo(含教学注释、集合名硬编码在常量区),必须改为配置驱动并纳入 `ToolExecutor` 治理 |
|
||
| 2 | **embedding 维度不匹配** | 适配器显式校验 `len(vector) == vector_dim`,不符时失败关闭;不静默降级为"无结果" |
|
||
| 3 | **降级路径漏掉有效期过滤** | 03 §6.2 明确要求降级仍执行 `published + active + 有效期`;写进单测断言 |
|
||
| 4 | **已下线知识仍被检索命中** | Milvus 命中后必须回查 `fin_knowledge_meta` 二次校验(步骤 3 第 2 条) |
|
||
| 5 | **pymilvus 异步 API 不确定** | 实现前先确认 `AsyncMilvusClient` 是否可用;不可用则用 `anyio.to_thread.run_sync` 包装同步客户端,**不得阻塞事件循环** |
|
||
| 6 | **引用 token 不签名** | 必须 HMAC 签名,否则用户可拼 token 越权读取他人知识引用 |
|
||
| 7 | **Milvus 建连时机** | 延迟导入 + 延迟建连(照 `EastmoneyAdapterFactory`),避免底座启动即依赖 Milvus |
|
||
| 8 | **集合未在 MySQL 建 ORM 映射** | `fin_knowledge_meta` 属 00 基线表,需补 ORM 映射;只读,不改表结构 |
|
||
|
||
---
|
||
|
||
## 7. 验收标准
|
||
|
||
```text
|
||
功能
|
||
[ ] 意图 faq / product_inquiry / policy_explain 分别命中对应集合,TopK 3/5/5
|
||
[ ] 返回结果携带 knowledge_id / collection / title / snippet / score
|
||
[ ] ToolExecutor 自动生成 SourceReference 并回填结果
|
||
|
||
降级
|
||
[ ] 停掉 Milvus,同一查询仍返回已发布且在有效期内的条目
|
||
[ ] 降级结果 degraded=True 且 degradation_reason 可辨识
|
||
[ ] embedding 维度不符时失败关闭(不返回错误结果)
|
||
|
||
合规与安全
|
||
[ ] 未发布(pending/approved/archived)知识不可检索
|
||
[ ] 已过期知识不可检索
|
||
[ ] 非白名单集合名被拒绝
|
||
[ ] 跨客户引用 token 返回 404
|
||
[ ] 检索调用写入 interaction_audit(由 ToolExecutor 统一完成)
|
||
|
||
质量
|
||
[ ] pytest 全绿;ruff 0;mypy strict 0
|
||
[ ] Milvus 不可用不影响底座启动
|
||
[ ] 知识集合路由可通过配置发布调整,无需改代码
|
||
```
|
||
|
||
---
|
||
|
||
## 8. 工作量与排期
|
||
|
||
| 阶段 | 内容 | 预估 |
|
||
|---|---|---|
|
||
| 步 1 | Embedding 出口 | 0.5 天 |
|
||
| 步 2 | Milvus 适配器 | 0.5 天 |
|
||
| 步 3 | 检索服务(含白名单、二次校验、重排) | 1 天 |
|
||
| 步 4 | 配置驱动路由 | 0.3 天 |
|
||
| 步 5 | MySQL 降级检索 | 0.5 天 |
|
||
| 步 6 | 工具注册与接线 | 0.3 天 |
|
||
| 步 7 | 知识引用占位修复(含 token 签名) | 0.5 天 |
|
||
| 步 8 | 测试(单元 + 契约 + 集成) | 1 天 |
|
||
| **合计** | | **约 4.6 人日** |
|
||
|
||
**建议分批**:步 1-3 为核心链路(2 天),完成后即可跑通检索;步 4-8 为接线与加固,可并行。
|
||
|
||
---
|
||
|
||
## 9. 前置条件(开工前必须确认;已于实施阶段落实)
|
||
|
||
| # | 条件 | 说明 |
|
||
|---|---|---|
|
||
| 1 | **可用的 embedding 模型端点** | `base_url` + `model_name` + `secret_ref=env:XXX`,且维度与集合一致(当前约定 1024) |
|
||
| 2 | **Milvus 实例与三集合已建** | 集合名、维度、度量方式(COSINE)、字段名需与第 4.2 节一致 |
|
||
| 3 | **知识数据已入库** | `fin_knowledge_meta.content_text` 已回填、`review_status='published'`、有效期字段可用(否则检索无数据,降级也查不到) |
|
||
| 4 | `pymilvus` 异步客户端可用性确认 | 决定步骤 2 的实现方式(原生异步 vs 线程包装) |
|
||
| 5 | **集合路由配置的发布责任人** | 步骤 4 的配置需走双人审核发布流程 |
|
||
|
||
---
|
||
|
||
## 10. 变更记录
|
||
|
||
| 版本 | 日期 | 变更 |
|
||
|---|---|---|
|
||
| v1.0 | 2026-09-09 | 首版:现状诊断、目标链路、数据契约、8 步实施方案、风险规避、验收标准、工作量与前置条件 |
|
||
| v1.0(状态复核) | 2026-09-11 | 只更新状态口径,不改方案正文:§2.2 缺失清单 7 项全部闭环;引用令牌校验、Milvus 读写适配器、`KnowledgeRetrievalService`、`query_knowledge` 工具均已落地;Milvus 检索已在 Phase 1 验收命中(score 0.7837) |
|