zcbot/web/routers/kb.py

159 lines
7.0 KiB
Python

"""个人知识库路由(纯文件 .kb/,照记忆机制范式,DESIGN §3.8)。
状态全在 FS(core/kb.py),端点只是薄壳;入库后台任务照定时执行器范式
(create_task + to_thread,per-(user,库) 锁在 core/kb_ingest.py 内去重)。
不设 HTTP 检索端点 —— agent 侧走 fs 工具 read/grep(kb_block 注入契约)。
"""
from __future__ import annotations
import asyncio
from uuid import UUID
from fastapi import Depends, File, HTTPException, UploadFile
from core.paths import ROOT
from ..schemas import KbCreateRequest
_kb_ingest_tasks: set = set() # 持引用防 GC(fire-and-forget task 的惯例)
def _spawn_kb_ingest(user_id: UUID, name: str) -> None:
from core.agent_builder import load_config, resolve_workspace
from core.kb_ingest import run_ingest
cfg = load_config()
ws = resolve_workspace(None)
t = asyncio.create_task(
asyncio.to_thread(run_ingest, ws, user_id, name, ROOT / cfg["models_dir"])
)
_kb_ingest_tasks.add(t)
t.add_done_callback(_kb_ingest_tasks.discard)
def register_kb_routes(app, *, require_user) -> None:
@app.get("/v1/kb", tags=["kb"])
def kb_list(user_id: UUID = Depends(require_user)):
"""所有库概览:[{name, doc_count, pending_count}]。前端入口一次拉满。"""
from core.agent_builder import resolve_workspace
from core.kb import list_kbs
return {"results": list_kbs(resolve_workspace(None), user_id)}
@app.post("/v1/kb", tags=["kb"])
def kb_create(body: KbCreateRequest, user_id: UUID = Depends(require_user)):
"""建库(幂等)。名字非法 → 400。"""
from core.agent_builder import resolve_workspace
from core.kb import create_kb
from core.kb_lock import KbBusyError
name = (body.name or "").strip()
try:
if create_kb(resolve_workspace(None), user_id, name) is None:
raise HTTPException(400, f"invalid kb name: {body.name!r}")
except KbBusyError:
raise HTTPException(409, "该库正在被其他操作修改")
return {"name": name}
@app.get("/v1/kb/{name}", tags=["kb"])
def kb_detail_ep(name: str, user_id: UUID = Depends(require_user)):
"""单库全貌:INDEX 条目 + 待入库列表 + 入库进度(前端轮询这里)。"""
from core.agent_builder import resolve_workspace
from core.kb import kb_detail
from core.kb_ingest import ingest_status
ws = resolve_workspace(None)
detail = kb_detail(ws, user_id, name)
if detail is None:
raise HTTPException(404, f"kb not found: {name!r}")
detail["ingest"] = ingest_status(user_id, name, ws)
return detail
@app.delete("/v1/kb/{name}", tags=["kb"])
def kb_delete(name: str, user_id: UUID = Depends(require_user)):
"""整库删除(原件 + docs + INDEX)。入库进行中 → 409(避免半截写盘)。"""
from core.agent_builder import resolve_workspace
from core.kb import delete_kb
from core.kb_lock import KbBusyError
try:
if not delete_kb(resolve_workspace(None), user_id, name):
raise HTTPException(404, f"kb not found: {name!r}")
except KbBusyError:
raise HTTPException(409, "该库正在入库或被其他操作修改")
return {"deleted": name}
@app.post("/v1/kb/{name}/upload", tags=["kb"])
async def kb_upload(
name: str,
files: list[UploadFile] = File(...),
user_id: UUID = Depends(require_user),
):
"""上传原件到 sources/ 并立即触发入库(上传即入库)。
同名覆盖 = 重新入库(core/kb.py::save_source 会摘掉旧 INDEX 条目重排队)。
磁盘配额 gate 与 /v1/files/upload 同款。
"""
from core.agent_builder import load_config as _load_cfg, resolve_workspace
from core.kb import kb_dir, save_sources
from core.kb_lock import KbBusyError
from core.storage.disk_quota import check_disk_quota, parse_bytes
_quotas_cfg = (_load_cfg().get("quotas") or {})
_limit = parse_bytes(_quotas_cfg.get("disk_bytes_per_user"))
if _limit is not None and _limit > 0:
_err = check_disk_quota(user_id, _limit)
if _err is not None:
raise HTTPException(413, _err)
ws = resolve_workspace(None)
d = kb_dir(ws, user_id, name)
if d is None or not d.is_dir():
raise HTTPException(404, f"kb not found: {name!r}")
uploads: list[tuple[str, bytes]] = []
for up in files or []:
raw_name = up.filename or ""
data = await up.read()
uploads.append((raw_name, data))
if not uploads:
raise HTTPException(400, "no files uploaded")
try:
saved_names = save_sources(ws, user_id, name, uploads)
except KbBusyError:
raise HTTPException(409, "该库正在入库或被其他操作修改")
if saved_names is None:
raise HTTPException(400, "存在非法文件名")
saved = [
{"name": saved_name, "size": len(data)}
for saved_name, (_raw_name, data) in zip(saved_names, uploads)
]
_spawn_kb_ingest(user_id, name)
return {"count": len(saved), "saved": saved}
@app.post("/v1/kb/{name}/ingest", tags=["kb"])
async def kb_ingest_trigger(name: str, user_id: UUID = Depends(require_user)):
"""手动触发入库(兜待入库缺口:上传时崩溃 / 单件失败重试)。幂等,已在跑则空转。"""
from core.agent_builder import resolve_workspace
from core.kb import kb_dir
if kb_dir(resolve_workspace(None), user_id, name) is None:
raise HTTPException(404, f"kb not found: {name!r}")
_spawn_kb_ingest(user_id, name)
return {"started": True}
@app.get("/v1/kb/{name}/docs/{filename}", tags=["kb"])
def kb_doc_read(name: str, filename: str, user_id: UUID = Depends(require_user)):
"""读单篇转换后 markdown(点开列表项时拉)。穿越校验收口 core/kb.py::read_doc。"""
from core.agent_builder import resolve_workspace
from core.kb import read_doc
content = read_doc(resolve_workspace(None), user_id, name, filename)
if content is None:
raise HTTPException(404, f"doc not found: {filename!r}")
return {"filename": filename, "content": content}
@app.delete("/v1/kb/{name}/docs/{filename}", tags=["kb"])
def kb_doc_delete(name: str, filename: str, user_id: UUID = Depends(require_user)):
"""删单篇(docs 文件 + INDEX 行 + source 原件,防原件被重新入库)。"""
from core.agent_builder import resolve_workspace
from core.kb import delete_doc
from core.kb_lock import KbBusyError
try:
if not delete_doc(resolve_workspace(None), user_id, name, filename):
raise HTTPException(404, f"doc not found: {filename!r}")
except KbBusyError:
raise HTTPException(409, "该库正在入库或被其他操作修改")
return {"deleted": filename}