deerflow-code/offline-backend-20260512/backend/app/gateway/routers/skill_knowledge.py
2026-09-07 18:24:55 +08:00

401 lines
17 KiB
Python

"""User-facing skill distillation APIs plus administrator maintenance APIs."""
from __future__ import annotations
import asyncio
from typing import Any, Literal
from croniter import croniter
from fastapi import APIRouter, HTTPException, Query, Request
from pydantic import BaseModel, Field, field_validator
from app.gateway.deps import get_optional_user_from_request
from app.gateway.skill_knowledge_job_executor import _skill_directory
from deerflow.config.system_settings import (
SkillKnowledgeAutoSyncSettings,
load_system_settings,
save_system_settings,
)
from deerflow.integrations.weknora.runtime import (
build_weknora_client,
get_resolved_llmwiki_runtime,
)
from deerflow.skill_knowledge import scan_skill
from deerflow.skills.storage import get_or_new_skill_storage
from deerflow.skills.types import SkillCategory
router = APIRouter(prefix="/api/skill-knowledge", tags=["skill-knowledge"])
async def _actor(request: Request) -> tuple[str, bool]:
user = await get_optional_user_from_request(request)
if user is None:
return "system", True
return str(user.id), getattr(user, "system_role", None) == "admin"
async def _require_admin(request: Request) -> str:
actor, is_admin = await _actor(request)
if not is_admin:
raise HTTPException(status_code=403, detail="技能知识归纳仅管理员可操作")
return actor
def _owned_by_actor(row: dict[str, Any] | None, actor: str, is_admin: bool) -> bool:
return bool(row and (is_admin or str(row.get("created_by") or "") == actor))
def _store(request: Request):
value = getattr(request.app.state, "skill_knowledge_store", None)
if value is None:
raise HTTPException(status_code=503, detail="技能知识归纳存储不可用")
return value
class SyncTarget(BaseModel):
target_type: Literal["weknora", "assistant"]
target_mode: Literal["wiki", "document", "faq", "conversation_archive", "readonly_sidecar", "assistant"]
target_id: str = Field(min_length=1, max_length=128)
target_name: str = Field(default="", max_length=255)
remote_business_key: str | None = Field(default=None, max_length=512)
readonly: bool = False
class CreateSyncJobRequest(BaseModel):
skill_names: list[str] = Field(min_length=1, max_length=200)
targets: list[SyncTarget] = Field(min_length=1, max_length=1)
force: bool = False
review_mode: Literal["required", "auto_high_confidence", "off"] = "off"
@field_validator("skill_names")
@classmethod
def validate_names(cls, values: list[str]) -> list[str]:
result = list(dict.fromkeys(value.strip() for value in values if value.strip()))
if not result:
raise ValueError("At least one skill is required")
return result
class ReviewRequest(BaseModel):
review_ids: list[str] = Field(min_length=1, max_length=500)
action: Literal["approve", "reject"]
edits: dict[str, dict[str, Any]] = Field(default_factory=dict)
reason: str | None = Field(default=None, max_length=1024)
class DetachRequest(BaseModel):
delete_contributions: bool
class AutoSyncSettingsUpdate(BaseModel):
watch_enabled: bool | None = None
schedule_enabled: bool | None = None
schedule_cron: str | None = Field(default=None, min_length=5, max_length=128)
debounce_seconds: int | None = Field(default=None, ge=10, le=86_400)
write_remote_enabled: bool | None = None
review_mode: Literal["required", "auto_high_confidence", "off"] | None = None
@field_validator("schedule_cron")
@classmethod
def validate_cron(cls, value: str | None) -> str | None:
if value is not None and not croniter.is_valid(value):
raise ValueError("无效的 cron 表达式")
return value
@router.get("/bindings")
async def list_bindings(
request: Request,
skill_names: list[str] | None = Query(default=None),
target_id: str | None = Query(default=None),
include_detached: bool = Query(default=False),
mine_only: bool = Query(default=False),
) -> dict[str, Any]:
actor, is_admin = await _actor(request)
rows = await _store(request).list_bindings(
skill_names=skill_names,
target_id=target_id,
include_detached=include_detached,
created_by=actor if mine_only or not is_admin else None,
)
return {"bindings": rows}
@router.get("/source-status")
async def source_status(request: Request, skill_names: list[str] = Query(...)) -> dict[str, Any]:
await _require_admin(request)
bindings = await _store(request).list_bindings(skill_names=skill_names)
by_skill: dict[str, dict[str, Any]] = {}
for name in list(dict.fromkeys(skill_names)):
try:
root = _skill_directory(request.app.state.config, name)
result = await asyncio.to_thread(scan_skill, name, root)
by_skill[name] = {"source_digest": result.source_digest, "counts": result.counts, "blocked_files": result.blocked_files}
except Exception as exc: # noqa: BLE001
by_skill[name] = {"error": str(exc)}
for binding in bindings:
status = by_skill.get(binding["skill_name"])
if status and status.get("source_digest") and status["source_digest"] != binding.get("synced_digest"):
binding["computed_status"] = "stale"
return {"skills": by_skill, "bindings": bindings}
@router.get("/auto-sync-settings", response_model=SkillKnowledgeAutoSyncSettings)
async def get_auto_sync_settings(request: Request) -> SkillKnowledgeAutoSyncSettings:
await _require_admin(request)
return load_system_settings().skill_knowledge_auto_sync
@router.put("/auto-sync-settings", response_model=SkillKnowledgeAutoSyncSettings)
async def update_auto_sync_settings(request: Request, body: AutoSyncSettingsUpdate) -> SkillKnowledgeAutoSyncSettings:
await _require_admin(request)
current = load_system_settings()
updated = current.skill_knowledge_auto_sync.model_copy(update=body.model_dump(exclude_unset=True))
saved = save_system_settings(current.model_copy(update={"skill_knowledge_auto_sync": updated}))
service = getattr(request.app.state, "skill_knowledge_auto_sync", None)
if service is not None:
service.nudge()
return saved.skill_knowledge_auto_sync
async def _set_binding_auto_sync(request: Request, binding_id: str, *, enabled: bool) -> dict[str, Any]:
await _require_admin(request)
row = await _store(request).set_binding_auto_sync(binding_id, enabled=enabled)
if row is None:
raise HTTPException(status_code=404, detail="绑定不存在")
service = getattr(request.app.state, "skill_knowledge_auto_sync", None)
if service is not None:
service.nudge()
return row
@router.post("/bindings/{binding_id}/enable-auto-sync")
async def enable_binding_auto_sync(request: Request, binding_id: str) -> dict[str, Any]:
return await _set_binding_auto_sync(request, binding_id, enabled=True)
@router.post("/bindings/{binding_id}/disable-auto-sync")
async def disable_binding_auto_sync(request: Request, binding_id: str) -> dict[str, Any]:
return await _set_binding_auto_sync(request, binding_id, enabled=False)
@router.post("/sync-jobs", status_code=202)
async def create_sync_job(request: Request, body: CreateSyncJobRequest) -> dict[str, Any]:
actor, is_admin = await _actor(request)
# Built-in and legacy skills need not have an ownership row. Validate
# against the same running filesystem used by the skills list/executor.
missing = []
favorite_names: set[str] = set()
custom_root: Any = None
if not is_admin:
favorite_names = set(await request.app.state.skill_store.list_favorite_skill_names(actor))
storage_root = get_or_new_skill_storage(
app_config=request.app.state.config
).get_skills_root_path().resolve()
custom_root = (storage_root / SkillCategory.CUSTOM.value).resolve()
for name in body.skill_names:
try:
skill_path = await asyncio.to_thread(_skill_directory, request.app.state.config, name)
if not is_admin:
visible = await request.app.state.skill_store.get_visible(name, actor)
is_public_filesystem_skill = not skill_path.resolve().is_relative_to(custom_root)
if visible is None and name not in favorite_names and not is_public_filesystem_skill:
missing.append(name)
except (FileNotFoundError, ValueError):
missing.append(name)
if missing:
raise HTTPException(status_code=404, detail=f"技能不存在:{', '.join(missing[:10])}")
targets: list[dict[str, Any]] = []
for target in body.targets:
data = target.model_dump()
if target.target_type != "weknora":
raise HTTPException(status_code=422, detail="不支持技能直接归纳到助手知识库,请先归纳到 WeKnora 普通知识库")
if target.target_mode != "wiki":
raise HTTPException(status_code=422, detail="技能归纳目标仅支持 WeKnora Wiki 知识库")
mapping = await request.app.state.llmwiki_store.get_authorized(
target.target_id,
actor,
write=True,
is_admin=is_admin,
)
if mapping is None:
raise HTTPException(status_code=404, detail=f"普通知识库不存在或当前用户不可写:{target.target_id}")
data["target_name"] = target.target_name or str(mapping.get("name") or target.target_id)
data["remote_business_key"] = target.remote_business_key or f"weknora:kb:{mapping['weknora_id']}"
targets.append(data)
job = await _store(request).create_job(
skill_names=body.skill_names,
targets=targets,
created_by=actor,
force=body.force,
review_mode=body.review_mode if is_admin else "off",
)
dispatcher = getattr(request.app.state, "skill_knowledge_dispatcher", None)
if dispatcher is not None:
dispatcher.nudge()
return job
@router.get("/sync-jobs")
async def list_sync_jobs(
request: Request,
limit: int = Query(default=50, ge=1, le=200),
mine_only: bool = Query(default=False),
) -> dict[str, Any]:
actor, is_admin = await _actor(request)
return {
"jobs": await _store(request).list_jobs(
limit=limit,
created_by=actor if mine_only or not is_admin else None,
)
}
@router.get("/sync-jobs/{job_id}")
async def get_sync_job(request: Request, job_id: str) -> dict[str, Any]:
actor, is_admin = await _actor(request)
row = await _store(request).get_job(job_id)
if not _owned_by_actor(row, actor, is_admin):
raise HTTPException(status_code=404, detail="归纳任务不存在")
return row
@router.post("/sync-jobs/{job_id}/cancel")
async def cancel_sync_job(request: Request, job_id: str) -> dict[str, bool]:
actor, is_admin = await _actor(request)
row = await _store(request).get_job(job_id)
if not _owned_by_actor(row, actor, is_admin):
raise HTTPException(status_code=404, detail="归纳任务不存在")
if not await _store(request).cancel_job(job_id):
raise HTTPException(status_code=404, detail="归纳任务不存在")
return {"canceled": True}
@router.post("/sync-items/{item_id}/retry", status_code=202)
async def retry_sync_item(request: Request, item_id: str) -> dict[str, Any]:
actor, is_admin = await _actor(request)
jobs = await _store(request).list_jobs(limit=200, created_by=None if is_admin else actor)
item = next(
(item for job in jobs for item in job.get("items", []) if item.get("id") == item_id),
None,
)
if item is None:
raise HTTPException(status_code=404, detail="归纳任务项不存在")
if item.get("target_type") == "weknora":
mapping = await request.app.state.llmwiki_store.get_authorized(
str(item["target_id"]), actor, write=True, is_admin=is_admin
)
if mapping is None:
raise HTTPException(status_code=404, detail="目标知识库不存在或当前用户不可写")
row = await _store(request).requeue_item(item_id)
if row is None:
raise HTTPException(status_code=404, detail="归纳任务项不存在")
request.app.state.skill_knowledge_dispatcher.nudge()
return row
@router.post("/bindings/{binding_id}/resync", status_code=202)
async def resync_binding(request: Request, binding_id: str) -> dict[str, Any]:
actor, is_admin = await _actor(request)
binding = await _store(request).get_binding(binding_id)
if not _owned_by_actor(binding, actor, is_admin) or not binding.get("active"):
raise HTTPException(status_code=404, detail="绑定不存在")
if binding.get("target_type") == "weknora":
mapping = await request.app.state.llmwiki_store.get_authorized(
str(binding["target_id"]), actor, write=True, is_admin=is_admin
)
if mapping is None:
raise HTTPException(status_code=404, detail="目标知识库不存在或当前用户不可写")
target = {
"target_type": binding["target_type"],
"target_mode": binding["target_mode"],
"target_id": binding["target_id"],
"target_name": binding["target_name"],
"remote_business_key": binding.get("remote_business_key"),
}
job = await _store(request).create_job(
skill_names=[binding["skill_name"]],
targets=[target],
created_by=actor,
force=True,
review_mode="auto_high_confidence",
)
request.app.state.skill_knowledge_dispatcher.nudge()
return job
@router.post("/bindings/{binding_id}/detach")
async def detach_binding(request: Request, binding_id: str, body: DetachRequest) -> dict[str, Any]:
actor, is_admin = await _actor(request)
store = _store(request)
binding = await store.get_binding(binding_id)
if not _owned_by_actor(binding, actor, is_admin):
raise HTTPException(status_code=404, detail="绑定不存在")
if await store.has_active_item(binding_id):
raise HTTPException(status_code=409, detail="该绑定仍有运行中的归纳项,请先取消任务")
cleanup_failed = False
if body.delete_contributions:
try:
if binding["target_type"] == "assistant":
sources = await request.app.state.assistant_knowledge_store.list_sources()
for source in sources:
if source.get("source_type") == "skill" and source.get("source_key") == binding["skill_name"]:
await request.app.state.assistant_knowledge_store.delete_source(source["id"])
elif binding["target_type"] == "weknora":
mapping = await request.app.state.llmwiki_store.get_authorized(
str(binding["target_id"]),
actor,
write=True,
is_admin=is_admin,
)
runtime = get_resolved_llmwiki_runtime(request.app.state.config)
if mapping is None or not runtime.weknora_enabled:
raise RuntimeError("远端普通知识库当前不可写")
client = build_weknora_client(runtime)
remote_kb_id = str(mapping["weknora_id"])
for artifact in await store.list_artifacts(binding_id):
if artifact["kind"] == "wiki_page" and artifact.get("slug"):
await client.delete_wiki_page(remote_kb_id, artifact["slug"])
elif artifact.get("remote_id") and artifact["kind"] in {
"raw_source",
"summary_document",
"faq_projection",
"conversation_archive",
"graph_projection",
}:
await client.delete_document(str(artifact["remote_id"]))
except Exception: # noqa: BLE001
cleanup_failed = True
row = await store.detach_binding(
binding_id,
delete_contributions=body.delete_contributions,
cleanup_failed=cleanup_failed,
)
if row is None:
raise HTTPException(status_code=404, detail="绑定不存在")
return row
@router.get("/reviews")
async def list_reviews(
request: Request,
job_id: str | None = Query(default=None),
item_id: str | None = Query(default=None),
status: str | None = Query(default="pending"),
) -> dict[str, Any]:
await _require_admin(request)
return {"reviews": await _store(request).list_reviews(job_id=job_id, item_id=item_id, status=status)}
@router.post("/reviews/actions")
async def review_actions(request: Request, body: ReviewRequest) -> dict[str, Any]:
actor = await _require_admin(request)
item_ids = await _store(request).review(body.review_ids, action=body.action, actor=actor, edits=body.edits, reason=body.reason)
resumed: list[str] = []
executor = request.app.state.skill_knowledge_executor
for item_id in item_ids:
if await _store(request).pending_review_count(item_id) == 0:
asyncio.create_task(executor.resume_reviewed_item(item_id))
resumed.append(item_id)
return {"updated": len(body.review_ids), "resumed_item_ids": resumed}