diff --git a/app/infrastructure/milvus_knowledge_writer.py b/app/infrastructure/milvus_knowledge_writer.py index ad1a903..340e3ff 100644 --- a/app/infrastructure/milvus_knowledge_writer.py +++ b/app/infrastructure/milvus_knowledge_writer.py @@ -4,30 +4,101 @@ `MilvusKnowledgeClient` / `KnowledgeRetrievalService`)不得导入本模块** —— 读写物理隔离: 检索进程永远不持有写客户端,向量库故障不能从写路径传染到问答主链路,反之亦然。 -幂等口径(Task 5 裁定 3):Milvus 的 `upsert` 按主键 `knowledge_id` **覆盖**同一实体的向量与 -标量字段,因此同一 `knowledge_id` 被重复投递时结果是"向量数仍为 1、内容是最后一次写入"。 -这正是重跑导入、或正文被 UPDATE 后重投时需要的行为 —— 本适配器**不做**任何"已存在就跳过" -的判断:跳过会让 Milvus 里留着旧正文对应的旧向量(检索命中旧答案)。 +## 字段名必须探测,不能硬编码 -连接是**惰性**的:`__init__` 不连 Milvus(Task 5 不碰真机,Task 6 才建集合), -首次写入时才 `import pymilvus` 并建立 `AsyncMilvusClient`;`pymilvus` 缺失/连不上统一 -转成 `RecoverableAgentError`,交 `OutboxWorker` 的退避重试与死信机制处理(本层不写重试逻辑)。 +同一批集合名(`fin_faq_collection` 等)在不同环境里是**两套不同的 schema**(实测确认): -装配位置:本适配器的实例由**组装层**(`app/service/agent/bootstrap.py`)创建后注入 -`build_knowledge_handlers(...)`,再注册进 `WorkerRuntime.dispatch_one` —— 这件事**还没做** -(Task 6 的独立待办;Task 5 不碰 `runtime.py`/`bootstrap.py`),所以目前只有 -`dispatch_knowledge_events(...)` 这条直接调用路径能用它。 +| 逻辑字段 | 一套环境 | 另一套环境 | +|---|---|---| +| 文档标识 | `knowledge_id` | `doc_id` | +| 正文 | `snippet` | `content` | +| 章节 / 可见性 / 来源文件 | 无 | `chapter` / `visibility` / `source_file` | + +Milvus 对不存在的字段直接报错(`Attempt to insert an unexpected field`),而这些集合都 +`enable_dynamic_field=False`,所以**写错一个键名整条 upsert 就失败**。检索侧早已改用运行时 +探测(`app/core/knowledge_schema.py`),写侧此前一直硬编码 `knowledge_id`/`snippet` —— +后果是:**只要环境不是这套名字,从接口上传的知识全部同步失败,而检索侧看不出异常** +(读得到老知识,新知识静默缺席)。2026-09-13 实测踩到:22 块新知识全部 +`RecoverableAgentError: 知识向量写入失败`,事件重试 3 次后进死信。 + +现在两侧共用 `resolve_schema` 的同一份映射表,调用方只用**逻辑字段名**, +由本模块映射到该集合的真实物理名;集合没有的逻辑字段(如另一套环境没有 `intent`) +**跳过而不是报错**。 + +## 缺字段也要补:非 nullable 的标量字段是必填 + +字段名对上了还不够。改了名字之后的第二次实测报的是另一件事: + + Insert missed an field `chapter` to collection without set nullable==true or set default_value + +即集合里存在、但调用方没有提供的**非 nullable 标量字段**,Milvus 在 insert/upsert 时要求 +必须给出 —— 缺一个就整条失败。所以本模块会把「集合有、这一行没给」的 VARCHAR 字段补成 +空串;`params.max_length` 是判断 VARCHAR 的稳定依据(无需 import pymilvus 的枚举)。 + +## 幂等口径 + +Milvus 的 `upsert` 按主键**覆盖**同一实体的向量与标量字段,因此同一知识被重复投递时结果是 +"向量数仍为 1、内容是最后一次写入"。这正是重跑导入、或正文被 UPDATE 后重投时需要的行为 —— +本适配器**不做**任何"已存在就跳过"的判断:跳过会让 Milvus 里留着旧正文对应的旧向量 +(检索命中旧答案)。 + +连接是**惰性**的:`__init__` 不连 Milvus,首次写入时才 `import pymilvus` 并建立 +`AsyncMilvusClient`;`pymilvus` 缺失/连不上统一转成 `RecoverableAgentError`, +交 `OutboxWorker` 的退避重试与死信机制处理(本层不写重试逻辑)。 """ +from collections.abc import Mapping, Sequence from typing import Any from app.core.errors import RecoverableAgentError +from app.core.knowledge_schema import CollectionSchema, SchemaCache, resolve_schema -#: 向量字段名必须与集合 schema 一致:`tools/setup_milvus_knowledge_collections.py` -#: 用 `embedding` 建字段,且集合 `enable_dynamic_field=False` —— 写成别的键(例如 `vector`) -#: 会让 `upsert` 直接失败,检索侧永远命中不到(真机已核实 schema 字段名为 `embedding`)。 +#: 向量字段名必须与集合 schema 一致。两套实测 schema 都叫 `embedding`, +#: 且集合 `enable_dynamic_field=False` —— 写成别的键(例如 `vector`)会让 `upsert` 直接失败, +#: 检索侧永远命中不到。 VECTOR_FIELD = "embedding" +#: 主键字段的**逻辑**名(物理名由探测决定:`doc_id` 或 `knowledge_id`)。 +PRIMARY_LOGICAL_FIELD = "doc_id" +#: 正文字段的逻辑名(物理名可能是 `content` 或 `snippet`)。 +CONTENT_LOGICAL_FIELD = "content" + +#: 「集合里有、这一行没给」的字段该补什么值。 +#: +#: 多数 VARCHAR 补空串即可,但**有语义的字段必须给对**:`visibility` 留空会让检索侧的 +#: `visibility == "public"` 过滤把整行排除掉。实测代价(2026-09-13):22 块知识全部 +#: 写进了 Milvus(`query` 能查到),却一条都检索不到(`search` 查不到)—— +#: 现象是"入库成功但客服永远答不出新知识",比写入失败更难查。 +FIELD_DEFAULTS: dict[str, str] = { + "visibility": "public", +} + + +def _varchar_fields(description: Any) -> frozenset[str]: + """从 `describe_collection` 结果里挑出 VARCHAR 字段名。 + + 判据用 `params.max_length`:Milvus 的 VarChar 字段必带它,而向量/数值字段不带。 + 这样不必 import `pymilvus` 的 `DataType` 枚举(本模块刻意对它惰性依赖)。 + """ + if not isinstance(description, Mapping): + return frozenset() + fields = description.get("fields") + if fields is None: + schema = description.get("schema") + fields = schema.get("fields") if isinstance(schema, Mapping) else None + if not isinstance(fields, Sequence): + return frozenset() + names: list[str] = [] + for item in fields: + if not isinstance(item, Mapping): + continue + name = item.get("name") + params = item.get("params") + if isinstance(name, str) and name and isinstance(params, Mapping): + if "max_length" in params: + names.append(name) + return frozenset(names) + class MilvusKnowledgeWriter: """知识向量的写边界:upsert(覆盖)/ delete,失败一律 `RecoverableAgentError`。""" @@ -36,6 +107,10 @@ class MilvusKnowledgeWriter: self._uri = uri self._token = token self._client: Any = None + self._schemas = SchemaCache() + #: 集合名 → (字段映射, 该集合的 VARCHAR 字段)。后者用于给「集合有、这一行没给」的 + #: 标量字段补空值,理由见模块 docstring。 + self._descriptions: dict[str, tuple[CollectionSchema, frozenset[str]]] = {} async def _ensure(self) -> Any: if self._client is None: @@ -49,6 +124,38 @@ class MilvusKnowledgeWriter: raise RecoverableAgentError("Milvus 写客户端初始化失败") from exc return self._client + async def _describe(self, collection: str) -> tuple[CollectionSchema, frozenset[str]]: + """探测(并缓存)该集合的字段映射与 VARCHAR 字段集合。 + + 与检索侧同一个 `resolve_schema`,但这里必须走**异步**调用:写侧用的是 + `AsyncMilvusClient`,而 `knowledge_schema.detect_schema` 是同步实现、只服务于 + 检索侧的 `MilvusClient`。探测失败不抛异常,返回不可用的 schema 由调用方判定。 + """ + cached = self._descriptions.get(collection) + if cached is not None: + return cached + client = await self._ensure() + try: + description = await client.describe_collection(collection_name=collection) + except Exception as exc: # 探测失败=该集合不可用,由调用方给出明确错误 + schema = CollectionSchema( + collection=collection, fields={}, physical_names=frozenset(), + missing_required=(PRIMARY_LOGICAL_FIELD, CONTENT_LOGICAL_FIELD), + error=f"{type(exc).__name__}: {exc}", + ) + result: tuple[CollectionSchema, frozenset[str]] = (schema, frozenset()) + else: + schema = resolve_schema(collection, description) + result = (schema, _varchar_fields(description)) + self._descriptions[collection] = result + self._schemas.put(schema) + return result + + async def schema_for(self, collection: str) -> CollectionSchema: + """探测(并缓存)该集合的字段映射。""" + schema, _ = await self._describe(collection) + return schema + async def upsert( self, *, @@ -57,8 +164,10 @@ class MilvusKnowledgeWriter: vector: list[float], fields: dict[str, Any], ) -> None: - """按 `knowledge_id` 覆盖写入一条向量(Milvus 主键 upsert,重复投递不产生重复向量)。 + """按主键覆盖写入一条向量(Milvus 主键 upsert,重复投递不产生重复向量)。 + `fields` 的键是**逻辑字段名**(`content`/`title`/`tags`/`version`/`intent`), + 由探测结果映射成物理名;集合没有的逻辑字段直接跳过。 `fields` 里**可以没有** `intent` 键:知识契约把 `intent` 定为稀疏标签, 无显式标签时由检索侧按集合名推断(见 `knowledge_vector_worker` 的约定说明)。 """ @@ -66,12 +175,28 @@ class MilvusKnowledgeWriter: raise RecoverableAgentError("knowledge_id 不能为空") if not vector: raise RecoverableAgentError("向量不能为空") + schema, varchar_fields = await self._describe(collection) + primary = schema.resolve(PRIMARY_LOGICAL_FIELD) + if primary is None or not schema.usable: + missing = list(schema.missing_required) or [schema.error or "未知原因"] + raise RecoverableAgentError( + f"知识集合 {collection} 缺少必要字段({missing}),无法写入向量" + ) + row: dict[str, Any] = {primary: knowledge_id, VECTOR_FIELD: vector} + for logical, value in fields.items(): + physical = schema.resolve(logical) + if physical is None: + continue # 该集合没有这个字段(如 `intent`),跳过而不是让整条写入失败 + row[physical] = value + # 集合里存在、但这一行没给的 VARCHAR 字段必须补值,否则 Milvus 报 + # `Insert missed an field ...`(非 nullable 且无默认值=必填)。主键已经填过, + # 这里跳过它免得覆盖掉真正的 id;有语义的字段按 FIELD_DEFAULTS 给对值。 + for physical in varchar_fields: + if physical != primary and physical not in row: + row[physical] = FIELD_DEFAULTS.get(physical, "") client = await self._ensure() try: - await client.upsert( - collection_name=collection, - data=[{"knowledge_id": knowledge_id, VECTOR_FIELD: vector, **fields}], - ) + await client.upsert(collection_name=collection, data=[row]) except RecoverableAgentError: raise except Exception as exc: diff --git a/app/worker/knowledge_vector_worker.py b/app/worker/knowledge_vector_worker.py index fa606de..f549d9b 100644 --- a/app/worker/knowledge_vector_worker.py +++ b/app/worker/knowledge_vector_worker.py @@ -105,9 +105,11 @@ Handler = Callable[[dict[str, Any]], Awaitable[None]] #: 会超限并使整条 upsert 失败(实测:政策文档标题 346 字节 > 256,报 #: `length of varchar field title exceeds max length`)。FAQ 语料标题短,掩盖了这个缺陷; #: 长标题的政策/产品文档一灌就炸。所以写入前必须逐个字段按字节截断。 +#: 键是**逻辑字段名**(见 `MilvusKnowledgeWriter`):物理名由集合 schema 运行时探测决定 +#: (两套环境分别叫 `content` 与 `snippet`),这里统一用 `content`。 VECTOR_FIELD_LIMITS: dict[str, int] = { "title": 256, - "snippet": 4000, + "content": 4000, "tags": 512, "version": 16, "intent": 32, @@ -117,7 +119,7 @@ VECTOR_FIELD_LIMITS: dict[str, int] = { def _fit_varchar(value: object, max_bytes: int) -> str: """把值转成字符串并按 **UTF-8 字节**截断到上限内(不切坏多字节字符)。 - 截断而非报错:`title`/`snippet` 只是检索辅助与展示字段,让整条向量同步因为字段过长 + 截断而非报错:`title`/`content` 只是检索辅助与展示字段,让整条向量同步因为字段过长 失败会阻塞知识入库;而正文的完整内容仍保存在 MySQL(权威源)。 """ text = "" if value is None else str(value) @@ -204,9 +206,10 @@ def build_knowledge_handlers( raise RecoverableAgentError("嵌入维度与集合定义不一致") fields: dict[str, Any] = { "title": _fit_varchar(row.title, VECTOR_FIELD_LIMITS["title"]), - # snippet 是检索返回给模型的正文;必须是**本次**读到的 content_text, - # 否则正文更新后重投会留下旧答案(见模块 docstring 的幂等口径)。 - "snippet": _fit_varchar(row.content_text, VECTOR_FIELD_LIMITS["snippet"]), + # `content` 是**逻辑**字段名,由 writer 按集合 schema 映射成 `content` 或 `snippet`。 + # 它必须是**本次**读到的 content_text,否则正文更新后重投会留下旧答案 + # (见模块 docstring 的幂等口径)。 + "content": _fit_varchar(row.content_text, VECTOR_FIELD_LIMITS["content"]), "tags": _fit_varchar(_tags_to_string(row.tags), VECTOR_FIELD_LIMITS["tags"]), "version": _fit_varchar( row.version or "", VECTOR_FIELD_LIMITS["version"] diff --git a/docs/43-场内基金产品手册(知识库入库版).md b/docs/43-场内基金产品手册(知识库入库版).md new file mode 100644 index 0000000..bc38ca0 --- /dev/null +++ b/docs/43-场内基金产品手册(知识库入库版).md @@ -0,0 +1,168 @@ +# 场内基金产品手册 + +> 本手册面向客户,说明本平台场内基金(ETF / LOF)的公开交易规则与产品资料。 +> 手册内容均为公开信息,不含任何个人账户数据。 + +## 一、交易规则 + +### 1.1 场内基金怎么买、最小买多少 + +问:场内基金怎么买?最小买多少?一次能买多少份? + +答:本平台的场内基金按"手"撮合,**1 手 = 100 份**,委托数量必须是 100 份的整数倍。 +最小交易金额以产品详情页披露的字段为准(平台产品库中登记为 100.00 元); +实际所需金额按成交价 × 份数计算,随市价变动。 +实际可买数量与所需金额请以交易页面的实时报价为准。 + +### 1.2 场内基金的费用有哪些 + +问:买场内基金要交哪些费用?费率是多少?管理费多少?托管费多少? + +答:场内基金涉及两类费用。 + +**(1)基金自身的费用**,从基金资产中按日计提、体现在净值里,不向您单独收取: +本平台 20 只场内基金的管理费率为 **0.15%~1.20%/年**、托管费率为 **0.05%~0.20%/年**。 +其中 ETF 通常较低(如沪深300ETF 管理费 0.15%/年、托管费 0.05%/年), +LOF 与主动管理型较高(如南方积极配置混合 1.20%/年、托管费 0.20%/年)。 +货币 ETF 另收销售服务费(如货币ETF南方 0.25%/年)。各产品的具体费率见下方费率一览。 + +**(2)交易费用**,即您买卖时产生的**券商佣金**,费率由您开户的券商决定, +本平台不代收、也不设定该费率。 + +### 1.3 报价的最小变动单位 + +问:场内基金报价最小变动多少? + +答:本平台场内基金的最小价格变动单位为 **0.001 元**。 + +### 1.4 风险等级与适当性 + +问:我能买什么风险等级的产品?风险等级不匹配能买吗? + +答:产品风险等级从 **R1 到 R5**,由低到高。平台按您的**有效风险测评结果**做匹配校验: +测评等级与产品等级不匹配时,下单会被拦截。 +请先完成风险承受能力测评,并以交易页面的校验结果为准。 +客服不能修改您的测评结果,也不能代为开通跨级权限。 + +### 1.5 场内基金与场外基金的区别 + +问:场内基金和场外基金有什么不一样? + +答:场内基金(ETF / LOF)在交易所挂牌,按市价、以"手"为单位实时撮合成交; +场外基金按当日净值申赎,通常 T+1 确认。 +本平台的模拟交易针对**场内**基金,下单即按市价全额成交。 + +## 二、产品清单 + +问:平台上有哪些场内基金?产品代码是什么?有哪些 ETF? + +答:本平台场内基金共 **20 只**,包含 13 只 ETF 与 7 只 LOF。 +下表净值为最新披露值,会随行情变动: + +| 代码 | 名称 | 交易所 | 类型 | 风险等级 | 净值 | 净值日期 | +|---|---|---|---|---|---|---| +| 159329 | 沙特ETF南方 | 深交所 | ETF | R5 | 0.9167 | 2026-09-10 | +| 159382 | 创业板人工智能ETF南方 | 深交所 | ETF | R4 | 2.5841 | 2026-09-11 | +| 159511 | 通信ETF南方 | 深交所 | ETF | R4 | 2.3755 | 2026-09-11 | +| 159615 | 恒生生物科技ETF南方 | 深交所 | ETF | R5 | 1.1027 | 2026-09-11 | +| 159687 | 亚太精选ETF南方 | 深交所 | ETF | R5 | 1.8891 | 2026-09-10 | +| 159700 | 科创债ETF南方 | 深交所 | ETF | R2 | 101.8009 | 2026-09-11 | +| 159948 | 创业板ETF南方 | 深交所 | ETF | R4 | 3.6895 | 2026-09-11 | +| 510300 | 沪深300ETF | 上交所 | ETF | R3 | 4.5794 | 2026-09-11 | +| 510500 | 中证500ETF南方 | 上交所 | ETF | R4 | 7.6027 | 2026-09-11 | +| 511070 | 公司债ETF南方 | 上交所 | ETF | R2 | 103.0145 | 2026-09-11 | +| 511810 | 货币ETF南方 | 上交所 | ETF | R1 | 0.2661 | 2026-09-11 | +| 588890 | 科创芯片ETF南方 | 上交所 | ETF | R4 | 1.2327 | 2026-09-11 | +| 160105 | 南方积极配置混合(LOF) | 深交所 | LOF | R3 | 1.2514 | 2026-09-11 | +| 160127 | 南方新兴消费增长股票(LOF)A | 深交所 | LOF | R4 | 0.8364 | 2026-09-11 | +| 160128 | 南方金利定开债券A | 深交所 | LOF | R2 | 1.0260 | 2026-09-11 | +| 160129 | 南方金利定开债券C | 深交所 | LOF | R2 | 1.0240 | 2026-09-11 | +| 160142 | 南方优势产业(LOF) | 深交所 | LOF | R3 | 1.0623 | 2026-09-11 | +| 160143 | 南方创业板2年定期开放混合 | 深交所 | LOF | R3 | 1.6105 | 2026-09-11 | +| 501018 | 南方原油A | 上交所 | LOF | R5 | 2.0750 | 2026-09-10 | +| 501062 | 南方瑞合定开混合(LOF) | 上交所 | LOF | R3 | 2.1048 | 2026-09-11 | + +宽基与债券类 ETF 的风险等级多为 R2~R3,行业主题与跨境类 ETF 多为 R4~R5。 +完整清单与实时净值请以产品列表页为准。 + +## 三、重点产品 + +### 3.1 沪深300ETF(510300) + +问:沪深300ETF 的起投金额是多少?南方沪深300ETF 怎么买?510300 的费率是多少? + +答:沪深300ETF 的产品代码为 **510300**,在**上交所**挂牌,属于 ETF,风险等级 **R3**。 +它按"手"交易,**1 手 = 100 份**;最小交易金额以产品详情页披露的字段为准(登记为 100.00 元)。 +该产品最新净值为 **4.5794 元**(2026-09-11), +管理费率 **0.15%/年**、托管费率 **0.05%/年**,从基金资产中计提、不单独向您收取。 +实际成交金额、可买数量与费用请以交易页面的实时报价为准。 + +**产品披露**:本产品为同指数参考产品,非本公司发行的基金,在本平台仅用于功能演示。 + +### 3.2 其他产品的查询 + +问:某只产品的风险等级和费率是多少?怎么查? + +答:每只产品的代码、名称、交易所、类型、风险等级、最新净值见上方产品清单, +管理费与托管费见下方费率一览;也可以在产品列表页搜索代码或名称后进入详情页查看。 + +## 四、费率一览 + +问:各家基金的管理费和托管费分别是多少? + +答:本平台 20 只场内基金的费率如下(单位:%/年,均为从基金资产中计提、不单独向客户收取): + +| 代码 | 名称 | 管理费率 | 托管费率 | 销售服务费率 | +|---|---|---|---|---| +| 159329 | 沙特ETF南方 | 0.50 | 0.10 | — | +| 159382 | 创业板人工智能ETF南方 | 0.50 | 0.10 | — | +| 159511 | 通信ETF南方 | 0.50 | 0.10 | — | +| 159615 | 恒生生物科技ETF南方 | 0.50 | 0.15 | — | +| 159687 | 亚太精选ETF南方 | 0.20 | 0.05 | — | +| 159700 | 科创债ETF南方 | 0.15 | 0.05 | — | +| 159948 | 创业板ETF南方 | 0.15 | 0.05 | — | +| 510300 | 沪深300ETF | 0.15 | 0.05 | — | +| 510500 | 中证500ETF南方 | 0.15 | 0.05 | — | +| 511070 | 公司债ETF南方 | 0.15 | 0.05 | — | +| 511810 | 货币ETF南方 | 0.30 | 0.05 | 0.25 | +| 588890 | 科创芯片ETF南方 | 0.50 | 0.10 | — | +| 160105 | 南方积极配置混合(LOF) | 1.20 | 0.20 | — | +| 160127 | 南方新兴消费增长股票(LOF)A | 1.20 | 0.20 | — | +| 160128 | 南方金利定开债券A | 0.50 | 0.15 | — | +| 160129 | 南方金利定开债券C | 0.50 | 0.15 | — | +| 160142 | 南方优势产业(LOF) | 1.20 | 0.20 | — | +| 160143 | 南方创业板2年定期开放混合 | 1.20 | 0.20 | — | +| 501018 | 南方原油A | 1.00 | 0.20 | — | +| 501062 | 南方瑞合定开混合(LOF) | 1.20 | 0.20 | — | + +"销售服务费率"为"—"表示该产品不收取销售服务费。 +除上表费用外,买卖时产生的券商佣金由您开户的券商决定,本平台不代收。 + +## 五、常见问答 + +### 5.1 ETF 能定投吗? + +问:ETF 能定投吗?可以自动扣款吗? + +答:本平台的场内基金按市价实时撮合成交,**不提供自动定投计划**。 +如需分批买入,可以自行分次下单,每次委托数量为 100 份的整数倍。 + +### 5.2 场内基金多久成交、多久到账? + +问:买入后多久成交? + +答:本平台为模拟交易,委托按市价**即时全额成交**,成交后持仓与资金变动实时更新。 + +### 5.3 产品的涨跌幅在哪里看? + +问:在哪里看基金今天的涨跌? + +答:产品列表页与详情页展示最新净值与当日涨跌;产品详情页另有历史净值走势图。 + +### 5.4 起投金额是固定的吗? + +问:起投金额是多少?会变吗? + +答:场内基金按"手"交易(1 手 = 100 份), +最小交易金额以产品详情页披露的字段为准(登记为 100.00 元); +实际所需金额随市价变动,请以交易页面的实时报价为准。 diff --git a/tests/unit/worker/test_knowledge_vector_worker.py b/tests/unit/worker/test_knowledge_vector_worker.py index a12e0c3..bfba1f3 100644 --- a/tests/unit/worker/test_knowledge_vector_worker.py +++ b/tests/unit/worker/test_knowledge_vector_worker.py @@ -250,8 +250,8 @@ async def test_re_sync_after_content_update_carries_new_text() -> None: {"knowledge_id": "11"} ) - assert writer.upserts[0]["fields"]["snippet"] != updated - assert writer.upserts[1]["fields"]["snippet"] == updated + assert writer.upserts[0]["fields"]["content"] != updated + assert writer.upserts[1]["fields"]["content"] == updated async def test_intent_keeps_contract_none_for_ordinary_knowledge() -> None: @@ -320,11 +320,34 @@ async def test_dispatch_publishes_exactly_one_event_per_call() -> None: assert session.commits == 0 +#: 两套**都真实存在**的集合 schema(2026-09-13 实测 describe_collection): +#: 本机现库是主键 `doc_id` + 正文 `content`,另一套环境是 `knowledge_id` + `snippet`。 +#: 写侧靠 `describe_collection` 把**逻辑**字段名映射成这些**物理**名 —— +#: 桩不提供它,探测就会失败、整条写入被拒(这正是修复前真机上的表现: +#: 22 块知识全部 `RecoverableAgentError`,事件进死信,而检索侧看不出异常)。 +SCHEMA_DOC_ID = ("doc_id", "title", "content", "chapter", "section", "tags", + "doc_no", "version", "visibility", "embedding") +SCHEMA_KNOWLEDGE_ID = ("knowledge_id", "title", "snippet", "tags", "version", "embedding") + + class _FakeMilvusClient: - def __init__(self) -> None: + def __init__(self, fields: tuple[str, ...] = SCHEMA_DOC_ID) -> None: self.upserts: list[dict[str, Any]] = [] self.deletes: list[dict[str, Any]] = [] self.closed = False + self._fields = fields + + async def describe_collection(self, *, collection_name: str) -> dict[str, Any]: + # 除向量字段外都带 `max_length`,与真实 VarChar 字段一致 —— 写侧靠它判断 + # 「哪些是必填标量字段、这一行没给就要补空串」。 + return { + "collection_name": collection_name, + "fields": [ + {"name": name} if name == "embedding" + else {"name": name, "params": {"max_length": 1024}} + for name in self._fields + ], + } async def upsert(self, *, collection_name: str, data: list[dict[str, Any]]) -> None: self.upserts.append({"collection_name": collection_name, "data": data}) @@ -336,8 +359,22 @@ class _FakeMilvusClient: self.closed = True -async def test_writer_upserts_one_row_keyed_by_knowledge_id() -> None: - client = _FakeMilvusClient() +@pytest.mark.parametrize( + ("present_fields", "id_key", "text_key"), + [ + (SCHEMA_DOC_ID, "doc_id", "content"), + (SCHEMA_KNOWLEDGE_ID, "knowledge_id", "snippet"), + ], +) +async def test_writer_maps_logical_fields_to_collection_schema( + present_fields: tuple[str, ...], id_key: str, text_key: str +) -> None: + """逻辑字段名 → 该集合的物理名:**两套 schema 都必须能写**。 + + 写侧此前硬编码 `knowledge_id`/`snippet`,遇到另一套 schema 时整条 upsert 直接失败; + 而且这条路径只在真机上暴露(桩没有 describe_collection),所以长期没被测出来。 + """ + client = _FakeMilvusClient(present_fields) writer = MilvusKnowledgeWriter("http://milvus:19530") writer._client = client @@ -345,20 +382,60 @@ async def test_writer_upserts_one_row_keyed_by_knowledge_id() -> None: collection="fin_faq_collection", knowledge_id="11", vector=[0.5] * 1024, - fields={"title": "标题", "snippet": "正文"}, + fields={"title": "标题", "content": "正文"}, ) - assert client.upserts == [{ - "collection_name": "fin_faq_collection", - "data": [{ - "knowledge_id": "11", - # 向量字段名必须与集合 schema 一致(`embedding`,不是 `vector`): - # 集合 `enable_dynamic_field=False`,写错键会让 upsert 直接失败。 - "embedding": [0.5] * 1024, - "title": "标题", - "snippet": "正文", - }], - }] + assert client.upserts[0]["collection_name"] == "fin_faq_collection" + row = client.upserts[0]["data"][0] + assert row[id_key] == "11" + # 向量字段名必须与集合 schema 一致(`embedding`,不是 `vector`): + # 集合 `enable_dynamic_field=False`,写错键会让 upsert 直接失败。 + assert row["embedding"] == [0.5] * 1024 + assert row["title"] == "标题" + assert row[text_key] == "正文" + # 集合有、这一行没给的 VARCHAR 字段补值:Milvus 对非 nullable 且无默认值的字段要求 + # 必须提供,缺一个整条 upsert 就失败(实测缺 `chapter` 报 `Insert missed an field`)。 + # 两套 schema 只有一套带这些字段,所以只断言"集合里有的那些"。 + for optional_field in ("chapter", "section"): + if optional_field in present_fields: + assert row[optional_field] == "" + # `visibility` 例外:留空会被检索侧 `visibility == "public"` 的过滤把整行排除, + # 表现为"入库成功却一条都检索不到"(实测 22 块知识全部如此),所以按 FIELD_DEFAULTS 填。 + if "visibility" in present_fields: + assert row["visibility"] == "public" + + +async def test_writer_skips_fields_the_collection_does_not_have() -> None: + """集合没有的逻辑字段**跳过而不是报错**(另一套 schema 里没有 `intent` 字段)。""" + client = _FakeMilvusClient(SCHEMA_KNOWLEDGE_ID) + writer = MilvusKnowledgeWriter("http://milvus:19530") + writer._client = client + + await writer.upsert( + collection="fin_faq_collection", + knowledge_id="11", + vector=[0.1] * 1024, + fields={"content": "正文", "intent": "chitchat"}, + ) + + row = client.upserts[0]["data"][0] + assert "intent" not in row + assert row["snippet"] == "正文" + + +async def test_writer_fails_closed_when_schema_is_unusable() -> None: + """探测不出必要字段时必须**报错**,不能静默写入检索不到的数据。""" + client = _FakeMilvusClient(("title", "embedding")) # 既无主键、也无正文 + writer = MilvusKnowledgeWriter("http://milvus:19530") + writer._client = client + + with pytest.raises(RecoverableAgentError): + await writer.upsert( + collection="fin_faq_collection", + knowledge_id="11", + vector=[0.1] * 1024, + fields={"content": "正文"}, + ) async def test_writer_delete_targets_single_id() -> None: @@ -387,6 +464,12 @@ async def test_writer_fails_closed_on_empty_input(kwargs: Mapping[str, Any]) -> async def test_writer_translates_backend_failure_to_recoverable() -> None: class Broken: + async def describe_collection(self, *, collection_name: str) -> dict[str, Any]: + return { + "collection_name": collection_name, + "fields": [{"name": name} for name in SCHEMA_DOC_ID], + } + async def upsert(self, **_kwargs: Any) -> None: raise RuntimeError("milvus down") diff --git a/tests/unit/worker/test_runtime_knowledge_wiring.py b/tests/unit/worker/test_runtime_knowledge_wiring.py index d4f8bec..5fdde39 100644 --- a/tests/unit/worker/test_runtime_knowledge_wiring.py +++ b/tests/unit/worker/test_runtime_knowledge_wiring.py @@ -62,6 +62,20 @@ class FakeMilvusClient: self.deletes: list[dict[str, Any]] = [] self._log = log + async def describe_collection(self, *, collection_name: str) -> dict[str, Any]: + """写侧靠它把**逻辑**字段名映射成物理名;桩不提供会走"探测失败"分支。 + + 这里给本机现库的 schema:主键 `doc_id`、正文 `content` + (另一套环境是 `knowledge_id`/`snippet`,两套都由 `knowledge_schema` 统一映射)。 + """ + return { + "collection_name": collection_name, + "fields": [{"name": name} for name in ( + "doc_id", "title", "content", "chapter", "section", + "tags", "doc_no", "version", "visibility", "embedding", + )], + } + async def upsert(self, *, collection_name: str, data: list[dict[str, Any]]) -> None: self._log.append("upsert") self.upserts.append({"collection_name": collection_name, "data": data}) @@ -180,11 +194,11 @@ async def test_dispatch_one_consumes_knowledge_sync_event( write = client.upserts[0] assert write["collection_name"] == "fin_faq_collection" row = write["data"][0] - assert row["knowledge_id"] == "11" + assert row["doc_id"] == "11" # 字段名必须是集合 schema 的 `embedding`:写成 `vector` 会让 upsert 失败、检索永远命中不到。 assert VECTOR_FIELD == "embedding" assert len(row[VECTOR_FIELD]) == 1024 - assert row["snippet"] == "交易日 15:00 前提交,T+1 确认份额。" + assert row["content"] == "交易日 15:00 前提交,T+1 确认份额。" # 嵌入的是知识行正文(不是事件 payload),端点走 customer_service/embedding。 assert embedder.texts == ["交易日 15:00 前提交,T+1 确认份额。"] # 事件本身被标记为已投递,不再 pending。