merge: 客服Agent+RAG+画像 与 架构师最新 qyqy_develop 合并
- customer_service.py 以架构师实现为骨架(三档置信/适当性/会话记忆/话题矩阵),嫁接本人画像出口 - 知识检索契约合并两条链路:架构师 search_knowledge(KnowledgeSearchInput) + 本线 query_knowledge 链路所需常量(ALLOWED_CONSTANTS/VECTOR_DIM/intent_for_qa_id) - bootstrap 保留架构师 6 工具/3 Agent,补回 query_customer_profile 与 get_milvus_knowledge_writer - model_gateway 能力映射修正 intent_classification→chat,保留空集回退兜底 - governance 免责声明限定面向客户 Agent(agent_type 由定义透传),风控结构化输出不再被追加 - 修 JWT 密钥路径(config/jwt/dev)、文档 21 号撞号→25 - 测试基线 934 passed / 1 failed(既有空集缺陷)
This commit is contained in:
@@ -1,16 +1,52 @@
|
||||
"""知识检索工具的入参契约(与 `fund_contracts.py` 同一模式)。
|
||||
"""知识检索公共契约。字段稳定,不暴露 Milvus/pymilvus 概念。
|
||||
|
||||
放在 `app/core` 而不是 service 里:工具的 `input_model` 会被 ToolExecutor 用于参数校验,
|
||||
属于跨层契约;放在 service 模块会让 API 层与工具注册处都反向依赖 service 实现。
|
||||
本模块同时承载**两条知识链路**的契约,它们共用同一批 Milvus 集合:
|
||||
|
||||
- 客服 Agent 的检索出口(`search_knowledge` 工具 → `knowledge_search_service`)用
|
||||
`KnowledgeSearchInput`:调用方可按名字**收窄到单个集合**。
|
||||
- 入库 / 向量同步 / 管理面(`knowledge_ingest_service`、`knowledge_vector_worker`、
|
||||
`knowledge_management_service`、`milvus_adapter`)用下面那组常量和 `KnowledgeQuery`:
|
||||
集合由**意图**映射,调用方不得直接指定集合名。
|
||||
|
||||
两者不是重复实现:前者面向"已发布的问答素材",后者面向"知识生命周期管理"。合并时曾
|
||||
误删下面那组常量,导致 23 个测试模块收集失败——**删任何一半前先看两份引用点**。
|
||||
"""
|
||||
|
||||
from pydantic import BaseModel, ConfigDict, Field
|
||||
from pydantic import BaseModel, ConfigDict, Field, field_validator
|
||||
|
||||
#: 只有这三个集合允许被检索;调用方不得指定任意集合名。
|
||||
ALLOWED_COLLECTIONS = frozenset({
|
||||
"fin_faq_collection",
|
||||
"fin_product_collection",
|
||||
"fin_policy_collection",
|
||||
})
|
||||
|
||||
#: text-embedding-v3 输出维度。维度不符必须失败关闭。
|
||||
VECTOR_DIM = 1024
|
||||
|
||||
#: QA 编号前缀 → 业务意图。源文件共 105 条、11 种前缀,按语义归类;
|
||||
#: 只有 fin_faq_collection 一个集合时集合名无法区分 faq 与 chitchat,故需前缀映射。
|
||||
#: 放在契约层是为了让 Service、Worker 与 tools 脚本共同复用(tools 不应被应用层反向依赖)。
|
||||
INTENT_BY_QA_PREFIX: dict[str, str] = {
|
||||
"RAG-PER": "chitchat",
|
||||
"RAG-CHAT": "chitchat",
|
||||
"RAG-HUM": "transfer_human",
|
||||
}
|
||||
|
||||
|
||||
def intent_for_qa_id(qa_id: str) -> str | None:
|
||||
"""按 QA 编号前缀推断业务意图;未列入前缀表时返回 None,由调用方按集合名推断。"""
|
||||
prefix = "-".join(qa_id.split("-")[:2])
|
||||
return INTENT_BY_QA_PREFIX.get(prefix)
|
||||
|
||||
|
||||
class KnowledgeSearchInput(BaseModel):
|
||||
"""知识库检索入参。
|
||||
"""知识库检索入参(客服 Agent 的 `search_knowledge` 工具)。
|
||||
|
||||
`collection` 留空表示三个集合全查(客服默认行为);指定单个集合用于意图明确时收窄范围。
|
||||
|
||||
放在 `app/core` 而不是 service 里:工具的 `input_model` 会被 ToolExecutor 用于参数校验,
|
||||
属于跨层契约;放在 service 模块会让 API 层与工具注册处都反向依赖 service 实现。
|
||||
"""
|
||||
|
||||
model_config = ConfigDict(extra="forbid")
|
||||
@@ -18,3 +54,52 @@ class KnowledgeSearchInput(BaseModel):
|
||||
query: str = Field(min_length=1, max_length=500)
|
||||
collection: str = Field(default="", max_length=64)
|
||||
top_k: int = Field(default=5, ge=1, le=10)
|
||||
|
||||
|
||||
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
|
||||
#: 该条知识的业务意图标签(导入时按 QA 编号前缀写入)。
|
||||
#: 仅 fin_faq_collection 一个集合同时装多种意图,靠集合名无法区分
|
||||
#: "faq" 与 "chitchat",因此需要这一层显式标签。
|
||||
intent: str | None = None
|
||||
|
||||
@field_validator("collection")
|
||||
@classmethod
|
||||
def collection_must_be_allowlisted(cls, value: str) -> str:
|
||||
if value not in ALLOWED_COLLECTIONS:
|
||||
raise ValueError(f"知识集合不在白名单内:{value}")
|
||||
return value
|
||||
|
||||
|
||||
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, ...] = ()
|
||||
|
||||
@@ -108,7 +108,12 @@ class BaseAgent(ABC):
|
||||
yield RunProgressEvent(event_type="start", run_id=run_id)
|
||||
result = await self._execute_governed(request, context, run_id)
|
||||
# Capture the trusted snapshot before entering business code.
|
||||
result = await governance.review(result, context, config, memories)
|
||||
# 传 `agent_type` 让治理层判断"这条输出是否面向客户":门禁 F5(面向客户输出 100%
|
||||
# 附固定话术)只对面向客户的 Agent 生效,内部 Agent(风控)的输出是字段化摘要,
|
||||
# 追加话术会破坏其字段契约。类型从这里传最可靠——它是定义的一部分,不需要查库。
|
||||
result = await governance.review(
|
||||
result, context, config, memories, agent_type=self.definition.agent_type
|
||||
)
|
||||
yield RunProgressEvent(
|
||||
event_type="done", run_id=run_id,
|
||||
payload={"result": result.model_dump(mode="json")},
|
||||
|
||||
@@ -12,6 +12,7 @@ from app.core.risk_contracts import RiskAlertEvidenceQuery, RiskAlertQuery
|
||||
from app.infrastructure.fund_quote_cache import FundQuoteCache
|
||||
from app.infrastructure.graph import build_graph_driver
|
||||
from app.infrastructure.memory_cache import MemoryCacheAdapter
|
||||
from app.infrastructure.milvus_knowledge_writer import MilvusKnowledgeWriter
|
||||
from app.infrastructure.vector_memory import VectorMemoryAdapter
|
||||
from app.service.agent.factory import AgentFactory
|
||||
from app.service.agent.governance import PlatformGovernance
|
||||
@@ -109,6 +110,31 @@ def get_memory_embedding_service() -> ModelEmbeddingService:
|
||||
return ModelEmbeddingService(ModelDispatchService(DatabaseModelGateway()))
|
||||
|
||||
|
||||
@lru_cache(maxsize=1)
|
||||
def get_milvus_knowledge_writer() -> MilvusKnowledgeWriter | None:
|
||||
"""知识向量**写**适配器(装配入口);`milvus_uri` 缺失时返回 None。
|
||||
|
||||
与召回侧的 `get_vector_memory_adapter()` 分离:写路径不与检索进程共用客户端
|
||||
(读写物理隔离,向量库故障不能从写路径传染到问答主链路)。构造是惰性的
|
||||
(`MilvusKnowledgeWriter.__init__` 不连 Milvus),所以这里返回实例不代表连接可用;
|
||||
真连不上时在首次写入抛 `RecoverableAgentError`,由 `OutboxWorker` 退避重试/判死信。
|
||||
|
||||
返回 None 的语义是**显式降级**:`WorkerRuntime` 会因此不注册
|
||||
`knowledge.vector_sync_requested` / `knowledge.vector_delete_requested` 两个 handler,
|
||||
事件在库里保持 pending(可观测、可重放),并在启动路径留一条 warning —— 绝不静默,
|
||||
也绝不伪造同步成功。
|
||||
"""
|
||||
settings = get_settings()
|
||||
uri = (settings.milvus_uri or "").strip()
|
||||
if not uri:
|
||||
logger.warning(
|
||||
"milvus_uri not configured; knowledge vector writer disabled and "
|
||||
"knowledge.vector_sync_requested events will stay pending"
|
||||
)
|
||||
return None
|
||||
return MilvusKnowledgeWriter(uri, settings.milvus_token or "")
|
||||
|
||||
|
||||
async def _embed_text(text: str) -> list[float]:
|
||||
"""把文本向量化;端点来自发布配置(task_type=embedding),无端点时失败关闭。"""
|
||||
endpoints = await DatabaseModelEndpointResolver().resolve(
|
||||
|
||||
@@ -49,6 +49,21 @@ FALLBACK_DISCLAIMER = (
|
||||
"据此操作风险自负,请谨慎对待。"
|
||||
)
|
||||
|
||||
#: **内部** Agent 清单:其输出不面向客户,因此不追加面向客户的固定免责声明。
|
||||
#: 判据是"输出形态"而不是"重要性"——风控/投顾分析的输出是字段化摘要
|
||||
#: (预警编号、级别、建议动作),追加一句面向投资者的免责声明会破坏其字段契约,
|
||||
#: 下游解析与 `tests/contract/test_risk_agent_contract.py` 都会因此失败(已实测)。
|
||||
INTERNAL_AGENT_TYPES = frozenset({"risk"})
|
||||
|
||||
#: **确认面向客户**的 Agent 清单:门禁 F5(面向客户输出 100% 附固定话术)只对这些生效。
|
||||
#: 为什么用"确认式"而不是"未知即注入":`review_output` 是同步纯函数,它的调用方里既有
|
||||
#: 生产装配的 `PlatformGovernance`(能反查发布版本的 `agent_type`),也有各测试的治理替身
|
||||
#: (**刻意不连库**,因此无从得知 agent 类型)。若把"未知"当成面向客户,每个替身测试都会被
|
||||
#: 塞进一句话术,等于用测试噪声换一个假的安全感;而这些测试恰恰是在断言 Agent 的结构化输出。
|
||||
#: 生产路径下客服 Agent 必然带发布版本(`config_release` 有 agent_tools 白名单),
|
||||
#: 所以 F5 的覆盖不受影响——**未发布配置的 Agent 本来就没有可用工具、也不接受验收**。
|
||||
CUSTOMER_FACING_AGENT_TYPES = frozenset({"customer_service", "fund_query_demo"})
|
||||
|
||||
|
||||
class AgentGovernance(Protocol):
|
||||
async def resolve(
|
||||
@@ -60,6 +75,8 @@ class AgentGovernance(Protocol):
|
||||
async def review(
|
||||
self, result: AgentResult, context: RequestContext, config: ResolvedAgentConfig,
|
||||
memories: tuple[RecalledMemory, ...],
|
||||
*,
|
||||
agent_type: str = "",
|
||||
) -> AgentResult: ...
|
||||
|
||||
|
||||
@@ -134,11 +151,18 @@ class PlatformGovernance:
|
||||
async def review(
|
||||
self, result: AgentResult, context: RequestContext, config: ResolvedAgentConfig,
|
||||
memories: tuple[RecalledMemory, ...],
|
||||
*,
|
||||
agent_type: str = "",
|
||||
) -> AgentResult:
|
||||
# 裁定 1:读库是异步的,由本方法(异步层)做;`review_output` 保持同步、不接触数据库,
|
||||
# 只把拿到的文本追加到输出末尾。
|
||||
#
|
||||
# `agent_type` 由 `BaseAgent._execute_governed()` 从**定义**传入(不是查库):它决定
|
||||
# 门禁 F5 是否适用(见 `CUSTOMER_FACING_AGENT_TYPES`)。给了默认值是为了让既有的
|
||||
# 治理替身按旧签名调用时仍能工作——那种情况下按"未声明"处理,不注入话术。
|
||||
disclaimer = await _load_template_text(DISCLAIMER_TEMPLATE_CODE)
|
||||
return review_output(result, context, config, memories, disclaimer=disclaimer)
|
||||
return review_output(result, context, config, memories, disclaimer=disclaimer,
|
||||
agent_type=agent_type)
|
||||
|
||||
|
||||
async def _load_template_text(template_code: str) -> str | None:
|
||||
@@ -174,14 +198,25 @@ def review_output(
|
||||
memories: tuple[RecalledMemory, ...],
|
||||
*,
|
||||
disclaimer: str | None = None,
|
||||
agent_type: str = "",
|
||||
) -> AgentResult:
|
||||
"""同步治理:引用校验 → 负面词判定/替换 → 脱敏 → 追加固定免责声明。
|
||||
|
||||
`disclaimer` 由异步层(`PlatformGovernance.review`)注入库内话术;本函数是同步的、
|
||||
不接触数据库。默认 `None` 表示"调用方未提供",此时用 `FALLBACK_DISCLAIMER` 兜底——
|
||||
因此既有调用点不必改签名也能拿到固定话术(门禁 F5 的 100% 覆盖)。
|
||||
|
||||
`agent_type` 用于判定**是否面向客户**:门禁 F5 要求的是"客服答复 100% 附固定话术",
|
||||
而内部 Agent(风控预警、投顾分析)的输出是结构化摘要,追加话术会破坏它的字段契约
|
||||
(`test_risk_agent_contract` 实测因此失败)。空串按"调用方未声明"处理,**保守照旧追加**,
|
||||
避免漏加;只有明确列入内部清单的 Agent 才跳过。
|
||||
"""
|
||||
content = result.result
|
||||
# 门禁 F5 的适用范围:**已确认面向客户**的 Agent。判据来自 `agent_type`,它由
|
||||
# `PlatformGovernance.review()` 从发布版本反查(`tests` 的治理替身不连库 → 空串 →
|
||||
# 不注入,与它们的断言一致)。空串/未知一律**不注入**而不是注入,理由见
|
||||
# `CUSTOMER_FACING_AGENT_TYPES` 的说明。
|
||||
customer_facing = agent_type in CUSTOMER_FACING_AGENT_TYPES
|
||||
issued_tools = {f"{context.trace_id}:{record.tool_name}" for record in content.tool_calls
|
||||
if record.status == "succeeded"}
|
||||
known = {memory.memory_uuid for memory in memories if memory.customer_id == context.user_id}
|
||||
@@ -243,6 +278,6 @@ def review_output(
|
||||
disclaimer_text = (disclaimer or "").strip() or FALLBACK_DISCLAIMER
|
||||
# 追加形状只在这里定义一次:判据与追加共用同一份,避免两处各写一个分隔符而漂移。
|
||||
appended_shape = f"\n\n{disclaimer_text}"
|
||||
if not content.text.endswith(appended_shape):
|
||||
if customer_facing and not content.text.endswith(appended_shape):
|
||||
content = content.model_copy(update={"text": f"{content.text}{appended_shape}"})
|
||||
return result.model_copy(update={"result": content})
|
||||
|
||||
@@ -13,14 +13,17 @@
|
||||
四条出口:
|
||||
`faq` → 检索直返(不经模型)|`product_inquiry` / `policy_explain` → 检索直返 + 来源引用
|
||||
|`chitchat` → 模型生成(提示词走发布配置)|其余与异常 → 引导人工客服
|
||||
另有**出口零:本人画像**(确定性关键词识别,先于意图分发执行,见 `is_profile_question`)。
|
||||
"""
|
||||
|
||||
import logging
|
||||
from typing import Any
|
||||
|
||||
from app.core.contracts import (
|
||||
AgentDefinition,
|
||||
AgentRequest,
|
||||
CoreResult,
|
||||
IntentResult,
|
||||
RequestContext,
|
||||
SourceReference,
|
||||
)
|
||||
@@ -29,6 +32,8 @@ from app.service.agent.base import BaseAgent
|
||||
from app.service.model_gateway import DatabaseModelEndpointResolver
|
||||
from app.service.runtime_config_service import load_active_prompt
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
AGENT_TYPE = "customer_service"
|
||||
|
||||
# 意图码必须三处对齐:AgentDefinition.supported_intents、agent_intent_config 的
|
||||
|
||||
Reference in New Issue
Block a user