141 lines
6.4 KiB
Python
141 lines
6.4 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
|
|
name = (body.name or "").strip()
|
|
if create_kb(resolve_workspace(None), user_id, name) is None:
|
|
raise HTTPException(400, f"invalid kb name: {body.name!r}")
|
|
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
|
|
detail = kb_detail(resolve_workspace(None), user_id, name)
|
|
if detail is None:
|
|
raise HTTPException(404, f"kb not found: {name!r}")
|
|
detail["ingest"] = ingest_status(user_id, name)
|
|
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_ingest import ingest_status
|
|
if ingest_status(user_id, name).get("running"):
|
|
raise HTTPException(409, "该库正在入库,等入库结束再删除")
|
|
if not delete_kb(resolve_workspace(None), user_id, name):
|
|
raise HTTPException(404, f"kb not found: {name!r}")
|
|
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_source
|
|
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}")
|
|
saved: list[dict] = []
|
|
for up in files or []:
|
|
raw_name = up.filename or ""
|
|
data = await up.read()
|
|
ok = save_source(ws, user_id, name, raw_name, data)
|
|
if ok is None:
|
|
raise HTTPException(400, f"invalid filename: {raw_name!r}")
|
|
saved.append({"name": ok, "size": len(data)})
|
|
if not saved:
|
|
raise HTTPException(400, "no files uploaded")
|
|
_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
|
|
if not delete_doc(resolve_workspace(None), user_id, name, filename):
|
|
raise HTTPException(404, f"doc not found: {filename!r}")
|
|
return {"deleted": filename}
|