72 KiB
客服 Agent + RAG 知识库实施计划(qyqy_develop 版)
For agentic workers: REQUIRED SUB-SKILL: 用
subagent-driven-development(推荐)或executing-plans逐任务实施。步骤用- [ ]复选框跟踪。本计划取代
2026-09-10-customer-service-agent-implementation.md(旧 12-Task 计划,为develop底座 + 只导 105 条 QA 设计,两个前提均已失效)。
Goal: 在 qyqy_develop 底座上交付"RAG 知识库 + 智能客服 Agent",对齐老师《需求文档-修改版》Phase 1 的 7 条验收标准,产出可答辩的端到端链路。
Architecture: 沿用底座既有 MVC+S 与统一 Agent 骨架。知识侧:文档解析 → 向量化 → Milvus 三集合(白名单)→ 读路径经 query_knowledge 只读工具暴露;文件存储走可替换的存储抽象(本地文件系统实现,MinIO 实现待实例就绪)。Agent 侧:CustomerServiceAgent 继承 BaseAgent,handle() 内先跑确定性安全路由(不依赖模型),再按置信度三档决定"直返/复述/兜底",全部输出经底座 governance.review_output()。
Tech Stack: Python 3.13、FastAPI、Pydantic v2、SQLAlchemy 2 async、Alembic、MySQL 8、Redis、pymilvus 2.6(AsyncMilvusClient)、DashScope text-embedding-v3(1024 维)、pytest + pytest-asyncio、ruff、mypy strict。
Spec:
docs/superpowers/analysis/2026-09-10-现状与差距分析.md(需求 vs 现状,含 10 张表核对)docs/superpowers/analysis/客服Agent专项设计方案-分析报告.md(22 条差异 + 9 条文档缺陷)docs/superpowers/analysis/2026-09-10-合规层种子数据缺失.md(合规种子列清单)docs/superpowers/specs/2026-09-10-knowledge-retrieval-infra-design.md(spec A,已修正)docs/superpowers/specs/2026-09-10-customer-service-agent-design.md(spec B)- 老师需求:
Desktop\金融\需求文档-修改版.html(优先级 1) - 底座接入:
docs/19-业务Agent接入实操(示例验证版).md、docs/16-Agent组员入门易懂版说明.md
Global Constraints
来自 AGENTS.md、docs/14、docs/19 与两份 spec,每个 Task 隐含包含:
- 数据库以
docs/00-新数据库基线设计.md为不可变基线;允许新增表/新增字段,禁止重命名、删除、复用已有字段,禁止改变已有字段类型、可空性、业务含义。 - 架构固定 MVC+S:Controller → Service → Repository → Model。业务 Agent 属 Service 层,不得直接建 Session、不得直接查 Repository/Model、不得直连 MySQL/Redis/Milvus/Neo4j。
- 业务 Agent 必须继承
BaseAgent并由AgentFactory创建,在app/service/agent/bootstrap.py的register_business_agents(factory)注册一行;不得覆盖execute/call_tool/generate_with_model/resolve_config等治理方法(BaseAgent.__init_subclass__会抛TypeError)。 - Agent 只能"提出"转人工(
CoreResult.transfer_required=True+transfer_reason),不得创建/分配/关闭工单(专项设计有 10 处要求"生成工单",与此冲突,见 §11-4,以本约束为准)。 - 工具白名单必须发布到
status='active'的config_release(namespace=agent_tools、config_key=<agent_type>:<intent>),否则工具失败关闭。代码里的allowed_tools是上限,配置只能收窄。 - 向量维度固定 1024;维度不符必须失败关闭,不得静默返回空结果。
- Milvus 集合白名单固定
fin_faq_collection/fin_product_collection/fin_policy_collection;集合名不得由调用方指定。metric_type必须与索引一致(COSINE)。 - Embedding 模型固定
text-embedding-v3;密钥只走secret_ref(格式^env:[A-Z][A-Z0-9_]{0,100}$),不得写入.py、数据库明文字段或日志。 - 合规红线:不代客交易;7 个零容忍负面词(保本/稳赚/无风险/保证收益/预期收益率/年化收益率/安全);面向客户输出末尾必须附固定免责声明;收益口径对外一律称「业绩比较基准」,不得出现「预期收益率」。
- 硬编码话术 vs 配置化:安全关键路径(P0/P1/P2 固定话术与联系方式)硬编码为代码常量,不查知识库、不依赖配置——这是刻意的取舍(spec B §5)。
handle()必须显式返回CoreResult.intent,否则会被底座自动分类结果覆盖。- 禁止任何 git 写操作(
add/commit/checkout/stash/reset):项目交接文档记录过用户不满自行提交的教训,且工作区有用户自己 staged 的文件。完成后在报告里给出建议的提交信息即可。 - 解释器固定
.\.venv\Scripts\python.exe(Python 3.13.5)。
环境前提(2026-09-10 实测)
| 项 | 状态 |
|---|---|
| 分支 | qyqy_develop @ 6516ccb |
| 数据库 | jr_agent,51 张表,audit_schema.py 通过 |
| Redis | 127.0.0.1:6379 可用 |
| Milvus | ❌ 未运行(19530 关闭)→ 需本地 Docker 起 Milvus |
| MinIO | ❌ 未安装未运行(9000/9001 关闭)→ 用存储抽象 + 本地实现 |
| embedding 端点 | ✅ embedding-primary 已 active(dashscope / text-embedding-v3 / capabilities=["embedding"]),1024 维真机验证过 |
| 工具白名单配置 | ❌ config_release 与 platform_config_item 均为空 → 必须先发布,否则工具被拒 |
| 测试账号 | 9001 客户 / 9002 风控 / 9003 管理员(含 model-endpoint:manage) |
| 测试基线 | 495 passed + 1 skipped + 1 failed(既有);ruff 全过;mypy 有 143 个分支自带既有错误(与本次无关) |
| 既有失败测试 | tests/unit/repository/test_fund_readonly_contract.py::test_fund_models_cover_all_fin_tables_with_expected_columns——已核实与本次改动无关(文件与 HEAD 逐字相同,且从不 import app.model.fund;属 import 顺序依赖的既有缺陷,本次不修) |
控制者已做的架构裁定(记录在案,可回退)
- MinIO 用可替换的存储抽象(
DocumentStorage协议),一期提供LocalDocumentStorage文件系统实现;MinioDocumentStorage留待实例就绪后补。- 理由:MinIO 未安装未运行,硬依赖会阻塞整个知识库链路;而需求要求的是"MinIO 连接与文件上传接口"的能力,抽象 + 本地实现能让上传/列表/删除/归档全链路先跑通可验收。
- 判断错的代价:答辩若被问"用的是不是真 MinIO",需现场说明实现已就绪、仅切换配置;若老师要求必须真 MinIO,补一个适配器即可(接口不变)。
- 不再导入
高频问答对.txt(用户明确说它只是参考文档)。一期知识内容用客服Agent知识库_QA问答对_v5_RAG发布候选版.txt(105 条)。- 风险:老师 Phase 1 验收点名的测试问题「基金申购后多久确认」不在 105 条内(已核实)。Task 3 会手工补 1 条该问答,使验收用例可命中。
- Milvus 集合字段名沿用 spec A 的
knowledge_id/title/snippet/tags/version/embedding,不采用专项设计 §4.2 的doc_id/content。- 理由:专项设计自身三套 schema 互相矛盾(§13.1),且
doc_id/content与已完成的KnowledgeHit契约不兼容;spec A 的 schema 三集合逐字统一。 - 判断错的代价:若答辩要求对齐专项设计字段名,需改 Milvus schema +
KnowledgeHit,成本约 1 个 Task。
- 理由:专项设计自身三套 schema 互相矛盾(§13.1),且
fin_faq_collection承载全部 105 条,fin_product_collection/fin_policy_collection只建不填(与 spec A 一致,与专项设计 §4.1 的门禁 F1 冲突——按 §11-1 处置)。
待用户确认(不阻塞开工,但需尽早给)
| # | 事项 | 影响 |
|---|---|---|
| 1 | Milvus 部署方式:本地 Docker 起 standalone,还是用某个现成实例? | Task 4/5/6 的真机联调 |
| 2 | MinIO 实例(地址/凭据)是否已有 | Task 6 的 MinIO 适配器与真机验证 |
| 3 | 免责声明文案最终用哪版(专项设计 §6.3 给了一版,胜宇一期边界基线也可能有) | Task 2 的种子数据 |
| 4 | 转人工率的统计口径(专项设计门禁 Q3 要求 ≤20%,但 Agent 不建工单) | 验收口径,Task 10 |
Task 1: 合规层种子数据(必须先做)
Files:
- Create:
tools/seed_compliance_baseline.py - Test:
tests/integration/test_compliance_seed_mysql.py
Interfaces:
- Consumes:
SessionFactory(app/infrastructure/db.py)、agent_negative_word/agent_reply_template表、sys_user.id=9003 - Produces: 7 条
agent_negative_word(+变体)与 6 条agent_reply_template记录,全部status='active'且reviewer_id/reviewed_at非空
为什么第一个做:
governance.pyL62-65 的查询条件是WHERE status='active' AND reviewer_id IS NOT NULL AND reviewed_at IS NOT NULL; 现库这两张表都是 0 行,所以 7 个零容忍词一个都不生效,客服 Agent 说"稳赚"系统不拦。 这是合规门禁 F4/F5 的前置,且改动面小、无依赖,适合先做。
- Step 1: 确认枚举与约束取值(不要臆造)
Run:
.\.venv\Scripts\python.exe -c "import pymysql;c=pymysql.connect(host='127.0.0.1',port=3306,user='root',password='123456',database='jr_agent');cur=c.cursor();cur.execute('SHOW CREATE TABLE agent_negative_word');print(cur.fetchone()[1]);cur.execute('SHOW CREATE TABLE agent_reply_template');print(cur.fetchone()[1])"
Expected: 输出两张表的完整 DDL,从中读出 match_type/severity/status/scene 的 CHECK 约束取值集。
把读出的取值记进报告,后续 SQL 只用这些值。
- Step 2: 写失败测试
创建 tests/integration/test_compliance_seed_mysql.py:
import pytest
from sqlalchemy import text
from app.infrastructure.db import SessionFactory
pytestmark = pytest.mark.integration
NEGATIVE_WORDS = ("保本", "稳赚", "无风险", "保证收益", "预期收益率", "年化收益率", "安全")
@pytest.mark.asyncio
async def test_all_seven_negative_words_are_active_and_reviewed() -> None:
async with SessionFactory() as session:
rows = (await session.execute(text(
"SELECT word_pattern FROM agent_negative_word "
"WHERE status='active' AND reviewer_id IS NOT NULL AND reviewed_at IS NOT NULL"
))).scalars().all()
patterns = " ".join(rows)
for word in NEGATIVE_WORDS:
assert word in patterns, f"负面词未生效:{word}"
@pytest.mark.asyncio
async def test_reply_templates_cover_all_six_scenes() -> None:
async with SessionFactory() as session:
rows = (await session.execute(text(
"SELECT template_code, scene FROM agent_reply_template "
"WHERE status='active' AND reviewer_id IS NOT NULL AND reviewed_at IS NOT NULL"
))).all()
scenes = {row.scene for row in rows}
assert len(rows) >= 6, f"话术模板不足 6 条:{len(rows)}"
assert len(scenes) >= 6, f"场景不足 6 类:{scenes}"
- Step 3: 跑测试确认失败
Run: .\.venv\Scripts\python.exe -m pytest tests/integration/test_compliance_seed_mysql.py -q
Expected: FAIL — 两处断言失败(表为空)
- Step 4: 写种子脚本
创建 tools/seed_compliance_baseline.py(列名与取值必须来自 Step 1 读出的 DDL):
"""合规层种子数据:7 个零容忍负面词 + 6 类固定话术。
背景:docs/02 §10.2/§10.3 要求这些数据,但迁移只建表、没有 seed,
导致 governance.py 的 WHERE 条件(status='active' AND reviewer_id/reviewed_at 非空)
一条都查不到,合规过滤静默失效。
幂等:按 rule_code / template_code upsert;可重复执行。
"""
from __future__ import annotations
import asyncio
import sys
from datetime import UTC, datetime
from pathlib import Path
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))
from sqlalchemy import text # noqa: E402
from app.infrastructure.db import SessionFactory # noqa: E402
ADMIN_ID = 9003
#: docs/02 §10.2 逐字要求的 7 个词 + 专项设计 §6.5 L2132 建议的变体
NEGATIVE_RULES: tuple[tuple[str, str, str, str], ...] = (
("NEG-001", "保本", "misleading", "资管新规后非保本"),
("NEG-002", "稳赚", "misleading", "误导性表述"),
("NEG-003", "无风险", "absolute", "绝对化表述"),
("NEG-004", "保证收益", "guarantee", "禁止刚兑承诺"),
("NEG-005", "预期收益率", "yield", "《理财销售办法》明文禁止"),
("NEG-006", "年化收益率", "yield", "监管处罚点名措辞"),
("NEG-007", "安全", "absolute", "绝对化表述"),
# 变体(专项设计 §6.5 评审建议)
("NEG-008", "零风险", "absolute", "绝对化表述变体"),
("NEG-009", "稳赚不赔", "misleading", "误导性表述变体"),
("NEG-010", "躺着赚", "misleading", "误导性表述变体"),
("NEG-011", "坐享收益", "misleading", "误导性表述变体"),
)
REPLY_TEMPLATES: tuple[tuple[str, str, str, str], ...] = (
("TPL_DISCLAIMER", "disclaimer", "固定免责声明",
"本内容仅为投资分析参考,不构成任何直接投资建议,不构成对任何产品的收益承诺,"
"据此操作风险自负,请谨慎对待。"),
("TPL_AI_NOTICE", "ai_notice", "AI 生成标识",
"本回答由 AI 生成,仅供参考。"),
("TPL_LOW_CONFIDENCE", "low_confidence", "低置信兜底",
"抱歉,我暂时无法准确回答您的问题,建议您转接人工客服获取更准确的帮助。"),
("TPL_COMPLIANCE_REFUSAL", "compliance_refusal", "合规拒答",
"根据监管要求,我不能对收益做出任何承诺。本产品为非保本浮动收益产品,请以产品说明书为准。"),
("TPL_TRANSFER_HUMAN", "transfer_human", "转人工提示",
"已为您转接人工客服,工作时间为工作日 09:00-18:00,客服电话 15936583816。"),
("TPL_SYSTEM_BUSY", "system_busy", "系统繁忙/模型故障",
"系统繁忙,暂时无法回答,请稍后重试或联系人工客服 15936583816。"),
)
async def seed() -> None:
now = datetime.now(UTC).replace(tzinfo=None)
async with SessionFactory() as session, session.begin():
for rule_code, word, category, reason in NEGATIVE_RULES:
await session.execute(text(
"INSERT INTO agent_negative_word"
" (rule_code, word_pattern, match_type, category, severity,"
" applicable_agents, safe_reply_template_code, status, version,"
" created_by, reviewer_id, reviewed_at, created_at, updated_at)"
" VALUES (:rule_code, :word, 'contains', :category, 'high',"
" :agents, 'TPL_COMPLIANCE_REFUSAL', 'active', 1,"
" :admin, :admin, :now, :now, :now)"
" ON DUPLICATE KEY UPDATE word_pattern=:word, category=:category,"
" status='active', reviewer_id=:admin, reviewed_at=:now, updated_at=:now"
), {
"rule_code": rule_code, "word": word, "category": category,
"agents": '["customer_service"]',
"admin": ADMIN_ID, "now": now,
})
for code, scene, title, content in REPLY_TEMPLATES:
await session.execute(text(
"INSERT INTO agent_reply_template"
" (template_code, scene, title, content_text, variables, locale,"
" version, status, active_key, created_by, reviewer_id, reviewed_at,"
" created_at, updated_at)"
" VALUES (:code, :scene, :title, :content, '{}', 'zh-CN',"
" 1, 'active', :active_key, :admin, :admin, :now, :now, :now)"
" ON DUPLICATE KEY UPDATE content_text=:content, title=:title,"
" status='active', reviewer_id=:admin, reviewed_at=:now, updated_at=:now"
), {
"code": code, "scene": scene, "title": title, "content": content,
"active_key": f"{scene}:zh-CN:1",
"admin": ADMIN_ID, "now": now,
})
async def main() -> None:
await seed()
print(f"seeded {len(NEGATIVE_RULES)} negative rules, {len(REPLY_TEMPLATES)} reply templates")
if __name__ == "__main__":
asyncio.run(main())
⚠️
category/severity的取值必须来自 Step 1 读到的 CHECK 约束。 上面写的misleading/absolute/guarantee/yield/high是占位假设; 若 DDL 里的取值集不同,以 DDL 为准并同步改本常量表。
- Step 5: 跑种子并验证
Run:
.\.venv\Scripts\python.exe tools\seed_compliance_baseline.py
.\.venv\Scripts\python.exe -m pytest tests/integration/test_compliance_seed_mysql.py -q
Expected: 脚本输出 seeded 11 negative rules, 6 reply templates;测试 PASS(2 passed)
- Step 6: 反证——证明拦截来自配置而非硬编码
Run:
.\.venv\Scripts\python.exe -c "
import asyncio
from sqlalchemy import text
from app.infrastructure.db import SessionFactory
from app.core.contracts import CoreResult, RequestContext, ResolvedAgentConfig
from app.service.agent.governance import review_output
async def main():
ctx = RequestContext(user_id='9001', trace_id='t', roles=('customer',))
cfg = ResolvedAgentConfig(config_version='v', prompt_version='p', model_endpoint='m')
bad = CoreResult(text='这只基金稳赚,保本无风险。')
try:
out = await review_output(bad, ctx, cfg, ())
print('输出:', out.text[:80])
except Exception as e:
print('被拦截:', type(e).__name__, e)
asyncio.run(main())
"
Expected: 输出被替换为安全话术或抛错——证明规则生效。把实际输出记进报告。
- Step 7: 提交前检查
Run:
.\.venv\Scripts\python.exe -m ruff check app tests tools
.\.venv\Scripts\python.exe -m pytest -q tests/unit tests/contract
Expected: ruff 通过;单测/契约测试无回归(基线 476 passed + 1 skipped)
Task 2: 免责声明与固定话术接入治理层
Files:
- Modify:
app/service/agent/governance.py - Test:
tests/unit/service/test_governance_disclaimer.py
Interfaces:
- Consumes: Task 1 落库的
agent_reply_template(必须按template_code解析,不要按scene) - Produces:
review_output()对所有面向客户的输出在末尾追加固定免责声明;常量DISCLAIMER_TEMPLATE_CODE = "TPL_DISCLAIMER"与代码兜底文案FALLBACK_DISCLAIMER
要求来源:专项设计 §0.3 + 门禁 F5「免责声明强制注入,面向客户输出 100% 附固定话术」。 当前
governance.py正常路径不注入免责声明(只在命中负面词时替换文本)。 硬编码兜底:话术从库里读,查不到时用代码常量兜底,符合"安全关键路径不依赖额外查询"的取舍。
Task 2 必须一并修的两个治理层缺陷(Task 1 审查发现,控制者已核实)
缺陷 1(安全,fail-open):governance.py L70-74 的判据是
if agents and definition.agent_type not in agents:
continue
applicable_agents 为 [] 或 NULL 时 if agents 为假 → 不做过滤 → 规则溢出到所有 Agent。
实测真值表:
applicable_agents |
客服加载 | 投顾加载 |
|---|---|---|
["customer_service"] |
11 | 0 |
[] / NULL |
11 | 11 ← fail-open 溢出 |
["advisor"] |
0 ← 静默失效 | 11 |
即"配置为空"不是失效而是扩大适用范围,与项目一贯的"失败关闭"原则相反。 修法:空数组/NULL 应视为不适用任何 Agent(fail-closed),并加反证测试。
缺陷 2(话术自绊,Task 1 已缓解但根因仍在):governance.py L131 命中负面词时替换的是
写死的句子「该内容需要人工核实。基金投资存在风险,本系统不代客交易。」,
不读 safe_reply_template_code——所以 Task 1 种的 6 条话术里,只有免责声明会被本 Task 用到,
其余 5 条是"白种"的。若本 Task 要把话术接进来,必须保证替换后的文本不再被负面词规则过滤一次
(否则 Task 1 已修复的"自绊"会以新形式复现)。同时注意:替换文本本身不得含任何禁用字面。
- Step 1: 写失败测试
创建 tests/unit/service/test_governance_disclaimer.py:
from app.core.contracts import AgentResult, CoreResult, RequestContext, ResolvedAgentConfig
from app.service.agent.governance import FALLBACK_DISCLAIMER, review_output
CONTEXT = RequestContext(user_id="9001", trace_id="t", roles=("customer",))
CONFIG = ResolvedAgentConfig(config_version="v", prompt_version="p", model_endpoint="m")
def _result(text: str) -> AgentResult:
return AgentResult(run_id="r", result=CoreResult(text=text))
def test_disclaimer_is_appended_to_normal_output() -> None:
result = review_output(_result("基金申购后 T+1 确认份额。"), CONTEXT, CONFIG, ())
assert FALLBACK_DISCLAIMER in result.result.text
assert result.result.text.startswith("基金申购后 T+1 确认份额。")
def test_disclaimer_is_not_duplicated() -> None:
once = review_output(_result("答案。"), CONTEXT, CONFIG, ())
twice = review_output(once, CONTEXT, CONFIG, ())
assert twice.result.text.count(FALLBACK_DISCLAIMER) == 1
def test_empty_applicable_agents_does_not_leak_rules_to_other_agents() -> None:
"""缺陷 1 的反证:applicable_agents 为空时不得把规则溢出到其他 Agent。"""
# 需要能注入自定义 negative_rules 的构造方式;若 review_output 的签名不便注入,
# 则改为直接测 PlatformGovernance.resolve 对空 applicable_agents 的行为(见集成测试)。
...
注意:
review_output是同步函数(不是 async),签名是review_output(result: AgentResult, context, config, memories),返回AgentResult。 断言文本要用result.result.text(不是result.text)。
-
Step 2: 跑测试确认失败 —
ImportError: cannot import name 'FALLBACK_DISCLAIMER' -
Step 3: 实现(新增
DISCLAIMER_TEMPLATE_CODE/FALLBACK_DISCLAIMER/_load_template_text; 在脱敏之后、返回之前追加免责声明,且不重复追加;同时把if agents and ...改为 fail-closed) -
Step 4: 跑本 Task 测试 + 既有 governance 测试,逐个判断既有断言是否需要更新(并在报告里列出)
-
Step 5: lint
Task 3: 知识种子解析与导入(105 条 + 补验收用例)
Files:
- Create:
tools/qa_source_parser.py - Create:
tools/import_knowledge_seed.py - Test:
tests/unit/tools/test_qa_source_parser.py
Interfaces:
- Consumes:
客服Agent知识库_QA问答对_v5_RAG发布候选版.txt(105 条,已审核 v5.8) - Produces:
class QaRecord(frozen dataclass):qa_id: str、question: str、synonyms: tuple[str, ...]、answer: strEXPECTED_RECORD_COUNT = 105、parse_qa_source(text, *, expected_count=None)、normalize_phrase()、phrase_hash()- CLI:
python tools/import_knowledge_seed.py --source <path> --created-by 9003
已核实的关键事实:源文件 105 条使用 12 种编号前缀(实测,逐前缀计数):
RAG-PER12 /RAG-PUB16 /RAG-RVW16 /RAG-CONFIG8 /RAG-CHAT4 /RAG-HUM4 /RAG-P12 /NF-SVC10 /NF-TRD11 /NF-CMP10 /NF-STS10 /NF-ACC2,合计 105。 只按RAG-解析会静默丢掉 43 条。 ⚠️ 修正记录:本计划初稿写"11 种前缀",漏了NF-ACC(2 条)——11 种最多只能凑出 103 条, 与文件声明的 105 条对不上,这个算术矛盾就是漏数的证据。Task 3 实施时实测发现并纠正。 实施教训:解析器因此不做前缀白名单(正则接受任意[Letters-…]形式),从根上避免漏数。非空列必须显式赋值(现库实测):
fin_knowledge_meta.knowledge_type/milvus_collection/content_text/status/created_at/updated_at;agent_faq_synonym.created_by(外键 → sys_user.id)/normalized_phrase/phrase_hash(= SHA-256(normalized_phrase),docs/02§7.3 L355)。验收缺口:老师 Phase 1 点名的「基金申购后多久确认」不在 105 条内——Step 6 会补一条。
- Step 1: 写失败测试
创建 tests/unit/tools/test_qa_source_parser.py:
import pytest
from tools.qa_source_parser import (
EXPECTED_RECORD_COUNT, parse_qa_source, phrase_hash, normalize_phrase,
)
SAMPLE = """
[RAG-PER-001]
问题:你是谁?
相似问法:你是人工吗?|你是真人客服吗?
回答:我是奶龙基金智能助手。
[NF-SVC-001]
问题:怎么联系你们?
相似问法:客服入口在哪?
回答:您可通过官网获取服务。
"""
def test_expected_count_is_105() -> None:
assert EXPECTED_RECORD_COUNT == 105
def test_parser_accepts_both_id_prefixes() -> None:
records = parse_qa_source(SAMPLE, expected_count=None)
assert [r.qa_id for r in records] == ["RAG-PER-001", "NF-SVC-001"]
def test_parser_splits_synonyms_on_pipe() -> None:
records = parse_qa_source(SAMPLE, expected_count=None)
assert records[0].synonyms == ("你是人工吗?", "你是真人客服吗?")
def test_parser_keeps_answer_verbatim() -> None:
records = parse_qa_source(SAMPLE, expected_count=None)
assert records[0].answer == "我是奶龙基金智能助手。"
def test_parser_rejects_record_without_synonyms() -> None:
with pytest.raises(ValueError, match="相似问法"):
parse_qa_source("[RAG-PER-001]\n问题:问?\n回答:答。\n", expected_count=None)
def test_parser_rejects_wrong_record_count_by_default() -> None:
with pytest.raises(ValueError, match="记录数不符"):
parse_qa_source(SAMPLE)
def test_phrase_hash_is_sha256_of_normalized_phrase() -> None:
import hashlib
value = "你是人工吗?"
assert phrase_hash(value) == hashlib.sha256(
normalize_phrase(value).encode("utf-8")).hexdigest()
assert len(phrase_hash(value)) == 64
- Step 2: 跑测试确认失败
Run: .\.venv\Scripts\python.exe -m pytest tests/unit/tools/test_qa_source_parser.py -q
Expected: FAIL — ModuleNotFoundError: No module named 'tools.qa_source_parser'
- Step 3: 写解析器
创建 tools/qa_source_parser.py:
"""解析已审核的 QA 源文件。
源文件共 105 条、11 种编号前缀(RAG-* 与 NF-*),每条固定三段:
[编号] / 问题: / 相似问法:(| 分隔)/ 回答:
只按 RAG- 前缀解析会静默丢掉 43 条 NF-* 记录。
"""
from __future__ import annotations
import re
from dataclasses import dataclass
from hashlib import sha256
EXPECTED_RECORD_COUNT = 105
_ID_PATTERN = re.compile(r"^\[([A-Za-z][A-Za-z0-9]*(?:-[A-Za-z0-9]+)+)\]$")
_FIELD_PREFIXES = ("问题:", "相似问法:", "回答:")
@dataclass(frozen=True)
class QaRecord:
qa_id: str
question: str
synonyms: tuple[str, ...]
answer: str
def normalize_phrase(phrase: str) -> str:
return " ".join(phrase.split()).strip().lower()
def phrase_hash(phrase: str) -> str:
return sha256(normalize_phrase(phrase).encode("utf-8")).hexdigest()
def _blocks(text: str) -> list[list[str]]:
blocks: list[list[str]] = []
current: list[str] | None = None
for raw in text.splitlines():
line = raw.strip()
if _ID_PATTERN.match(line):
if current is not None:
blocks.append(current)
current = [line]
continue
if current is None or not line:
continue
if line.startswith(_FIELD_PREFIXES):
current.append(line)
if current is not None:
blocks.append(current)
return blocks
def parse_qa_source(
text: str, *, expected_count: int | None = EXPECTED_RECORD_COUNT
) -> list[QaRecord]:
records: list[QaRecord] = []
seen: set[str] = set()
for block in _blocks(text):
match = _ID_PATTERN.match(block[0])
assert match is not None
qa_id = match.group(1)
if qa_id in seen:
raise ValueError(f"重复的 qa_id:{qa_id}")
seen.add(qa_id)
fields: dict[str, str] = {}
for line in block[1:]:
for prefix in _FIELD_PREFIXES:
if line.startswith(prefix):
fields[prefix] = line[len(prefix):].strip()
break
question = fields.get("问题:", "")
synonyms_raw = fields.get("相似问法:", "")
answer = fields.get("回答:", "")
if not question:
raise ValueError(f"{qa_id} 缺少问题")
if not synonyms_raw:
raise ValueError(f"{qa_id} 缺少相似问法")
if not answer:
raise ValueError(f"{qa_id} 缺少回答")
synonyms = tuple(dict.fromkeys(
item.strip() for item in synonyms_raw.split("|") if item.strip()))
if not synonyms:
raise ValueError(f"{qa_id} 相似问法为空")
records.append(QaRecord(qa_id, question, synonyms, answer))
if expected_count is not None and len(records) != expected_count:
raise ValueError(f"记录数不符:期望 {expected_count},实际 {len(records)}")
return records
- Step 4: 跑测试 + 对真实源文件验证
Run:
.\.venv\Scripts\python.exe -m pytest tests/unit/tools/test_qa_source_parser.py -q
.\.venv\Scripts\python.exe -c "from pathlib import Path; from tools.qa_source_parser import parse_qa_source; p=Path(r'C:\Users\Windows\Desktop\DSH开发智能助手\胜宇前期资料\客服Agent知识库_QA问答对_v5_RAG发布候选版.txt'); r=parse_qa_source(p.read_text(encoding='utf-8')); print('records:', len(r)); print('synonyms total:', sum(len(x.synonyms) for x in r))"
Expected: 测试 7 passed;records: 105;synonyms total: 为正整数
- Step 5: 写导入脚本
创建 tools/import_knowledge_seed.py:
"""把已审核 QA 一次性导入 fin_knowledge_meta 与 agent_faq_synonym。
幂等:按 tags.qa_id 定位既有记录并更新;随后写入向量同步 Outbox 事件。
本脚本不直接调 Milvus(写入由 Outbox Worker 负责)。
"""
from __future__ import annotations
import argparse
import asyncio
import json
import sys
from datetime import UTC, date, datetime, timedelta
from pathlib import Path
from uuid import uuid4
ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))
from sqlalchemy import text # noqa: E402
from app.infrastructure.db import SessionFactory # noqa: E402
from tools.qa_source_parser import ( # noqa: E402
QaRecord, normalize_phrase, parse_qa_source, phrase_hash,
)
FAQ_COLLECTION = "fin_faq_collection"
KNOWLEDGE_TYPE = "faq"
SOURCE_VERSION = "v5.8"
VECTOR_SYNC_EVENT = "knowledge.vector_sync_requested"
#: 老师 Phase 1 验收点名的测试问题不在 105 条内,此处补齐使验收用例可命中。
ACCEPTANCE_SUPPLEMENT = QaRecord(
qa_id="SUP-001",
question="基金申购后多久确认?",
synonyms=("申购多久确认份额?", "买基金几天能确认?", "基金申购确认时间"),
answer=(
"交易日 15:00 前提交的申购申请,T+1 日确认份额(QDII 基金为 T+2 日);"
"15:00 后提交则顺延至下一交易日。非交易日(周末及法定节假日)提交的申请顺延至"
"下一交易日处理。确认后即可查看持仓。具体以产品说明书和交易页面显示为准。"
),
)
async def _upsert_meta(session, record: QaRecord, source_file: str,
created_by: int, now: datetime) -> int:
tags = json.dumps(
{"qa_id": record.qa_id, "phase": "phase_1", "source_version": SOURCE_VERSION},
ensure_ascii=False)
existing = await session.scalar(text(
"SELECT id FROM fin_knowledge_meta"
" WHERE JSON_UNQUOTE(JSON_EXTRACT(tags,'$.qa_id'))=:qa_id LIMIT 1"
), {"qa_id": record.qa_id})
if existing is None:
await session.execute(text(
"INSERT INTO fin_knowledge_meta (knowledge_type, title, source_file,"
" milvus_collection, version, content_text, tags, reviewer_id,"
" review_status, status, created_at, updated_at)"
" VALUES (:kt, :title, :src, :coll, :ver, :body, :tags, :who,"
" 'published', 'active', :now, :now)"
), {"kt": KNOWLEDGE_TYPE, "title": record.question, "src": source_file,
"coll": FAQ_COLLECTION, "ver": SOURCE_VERSION, "body": record.answer,
"tags": tags, "who": created_by, "now": now})
return int(await session.scalar(text("SELECT LAST_INSERT_ID()")))
await session.execute(text(
"UPDATE fin_knowledge_meta SET title=:title, content_text=:body,"
" source_file=:src, version=:ver, review_status='published', status='active',"
" tags=:tags, reviewer_id=:who, updated_at=:now WHERE id=:id"
), {"title": record.question, "body": record.answer, "src": source_file,
"ver": SOURCE_VERSION, "tags": tags, "who": created_by,
"now": now, "id": int(existing)})
return int(existing)
async def _replace_synonyms(session, knowledge_id: int, record: QaRecord,
created_by: int, now: datetime) -> None:
await session.execute(text(
"DELETE FROM agent_faq_synonym WHERE knowledge_id=:id AND source_type='import'"
), {"id": knowledge_id})
seen: set[str] = set()
for phrase in (record.question, *record.synonyms):
digest = phrase_hash(phrase)
if digest in seen:
continue
seen.add(digest)
await session.execute(text(
"INSERT INTO agent_faq_synonym (knowledge_id, phrase, normalized_phrase,"
" phrase_hash, language_code, source_type, hit_count, status, created_by,"
" reviewer_id, reviewed_at, created_at, updated_at)"
" VALUES (:kid, :phrase, :norm, :hash, 'zh-CN', 'import', 0, 'approved',"
" :who, :who, :now, :now, :now)"
), {"kid": knowledge_id, "phrase": phrase, "norm": normalize_phrase(phrase),
"hash": digest, "who": created_by, "now": now})
async def _enqueue_vector_sync(session, knowledge_id: int, now: datetime) -> None:
await session.execute(text(
"INSERT INTO domain_event_outbox (event_id, event_type, aggregate_type,"
" aggregate_id, trace_id, payload, status, retry_count, occurred_at,"
" created_at, updated_at)"
" VALUES (:eid, :etype, 'knowledge_meta', :agg, :trace, :payload, 'pending',"
" 0, :now, :now, :now)"
), {"eid": str(uuid4()), "etype": VECTOR_SYNC_EVENT, "agg": str(knowledge_id),
"trace": str(uuid4()),
"payload": json.dumps({"knowledge_id": str(knowledge_id)}), "now": now})
async def import_records(records: list[QaRecord], *, source_file: str,
created_by: int, now: datetime) -> int:
async with SessionFactory() as session, session.begin():
for record in records:
kid = await _upsert_meta(session, record, source_file, created_by, now)
await _replace_synonyms(session, kid, record, created_by, now)
await _enqueue_vector_sync(session, kid, now)
return len(records)
async def main() -> None:
parser = argparse.ArgumentParser(description="导入已审核 QA 知识库")
parser.add_argument("--source", required=True, type=Path)
parser.add_argument("--created-by", required=True, type=int)
parser.add_argument("--dry-run", action="store_true")
args = parser.parse_args()
records = parse_qa_source(args.source.read_text(encoding="utf-8"))
records.append(ACCEPTANCE_SUPPLEMENT)
print(f"parsed {len(records)} records (incl. 1 acceptance supplement)")
if args.dry_run:
print("dry-run: 未写库")
return
count = await import_records(
records, source_file=args.source.name, created_by=args.created_by,
now=datetime.now(UTC).replace(tzinfo=None))
print(f"imported {count} records; vector sync events queued")
if __name__ == "__main__":
asyncio.run(main())
说明:
effective_date/expire_date留空表示长期有效(spec A 口径)。
- Step 6: 导入并验证
Run:
.\.venv\Scripts\python.exe tools\import_knowledge_seed.py --source "C:\Users\Windows\Desktop\DSH开发智能助手\胜宇前期资料\客服Agent知识库_QA问答对_v5_RAG发布候选版.txt" --created-by 9003
.\.venv\Scripts\python.exe -c "import pymysql;c=pymysql.connect(host='127.0.0.1',port=3306,user='root',password='123456',database='jr_agent');cur=c.cursor();cur.execute(\"SELECT COUNT(*) FROM fin_knowledge_meta WHERE review_status='published' AND status='active'\");print('published meta:',cur.fetchone()[0]);cur.execute('SELECT COUNT(*) FROM agent_faq_synonym');print('synonyms:',cur.fetchone()[0]);cur.execute(\"SELECT content_text FROM fin_knowledge_meta WHERE title='基金申购后多久确认?'\");r=cur.fetchone();print('acceptance supplement:', 'FOUND' if r else 'MISSING')"
Expected: imported 106 records;published meta: 106;synonyms: 为正整数;acceptance supplement: FOUND
- Step 7: 提交前检查
Run: .\.venv\Scripts\python.exe -m ruff check app tests tools
Expected: 通过(tools/*.py 已放宽 E501/I001/B007)
Task 4: 文档解析与可替换的文档存储
Files:
- Create:
app/infrastructure/document_storage.py - Create:
app/service/document_parser.py - Test:
tests/unit/infrastructure/test_document_storage.py、tests/unit/service/test_document_parser.py
Interfaces:
- Consumes: 无
- Produces:
class DocumentStorage(Protocol):async save(self, *, key: str, content: bytes, content_type: str) -> str(返回存储路径)、async read(self, *, key: str) -> bytes、async delete(self, *, key: str) -> None、async archive(self, *, key: str) -> Noneclass LocalDocumentStorage:__init__(self, root: Path),实现上述协议class ParsedChunk(frozen dataclass):text: str、heading_path: tuple[str, ...]、index: intclass DocumentParser:chunk_size: int = 512、chunk_overlap: int = 64;parse(self, *, filename: str, content: bytes) -> list[ParsedChunk]SUPPORTED_EXTENSIONS = frozenset({".txt", ".md", ".docx"})- 异常:
UnsupportedDocumentError
需求来源:修改版 F1.1「MinIO 连接与文件上传接口」、F1.2「支持 txt/md/docx,512 token/块, overlap 64,保留标题层级作为 metadata」。
架构裁定:MinIO 未安装未运行,故用协议 + 本地实现先跑通全链路;
MinioDocumentStorage待实例就绪后在 Task 6 补(接口不变,业务代码零改动)。docx 解析:需要新增
python-docx依赖。若沙箱拒绝安装,改为降级为"仅 txt/md 可用、 docx 返回UnsupportedDocumentError"并在报告里标注为已知限制,不要伪造实现。
- Step 1: 写失败测试
创建 tests/unit/infrastructure/test_document_storage.py:
from pathlib import Path
import pytest
from app.infrastructure.document_storage import LocalDocumentStorage
@pytest.mark.asyncio
async def test_save_read_delete_round_trip(tmp_path: Path) -> None:
storage = LocalDocumentStorage(tmp_path)
key = "kb/2026/faq.txt"
path = await storage.save(key=key, content=b"hello", content_type="text/plain")
assert path == key
assert await storage.read(key=key) == b"hello"
await storage.delete(key=key)
with pytest.raises(FileNotFoundError):
await storage.read(key=key)
@pytest.mark.asyncio
async def test_archive_moves_without_deleting(tmp_path: Path) -> None:
storage = LocalDocumentStorage(tmp_path)
await storage.save(key="a.txt", content=b"x", content_type="text/plain")
await storage.archive(key="a.txt")
assert (tmp_path / "archive" / "a.txt").read_bytes() == b"x"
assert not (tmp_path / "a.txt").exists()
创建 tests/unit/service/test_document_parser.py:
import pytest
from app.service.document_parser import (
SUPPORTED_EXTENSIONS, DocumentParser, UnsupportedDocumentError,
)
def test_supported_extensions_match_requirement() -> None:
assert SUPPORTED_EXTENSIONS == frozenset({".txt", ".md", ".docx"})
def test_txt_is_split_with_overlap_and_indexed() -> None:
parser = DocumentParser(chunk_size=10, chunk_overlap=2)
text = "\n\n".join(f"段落{i}" + "内容" * 20 for i in range(4))
chunks = parser.parse(filename="a.txt", content=text.encode("utf-8"))
assert len(chunks) > 1
assert [c.index for c in chunks] == list(range(len(chunks)))
assert all(c.text.strip() for c in chunks)
def test_markdown_keeps_heading_path_as_metadata() -> None:
parser = DocumentParser(chunk_size=200, chunk_overlap=0)
text = "# 第一章\n\n## 第一节\n\n这是正文内容,需要足够长以便成块。\n"
chunks = parser.parse(filename="a.md", content=text.encode("utf-8"))
assert chunks
assert chunks[0].heading_path == ("第一章", "第一节")
def test_unsupported_extension_is_rejected() -> None:
parser = DocumentParser()
with pytest.raises(UnsupportedDocumentError):
parser.parse(filename="a.pdf", content=b"%PDF-1.4")
- Step 2: 跑测试确认失败
Run: .\.venv\Scripts\python.exe -m pytest tests/unit/infrastructure/test_document_storage.py tests/unit/service/test_document_parser.py -q
Expected: FAIL — 两个模块都不存在
- Step 3: 实现存储抽象
创建 app/infrastructure/document_storage.py:
"""文档存储抽象。
一期用本地文件系统实现,MinIO 实现待实例就绪后补(接口不变)。
key 使用相对路径语义(如 `kb/2026/faq.txt`),不得包含 `..`。
"""
from pathlib import Path
from typing import Protocol
class DocumentStorage(Protocol):
async def save(self, *, key: str, content: bytes, content_type: str) -> str: ...
async def read(self, *, key: str) -> bytes: ...
async def delete(self, *, key: str) -> None: ...
async def archive(self, *, key: str) -> None: ...
class LocalDocumentStorage:
"""本地文件系统实现;仅用于开发与环境未就绪时的真机替代。"""
def __init__(self, root: Path) -> None:
self._root = Path(root)
def _resolve(self, key: str) -> Path:
if ".." in Path(key).parts:
raise ValueError("key 不得包含上级目录引用")
return self._root / key
async def save(self, *, key: str, content: bytes, content_type: str) -> str:
del content_type # 本地实现不记录 MIME,MinIO 实现需要
path = self._resolve(key)
path.parent.mkdir(parents=True, exist_ok=True)
path.write_bytes(content)
return key
async def read(self, *, key: str) -> bytes:
return self._resolve(key).read_bytes()
async def delete(self, *, key: str) -> None:
self._resolve(key).unlink(missing_ok=True)
async def archive(self, *, key: str) -> None:
source = self._resolve(key)
if not source.exists():
return
target = self._root / "archive" / key
target.parent.mkdir(parents=True, exist_ok=True)
source.replace(target)
- Step 4: 实现文档解析
创建 app/service/document_parser.py:
"""文档解析与切分。
需求(修改版 F1.2):支持 txt/md/docx;512 token/块、overlap 64;保留标题层级作为 metadata。
本实现按字符近似 token(中文场景 1 字≈1 token 的保守估计),不引入 tokenizer 依赖。
"""
from dataclasses import dataclass
from pathlib import PurePosixPath
SUPPORTED_EXTENSIONS = frozenset({".txt", ".md", ".docx"})
class UnsupportedDocumentError(ValueError):
"""不支持的文件类型。"""
@dataclass(frozen=True)
class ParsedChunk:
text: str
heading_path: tuple[str, ...]
index: int
class DocumentParser:
def __init__(self, *, chunk_size: int = 512, chunk_overlap: int = 64) -> None:
if chunk_size <= 0:
raise ValueError("chunk_size 必须为正")
if not 0 <= chunk_overlap < chunk_size:
raise ValueError("chunk_overlap 必须满足 0 <= overlap < chunk_size")
self.chunk_size = chunk_size
self.chunk_overlap = chunk_overlap
def parse(self, *, filename: str, content: bytes) -> list[ParsedChunk]:
suffix = PurePosixPath(filename).suffix.lower()
if suffix not in SUPPORTED_EXTENSIONS:
raise UnsupportedDocumentError(f"不支持的文件类型:{suffix}")
if suffix == ".docx":
text = self._read_docx(content)
headings: list[str] = []
else:
text = content.decode("utf-8", errors="replace")
headings = [line.lstrip("#").strip()
for line in text.splitlines() if line.startswith("#")]
return self._split(text, tuple(headings))
@staticmethod
def _read_docx(content: bytes) -> str:
import io
from docx import Document # 需要 python-docx
document = Document(io.BytesIO(content))
return "\n\n".join(p.text for p in document.paragraphs if p.text.strip())
def _split(self, text: str, heading_path: tuple[str, ...]) -> list[ParsedChunk]:
paragraphs = [p.strip() for p in text.split("\n\n") if p.strip()]
chunks: list[ParsedChunk] = []
buffer = ""
for paragraph in paragraphs:
candidate = f"{buffer}\n\n{paragraph}".strip() if buffer else paragraph
if len(candidate) <= self.chunk_size:
buffer = candidate
continue
if buffer:
chunks.append(ParsedChunk(buffer, heading_path, len(chunks)))
buffer = buffer[-self.chunk_overlap:] if self.chunk_overlap else ""
candidate = f"{buffer}\n\n{paragraph}".strip() if buffer else paragraph
while len(candidate) > self.chunk_size:
chunks.append(
ParsedChunk(candidate[: self.chunk_size], heading_path, len(chunks)))
candidate = candidate[self.chunk_size - self.chunk_overlap:]
buffer = candidate
if buffer:
chunks.append(ParsedChunk(buffer, heading_path, len(chunks)))
return chunks
- Step 5: 安装 docx 依赖并跑测试
Run:
.\.venv\Scripts\python.exe -m pip install python-docx
.\.venv\Scripts\python.exe -m pytest tests/unit/infrastructure/test_document_storage.py tests/unit/service/test_document_parser.py -q
Expected: 6 passed。
若 pip 被沙箱拒绝:把 _read_docx 改为抛 UnsupportedDocumentError("docx 解析依赖未安装"),
把 docx 测试改为断言该异常,并在报告里标注为已知限制。
- Step 6: 提交前检查
Run: .\.venv\Scripts\python.exe -m ruff check app tests 与 .\.venv\Scripts\python.exe -m mypy app
Expected: ruff 通过;mypy 不新增错误
Task 5: 知识入库服务与 Outbox 向量同步
Files:
- Create:
app/service/knowledge_ingest_service.py - Create:
app/infrastructure/milvus_knowledge_writer.py - Create:
app/worker/knowledge_vector_worker.py - Test:
tests/unit/service/test_knowledge_ingest_service.py、tests/unit/worker/test_knowledge_vector_worker.py
Interfaces:
- Consumes: Task 3 的
fin_knowledge_meta;Task 4 的DocumentParser/DocumentStorage;底座的ModelEmbeddingService(单文本embed(endpoints, text) -> EmbeddingExecution,.vector)、DatabaseModelEndpointResolver - Produces:
class KnowledgeIngestService:__init__(self, *, session, parser, storage, embedder, endpoint_resolver);async ingest(self, *, filename: str, content: bytes, knowledge_type: str, collection: str, created_by: int) -> int(返回knowledge_id)class MilvusKnowledgeWriter:async upsert(self, *, collection, knowledge_id, vector, fields) -> None、async delete(self, *, collection, knowledge_id) -> NoneVECTOR_SYNC_EVENT = "knowledge.vector_sync_requested"、VECTOR_DELETE_EVENT = "knowledge.vector_delete_requested"、build_knowledge_handlers(...)
要点:
- 底座已有的是单文本
embed()(docs/19与bootstrap.py确认),不是批量; 旧计划的"批量 embed"设计作废。- 维度必须等于 1024,不等则失败关闭(抛异常,交 Outbox 重试)。
- 写路径与读路径物理隔离:检索侧永远不持有写客户端。
- 重试完全复用
OutboxWorker(指数退避、5 次上限、死信),不新写重试逻辑。
- Step 1: 写失败测试
创建 tests/unit/worker/test_knowledge_vector_worker.py:
import pytest
from app.core.errors import RecoverableAgentError
from app.worker.knowledge_vector_worker import (
VECTOR_DELETE_EVENT, VECTOR_SYNC_EVENT, build_knowledge_handlers,
)
class FakeRow:
id = 11
title = "基金申购后多久确认?"
content_text = "交易日 15:00 前提交,T+1 确认份额。"
milvus_collection = "fin_faq_collection"
version = "v5.8"
tags = {"qa_id": "SUP-001"}
class FakeSession:
def __init__(self, row): self._row = row
async def scalar(self, statement): return self._row
class FakeWriter:
def __init__(self): self.upserts, self.deletes = [], []
async def upsert(self, **kw): self.upserts.append(kw)
async def delete(self, **kw): self.deletes.append(kw)
class FakeEmbedder:
def __init__(self, dim=1024): self._dim = dim
async def embed(self, endpoints, text, *, max_attempts=2):
class E: vector = [0.5] * self._dim
return E()
class FakeResolver:
async def resolve(self, *, agent_type, task_type): return [object()]
def _handlers(row, writer, dim=1024):
return build_knowledge_handlers(
FakeSession(row), writer=writer, embedder=FakeEmbedder(dim),
endpoint_resolver=FakeResolver())
@pytest.mark.asyncio
async def test_sync_upserts_with_expected_fields() -> None:
writer = FakeWriter()
await _handlers(FakeRow(), writer)[VECTOR_SYNC_EVENT]({"knowledge_id": "11"})
call = writer.upserts[0]
assert call["collection"] == "fin_faq_collection"
assert call["knowledge_id"] == "11"
assert len(call["vector"]) == 1024
assert call["fields"]["title"] == "基金申购后多久确认?"
assert call["fields"]["tags"] == "SUP-001"
@pytest.mark.asyncio
async def test_sync_fails_closed_on_wrong_dimension() -> None:
with pytest.raises(RecoverableAgentError):
await _handlers(FakeRow(), FakeWriter(), dim=768)[VECTOR_SYNC_EVENT](
{"knowledge_id": "11"})
@pytest.mark.asyncio
async def test_sync_fails_when_row_missing() -> None:
with pytest.raises(RecoverableAgentError):
await _handlers(None, FakeWriter())[VECTOR_SYNC_EVENT]({"knowledge_id": "404"})
@pytest.mark.asyncio
async def test_delete_calls_writer() -> None:
writer = FakeWriter()
await _handlers(FakeRow(), writer)[VECTOR_DELETE_EVENT]({"knowledge_id": "11"})
assert writer.deletes == [
{"collection": "fin_faq_collection", "knowledge_id": "11"}]
- Step 2: 跑测试确认失败
Run: .\.venv\Scripts\python.exe -m pytest tests/unit/worker/test_knowledge_vector_worker.py -q
Expected: FAIL — 模块不存在
- Step 3: 实现写适配器与 Worker
创建 app/infrastructure/milvus_knowledge_writer.py:
"""知识写路径 Milvus 适配器。
只被 Outbox Worker 使用;检索路径不导入本模块(读写物理隔离)。
"""
from typing import Any
from app.core.errors import RecoverableAgentError
class MilvusKnowledgeWriter:
def __init__(self, uri: str, token: str = "") -> None:
self._uri, self._token, self._client = uri, token, 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 upsert(self, *, collection: str, knowledge_id: str,
vector: list[float], fields: dict[str, Any]) -> None:
if not knowledge_id or not vector:
raise RecoverableAgentError("knowledge_id 与向量都不能为空")
client = await self._ensure()
try:
await client.upsert(collection_name=collection, data=[
{"knowledge_id": knowledge_id, "vector": vector, **fields}])
except Exception as exc:
raise RecoverableAgentError("知识向量写入失败") from exc
async def delete(self, *, collection: str, knowledge_id: str) -> None:
if not knowledge_id:
raise RecoverableAgentError("knowledge_id 不能为空")
client = await self._ensure()
try:
await client.delete(collection_name=collection, ids=[knowledge_id])
except Exception as exc:
raise RecoverableAgentError("知识向量删除失败") from exc
async def close(self) -> None:
if self._client is not None:
await self._client.close()
self._client = None
创建 app/worker/knowledge_vector_worker.py:
"""知识向量同步 Outbox handler。
handler 必须使用传入的 session 做本地读写,且**不得 commit**(OutboxWorker 契约)。
重试/退避/死信复用 OutboxWorker 既有机制。
"""
from collections.abc import Awaitable, Callable
from typing import Any
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.errors import RecoverableAgentError
from app.core.knowledge_contracts import VECTOR_DIM
from app.model.knowledge import KnowledgeMeta
from app.worker.outbox_worker import OutboxWorker
VECTOR_SYNC_EVENT = "knowledge.vector_sync_requested"
VECTOR_DELETE_EVENT = "knowledge.vector_delete_requested"
def _tags_to_string(tags: Any) -> str:
if isinstance(tags, dict):
return ",".join(str(v) for v in tags.values())
if isinstance(tags, (list, tuple)):
return ",".join(str(v) for v in tags)
return ""
def build_knowledge_handlers(
session: AsyncSession, *, writer: Any, embedder: Any, endpoint_resolver: Any,
vector_dim: int = VECTOR_DIM,
) -> dict[str, Callable[[dict[str, Any]], Awaitable[None]]]:
async def _load(knowledge_id: str) -> KnowledgeMeta:
if not knowledge_id.isdigit():
raise RecoverableAgentError("knowledge_id 非法")
row = await session.scalar(
select(KnowledgeMeta).where(KnowledgeMeta.id == int(knowledge_id)))
if row is None:
raise RecoverableAgentError("知识记录不存在")
return row
async def sync(payload: dict[str, Any]) -> None:
from app.core.knowledge_contracts import intent_for_qa_id
row = await _load(str(payload.get("knowledge_id", "")))
endpoints = await endpoint_resolver.resolve(
agent_type="customer_service", task_type="embedding")
execution = await embedder.embed(endpoints, row.content_text)
vector = list(execution.vector)
if len(vector) != vector_dim:
raise RecoverableAgentError("嵌入维度与集合定义不一致")
tags = row.tags if isinstance(row.tags, dict) else {}
await writer.upsert(
collection=row.milvus_collection, knowledge_id=str(row.id),
vector=vector,
fields={"title": row.title, "snippet": row.content_text,
"tags": _tags_to_string(row.tags), "version": row.version or "",
"intent": intent_for_qa_id(str(tags.get("qa_id", ""))) or ""})
async def remove(payload: dict[str, Any]) -> None:
row = await _load(str(payload.get("knowledge_id", "")))
await writer.delete(collection=row.milvus_collection, knowledge_id=str(row.id))
return {VECTOR_SYNC_EVENT: sync, VECTOR_DELETE_EVENT: remove}
async def dispatch_knowledge_events(
session: AsyncSession, *, writer: Any, embedder: Any, endpoint_resolver: Any,
vector_dim: int = VECTOR_DIM,
) -> bool:
handlers = build_knowledge_handlers(
session, writer=writer, embedder=embedder,
endpoint_resolver=endpoint_resolver, vector_dim=vector_dim)
return await OutboxWorker(session, handlers).publish_one()
- Step 4: 实现入库服务
创建 app/service/knowledge_ingest_service.py(要点:写文件 → 建元数据 → 拼 chunks → 逐块写 fin_knowledge_meta → 投递 Outbox 事件):
"""知识入库:文件落存储 → 元数据落库 → 切分 → 向量同步事件。
Service 层;不直接连 Milvus(由 Outbox Worker 负责),不直接拼 SQL 之外的存储细节。
"""
import hashlib
import json
from datetime import UTC, datetime
from pathlib import Path
from uuid import uuid4
from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.errors import ValidationAgentError
from app.service.document_parser import DocumentParser
VECTOR_SYNC_EVENT = "knowledge.vector_sync_requested"
ALLOWED_KNOWLEDGE_TYPES = frozenset({"faq", "product", "policy"})
TYPE_TO_COLLECTION = {
"faq": "fin_faq_collection",
"product": "fin_product_collection",
"policy": "fin_policy_collection",
}
class KnowledgeIngestService:
def __init__(self, *, session: AsyncSession, parser: DocumentParser,
storage: object, created_by: int) -> None:
self._session = session
self._parser = parser
self._storage = storage
self._created_by = created_by
async def ingest(self, *, filename: str, content: bytes,
knowledge_type: str) -> list[int]:
if knowledge_type not in ALLOWED_KNOWLEDGE_TYPES:
raise ValidationAgentError("knowledge_type 非法")
collection = TYPE_TO_COLLECTION[knowledge_type]
digest = hashlib.sha256(content).hexdigest()
key = f"kb/{knowledge_type}/{digest[:16]}-{Path(filename).name}"
await self._storage.save(key=key, content=content,
content_type="application/octet-stream")
chunks = self._parser.parse(filename=filename, content=content)
now = datetime.now(UTC).replace(tzinfo=None)
ids: list[int] = []
for chunk in chunks:
await self._session.execute(text(
"INSERT INTO fin_knowledge_meta (knowledge_type, title, source_file,"
" minio_path, milvus_collection, version, content_text, tags,"
" reviewer_id, review_status, status, created_at, updated_at)"
" VALUES (:kt, :title, :src, :path, :coll, 'v1', :body, :tags,"
" :who, 'published', 'active', :now, :now)"
), {"kt": knowledge_type, "title": chunk.text[:120], "src": filename,
"path": key, "coll": collection, "body": chunk.text,
"tags": json.dumps({"heading_path": list(chunk.heading_path),
"chunk_index": chunk.index}, ensure_ascii=False),
"who": self._created_by, "now": now})
knowledge_id = int(await self._session.scalar(text("SELECT LAST_INSERT_ID()")))
ids.append(knowledge_id)
await self._session.execute(text(
"INSERT INTO domain_event_outbox (event_id, event_type,"
" aggregate_type, aggregate_id, trace_id, payload, status,"
" retry_count, occurred_at, created_at, updated_at)"
" VALUES (:eid, :etype, 'knowledge_meta', :agg, :trace, :payload,"
" 'pending', 0, :now, :now, :now)"
), {"eid": str(uuid4()), "etype": VECTOR_SYNC_EVENT,
"agg": str(knowledge_id), "trace": str(uuid4()),
"payload": json.dumps({"knowledge_id": str(knowledge_id)}), "now": now})
return ids
- Step 5: 跑测试
Run: .\.venv\Scripts\python.exe -m pytest tests/unit/worker/test_knowledge_vector_worker.py -q
Expected: 4 passed
- Step 6: 提交前检查
Run:
.\.venv\Scripts\python.exe -m ruff check app tests
.\.venv\Scripts\python.exe -m pytest -q tests/unit tests/contract
Expected: ruff 通过;既有测试无回归
Task 6: Milvus 集合与检索服务
Files:
- Create:
tools/setup_milvus_knowledge_collections.py - Create:
app/infrastructure/milvus_adapter.py - Create:
app/service/knowledge_retrieval_service.py - Test:
tests/unit/service/test_knowledge_retrieval_service.py
Interfaces:
- Consumes:
KnowledgeMeta(ORM,Task 3/5 写入)、MilvusKnowledgeClient - Produces:
- 建集合 CLI:
python tools/setup_milvus_knowledge_collections.py class MilvusKnowledgeClient:async search(self, *, collection, vector, top_k, output_fields=...) -> list[dict]class KnowledgeRuntimeConfig:DEFAULT_ROUTES = {"faq": ("fin_faq_collection", 3), "product_inquiry": ("fin_product_collection", 5), "policy_explain": ("fin_policy_collection", 5)}class KnowledgeRetrievalService:async search(self, query, *, embedding_endpoints, embedder) -> KnowledgeSearchResult
- 建集合 CLI:
schema(沿用 spec A,三集合逐字相同):
knowledge_id(VARCHAR64,主键) /title(256) /snippet(4000) /tags(512) /version(16) /intent(32) /embedding(FLOAT_VECTOR 1024), 索引AUTOINDEX+metric_type="COSINE"。生效日期过滤不在 Milvus 侧:走 MySQL 二次校验(专项设计在此处有缺陷——产品/政策集合 无
expire_date字段却无条件拼该过滤,见分析报告 §13.2)。
- Step 1: 写建集合脚本
创建 tools/setup_milvus_knowledge_collections.py(幂等:has_collection 判断 + 建索引 + load)。执行前必须先人工核实 Milvus 中是否已存在同名集合:
Run: .\.venv\Scripts\python.exe -c "import asyncio;from pymilvus import AsyncMilvusClient;print(asyncio.run((lambda c: c.list_collections())(AsyncMilvusClient(uri='http://127.0.0.1:19530'))))"
Expected: 列出既有集合;若已有 fin_faq_collection 且字段结构不同,停下来报告,不要覆盖。
- Step 2: 写检索服务失败测试
创建 tests/unit/service/test_knowledge_retrieval_service.py,覆盖:
faq意图路由到fin_faq_collection且 top_k=3- 白名单外集合被拒(
ForbiddenAgentError) - 维度不等于 1024 → 失败关闭
- Milvus 失败 → 降级 MySQL LIKE 且
degraded=True - 降级路径仍执行「已发布+有效期」过滤
- Step 3-5: 实现 → 跑测试 → lint
(实现要点:_assert_collections_allowed 在任何网络调用之前执行;
_to_hits 的 intent 取值遵循 spec A 第 10 节裁定——有标签用标签、无标签用 QA 前缀回退、都没有则 None。)
- Step 6: 真机联调(需 Milvus 就绪)
Run:
.\.venv\Scripts\python.exe tools\setup_milvus_knowledge_collections.py
.\.venv\Scripts\python.exe -m app.worker --once
.\.venv\Scripts\python.exe -c "import pymysql;c=pymysql.connect(host='127.0.0.1',port=3306,user='root',password='123456',database='jr_agent');cur=c.cursor();cur.execute(\"SELECT status, COUNT(*) FROM domain_event_outbox GROUP BY status\");print(cur.fetchall())"
Expected: 三集合创建成功;106 条事件的 domain_event_outbox.status 变为 published。
Task 7: query_knowledge 工具注册与配置发布
Files:
- Create:
app/service/knowledge_tool.py - Modify:
app/service/agent/bootstrap.py - Test:
tests/unit/service/test_knowledge_tool.py
Interfaces:
- Produces:
TOOL_NAME = "query_knowledge"、REQUIRED_PERMISSION = "knowledge:query"、ALLOWED_ROLES = ("customer", "operator", "advisor", "risk_operator", "admin");bootstrap 注册ToolDefinition(...)
关键:工具注册只声明上限。必须发布配置(
config_release+platform_config_item) 才会生效;现库config_release为空,所以这一步必须包含发布动作,否则端到端必然AGENT_PERMISSION_DENIED。 发布流程见docs/19§3.1(create → submit → review → activate,单管理员可自审)。
- Step 1: 写失败测试(工具声明自洽、read_only、出参可序列化)
- Step 2: 跑测试确认失败
- Step 3: 实现工具 + bootstrap 注册
- Step 4: 跑测试
- Step 5: 发布配置并验证生效
Run: 参照 tools/demo_agent_e2e.py 的方式发布 agent_tools / customer_service:faq → {"allowed_tools": ["query_knowledge"]},然后核对 SQL:
SELECT i.release_id, i.namespace, i.config_key, i.value_json
FROM platform_config_item i JOIN config_release r ON r.id=i.release_id
WHERE r.status='active' AND i.namespace='agent_tools';
Expected: 唯一 active 版本,config_key='customer_service:faq',value_json={"allowed_tools":["query_knowledge"]}
Task 8: 客服 Agent 安全路由(确定性规则)
Files:
- Create:
app/core/customer_service_rules.py - Test:
tests/unit/core/test_customer_service_rules.py
Interfaces:
- Produces:
class SafetyRoute(frozen dataclass:priority/intent/needs_clarification/reply/transfer_required/transfer_reason);route_message(message) -> SafetyRoute | None(None= 落入知识检索路径);常量P0_KEYWORDS/P1_KEYWORDS/P2_KEYWORDS、P0_REPLY/P1_REPLY/P2_REPLY/P4_REPLY、CONTACT_PHONE = "15936583816"
硬约束:P0/P1/P2 话术与联系方式硬编码,不查知识库(spec B §5)。 编号冲突提醒:本模块的 P0-P4 是安全路由优先级(spec B 语义), 与专项设计 §5.1 的「转人工分级 P0/P1/P2」同名不同义,注释里必须写明,避免后续混用。
必须覆盖的测试用例(分析报告 §3.1 发现的矛盾): 用户会问出含零容忍词的句子(如"有什么年化5%以上的理财"), Agent 必须走合规拒答/引导而不是当普通产品咨询回答。
- Step 1-4: 写测试 → 实现 → 跑通(覆盖 P0 诈骗/验证码、P1 本人账户数据、P2 转人工诉求、含「年化收益率」的提问、公开问题返回
None) - Step 5: lint
Task 9: 客服 Agent 本体与注册
Files:
- Create:
app/service/agent/implementations/customer_service.py - Modify:
app/service/agent/bootstrap.py - Test:
tests/unit/service/test_customer_service_agent.py
Interfaces:
- Produces:
class CustomerServiceAgent(BaseAgent),agent_type="customer_service"、allowed_roles=("customer",)、allowed_portals=("api",)、allowed_tools=("query_knowledge",)、supported_intents=("faq","product_inquiry","policy_explain","chitchat","transfer_human")
硬约束:不得覆盖治理方法;不得直连数据库;
handle()必须显式设CoreResult.intent; 置信度三档:高=直返、中=回答+「以上信息可能不完整…」、低=兜底话术+建议转人工(专项设计 §2.3);faq命中 1 条直返原文不调模型,多条才允许generate_with_model()且 Prompt 强约束"不得新增事实"。
- Step 1-5: 写测试 → 实现 → 注册 → 跑通 → lint
- Step 6: 端到端验证
Run: 参照 docs/19 §4 的 POST /api/v1/agent-runs 链路,用客户身份 9001;
先停常驻 Worker(它会抢 run 导致假失败)。
Expected: status=succeeded,result.tool_calls 含 query_knowledge,输出末尾含免责声明。
Task 10: 验收对齐与交付物
Files:
-
Create:
docs/superpowers/analysis/2026-09-10-Phase1验收报告.md -
逐条验证老师 Phase 1 的 7 条验收标准,每条给出可复现命令 + 实际输出
-
运行双轨评测集:业务准确率(≥80% 必达、90% 目标)与合规通过率(违规 0)
-
产出答辩所需交付物:API 接口文档更新、每日会议纪要、演示脚本(3-5 个演示问题)
-
全量回归:
pytest -q、ruff check app tests、mypy app、tools/audit_schema.py
自查记录
1. Spec 覆盖:需求 Phase 1 的 16 条逐条对应 Task(F1.1 → T4/T5/T6/T7;F1.2 → T3/T4/T5/T6; F1.3 → T8/T9);合规门禁 F4/F5 → T1/T2;接口 → T7/T9;验收 → T10。
2. 占位符扫描:Task 6-10 的实现细节以要点+接口形式给出,未逐行展开代码—— 这是本计划的已知不足:Task 1-5 有可直接落盘的完整代码,Task 6-10 需要实施者在 读完 Spec 后补全实现。建议实施 Task 6 前先补写该 Task 的完整代码块。
3. 类型一致性:KnowledgeHit/KnowledgeQuery/KnowledgeSearchResult 沿用 Task 1(旧会话)
已落盘的契约;MilvusKnowledgeWriter.upsert/delete 签名与 Task 5 handler 调用一致;
事件常量 knowledge.vector_sync_requested 在 T3(导入脚本)与 T5(handler)字面量一致。
4. 外部依赖:Milvus(T6 真机)、MinIO(T4 抽象已就绪,适配器待实例)、python-docx(T4 依赖)、
config_release 发布(T7 必需)。
子项目 C:游客版智能客服(不在一期范围,答辩后再做)
用户 2026-09-10 提出要做"针对游客的智能客服 agent"。本文档把它独立成子项目 C, 不混进上面 10 个 Task。理由:老师《需求文档-修改版》Phase 1 只要求"已登录客户"这条链路; 且它必须触碰身份认证(
AGENTS.md列为高风险区域),需要单独一轮设计与评审。
C.1 已核实的卡点
读 app/api/dependencies/auth.py 与 app/core/security.py 确认:
- 全代码库没有匿名访问路径:
build_request_context强制要求 Bearer JWT,credentials is None直接抛UnauthorizedAgentError(401AUTHENTICATION_REQUIRED)。 JwtAuthenticator.authenticate()校验sub必须是纯数字且0 < int(sub) <= 2^64-1、 长度 ≤20 位——无法用"guest-abc"这类字符串主体。IdentityService.resolve()会拿该 id 去查真实用户并读角色权限——虚拟数字 id 查不到用户仍会 401。- 现有
app/api/controllers/public_platform.py名字虽叫 public,每个端点都依赖build_request_context, 并非匿名入口。
→ 结论:"给游客发个 JWT"这条捷径走不通,必须设计独立的匿名身份机制。
C.2 建议方案(待单独 brainstorm 确认)
| 设计点 | 方案 | 理由 |
|---|---|---|
| 身份签发 | 新增独立的匿名身份接口(公开站点调用,发短期匿名 token),不改动 JwtAuthenticator 的既有校验逻辑 |
避开高风险区,不破坏已验收的鉴权链路 |
| 角色 | 新增 guest 角色 |
与 customer 隔离,权限可单独收窄 |
| Agent | 独立 AgentDefinition:agent_type="guest_service"、allowed_roles=("guest",)、allowed_portals=("public",)、allowed_tools=("query_knowledge",)、supported_intents=("faq","product_inquiry","policy_explain","chitchat","transfer_human") |
复用同一个 CustomerServiceAgent 的规则与话术,只换定义 |
| 安全路由 | 完整复用 Task 8 的确定性 P0-P2 规则 | 游客最需要防的正是"被问到个人数据/诈骗",P1 分支照样拦 |
| 记忆 | 禁止召回(PlatformGovernance 的 recall_factory 对 guest 返回空) |
游客无客户归属,不得有任何记忆数据 |
| 工具白名单 | 必须单独发布 agent_tools / guest_service:faq → {"allowed_tools": ["query_knowledge"]} |
否则工具失败关闭 |
| 会话 | 独立存储 + 过期时间,不与登录用户会话混同 | 避免越权读到他人会话 |
| 转人工 | 只置 transfer_required,且话术走"引导联系官方渠道"而非"已为您转接" |
游客无工单归属 |
C.3 启动前必须先想清楚的问题(留给子项目 C 的设计轮)
- 匿名身份如何签发与校验?token 有效期多长?如何防止刷取?
guest角色的权限边界与customer的差集,是否需要在sys_role/sys_permission种子数据中固化?- 游客会话在用户登录后是否需要平滑过渡为正式会话?若需要,会话归属如何迁移?
- 游客能否触发转人工?若能,转人工申请的归属主体是谁(无
customer_id)?