"""知识入库服务:文件落存储 → 元数据落库(切块逐行)→ 投递向量同步事件。 ## 分层边界 Service 层(MVC+S):使用**传入的** `AsyncSession`,不自建 Session、不直连 MySQL/Redis/Milvus, 也不调用 embedding —— Milvus 的写入与向量化全部由 `knowledge.vector_sync_requested` 事件的 Outbox handler(`app/worker/knowledge_vector_worker.py`)完成。这里只把"这件事需要做"写进 `domain_event_outbox`,与知识行处于**同一个事务**(调用方提交才一起生效),因此不存在 "知识入库了但向量同步事件丢了"的中间态。 构造器收 `embedder` / `endpoint_resolver` 不是摆设:它们是**后续对账自愈**要用的依赖 (裁定 3 —— 自愈补投的事件最终要靠它们落地向量),由调用方在组合根一次性注入, 避免 Service 以后为了自愈去偷偷 new 一个模型网关。 ## 为什么 `ingest` 返回 `list[int]` 而不是单个 `int` 一份文档会被切分成**多个 chunk**(`DocumentParser`:512 字符/块、overlap 64,且 overlap 会在 相邻块间重复出现,见该模块 docstring),每个 chunk 是 `fin_knowledge_meta` 里**独立的一行** (独立 `knowledge_id`、独立向量)。返回单个 int 无法表达"这次入库产生了哪些知识", 调用方也没法为多行知识建同义问法或回滚。所以返回**按 chunk 顺序排列的 `knowledge_id` 列表**; 空文档返回空列表。 ## 尚未接入生产 worker(**已知缺口,不在 Task 5 范围**) 本服务投出的事件(`knowledge.vector_sync_requested`)**目前没有生产消费者**: `WorkerRuntime.dispatch_one`(`app/worker/runtime.py`)的 handlers 是硬编码白名单,里面没有这两个 知识事件,需由组装层(`app/service/agent/bootstrap.py`,Task 7/9 范围)调用 `app/worker/knowledge_vector_worker.py` 的 `build_knowledge_handlers(...)` 并**注册进 `WorkerRuntime.dispatch_one` 的 handlers 字典**才能被领取。在那之前事件只会堆在 `domain_event_outbox`(现库已有 408 条),"入库成功"只代表**知识行与事件同一事务落库**, **不代表向量已同步**。 ## 校验口径 - `knowledge_type` 必须属于白名单; - 集合名**不得由调用方随意指定**:只允许取 `TYPE_TO_COLLECTION[knowledge_type]`, 显式传入就必须**等于**它,否则 `ValidationAgentError`(全库只有三个知识集合, 写错集合等于让检索永远查不到这条知识); - `created_by` 解析顺序:`ingest(created_by=...)` > 构造器 `created_by`;两者都没有则 `ValidationAgentError`(`fin_knowledge_meta.reviewer_id` 是"谁导入的"唯一线索,不能静默写空)。 """ import hashlib import json from datetime import UTC, datetime from pathlib import Path from typing import Any from uuid import uuid4 from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncSession from app.core.errors import ValidationAgentError from app.core.knowledge_contracts import ALLOWED_COLLECTIONS from app.service.document_parser import DocumentParser # 裁定 5:事件常量只有一份真相,在 worker 模块里;这里 import 而不是另写字面量。 from app.worker.knowledge_vector_worker import KNOWLEDGE_AGGREGATE_TYPE, VECTOR_SYNC_EVENT #: 允许入库的知识类型与它们**唯一**对应的 Milvus 集合(`ALLOWED_COLLECTIONS` 是契约层白名单)。 TYPE_TO_COLLECTION: dict[str, str] = { "faq": "fin_faq_collection", "product": "fin_product_collection", "policy": "fin_policy_collection", } ALLOWED_KNOWLEDGE_TYPES = frozenset(TYPE_TO_COLLECTION) #: 入库默认版本号(同一份文档重导不会改变它;版本治理由后续任务的审核流程负责)。 DEFAULT_VERSION = "v1" #: 标题列是 `String(256)`,chunk 正文可能远超;标题取正文前若干字符即可。 TITLE_MAX_LENGTH = 120 INSERT_KNOWLEDGE_META = 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 (:knowledge_type, :title, :source_file, :minio_path, :milvus_collection, :version," " :content_text, :tags, :reviewer_id, 'published', 'active', :now, :now)" ) SELECT_LAST_INSERT_ID = text("SELECT LAST_INSERT_ID()") INSERT_OUTBOX_EVENT = 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 (:event_id, :event_type, :aggregate_type, :aggregate_id, :trace_id, :payload," " 'pending', 0, :now, :now, :now)" ) class KnowledgeIngestService: def __init__( self, *, session: AsyncSession, parser: DocumentParser, storage: Any, embedder: Any, endpoint_resolver: Any, created_by: int | None = None, ) -> None: self._session = session self._parser = parser self._storage = storage self._embedder = embedder self._endpoint_resolver = endpoint_resolver self._created_by = created_by async def ingest( self, *, filename: str, content: bytes, knowledge_type: str, collection: str | None = None, created_by: int | None = None, ) -> list[int]: """把一份文档切成多块写入知识表,并为每一块投一条向量同步事件。 返回**该文档产生的 `knowledge_id` 列表**(按 chunk 顺序;空文档返回 `[]`)。 不 commit:事务归调用方(通常是与审核/审计记录一起提交,或整体回滚)。 """ expected_collection = self._collection_for(knowledge_type, collection) reviewer_id = self._resolve_created_by(created_by) key = self._storage_key(knowledge_type, filename, content) await self._storage.save( key=key, content=content, content_type="application/octet-stream" ) # 切分在落存储之后、写库之前:解析失败(UnsupportedDocumentError)时库里不留半截知识, # 只在存储里留一个未引用的对象,比留一堆无事件的孤儿知识行安全。 chunks = self._parser.parse(filename=filename, content=content) now = datetime.now(UTC).replace(tzinfo=None) knowledge_ids: list[int] = [] for chunk in chunks: knowledge_id = await self._insert_knowledge_row( chunk=chunk, filename=filename, knowledge_type=knowledge_type, collection=expected_collection, storage_key=key, reviewer_id=reviewer_id, now=now, ) knowledge_ids.append(knowledge_id) await self._enqueue_vector_sync(knowledge_id, now=now) return knowledge_ids def _collection_for(self, knowledge_type: str, collection: str | None) -> str: """白名单校验:类型必须已知,显式给定的集合必须等于该类型的唯一集合。""" if knowledge_type not in ALLOWED_KNOWLEDGE_TYPES: raise ValidationAgentError(f"knowledge_type 非法:{knowledge_type}") expected = TYPE_TO_COLLECTION[knowledge_type] if collection is not None and collection != expected: raise ValidationAgentError( f"集合名只能是 {expected},不接受调用方指定的 {collection}" ) if expected not in ALLOWED_COLLECTIONS: # pragma: no cover - 契约层与本地映射的护栏 raise ValidationAgentError(f"集合不在契约白名单内:{expected}") return expected def _resolve_created_by(self, created_by: int | None) -> int: resolved = created_by if created_by is not None else self._created_by if resolved is None: raise ValidationAgentError("缺少 created_by:入库必须记录导入人") return resolved @staticmethod def _storage_key(knowledge_type: str, filename: str, content: bytes) -> str: """`kb//<内容前 16 位 sha256>-<文件名>`:同一份文档重导落到同一个 key。""" digest = hashlib.sha256(content).hexdigest() return f"kb/{knowledge_type}/{digest[:16]}-{Path(filename).name}" async def _insert_knowledge_row( self, *, chunk: Any, filename: str, knowledge_type: str, collection: str, storage_key: str, reviewer_id: int, now: datetime, ) -> int: """插入一行知识并返回它的自增 `id`。 ## 与 `LAST_INSERT_ID()` 的硬耦合(已知脆弱点,**不要**在中间插语句) `LAST_INSERT_ID()` 是 **MySQL 连接级的隐式状态**,不是"上一次 INSERT 的返回值": 它只在**同一个连接**上、且**中间没有其它 INSERT** 时才等于本行 id。 因此本方法把 `INSERT` 与紧随其后的 `SELECT LAST_INSERT_ID()` 放在**同一方法体内成对**执行, 中间的语句数是 **0**,并删掉了 `sqlalchemy.text` 之外的任何额外查询 (`test_last_insert_id_is_read_once_per_chunk_in_the_same_statement_order` 把"成对且相邻"的语句序列固定下来)。调用方在本方法返回前不得插入别的写操作。 **为什么不用"显式回读"替代**(本方案评估过、明确否决): `fin_knowledge_meta` 的唯一键只有主键 `id`(DDL 见 `docs/evidence/20260909-before-constraint-fix.sql`: `PRIMARY KEY (id)` + `KEY idx_knowledge_status(status)`), `(knowledge_type, source_file, minio_path)` 上**没有唯一约束**。回读只能靠这几列近似匹配, 而同一份文档用同一 key 重导会产出**完全相同的行** → 回读 `LIMIT 1` 可能拿到上次导入的 旧行 id(`ORDER BY id` 也救不了:新行 id 更大,反而更易混淆),会静默写坏 outbox 事件。 "先算好外部 id 再插入"同样不可行:`id` 是 `bigint unsigned AUTO_INCREMENT`,要显式指定就得 自己分配全局唯一整数,等于引入新表/新序列(触碰基线),且要额外处理并发分配的唯一性。 所以在本 Schema 下 `LAST_INSERT_ID()` 是**唯一**无需改基线的可行做法。 ## 失败路径:两个分支的文案必须区分真相 - `rowcount == 0`(INSERT 没生效):此时 `LAST_INSERT_ID()` 保持**上一次**的值、 新连接上是 0,所以真相是"插入没生效", **不能**报成"取不到 id"(那会把排障引向错误方向)。 - `LAST_INSERT_ID() == 0` 而 INSERT 影响了 1 行:说明连接上还没有过自增插入就拿到了 0, 文案直接点出 `LAST_INSERT_ID=0`,并提示不要对它取整。 """ result = await self._session.execute(INSERT_KNOWLEDGE_META, { "knowledge_type": knowledge_type, "title": chunk.text[:TITLE_MAX_LENGTH], "source_file": filename, "minio_path": storage_key, "milvus_collection": collection, "version": DEFAULT_VERSION, "content_text": chunk.text, # `fin_knowledge_meta.tags` 是 JSON 列:这里写 json 字符串(与 # `tools/import_knowledge_seed.py` 同口径),读回来是 dict。 "tags": json.dumps( {"heading_path": list(chunk.heading_path), "chunk_index": chunk.index}, ensure_ascii=False, ), "reviewer_id": reviewer_id, "now": now, }) # `rowcount` 可能是 -1(驱动不报告行数),所以只在"明确为 0"时判失败,不把 -1 当失败。 affected = getattr(result, "rowcount", None) if isinstance(affected, int) and affected == 0: raise ValidationAgentError( "知识行插入未生效(INSERT 影响 0 行),未继续取 LAST_INSERT_ID()" ) inserted = await self._session.scalar(SELECT_LAST_INSERT_ID) if inserted is None: # pragma: no cover - 只在非 MySQL 后端上发生 raise ValidationAgentError("插入知识行后取不到 LAST_INSERT_ID()(返回 NULL)") if int(inserted) == 0: # pragma: no cover - 真机上的异常连接状态 raise ValidationAgentError( "插入知识行后 LAST_INSERT_ID=0:取不到本次自增 id,不能拿它当 knowledge_id" ) return int(inserted) async def _enqueue_vector_sync(self, knowledge_id: int, *, now: datetime) -> None: """投一条向量同步事件;`aggregate_id` 与 payload 都用 `str(knowledge_id)`。 `domain_event_outbox` 上没有唯一约束,所以**重复导入会产生多条同 `aggregate_id` 的事件**: 这是已知且被接受的 —— 消费侧按 `aggregate_id` 覆盖写入而非跳过 (见 `knowledge_vector_worker` 的幂等口径),因此重复投递不会在 Milvus 里留下重复向量。 """ await self._session.execute(INSERT_OUTBOX_EVENT, { "event_id": str(uuid4()), "event_type": VECTOR_SYNC_EVENT, "aggregate_type": KNOWLEDGE_AGGREGATE_TYPE, "aggregate_id": str(knowledge_id), "trace_id": str(uuid4()), "payload": json.dumps({"knowledge_id": str(knowledge_id)}), "now": now, })