Files
group_fqcd_jr/app/service/knowledge_management_service.py
lzf_0626 58b73ff594 知识库三项收口:向量-元数据对账 + 导入侧幂等 + 过期行向量清理入口
① 只读对账 tools/reconcile_knowledge_vectors.py
   按集合列出:孤儿向量 / 死向量 / 缺向量 / 重复正文 / 低信息量碎片 / 纯标题。
   关键口径:非数字 id(FAQ-0013 这类语义 id)是灌库脚本有意写进 Milvus 的,
   单独归类、不建议删;向量数取自 query 实际行数,不用 get_collection_stats
   (后者含已软删未 compaction 的行)。

② 导入侧幂等:同 source_file + 集合重传 = 覆盖上一版
   app/service/knowledge_ingest_service.py 新增 _supersede_previous_version:
   把上一版 active 行置为 expired,并逐行投 knowledge.vector_delete_requested
   (与本次入库同事务)。写入侧只认 active 而检索侧不看 status,旧向量不清掉
   会继续参与排序、和同题活块抢答。
   顺带修掉一个真 bug:改为先判 chunks 非空再下线 —— 否则传一份解析出 0 块的
   文档会把上一版下架、新版一行没写,这份文档在检索侧凭空消失。

③ 清理入口:POST /api/v1/knowledge/{knowledge_id}/vector-cleanups
   给历史上"被别的途径置为 expired、从未投过删除事件"的行补投向量清理。
   DELETE 对已过期行返回 404 的口径保持不变(重复删除静默成功会让调用方
   分不清"这次真下线了"和"早就过期了"),因此新开一个语义明确的端点:
   不存在 404 / 仍是 active 422(请改用 DELETE)/ 已 expired 200 并回传事件名。
   配套 tools/purge_expired_knowledge_vectors.py(默认 dry-run)批量驱动该端点。

文档:docs/演示用/知识库向量对账与清理-2026-09-15.md(含真机验证输出),
并对 docs/演示用/知识库问答诊断-2026-09-14.md 做两处更正 —— 实测孤儿向量 0 条、
那 175 行历史副本从来没有向量(不参与排序),当时的差额来自 get_collection_stats
把已软删行算进去。

新发现(未修,需业务拍板):661 条向量里 451 条正文不到 40 字,是灌库时把
markdown 表格/标题切碎产生的碎片。「风险评估问卷怎么评分」实测前 4 名是 4 条
一模一样的 19 字碎片(gap 0.0024),真正 2828 字的答案排第 5 → 客服必然转人工。
属灌库切分缺陷,补内容救不了,也不应靠放宽 MIN_GAP 解决。

验证:pytest tests/unit tests/contract → 1500 passed, 2 skipped, 0 failed;
mypy app → 3 个错全在组员文件中(与本次改动无关);ruff 本次改动文件 0 错。
真机端到端:重传 → 旧行 expired + 删除事件 published + 旧向量已从 Milvus 删除;
两个问句回归仍正常回答(r1到r5 gap 0.0766;申购确认 0.8453)。
2026-09-15 09:06:57 +08:00

351 lines
17 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""知识库管理服务(Task 11):上传 / 查询 / 删除文档 / 补投向量清理。
## 为什么单独一个 Service,而不是塞进 `KnowledgeIngestService`
`KnowledgeIngestService.ingest(...)` 是**切分与入库的活动链**(一个 chunk 一行知识 + 一条
向量同步事件),但它**不 commit**(事务归调用方)、不认识权限、不认识删除语义。接口层需要的
是"一个文档维度的管理用例":鉴权 → 开事务 → 复用入库链 → 提交。那一层就是本模块,
`KnowledgeIngestService` 一行不改(它的 `LAST_INSERT_ID()` 成对约束与事件形状都保持在原处)。
## 事务边界(谁 commit)
本模块**拥有事务**:`async with session.begin()` 包住整段写入并提交;入库链中途抛错
(不支持的文件类型、非法 knowledge_type)时事务整体回滚,不会留下半截知识行。
`ingest` 内部投的 `knowledge.vector_sync_requested` 与知识行同事务,因此"知识入库了但
向量同步事件丢了"这个中间态在接口层同样不成立。
## 删除语义(老师原文:标记 + 投删除事件)
一次删除在**同一个事务**里做两件事:`fin_knowledge_meta.status = 'expired'` 与投
`knowledge.vector_delete_requested`(消费侧见 `app/worker/knowledge_vector_worker.py`
的 `remove`)。**不做物理删除**:知识行是审计与对账的锚点,删掉就再也说不清"这条知识
什么时候被谁下线过"。文件字节的归档(`LocalDocumentStorage.archive()`)在提交之后做,
且**尽力而为**:归档失败只记日志,不回滚已经提交的删除 —— 否则会出现"库里还是 active、
存储里已经归档"的不一致,比"库已 expired、文件还在"危险得多(后者只浪费磁盘)。
删除只覆盖"**本次调用**造成的 active→expired"。历史上已经过期、但从未投过删除事件的行
(早期脚本改库、导入侧幂等上线前的重复副本)需要 `cleanup_vector` 补投一次 —— 为什么不能
让 `delete_document` 顺带兼容它们:那会让"重复删除"变成静默成功,调用方再也分不清
"这一次真的下线了"和"它早就过期了,我什么都没改"。
## 权限
四个端点统一要求 `knowledge:manage`(admin 角色持有;权限码在 `tools/seed_test_rbac.py`
的 `PERMISSIONS` 里,id 9019)。复用既有 `knowledge:query`(9018)会让**客户角色**看见
整库文档清单与正文摘要,知识库的后台维护面不该向客户开放。
"""
import base64
import binascii
import logging
from collections.abc import Callable
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
from uuid import uuid4
from sqlalchemy import select, update
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.contracts import RequestContext
from app.core.errors import GenericResourceNotFoundError, ValidationAgentError
from app.infrastructure.db import SessionFactory
from app.model.knowledge import KnowledgeMeta
from app.model.platform import DomainEventOutbox
from app.service.authorization_service import AuthorizationService
from app.service.document_parser import DocumentParser
from app.service.knowledge_ingest_service import (
ALLOWED_KNOWLEDGE_TYPES,
KnowledgeIngestService,
)
from app.worker.knowledge_vector_worker import (
KNOWLEDGE_AGGREGATE_TYPE,
VECTOR_DELETE_EVENT,
)
logger = logging.getLogger(__name__)
#: 管理面的唯一权限码(见模块 docstring)。
REQUIRED_PERMISSION = "knowledge:manage"
#: 删除后的状态取值。`expired` 表示已过期/已下线的知识,与
#: `KnowledgeReferenceService` 的 `ACTIVE_STATUS = "active"` 相互排斥。
EXPIRED_STATUS = "expired"
#: 列表返回的正文预览长度:管理面需要能分辨同名文档,但接口不该整篇回吐正文
#: (正文是检索侧素材,全量回吐等于绕开引用接口的脱敏口径)。
CONTENT_PREVIEW_LENGTH = 200
#: 上传时 `filename` 的最大长度(`fin_knowledge_meta.source_file` 是 `String(256)`)。
FILENAME_MAX_LENGTH = 256
#: 按"当次事务的 session"构造入库链的工厂:事务归本服务,入库链**不持有** session,
#: 否则一个 Service 实例被两次请求复用时会串事务。
IngestFactory = Callable[[AsyncSession], Any]
#: 本地文档存储根目录(MinIO 无实例时的一期实现,见 `app/infrastructure/document_storage.py`)。
#: 只影响字节落盘位置,不影响 `fin_knowledge_meta.minio_path` 里记录的相对 key。
KNOWLEDGE_STORAGE_ROOT = Path("data") / "knowledge_documents"
def decode_content(content_base64: str) -> bytes:
"""严格 base64 解码:非法输入直接失败,不做"忽略错误字符"的宽容解码。
宽容解码(`validate=False`)会把 `!!!` 之类解码成空字节,最终表现为
"上传成功但文档入库为空"——错误被推到很远的地方才暴露。
"""
try:
return base64.b64decode(content_base64, validate=True)
except (binascii.Error, ValueError) as exc:
raise ValidationAgentError("content_base64 不是合法的 base64 内容") from exc
class KnowledgeManagementService:
"""知识文档管理用例:上传(切分+入库)、列表查询、删除(标记+投递向量删除)。"""
def __init__(
self,
*,
ingest_factory: IngestFactory,
session_factory: Callable[[], Any] | None = None,
storage: Any = None,
archive_on_delete: bool = True,
) -> None:
self._ingest_factory = ingest_factory
self._session_factory: Callable[[], Any] = session_factory or SessionFactory
self._storage = storage
self._archive_on_delete = archive_on_delete
async def upload(
self,
context: RequestContext,
*,
filename: str,
content: bytes,
knowledge_type: str,
) -> dict[str, Any]:
"""上传一份文档并自动入库,返回本次产生的 `knowledge_id` 列表。
`created_by` **只**来自认证上下文(`context.user_id`),方法不提供该参数:
允许调用方传就等于允许伪造导入人,而 `fin_knowledge_meta.reviewer_id` 是
"谁导入了这条知识"的唯一线索。
"""
await AuthorizationService.require(context, REQUIRED_PERMISSION)
self._assert_filename(filename)
async with self._session_factory() as session, session.begin():
knowledge_ids = await self._ingest_factory(session).ingest(
filename=filename,
content=content,
knowledge_type=knowledge_type,
created_by=int(context.user_id),
)
return {
"knowledge_ids": list(knowledge_ids),
"filename": filename,
"knowledge_type": knowledge_type,
"created_by": context.user_id,
"chunk_count": len(knowledge_ids),
}
async def list_documents(
self,
context: RequestContext,
*,
limit: int,
offset: int,
knowledge_type: str | None = None,
) -> dict[str, Any]:
"""列出未过期(`status != 'expired'`)的知识行,按 id 倒序。
默认口径是**只返回未过期行**:被删除的知识仍在库里(删除是标记而非物理删除),
管理面默认视图不应该再列出它们。
"""
await AuthorizationService.require(context, REQUIRED_PERMISSION)
if knowledge_type is not None and knowledge_type not in ALLOWED_KNOWLEDGE_TYPES:
raise ValidationAgentError(f"knowledge_type 非法:{knowledge_type}")
statement = (
select(KnowledgeMeta)
.where(KnowledgeMeta.status != EXPIRED_STATUS)
.order_by(KnowledgeMeta.id.desc())
.limit(limit)
.offset(offset)
)
if knowledge_type is not None:
statement = statement.where(KnowledgeMeta.knowledge_type == knowledge_type)
async with self._session_factory() as session:
rows = list((await session.scalars(statement)).all())
return {"items": [self._to_item(row) for row in rows], "count": len(rows)}
async def delete_document(
self, context: RequestContext, knowledge_id: int
) -> dict[str, Any]:
"""删除一份文档:标记 `expired` + 投递向量删除事件(+ 尽力归档文件字节)。
不存在的 id(以及已经删除过的 id)一律抛 `SESSION_NOT_FOUND` 语义的 404,
**不静默成功**:静默成功会让调用方以为向量已清理,而实际上什么都没发生。
"""
await AuthorizationService.require(context, REQUIRED_PERMISSION)
async with self._session_factory() as session, session.begin():
row = await session.get(KnowledgeMeta, knowledge_id)
if row is None or row.status == EXPIRED_STATUS:
raise GenericResourceNotFoundError("知识文档不存在")
storage_key = row.minio_path
now = datetime.now(UTC).replace(tzinfo=None)
await session.execute(
update(KnowledgeMeta)
.where(KnowledgeMeta.id == knowledge_id)
.values(status=EXPIRED_STATUS, updated_at=now)
)
self._enqueue_vector_delete(
session, knowledge_id, trace_id=context.trace_id, now=now
)
await self._archive(storage_key)
return {
"knowledge_id": knowledge_id,
"status": EXPIRED_STATUS,
"vector_delete_event": VECTOR_DELETE_EVENT,
}
async def cleanup_vector(
self, context: RequestContext, knowledge_id: int
) -> dict[str, Any]:
"""给**已经过期**的历史行补投一次向量删除事件(清掉向量侧残留)。
## 为什么需要它
`delete_document` 只在"active → expired"这**一次状态转换**上投事件。历史上有一批行
是**别的途径**变成 `expired` 的(早期脚本直接改库、重复入库被覆盖、导入侧幂等上线前的
重复副本),它们**从来没投过**删除事件,向量还留在 Milvus 里继续参与排序;
而 `delete_document` 对这些行**刻意返回 404**("不静默成功"),于是管理端口上
根本没有一个能清理它们的动作 —— 这是行为与运维需求之间的缺口,本方法就是补这个缺口。
## 口径(三条,都是"宁可报错也不假装成功")
1. id 不存在 → 404(`SESSION_NOT_FOUND` 语义),与删除接口一致;
2. 行**仍是 active** → 校验失败(422):在用的文档不该走"只清向量不标状态"这条路,
那会让检索侧读不到正文而知识行看起来还好好的,请改用 DELETE(它同时做两件事);
3. 成功 → 投出一条**新的** `knowledge.vector_delete_requested`(`event_id` 是新 uuid4,
消费侧删除幂等,重复清理同一行只是多删一次空集,不会报错)。
返回 `status`(当前状态)与 `vector_delete_event`(投出的事件名),
调用方据此能确认"事件确实被投出去了",而不是只看一个 200。
"""
await AuthorizationService.require(context, REQUIRED_PERMISSION)
async with self._session_factory() as session, session.begin():
row = await session.get(KnowledgeMeta, knowledge_id)
if row is None:
raise GenericResourceNotFoundError("知识文档不存在")
status = str(row.status)
if status != EXPIRED_STATUS:
raise ValidationAgentError(
f"知识文档当前状态为 {status},只有已过期/已下线的文档才需要补投向量清理;"
"在用的文档请改用 DELETE 接口(它同时标记过期并投出删除事件)"
)
now = datetime.now(UTC).replace(tzinfo=None)
self._enqueue_vector_delete(
session, knowledge_id, trace_id=context.trace_id, now=now
)
logger.info(
"补投向量清理事件:knowledge_id=%s status=%s trace_id=%s",
knowledge_id, status, context.trace_id,
)
return {
"knowledge_id": knowledge_id,
"status": status,
"vector_delete_event": VECTOR_DELETE_EVENT,
}
# --- 内部 ---------------------------------------------------------------
@staticmethod
def _assert_filename(filename: str) -> None:
"""文件名不得为空、不得超列宽、不得带路径:`source_file` 只存文件名本身。
带路径的文件名会被 `KnowledgeIngestService._storage_key` 剥成 basename 落存储,
但 `source_file` 列的取值同样要干净,否则管理面列表里会出现调用方的本地路径。
"""
if not filename.strip():
raise ValidationAgentError("filename 不得为空")
if len(filename) > FILENAME_MAX_LENGTH:
raise ValidationAgentError(f"filename 超过 {FILENAME_MAX_LENGTH} 字符")
if filename != Path(filename).name or "\\" in filename:
raise ValidationAgentError("filename 只能是文件名,不得包含路径")
@staticmethod
def _enqueue_vector_delete(
session: AsyncSession, knowledge_id: int, *, trace_id: str, now: datetime
) -> None:
"""投一条向量删除事件;`aggregate_id` 与 payload 都是 `str(knowledge_id)`(投递同口径)。
`event_id` 是唯一键,显式给 uuid4;`aggregate_type` 沿用投递侧常量,
保证消费侧按同一聚合维度幂等覆盖。
"""
session.add(DomainEventOutbox(
event_id=str(uuid4()),
event_type=VECTOR_DELETE_EVENT,
aggregate_type=KNOWLEDGE_AGGREGATE_TYPE,
aggregate_id=str(knowledge_id),
trace_id=trace_id or str(uuid4()),
payload={"knowledge_id": str(knowledge_id)},
status="pending",
retry_count=0,
occurred_at=now,
created_at=now,
updated_at=now,
))
async def _archive(self, storage_key: str | None) -> None:
"""提交之后把文件字节移入 `archive/`;失败只记日志(见模块 docstring 的取舍)。"""
if not self._archive_on_delete or not storage_key or self._storage is None:
return
try:
await self._storage.archive(key=storage_key)
except Exception: # noqa: BLE001 - 归档尽力而为,任何失败都不影响已提交的删除
logger.warning("知识文档归档失败,已提交的删除不回滚 storage_key=%s", storage_key)
@staticmethod
def _to_item(row: KnowledgeMeta) -> dict[str, Any]:
content = row.content_text or ""
tags = row.tags if isinstance(row.tags, dict) else {}
return {
"knowledge_id": row.id,
"knowledge_type": row.knowledge_type,
"title": row.title,
"source_file": row.source_file,
"collection": row.milvus_collection,
"version": row.version,
"status": row.status,
"review_status": row.review_status,
"created_by": row.reviewer_id,
"tags": tags,
"content_preview": content[:CONTENT_PREVIEW_LENGTH],
"content_length": len(content),
}
def build_knowledge_management_service() -> KnowledgeManagementService:
"""组合根:把接口层需要的依赖一次装好(Controller 只调它,不自己拼装配)。
`embedder` / `endpoint_resolver` 在 `KnowledgeIngestService` 里是**给后续对账自愈预留的**
依赖(入库链自己不调用 embedding,向量化由 Outbox Worker 做),因此直接用组装层
同一套工厂函数(`app.service.agent.bootstrap`)构造,不另写一套模型网关装配。
"""
from app.infrastructure.document_storage import LocalDocumentStorage
from app.service.agent.bootstrap import get_memory_embedding_service
from app.service.model_gateway import DatabaseModelEndpointResolver
def ingest_for(session: AsyncSession) -> KnowledgeIngestService:
return KnowledgeIngestService(
session=session,
parser=DocumentParser(),
storage=LocalDocumentStorage(KNOWLEDGE_STORAGE_ROOT),
embedder=get_memory_embedding_service(),
endpoint_resolver=DatabaseModelEndpointResolver(),
)
return KnowledgeManagementService(
ingest_factory=ingest_for,
session_factory=SessionFactory,
storage=LocalDocumentStorage(KNOWLEDGE_STORAGE_ROOT),
)