3456 lines
151 KiB
Python
3456 lines
151 KiB
Python
import asyncio
|
||
import json
|
||
import logging
|
||
import mimetypes
|
||
import os
|
||
import re
|
||
import shutil
|
||
import subprocess
|
||
import sys
|
||
import tempfile
|
||
import zipfile
|
||
from datetime import UTC, datetime
|
||
from pathlib import Path
|
||
from typing import Any
|
||
|
||
from fastapi import APIRouter, Body, Depends, File, HTTPException, Query, Request, UploadFile
|
||
from fastapi.responses import FileResponse
|
||
from pydantic import BaseModel, Field
|
||
from sqlalchemy import select
|
||
|
||
from app.gateway.deps import get_agent_store, get_config, get_skill_store, get_tag_store
|
||
from app.gateway.path_utils import aresolve_thread_virtual_path
|
||
from deerflow.agents.lead_agent.prompt import refresh_skills_system_prompt_cache_async
|
||
from deerflow.config.app_config import AppConfig
|
||
from deerflow.config.extensions_config import ExtensionsConfig, SkillDisplayConfig, SkillStateConfig, get_extensions_config, reload_extensions_config
|
||
from deerflow.persistence.engine import get_engine, get_session_factory
|
||
from deerflow.persistence.skills import SkillStore
|
||
from deerflow.persistence.skills.model import SkillDisplayAdapterRow
|
||
from deerflow.persistence.tags import TagStore
|
||
from deerflow.runtime.user_context import get_effective_user_id
|
||
from deerflow.skills import Skill
|
||
from deerflow.skills.cascade import cascade_skill_unpublished_or_deleted
|
||
from deerflow.skills.installer import SkillAlreadyExistsError
|
||
from deerflow.skills.security_scanner import scan_skill_content
|
||
from deerflow.skills.storage import SkillStorage, get_or_new_skill_storage
|
||
from deerflow.skills.types import SKILL_MD_FILE, SkillCategory
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
router = APIRouter(prefix="/api", tags=["skills"])
|
||
_UNSET = object()
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# Per-user authz helpers (mirrors routers/agents.py)
|
||
# ---------------------------------------------------------------------------
|
||
|
||
|
||
def _current_user_id(request: Request) -> str:
|
||
user = getattr(request.state, "user", None)
|
||
if user is not None:
|
||
return str(user.id)
|
||
return get_effective_user_id()
|
||
|
||
|
||
def _is_admin_user(request: Request) -> bool:
|
||
user = getattr(request.state, "user", None)
|
||
return getattr(user, "system_role", None) == "admin"
|
||
|
||
|
||
def _email_local_part(email: str | None) -> str | None:
|
||
"""Return the username portion (before ``@``) of an email address."""
|
||
if not email:
|
||
return None
|
||
at = email.find("@")
|
||
return email[:at] if at > 0 else email
|
||
|
||
|
||
async def _attach_owner_usernames(responses: list["SkillResponse"]) -> list["SkillResponse"]:
|
||
"""Resolve each custom skill's ``owner_user_id`` to a friendly username.
|
||
|
||
One lookup per distinct owner id; tolerates deleted accounts and an
|
||
uninitialised auth provider by leaving ``owner_username`` as None.
|
||
"""
|
||
owner_ids = {r.owner_user_id for r in responses if r.owner_user_id}
|
||
if not owner_ids:
|
||
return responses
|
||
try:
|
||
from app.gateway.deps import get_local_provider
|
||
|
||
provider = get_local_provider()
|
||
except Exception:
|
||
return responses
|
||
names: dict[str, str] = {}
|
||
for uid in owner_ids:
|
||
try:
|
||
user = await provider.get_user(uid)
|
||
except Exception:
|
||
user = None
|
||
local = _email_local_part(getattr(user, "email", None)) if user is not None else None
|
||
if local:
|
||
names[uid] = local
|
||
for resp in responses:
|
||
if resp.owner_user_id:
|
||
resp.owner_username = names.get(resp.owner_user_id)
|
||
return responses
|
||
|
||
|
||
def _can_view_skill_source(
|
||
category: SkillCategory,
|
||
owner_user_id: str | None,
|
||
user_id: str,
|
||
*,
|
||
is_admin: bool,
|
||
) -> bool:
|
||
"""Whether the raw SKILL.md source is visible to a user.
|
||
|
||
Admins may view any skill's source; otherwise only the publisher (the
|
||
owner of a custom skill) qualifies. Built-in / public skills have no human
|
||
owner, so for non-admins they are never source-visible.
|
||
"""
|
||
if is_admin:
|
||
return True
|
||
return category == SkillCategory.CUSTOM and owner_user_id is not None and owner_user_id == user_id
|
||
|
||
|
||
async def _get_owned_skill_or_raise(store: SkillStore, name: str, user_id: str) -> dict[str, Any]:
|
||
owned = await store.get_owned(name, user_id)
|
||
if owned is not None:
|
||
return owned
|
||
visible = await store.get_visible(name, user_id)
|
||
if visible is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{name}' not found")
|
||
raise HTTPException(status_code=403, detail="You can only modify skills that you own")
|
||
|
||
|
||
async def _get_editable_skill_or_raise(
|
||
store: SkillStore,
|
||
name: str,
|
||
user_id: str,
|
||
*,
|
||
is_admin: bool,
|
||
) -> tuple[dict[str, Any], str]:
|
||
"""Return (record, edit_mode) where edit_mode is 'owner' or 'admin'.
|
||
|
||
Admins can edit/delete any skill (including others' published ones, since
|
||
the spec lets them force-takedown). Non-admins are limited to their own.
|
||
"""
|
||
owned = await store.get_owned(name, user_id)
|
||
if owned is not None:
|
||
return owned, "owner"
|
||
record = await store.get_any(name)
|
||
if record is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{name}' not found")
|
||
if is_admin:
|
||
return record, "admin"
|
||
raise HTTPException(status_code=403, detail="You can only modify skills that you own")
|
||
|
||
|
||
class SkillResponse(BaseModel):
|
||
"""Response model for skill information."""
|
||
|
||
name: str = Field(..., description="Name of the skill")
|
||
name_zh: str | None = Field(default=None, description="用户配置或大模型生成的中文名称")
|
||
description: str = Field(..., description="Description of what the skill does")
|
||
license: str | None = Field(None, description="License information")
|
||
category: SkillCategory = Field(..., description="Category of the skill (public or custom)")
|
||
enabled: bool = Field(default=True, description="Whether this skill is enabled")
|
||
always_on: bool = Field(default=False, description="常驻(免压缩):开启技能压缩后仍全文注入系统提示词")
|
||
use_count: int = Field(default=0, description="Number of times the skill has been used")
|
||
view_count: int = Field(default=0, description="Number of times the skill has been viewed")
|
||
patch_count: int = Field(default=0, description="Number of times the skill has been patched")
|
||
last_activity_at: datetime | None = Field(default=None, description="Most recent activity timestamp")
|
||
created_at: datetime | None = Field(default=None, description="Skill creation timestamp when known")
|
||
state: str = Field(default="active", description="Lifecycle state: active, stale, or archived")
|
||
pinned: bool = Field(default=False, description="Whether this skill is pinned (excluded from curation)")
|
||
featured: bool = Field(default=False, description="Whether this skill is pinned to the top of the list (sort only)")
|
||
source: str = Field(default="upload", description="Origin: 'agent' for agent-created, 'upload' for user-uploaded")
|
||
owner_user_id: str | None = Field(default=None, description="Custom skill owner (null = built-in or legacy)")
|
||
owner_username: str | None = Field(default=None, description="Publisher username (email local part). Null for built-in skills.")
|
||
published: bool = Field(default=False, description="Whether the owner has made this skill visible to other users")
|
||
square_id: str = Field(default="", description="Publish square (发布广场) the skill is filed under once published. Empty string = the default square.")
|
||
builtin: bool = Field(default=False, description="True when owner_user_id is null (no human owner)")
|
||
favorited: bool = Field(default=False, description="True when this skill is in the caller's 我的技能 collection (e.g. granted by their 岗位/position). Owned skills are detected separately via owner_user_id.")
|
||
detail: str | None = Field(default=None, description="Human-authored long-form detail text (Markdown), editable by owner/admin")
|
||
display: dict[str, Any] | None = Field(default=None, description="Optional frontend display adapter")
|
||
display_sample: dict[str, Any] | None = Field(default=None, alias="displaySample", description="Persisted sample output used by the display adapter editor")
|
||
tags: list[dict] = Field(default_factory=list, description="Tags assigned to this skill")
|
||
|
||
|
||
class SkillsListResponse(BaseModel):
|
||
"""Response model for listing all skills."""
|
||
|
||
skills: list[SkillResponse]
|
||
|
||
|
||
class SkillUpdateRequest(BaseModel):
|
||
"""Request model for updating a skill."""
|
||
|
||
enabled: bool = Field(..., description="Whether to enable or disable the skill")
|
||
|
||
|
||
class SkillAlwaysOnUpdateRequest(BaseModel):
|
||
"""Request model for updating a skill's 常驻(免压缩) flag."""
|
||
|
||
always_on: bool = Field(..., description="Whether this skill bypasses skill compression (full-text injected)")
|
||
|
||
|
||
class SkillDisplayUpdateRequest(BaseModel):
|
||
"""Request model for updating a skill's frontend display adapter."""
|
||
|
||
display: SkillDisplayConfig | None = Field(default=None, description="Display adapter configuration; null clears it")
|
||
sample: dict[str, Any] | None = Field(default=None, description="Optional sample payload to persist for future editor sessions")
|
||
|
||
|
||
class SkillDisplaySampleResponse(BaseModel):
|
||
"""Sample structured output used to infer a display adapter."""
|
||
|
||
skill_name: str
|
||
sample: dict[str, Any]
|
||
|
||
|
||
class PersistedSkillDisplaySampleResponse(BaseModel):
|
||
"""Persisted sample output loaded from the database only."""
|
||
|
||
skill_name: str
|
||
sample: dict[str, Any] | None = None
|
||
|
||
|
||
class SkillInstallRequest(BaseModel):
|
||
"""Request model for installing a skill from a .skill file."""
|
||
|
||
thread_id: str = Field(..., description="The thread ID where the .skill file is located")
|
||
path: str = Field(..., description="Virtual path to the .skill file (e.g., mnt/user-data/outputs/my-skill.skill)")
|
||
skill_name: str | None = Field(default=None, description="Optional custom install name. When set, rewrites the archive SKILL.md frontmatter name before install.")
|
||
target_category: SkillCategory = Field(default=SkillCategory.CUSTOM, description="Destination category: 'public' (admin only) or 'custom'")
|
||
|
||
|
||
class SkillInstallResponse(BaseModel):
|
||
"""Response model for skill installation."""
|
||
|
||
success: bool = Field(..., description="Whether the installation was successful")
|
||
skill_name: str = Field(..., description="Name of the installed skill")
|
||
message: str = Field(..., description="Installation result message")
|
||
overwritten: bool = Field(default=False, description="True if an existing skill was replaced")
|
||
category: str = Field(default="custom", description="Destination category: 'public' (admin) or 'custom'")
|
||
|
||
|
||
class SkillInstallPreviewRequest(BaseModel):
|
||
"""Request model for preflighting a .skill artifact install."""
|
||
|
||
thread_id: str = Field(..., description="The thread ID where the .skill file is located")
|
||
path: str = Field(..., description="Virtual path to the .skill file")
|
||
target_category: SkillCategory = Field(default=SkillCategory.CUSTOM, description="Destination category: 'public' (admin only) or 'custom'")
|
||
|
||
|
||
class SkillInstallPreviewResponse(BaseModel):
|
||
"""Response model for skill install preflight."""
|
||
|
||
ok: bool = Field(..., description="Whether the archive could be parsed")
|
||
skill_name: str = Field(..., description="Skill name parsed from the archive frontmatter")
|
||
description: str = Field(default="", description="Parsed description from the archive frontmatter")
|
||
license: str | None = Field(default=None, description="Parsed license from the archive frontmatter")
|
||
allowed_tools: list[str] = Field(default_factory=list, description="Allowed tools declared in the archive")
|
||
can_install_as_is: bool = Field(default=True, description="Whether the current user may install the archive without renaming")
|
||
owned_by_me: bool = Field(default=False, description="True when the existing conflicting skill is already owned by the caller")
|
||
owned_by_other: bool = Field(default=False, description="True when another user already owns the conflicting custom skill")
|
||
builtin_conflict: bool = Field(default=False, description="True when the parsed name collides with a built-in/public skill")
|
||
owner_username: str | None = Field(default=None, description="Friendly owner username for owned_by_other conflicts")
|
||
owner_user_id: str | None = Field(default=None, description="Owner user id for custom conflicts")
|
||
message: str = Field(default="", description="Human-readable install guidance for the UI")
|
||
warnings: list[str] = Field(default_factory=list, description="Non-fatal warnings from validation")
|
||
|
||
|
||
class SkillUploadValidateResponse(BaseModel):
|
||
"""Response model for the pre-install validation endpoint."""
|
||
|
||
ok: bool = Field(..., description="Whether the archive passed cheap validation")
|
||
skill_name: str = Field(..., description="Skill name parsed from frontmatter")
|
||
description: str = Field(default="", description="Skill description parsed from frontmatter")
|
||
license: str | None = Field(default=None, description="Skill license, if declared")
|
||
allowed_tools: list[str] = Field(default_factory=list, description="Tools the skill may invoke")
|
||
conflict: bool = Field(default=False, description="True if a skill with the same name already exists")
|
||
conflict_kind: str | None = Field(default=None, description="'public' or 'custom' when conflict=true")
|
||
can_install_as_is: bool = Field(default=True, description="Whether the archive can be installed without overwrite")
|
||
can_overwrite: bool = Field(default=False, description="Whether the caller may overwrite the conflicting skill")
|
||
owned_by_me: bool = Field(default=False, description="True when the caller owns the conflicting custom skill")
|
||
owned_by_other: bool = Field(default=False, description="True when another user owns the conflicting custom skill")
|
||
builtin_conflict: bool = Field(default=False, description="True when the name conflicts with a built-in/public skill")
|
||
owner_username: str | None = Field(default=None, description="Friendly username for a conflicting skill owned by another user")
|
||
owner_user_id: str | None = Field(default=None, description="Owner user id for a conflicting custom skill")
|
||
message: str = Field(default="", description="Human-readable upload guidance")
|
||
warnings: list[str] = Field(default_factory=list, description="Non-fatal warnings to surface to the user")
|
||
|
||
|
||
class SkillReconcileResponse(BaseModel):
|
||
"""Response model for reconciling custom skills found on disk."""
|
||
|
||
adopted: list[str] = Field(default_factory=list, description="Skills found on disk that had no DB row; now owned by the caller")
|
||
already_registered: list[str] = Field(default_factory=list, description="Skills found on disk that already had a DB row (unchanged)")
|
||
|
||
|
||
class CustomSkillContentResponse(SkillResponse):
|
||
content: str = Field(..., description="Raw SKILL.md content")
|
||
|
||
|
||
class CustomSkillUpdateRequest(BaseModel):
|
||
content: str = Field(..., description="Replacement SKILL.md content")
|
||
|
||
|
||
class SkillDuplicateRequest(BaseModel):
|
||
"""Request model for duplicating an existing skill into a new custom skill."""
|
||
|
||
new_name: str = Field(..., description="Name for the new custom skill copy (lowercase letters/digits, '-'/'_' separators)")
|
||
|
||
|
||
class CustomSkillHistoryResponse(BaseModel):
|
||
history: list[dict]
|
||
|
||
|
||
class SkillRollbackRequest(BaseModel):
|
||
history_index: int = Field(default=-1, description="History entry index to restore from, defaulting to the latest change.")
|
||
|
||
|
||
def _skill_to_response(
|
||
skill: Skill,
|
||
ownership: dict[str, Any] | None = None,
|
||
*,
|
||
display_override: dict[str, Any] | None | object = _UNSET,
|
||
display_sample: dict[str, Any] | None = None,
|
||
favorited: bool = False,
|
||
) -> SkillResponse:
|
||
"""Convert a Skill object to a SkillResponse.
|
||
|
||
``ownership`` is the skill's record from ``SkillStore`` (or None for public
|
||
built-in skills which never have a DB row). When present, contributes the
|
||
``owner_user_id`` / ``published`` / ``builtin`` fields used for per-user
|
||
visibility and the frontend's owner badges.
|
||
"""
|
||
from deerflow.skills.usage import get_entry, latest_activity_at
|
||
|
||
entry = get_entry(skill.name)
|
||
created_at_str = entry.get("created_at")
|
||
last_at_str = latest_activity_at(skill.name)
|
||
created_at: datetime | None = None
|
||
last_at: datetime | None = None
|
||
if created_at_str:
|
||
try:
|
||
created_at = datetime.fromisoformat(created_at_str)
|
||
except Exception:
|
||
pass
|
||
if created_at is None:
|
||
try:
|
||
created_at = datetime.fromtimestamp(skill.skill_file.stat().st_mtime, tz=UTC)
|
||
except Exception:
|
||
pass
|
||
if last_at_str:
|
||
try:
|
||
last_at = datetime.fromisoformat(last_at_str)
|
||
except Exception:
|
||
pass
|
||
|
||
owner_user_id: str | None = None
|
||
published = False
|
||
builtin = skill.category != SkillCategory.CUSTOM
|
||
detail: str | None = None
|
||
name_zh: str | None = None
|
||
square_id = ""
|
||
# 常驻(免压缩) is stored in the DB (SkillRow.always_on), surfaced via the
|
||
# ownership record — never read from the on-disk extensions_config.json.
|
||
always_on = False
|
||
if ownership is not None:
|
||
owner_user_id = ownership.get("owner_user_id")
|
||
published = bool(ownership.get("published", False))
|
||
builtin = bool(ownership.get("builtin", owner_user_id is None))
|
||
detail = ownership.get("detail")
|
||
name_zh = ownership.get("name_zh")
|
||
square_id = str(ownership.get("square_id") or "")
|
||
always_on = bool(ownership.get("always_on", False))
|
||
|
||
display = None if display_override is _UNSET else display_override
|
||
if display_override is _UNSET:
|
||
try:
|
||
skill_config = get_extensions_config().skills.get(skill.name)
|
||
if skill_config and skill_config.display:
|
||
display = skill_config.display.model_dump(by_alias=True)
|
||
except Exception:
|
||
logger.debug("Failed to read display adapter for %s", skill.name, exc_info=True)
|
||
|
||
return SkillResponse(
|
||
name=skill.name,
|
||
name_zh=name_zh,
|
||
description=skill.description,
|
||
license=skill.license,
|
||
category=skill.category,
|
||
enabled=skill.enabled,
|
||
always_on=always_on,
|
||
use_count=entry.get("use_count", 0),
|
||
view_count=entry.get("view_count", 0),
|
||
patch_count=entry.get("patch_count", 0),
|
||
last_activity_at=last_at,
|
||
created_at=created_at,
|
||
state=entry.get("state", "active"),
|
||
pinned=entry.get("pinned", False),
|
||
featured=entry.get("featured", False),
|
||
source=entry.get("source", "upload"),
|
||
owner_user_id=owner_user_id,
|
||
published=published,
|
||
square_id=square_id,
|
||
builtin=builtin,
|
||
favorited=favorited,
|
||
detail=detail,
|
||
display=display,
|
||
displaySample=display_sample,
|
||
)
|
||
|
||
|
||
def _dump_extensions_config(config_obj: ExtensionsConfig) -> dict[str, Any]:
|
||
return {
|
||
"mcpServers": {name: server.model_dump(by_alias=True) for name, server in config_obj.mcp_servers.items()},
|
||
"skills": {name: skill_config.model_dump(by_alias=True, exclude_none=True) for name, skill_config in config_obj.skills.items()},
|
||
}
|
||
|
||
|
||
async def _ensure_skill_display_table() -> bool:
|
||
"""Create the display adapter table when DB persistence is enabled."""
|
||
engine = get_engine()
|
||
sf = get_session_factory()
|
||
if engine is None or sf is None:
|
||
return False
|
||
async with engine.begin() as conn:
|
||
await conn.run_sync(SkillDisplayAdapterRow.__table__.create, checkfirst=True)
|
||
return True
|
||
|
||
|
||
def _json_text_loads(value: str | None) -> dict[str, Any] | None:
|
||
if not value:
|
||
return None
|
||
try:
|
||
parsed = json.loads(value)
|
||
return parsed if isinstance(parsed, dict) else None
|
||
except Exception:
|
||
logger.warning("Failed to parse persisted skill display JSON", exc_info=True)
|
||
return None
|
||
|
||
|
||
def _json_text_dumps(value: dict[str, Any] | None) -> str | None:
|
||
if value is None:
|
||
return None
|
||
return json.dumps(value, ensure_ascii=False)
|
||
|
||
|
||
def _display_override_from_adapter(adapter: dict[str, Any] | None) -> dict[str, Any] | None | object:
|
||
return adapter.get("display") if adapter is not None else _UNSET
|
||
|
||
|
||
async def _get_skill_display_adapter(skill_name: str) -> dict[str, Any] | None:
|
||
if not await _ensure_skill_display_table():
|
||
return None
|
||
sf = get_session_factory()
|
||
if sf is None:
|
||
return None
|
||
async with sf() as session:
|
||
row = await session.get(SkillDisplayAdapterRow, skill_name)
|
||
if row is None:
|
||
return None
|
||
return {
|
||
"display": _json_text_loads(row.display_config),
|
||
"sample": _json_text_loads(row.sample_payload),
|
||
}
|
||
|
||
|
||
async def _list_skill_display_adapters(skill_names: list[str]) -> dict[str, dict[str, Any]]:
|
||
if not skill_names or not await _ensure_skill_display_table():
|
||
return {}
|
||
sf = get_session_factory()
|
||
if sf is None:
|
||
return {}
|
||
async with sf() as session:
|
||
result = await session.execute(
|
||
select(SkillDisplayAdapterRow).where(SkillDisplayAdapterRow.skill_name.in_(skill_names))
|
||
)
|
||
return {
|
||
row.skill_name: {
|
||
"display": _json_text_loads(row.display_config),
|
||
"sample": _json_text_loads(row.sample_payload),
|
||
}
|
||
for row in result.scalars()
|
||
}
|
||
|
||
|
||
async def _upsert_skill_display_adapter(
|
||
skill_name: str,
|
||
*,
|
||
display: dict[str, Any] | None | object = ...,
|
||
sample: dict[str, Any] | None | object = ...,
|
||
) -> None:
|
||
if not await _ensure_skill_display_table():
|
||
return
|
||
sf = get_session_factory()
|
||
if sf is None:
|
||
return
|
||
async with sf() as session:
|
||
row = await session.get(SkillDisplayAdapterRow, skill_name)
|
||
if row is None:
|
||
row = SkillDisplayAdapterRow(skill_name=skill_name)
|
||
session.add(row)
|
||
if display is not ...:
|
||
row.display_config = _json_text_dumps(display) # type: ignore[arg-type]
|
||
if sample is not ...:
|
||
row.sample_payload = _json_text_dumps(sample) # type: ignore[arg-type]
|
||
row.updated_at = datetime.now()
|
||
await session.commit()
|
||
|
||
|
||
_SKILL_SAMPLE_SYSTEM_PROMPT = (
|
||
"You infer realistic JSON return shapes for AI agent skills. "
|
||
"Given a skill's name, description, and full SKILL.md, output a SINGLE JSON object that "
|
||
"represents a plausible successful response from this skill's primary action. "
|
||
"Rules:\n"
|
||
"- Output JSON ONLY. No prose. No markdown fences. No comments.\n"
|
||
"- The top level must be a JSON object (not an array).\n"
|
||
"- Include a results/items/list/records array (pick the most natural key) with 2-3 example entries.\n"
|
||
"- Use realistic, idiomatic field names (e.g. title, content, source, publish_time, score, url, id) that match the skill's domain.\n"
|
||
"- Each entry should have enough fields to drive a citation/list/table UI: a title, a short snippet/content, a source, a time/date, and an id or url when relevant.\n"
|
||
"- Numbers stay numbers, dates stay ISO strings, scores between 0 and 1 when applicable.\n"
|
||
"- Do not invent fields that contradict the skill description."
|
||
)
|
||
|
||
|
||
def _extract_json_from_llm_output(raw: str) -> dict[str, Any] | None:
|
||
"""Best-effort extraction of a JSON object from an LLM response."""
|
||
text = (raw or "").strip()
|
||
if not text:
|
||
return None
|
||
fence = re.match(r"^```(?:json)?\s*(.*?)\s*```$", text, re.DOTALL | re.IGNORECASE)
|
||
if fence:
|
||
text = fence.group(1).strip()
|
||
try:
|
||
parsed = json.loads(text)
|
||
return parsed if isinstance(parsed, dict) else None
|
||
except json.JSONDecodeError:
|
||
pass
|
||
match = re.search(r"\{.*\}", text, re.DOTALL)
|
||
if not match:
|
||
return None
|
||
try:
|
||
parsed = json.loads(match.group(0))
|
||
return parsed if isinstance(parsed, dict) else None
|
||
except json.JSONDecodeError:
|
||
return None
|
||
|
||
|
||
def _skill_sample_python_candidates() -> list[str]:
|
||
candidates = [
|
||
os.environ.get("DEERFLOW_SKILL_SAMPLE_PYTHON"),
|
||
r"C:\Users\ji\anaconda3\envs\open_manus\python.exe",
|
||
sys.executable,
|
||
"python",
|
||
]
|
||
result: list[str] = []
|
||
for candidate in candidates:
|
||
if not candidate or candidate in result:
|
||
continue
|
||
result.append(candidate)
|
||
return result
|
||
|
||
|
||
def _find_skill_sample_script(skill: Skill) -> Path | None:
|
||
skill_dir = skill.skill_file.parent
|
||
preferred = [
|
||
skill_dir / "scripts" / "search.py",
|
||
skill_dir / "search.py",
|
||
skill_dir / "scripts" / "run.py",
|
||
skill_dir / "run.py",
|
||
]
|
||
for path in preferred:
|
||
if path.exists() and path.is_file():
|
||
return path
|
||
scripts_dir = skill_dir / "scripts"
|
||
if scripts_dir.exists():
|
||
for path in sorted(scripts_dir.glob("*.py")):
|
||
if path.is_file():
|
||
return path
|
||
return None
|
||
|
||
|
||
async def _run_skill_python_sample(skill: Skill) -> dict[str, Any] | None:
|
||
"""Run a skill's Python script to capture a real JSON sample.
|
||
|
||
This is used by the display adapter editor. It intentionally runs the same
|
||
script the agent would run, instead of guessing from SKILL.md examples.
|
||
"""
|
||
script = _find_skill_sample_script(skill)
|
||
if script is None:
|
||
return None
|
||
|
||
query = os.environ.get("DEERFLOW_SKILL_SAMPLE_QUERY") or "参考文献模式测试"
|
||
env = os.environ.copy()
|
||
env.setdefault("PYTHONIOENCODING", "utf-8")
|
||
env.setdefault("DEERFLOW_MOCK_KNOWLEDGE_SEARCH", "1")
|
||
env.setdefault("KNOWLEDGE_SEARCH_V2_FIXED", "1")
|
||
|
||
for python_exe in _skill_sample_python_candidates():
|
||
try:
|
||
completed = await asyncio.to_thread(
|
||
subprocess.run,
|
||
[python_exe, str(script), query],
|
||
cwd=str(script.parent),
|
||
env=env,
|
||
capture_output=True,
|
||
text=True,
|
||
encoding="utf-8",
|
||
errors="replace",
|
||
timeout=60,
|
||
)
|
||
except (subprocess.TimeoutExpired, FileNotFoundError, OSError) as exc:
|
||
logger.warning("Skill sample runner failed for %s via %s: %s", skill.name, python_exe, exc)
|
||
continue
|
||
except Exception:
|
||
logger.warning("Skill sample runner crashed for %s via %s", skill.name, python_exe, exc_info=True)
|
||
continue
|
||
|
||
out_text = (completed.stdout or "").strip()
|
||
err_text = (completed.stderr or "").strip()
|
||
if completed.returncode != 0:
|
||
logger.warning(
|
||
"Skill sample script exited %s for %s via %s: %s",
|
||
completed.returncode,
|
||
skill.name,
|
||
python_exe,
|
||
err_text[-1000:],
|
||
)
|
||
continue
|
||
parsed = _extract_json_from_llm_output(out_text)
|
||
if parsed:
|
||
parsed.setdefault("__sampleSource", "python-script")
|
||
parsed.setdefault("__sampleScript", str(script))
|
||
return parsed
|
||
logger.warning(
|
||
"Skill sample script returned non-JSON for %s via %s: stdout=%s stderr=%s",
|
||
skill.name,
|
||
python_exe,
|
||
out_text[-1000:],
|
||
err_text[-1000:],
|
||
)
|
||
return None
|
||
|
||
|
||
async def _generate_skill_display_sample(skill: Skill, app_config: AppConfig) -> dict[str, Any]:
|
||
"""Generate a realistic return-shape sample for ``skill`` via the LLM.
|
||
|
||
Reads SKILL.md, asks the configured curator model (or main model) to infer
|
||
a plausible JSON response. Falls back to :func:`_mock_skill_display_sample`
|
||
when the model is unconfigured, the call times out, or the output cannot
|
||
be parsed as a JSON object.
|
||
"""
|
||
script_sample = await _run_skill_python_sample(skill)
|
||
if script_sample:
|
||
return script_sample
|
||
|
||
try:
|
||
skill_content = skill.skill_file.read_text(encoding="utf-8")
|
||
except OSError as exc:
|
||
logger.debug("Cannot read SKILL.md for %s (%s); falling back to description only", skill.name, exc)
|
||
skill_content = ""
|
||
# Cap the SKILL.md body so very large workflow docs don't blow the context.
|
||
max_chars = 6000
|
||
if len(skill_content) > max_chars:
|
||
skill_content = skill_content[:max_chars] + "\n…(truncated)"
|
||
|
||
user_prompt = (
|
||
f"Skill name: {skill.name}\n"
|
||
f"Short description: {skill.description}\n\n"
|
||
f"SKILL.md:\n-----\n{skill_content or '(SKILL.md unavailable; rely on the short description above.)'}\n-----\n\n"
|
||
"Return one JSON object that matches what this skill would actually produce on success."
|
||
)
|
||
|
||
try:
|
||
from deerflow.models import create_chat_model
|
||
|
||
# Reuse the curator/security-scan model setting; falls back to the main model when unset.
|
||
model_name = app_config.skill_evolution.curator.model_name if app_config.skill_evolution and app_config.skill_evolution.curator else None
|
||
model = create_chat_model(name=model_name, thinking_enabled=False, app_config=app_config)
|
||
response = await asyncio.wait_for(
|
||
model.ainvoke(
|
||
[
|
||
{"role": "system", "content": _SKILL_SAMPLE_SYSTEM_PROMPT},
|
||
{"role": "user", "content": user_prompt},
|
||
],
|
||
config={"run_name": "skill_display_sample"},
|
||
),
|
||
timeout=45,
|
||
)
|
||
parsed = _extract_json_from_llm_output(str(getattr(response, "content", "") or ""))
|
||
if parsed:
|
||
return parsed
|
||
logger.warning("Skill display sample model returned non-JSON for %s; falling back to static mock", skill.name)
|
||
except asyncio.TimeoutError:
|
||
logger.warning("Skill display sample LLM call timed out for %s; falling back to static mock", skill.name)
|
||
except Exception:
|
||
logger.warning("Skill display sample LLM call failed for %s; falling back to static mock", skill.name, exc_info=True)
|
||
return _generic_skill_display_sample(skill.name)
|
||
|
||
|
||
def _generic_skill_display_sample(skill_name: str) -> dict[str, Any]:
|
||
return {
|
||
"success": True,
|
||
"query": "sample query",
|
||
"results": [
|
||
{
|
||
"id": f"{skill_name}-sample-001",
|
||
"title": "Sample result",
|
||
"summary": "Sample structured output used to infer field mapping.",
|
||
"source": skill_name,
|
||
"updatedAt": "2026-05-25",
|
||
}
|
||
],
|
||
}
|
||
|
||
|
||
def _mock_skill_display_sample(skill_name: str) -> dict[str, Any]:
|
||
if skill_name == "knowledge-search-v2":
|
||
return {
|
||
"success": True,
|
||
"query": "参考文献模式测试 知识库固定测试",
|
||
"count": 4,
|
||
"source": "fixed-display-sample",
|
||
"results": [
|
||
{
|
||
"index": 1,
|
||
"recUuid": "ksv2-fixed-001",
|
||
"relevance_score": 0.982,
|
||
"title": "参考文献模式测试:参考文献模式固定测试数据一",
|
||
"source_category": "知识库/测试资料",
|
||
"source": "knowledge-search-v2/fixed/source-001",
|
||
"publish_time": "2026-05-25",
|
||
"content_preview": "这是 knowledge-search-v2 返回的第一条固定测试资料,用于验证右侧参考文献列表、正文编号以及 recUuid 模板跳转是否正常。",
|
||
"content_length": 168,
|
||
},
|
||
{
|
||
"index": 2,
|
||
"recUuid": "ksv2-fixed-002",
|
||
"relevance_score": 0.936,
|
||
"title": "参考文献模式测试:结构化字段映射验证",
|
||
"source_category": "知识库/字段映射",
|
||
"source": "knowledge-search-v2/fixed/source-002",
|
||
"publish_time": "2026-05-24",
|
||
"content_preview": "本条用于验证 title、source_category、source、publish_time、content_preview、relevance_score、recUuid 等字段能被前端正确识别并渲染。",
|
||
"content_length": 154,
|
||
},
|
||
{
|
||
"index": 3,
|
||
"recUuid": "ksv2-fixed-003",
|
||
"relevance_score": 0.887,
|
||
"title": "参考文献模式测试:多引用排序验证",
|
||
"source_category": "知识库/排序验证",
|
||
"source": "knowledge-search-v2/fixed/source-003",
|
||
"publish_time": "2026-05-23",
|
||
"content_preview": "本条用于验证多条参考文献的编号顺序是否与 results 数组顺序一致,并验证点击正文编号时能定位到右侧对应参考文献。",
|
||
"content_length": 132,
|
||
},
|
||
{
|
||
"index": 4,
|
||
"recUuid": "",
|
||
"relevance_score": 0.801,
|
||
"title": "参考文献模式测试:无 recUuid 跳转隐藏验证",
|
||
"source_category": "知识库/边界场景",
|
||
"source": "knowledge-search-v2/fixed/source-004",
|
||
"publish_time": "2026-05-22",
|
||
"content_preview": "本条故意把 recUuid 置为空,用于验证右侧弹框仍展示详情,但不展示外部跳转按钮。",
|
||
"content_length": 98,
|
||
},
|
||
],
|
||
}
|
||
return {
|
||
"success": True,
|
||
"query": "sample query",
|
||
"results": [
|
||
{
|
||
"id": "sample-001",
|
||
"title": "Sample result",
|
||
"summary": "Sample structured output used to infer field mapping.",
|
||
"source": "sample",
|
||
"updatedAt": "2026-05-25",
|
||
}
|
||
],
|
||
}
|
||
|
||
|
||
async def _build_ownership_index(skill_store: SkillStore, user_id: str) -> dict[str, dict[str, Any]]:
|
||
"""Index visible-to-user skill records by name. Internal call cache for list endpoints."""
|
||
records = await skill_store.list_visible(user_id)
|
||
return {record["name"]: record for record in records}
|
||
|
||
|
||
async def _attach_and_filter_tags(
|
||
tag_store: TagStore,
|
||
responses: list[SkillResponse],
|
||
search: str | None,
|
||
tag_ids: list[str] | None,
|
||
) -> list[SkillResponse]:
|
||
"""Bulk-attach tags to skill responses, then apply name/tag filters."""
|
||
assignments = await tag_store.list_assignments("skill", [r.name for r in responses])
|
||
for resp in responses:
|
||
resp.tags = assignments.get(resp.name, [])
|
||
if search:
|
||
needle = search.strip().lower()
|
||
# Match the English name OR the Chinese display name so 中文 search works.
|
||
responses = [r for r in responses if needle in r.name.lower() or (r.name_zh or "").lower().find(needle) >= 0]
|
||
if tag_ids:
|
||
wanted = set(tag_ids)
|
||
responses = [r for r in responses if any(t.get("id") in wanted for t in r.tags)]
|
||
return responses
|
||
|
||
|
||
def _sort_skill_responses(responses: list[SkillResponse], sort: str) -> list[SkillResponse]:
|
||
reverse = sort not in {"time_asc", "activity_asc", "name_asc", "usage_asc"}
|
||
if sort in {"name_asc", "name_desc"}:
|
||
return sorted(responses, key=lambda r: (r.name or "").lower(), reverse=reverse)
|
||
if sort in {"usage_desc", "usage_asc"}:
|
||
return sorted(responses, key=lambda r: int(r.use_count or 0), reverse=reverse)
|
||
if sort in {"activity_desc", "activity_asc"}:
|
||
return sorted(responses, key=lambda r: str(r.last_activity_at or r.created_at or ""), reverse=reverse)
|
||
return sorted(responses, key=lambda r: str(r.created_at or r.last_activity_at or ""), reverse=reverse)
|
||
|
||
|
||
@router.get(
|
||
"/skills",
|
||
response_model=SkillsListResponse,
|
||
summary="List All Skills",
|
||
description="Retrieve a list of all available skills from both public and custom directories.",
|
||
)
|
||
async def list_skills(
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
tag_store: TagStore = Depends(get_tag_store),
|
||
search: str | None = Query(default=None, description="Case-insensitive fuzzy match on skill name or Chinese name (name_zh)"),
|
||
tag_ids: list[str] | None = Query(default=None, description="Keep skills carrying any of these tags (OR)"),
|
||
square: str | None = Query(default=None, description="Keep only skills filed under this publish square (发布广场). Omit for all squares."),
|
||
scope: str | None = Query(default=None, description="``mine`` filters to the skills the caller's 岗位 (position) makes visible."),
|
||
sort: str = Query(
|
||
default="time_desc",
|
||
pattern="^(time_desc|time_asc|activity_desc|activity_asc|usage_desc|usage_asc|name_asc|name_desc)$",
|
||
description="Sort order. Default is newest created/modified first.",
|
||
),
|
||
admin: bool = Query(
|
||
default=False,
|
||
description=(
|
||
"Admin-only override. When ``true`` and the caller is a system admin, "
|
||
"every custom skill is returned (including unpublished skills owned by "
|
||
"other users) regardless of the usual ownership filter."
|
||
),
|
||
),
|
||
) -> SkillsListResponse:
|
||
try:
|
||
user_id = _current_user_id(request)
|
||
if admin and not _is_admin_user(request):
|
||
raise HTTPException(status_code=403, detail="Only system admins can list all skills")
|
||
all_skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
if admin:
|
||
ownership_records = await skill_store.list_all()
|
||
ownership = {record["name"]: record for record in ownership_records}
|
||
else:
|
||
ownership = await _build_ownership_index(skill_store, user_id)
|
||
# Skills the caller's 岗位 (position) granted into their "我的技能"
|
||
# collection. Admins browse everything, so the badge is moot for them.
|
||
favorite_names: set[str] = set()
|
||
if not admin:
|
||
try:
|
||
favorite_names = await skill_store.list_favorite_skill_names(user_id)
|
||
except Exception:
|
||
logger.debug("skill favorites lookup failed", exc_info=True)
|
||
favorite_names = set()
|
||
display_adapters = await _list_skill_display_adapters([skill.name for skill in all_skills])
|
||
responses: list[SkillResponse] = []
|
||
for skill in all_skills:
|
||
adapter = display_adapters.get(skill.name)
|
||
# Public/built-in skills (no DB row) are visible to everyone.
|
||
# Custom skills must have an ownership record in the visible set.
|
||
if skill.category != SkillCategory.CUSTOM:
|
||
# Built-in skills usually have no DB row, but an admin-authored
|
||
# ``detail`` annotation materializes one — surface it when present.
|
||
responses.append(
|
||
_skill_to_response(
|
||
skill,
|
||
ownership.get(skill.name),
|
||
display_override=_display_override_from_adapter(adapter),
|
||
display_sample=adapter.get("sample") if adapter else None,
|
||
favorited=skill.name in favorite_names,
|
||
)
|
||
)
|
||
continue
|
||
record = ownership.get(skill.name)
|
||
if record is None:
|
||
if admin:
|
||
# Admin sees every custom skill on disk even without an
|
||
# ownership row (legacy / orphaned), so they can clean it up.
|
||
record = None
|
||
elif skill.name in favorite_names:
|
||
# A position granted this skill even though it isn't
|
||
# otherwise visible to the user (e.g. an unpublished custom
|
||
# skill). Surface it so it lands in their 我的技能.
|
||
record = await skill_store.get_any(skill.name)
|
||
else:
|
||
continue
|
||
responses.append(
|
||
_skill_to_response(
|
||
skill,
|
||
record,
|
||
display_override=_display_override_from_adapter(adapter),
|
||
display_sample=adapter.get("sample") if adapter else None,
|
||
favorited=skill.name in favorite_names,
|
||
)
|
||
)
|
||
responses = await _attach_and_filter_tags(tag_store, responses, search, tag_ids)
|
||
responses = await _attach_owner_usernames(responses)
|
||
square_id = square.strip() if square and square.strip() else None
|
||
if square_id:
|
||
from deerflow.config.system_settings import load_system_settings, resolve_default_square_id
|
||
|
||
square_is_default = square_id == resolve_default_square_id(load_system_settings().publish_squares)
|
||
responses = [
|
||
r
|
||
for r in responses
|
||
if (r.square_id or "") == square_id or (square_is_default and not (r.square_id or ""))
|
||
]
|
||
# 岗位 (position) no longer *restricts* the skill list: a position grants
|
||
# skills by materializing them into the member's 我的技能 (the ``favorited``
|
||
# flag above + the ``skill_favorites`` table, synced by PositionSyncService),
|
||
# mirroring how agents are added to 我的智能体. The browse-all 广场 therefore
|
||
# always shows every visible skill — the ``scope`` param is kept for API
|
||
# compatibility but intentionally does not filter here.
|
||
responses = _sort_skill_responses(responses, sort)
|
||
return SkillsListResponse(skills=responses)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error(f"Failed to load skills: {e}", exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to load skills: {str(e)}")
|
||
|
||
|
||
_MAX_SKILL_UPLOAD_BYTES = 512 * 1024 * 1024
|
||
_ZIP_MAGIC = b"PK\x03\x04"
|
||
|
||
|
||
async def _persist_uploaded_skill(file: UploadFile) -> Path:
|
||
"""Stream an uploaded skill archive to a temp file. Caller must unlink it.
|
||
|
||
Validates the leading bytes look like a ZIP archive and enforces the same
|
||
512 MB cap that the extractor uses, so we reject obvious garbage early.
|
||
"""
|
||
if file.filename is None:
|
||
raise HTTPException(status_code=400, detail="No file provided")
|
||
suffix = Path(file.filename).suffix.lower()
|
||
if suffix not in {".zip", ".skill"}:
|
||
raise HTTPException(status_code=400, detail="Skill upload must be a .zip or .skill file")
|
||
|
||
fd, tmp_name = tempfile.mkstemp(suffix=".skill", prefix="skill-upload-")
|
||
tmp_path = Path(tmp_name)
|
||
total = 0
|
||
head: bytes = b""
|
||
try:
|
||
with open(fd, "wb") as out:
|
||
while True:
|
||
chunk = await file.read(1024 * 1024)
|
||
if not chunk:
|
||
break
|
||
if not head:
|
||
head = chunk[: len(_ZIP_MAGIC)]
|
||
total += len(chunk)
|
||
if total > _MAX_SKILL_UPLOAD_BYTES:
|
||
raise HTTPException(status_code=400, detail="Uploaded skill archive exceeds the 512 MB limit")
|
||
out.write(chunk)
|
||
except HTTPException:
|
||
tmp_path.unlink(missing_ok=True)
|
||
raise
|
||
except Exception:
|
||
tmp_path.unlink(missing_ok=True)
|
||
raise
|
||
|
||
if not head.startswith(_ZIP_MAGIC):
|
||
tmp_path.unlink(missing_ok=True)
|
||
raise HTTPException(status_code=400, detail="Uploaded file is not a valid ZIP archive")
|
||
return tmp_path
|
||
|
||
|
||
@router.post(
|
||
"/skills/validate",
|
||
response_model=SkillUploadValidateResponse,
|
||
summary="Validate Skill Archive (No Install)",
|
||
description="Inspect an uploaded .zip/.skill archive and report parsed metadata plus ownership-aware conflicts. Does NOT install or run the LLM security scan.",
|
||
)
|
||
async def validate_skill_upload(
|
||
request: Request,
|
||
file: UploadFile = File(..., description="The .zip or .skill archive to validate"),
|
||
target_category: SkillCategory = SkillCategory.CUSTOM,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillUploadValidateResponse:
|
||
tmp_path = await _persist_uploaded_skill(file)
|
||
try:
|
||
result = get_or_new_skill_storage(app_config=config).validate_skill_archive(tmp_path)
|
||
(
|
||
can_install_as_is,
|
||
can_overwrite,
|
||
owned_by_me,
|
||
owned_by_other,
|
||
builtin_conflict,
|
||
owner_user_id,
|
||
owner_username,
|
||
message,
|
||
) = await _preview_install_target(
|
||
request=request,
|
||
skill_store=skill_store,
|
||
skill_name=str(result.get("skill_name") or ""),
|
||
target_category=target_category,
|
||
conflict_kind=result.get("conflict_kind"),
|
||
)
|
||
effective_conflict_kind = result.get("conflict_kind")
|
||
if effective_conflict_kind is None:
|
||
if owned_by_me or owned_by_other:
|
||
effective_conflict_kind = "custom"
|
||
elif builtin_conflict:
|
||
effective_conflict_kind = "public"
|
||
return SkillUploadValidateResponse(
|
||
**{
|
||
**result,
|
||
"conflict": bool(result.get("conflict") or effective_conflict_kind),
|
||
"conflict_kind": effective_conflict_kind,
|
||
"can_install_as_is": can_install_as_is,
|
||
"can_overwrite": can_overwrite,
|
||
"owned_by_me": owned_by_me,
|
||
"owned_by_other": owned_by_other,
|
||
"builtin_conflict": builtin_conflict,
|
||
"owner_username": owner_username,
|
||
"owner_user_id": owner_user_id,
|
||
"message": message,
|
||
}
|
||
)
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to validate uploaded skill: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to validate skill: {e}")
|
||
finally:
|
||
tmp_path.unlink(missing_ok=True)
|
||
|
||
|
||
@router.post(
|
||
"/skills/install-upload",
|
||
response_model=SkillInstallResponse,
|
||
summary="Install Skill from Upload",
|
||
description="Install a skill directly from an uploaded .zip/.skill archive. Set ?overwrite=true to replace the caller's own custom skill, or an existing public skill when installing as an admin.",
|
||
)
|
||
async def install_skill_upload(
|
||
request: Request,
|
||
file: UploadFile = File(..., description="The .zip or .skill archive to install"),
|
||
overwrite: bool = False,
|
||
target_category: SkillCategory = SkillCategory.CUSTOM,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillInstallResponse:
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
# Only admins may install to the public (built-in) directory.
|
||
if target_category == SkillCategory.PUBLIC and not is_admin:
|
||
raise HTTPException(status_code=403, detail="Only admin users can install skills to the public directory")
|
||
tmp_path = await _persist_uploaded_skill(file)
|
||
try:
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
# Re-run the same ownership-aware preflight used by /skills/validate.
|
||
# This closes the time-of-check/time-of-use gap and prevents a crafted
|
||
# request (including an admin installing into custom/) from replacing a
|
||
# different user's skill.
|
||
validation = storage.validate_skill_archive(tmp_path)
|
||
preview_name = str(validation.get("skill_name") or "")
|
||
(
|
||
can_install_as_is,
|
||
can_overwrite,
|
||
_owned_by_me,
|
||
_owned_by_other,
|
||
_builtin_conflict,
|
||
_owner_user_id,
|
||
_owner_username,
|
||
conflict_message,
|
||
) = await _preview_install_target(
|
||
request=request,
|
||
skill_store=skill_store,
|
||
skill_name=preview_name,
|
||
target_category=target_category,
|
||
conflict_kind=validation.get("conflict_kind"),
|
||
)
|
||
if (overwrite and not can_install_as_is and not can_overwrite) or (not overwrite and not can_install_as_is):
|
||
raise HTTPException(status_code=409, detail=conflict_message)
|
||
result = await storage.ainstall_skill_from_archive(tmp_path, overwrite=overwrite, require_skill_suffix=False, target_category=target_category)
|
||
try:
|
||
from deerflow.skills.usage import mark_uploaded
|
||
mark_uploaded(result["skill_name"])
|
||
except Exception:
|
||
pass
|
||
# For public installs (admin), register as built-in with no owner.
|
||
# For custom installs, register ownership under the uploading user.
|
||
if target_category == SkillCategory.PUBLIC:
|
||
await skill_store.ensure_legacy(
|
||
{"name": result["skill_name"], "owner_user_id": None, "published": True}
|
||
)
|
||
else:
|
||
# ``ensure_legacy`` is idempotent and only patches a NULL owner; it
|
||
# never overwrites an existing owner, so the pre-parse check above
|
||
# is what really guards cross-user steals.
|
||
await skill_store.ensure_legacy(
|
||
{"name": result["skill_name"], "owner_user_id": user_id, "published": False}
|
||
)
|
||
await refresh_skills_system_prompt_cache_async()
|
||
return SkillInstallResponse(**result)
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except SkillAlreadyExistsError as e:
|
||
raise HTTPException(status_code=409, detail=str(e))
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to install uploaded skill: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to install skill: {e}")
|
||
finally:
|
||
tmp_path.unlink(missing_ok=True)
|
||
|
||
|
||
async def _preview_install_target(
|
||
*,
|
||
request: Request,
|
||
skill_store: SkillStore,
|
||
skill_name: str,
|
||
target_category: SkillCategory,
|
||
conflict_kind: str | None = None,
|
||
) -> tuple[bool, bool, bool, bool, bool, str | None, str | None, str]:
|
||
"""Return installability / ownership info for a parsed skill name."""
|
||
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
existing = await skill_store.get_any(skill_name)
|
||
owner_user_id = existing.get("owner_user_id") if existing else None
|
||
owner_username: str | None = None
|
||
if owner_user_id:
|
||
try:
|
||
attached = await _attach_owner_usernames([SkillResponse(name=skill_name, description="", license=None, category=SkillCategory.CUSTOM, owner_user_id=owner_user_id)])
|
||
owner_username = attached[0].owner_username if attached else None
|
||
except Exception:
|
||
owner_username = None
|
||
|
||
custom_conflict = conflict_kind == "custom" or bool(existing and owner_user_id is not None)
|
||
builtin_conflict = conflict_kind == "public" or bool(existing and owner_user_id is None and conflict_kind != "custom")
|
||
owned_by_me = bool(custom_conflict and owner_user_id == user_id)
|
||
owned_by_other = bool(custom_conflict and owner_user_id not in (None, user_id))
|
||
unattributed_custom_conflict = bool(custom_conflict and owner_user_id is None and not builtin_conflict)
|
||
|
||
if target_category == SkillCategory.PUBLIC and not is_admin:
|
||
return False, False, owned_by_me, owned_by_other, builtin_conflict, owner_user_id, owner_username, "Only admin users can install skills to the public directory"
|
||
|
||
if target_category == SkillCategory.PUBLIC:
|
||
if custom_conflict:
|
||
who = owner_username or owner_user_id
|
||
if owned_by_other and who:
|
||
message = f"Skill '{skill_name}' has already been uploaded by user '{who}' and cannot be overwritten."
|
||
elif owned_by_me:
|
||
message = f"You already own custom skill '{skill_name}'. Switch the target to custom to overwrite it."
|
||
else:
|
||
message = f"Skill '{skill_name}' already exists as a custom skill and its owner could not be determined."
|
||
return False, False, owned_by_me, owned_by_other, builtin_conflict, owner_user_id, owner_username, message
|
||
if builtin_conflict:
|
||
return False, True, False, False, True, owner_user_id, owner_username, f"Built-in/public skill '{skill_name}' already exists and can be overwritten by an admin."
|
||
return True, False, False, False, False, owner_user_id, owner_username, f"Skill '{skill_name}' can be installed."
|
||
|
||
if builtin_conflict:
|
||
return False, False, owned_by_me, owned_by_other, True, owner_user_id, owner_username, f"Skill '{skill_name}' already exists as a built-in/public skill. Please install it under a different name."
|
||
if owned_by_other:
|
||
who = owner_username or owner_user_id or "another user"
|
||
return False, False, False, True, False, owner_user_id, owner_username, f"Skill '{skill_name}' has already been uploaded by user '{who}' and cannot be overwritten."
|
||
if owned_by_me:
|
||
return False, True, True, False, False, owner_user_id, owner_username, f"You have already uploaded skill '{skill_name}'. You may upload it again and overwrite your existing copy."
|
||
if unattributed_custom_conflict:
|
||
return False, False, False, False, False, owner_user_id, owner_username, f"Skill '{skill_name}' already exists, but its owner could not be determined. It cannot be overwritten safely."
|
||
return True, False, False, False, False, owner_user_id, owner_username, f"Skill '{skill_name}' can be installed."
|
||
|
||
|
||
async def _build_renamed_skill_archive(
|
||
archive_path: str | Path,
|
||
*,
|
||
new_name: str,
|
||
require_skill_suffix: bool,
|
||
) -> Path:
|
||
"""Create a temporary archive whose SKILL.md frontmatter name is rewritten."""
|
||
|
||
from deerflow.skills.installer import resolve_skill_dir_from_archive, safe_extract_skill_archive
|
||
|
||
SkillStorage.validate_skill_name(new_name)
|
||
tmp_root = Path(tempfile.mkdtemp(prefix="skill-install-rename-"))
|
||
try:
|
||
path = Path(archive_path)
|
||
if require_skill_suffix and path.suffix != ".skill":
|
||
raise ValueError("File must have .skill extension")
|
||
with zipfile.ZipFile(path, "r") as zf:
|
||
safe_extract_skill_archive(zf, tmp_root)
|
||
skill_dir = resolve_skill_dir_from_archive(tmp_root)
|
||
skill_md = skill_dir / SKILL_MD_FILE
|
||
skill_md.write_text(
|
||
_rewrite_skill_frontmatter_name(skill_md.read_text(encoding="utf-8"), new_name),
|
||
encoding="utf-8",
|
||
)
|
||
renamed_archive = tmp_root / f"{new_name}.skill"
|
||
with zipfile.ZipFile(renamed_archive, "w", compression=zipfile.ZIP_DEFLATED) as zf:
|
||
for path in sorted(tmp_root.rglob("*")):
|
||
if not path.is_file() or path == renamed_archive:
|
||
continue
|
||
zf.write(path, path.relative_to(tmp_root))
|
||
return renamed_archive
|
||
except Exception:
|
||
shutil.rmtree(tmp_root, ignore_errors=True)
|
||
raise
|
||
|
||
|
||
async def _install_archive_with_optional_name(
|
||
*,
|
||
storage: SkillStorage,
|
||
archive_path: str | Path,
|
||
requested_name: str | None,
|
||
target_category: SkillCategory,
|
||
overwrite: bool = False,
|
||
require_skill_suffix: bool = True,
|
||
) -> SkillInstallResponse:
|
||
normalized_requested = SkillStorage.validate_skill_name(requested_name) if requested_name else None
|
||
temp_root: Path | None = None
|
||
install_source = Path(archive_path)
|
||
if normalized_requested:
|
||
renamed_archive = await _build_renamed_skill_archive(
|
||
archive_path,
|
||
new_name=normalized_requested,
|
||
require_skill_suffix=require_skill_suffix,
|
||
)
|
||
temp_root = renamed_archive.parent
|
||
install_source = renamed_archive
|
||
try:
|
||
result = await storage.ainstall_skill_from_archive(
|
||
install_source,
|
||
overwrite=overwrite,
|
||
require_skill_suffix=require_skill_suffix,
|
||
target_category=target_category,
|
||
)
|
||
return SkillInstallResponse(**result)
|
||
finally:
|
||
if temp_root is not None:
|
||
shutil.rmtree(temp_root, ignore_errors=True)
|
||
|
||
|
||
@router.post(
|
||
"/skills/install-preview",
|
||
response_model=SkillInstallPreviewResponse,
|
||
summary="Preview Skill Install",
|
||
description="Parse a .skill artifact from a thread and return ownership/conflict info before install.",
|
||
)
|
||
async def preview_skill_install(
|
||
body: SkillInstallPreviewRequest,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillInstallPreviewResponse:
|
||
try:
|
||
skill_file_path = await aresolve_thread_virtual_path(body.thread_id, body.path)
|
||
validation = get_or_new_skill_storage(app_config=config).validate_skill_archive(skill_file_path, require_skill_suffix=False)
|
||
skill_name = str(validation.get("skill_name") or "")
|
||
can_install_as_is, _can_overwrite, owned_by_me, owned_by_other, builtin_conflict, owner_user_id, owner_username, message = await _preview_install_target(
|
||
request=request,
|
||
skill_store=skill_store,
|
||
skill_name=skill_name,
|
||
target_category=body.target_category,
|
||
conflict_kind=validation.get("conflict_kind"),
|
||
)
|
||
return SkillInstallPreviewResponse(
|
||
ok=True,
|
||
skill_name=skill_name,
|
||
description=str(validation.get("description") or ""),
|
||
license=validation.get("license"),
|
||
allowed_tools=[str(item) for item in (validation.get("allowed_tools") or [])],
|
||
can_install_as_is=can_install_as_is,
|
||
owned_by_me=owned_by_me,
|
||
owned_by_other=owned_by_other,
|
||
builtin_conflict=builtin_conflict,
|
||
owner_username=owner_username,
|
||
owner_user_id=owner_user_id,
|
||
message=message,
|
||
warnings=[str(item) for item in (validation.get("warnings") or [])],
|
||
)
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to preview skill install: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to preview skill install: {e}")
|
||
|
||
|
||
@router.post(
|
||
"/skills/install",
|
||
response_model=SkillInstallResponse,
|
||
summary="Install Skill",
|
||
description="Install a skill from a .skill file (ZIP archive) located in the thread's user-data directory.",
|
||
)
|
||
async def install_skill(
|
||
body: SkillInstallRequest,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillInstallResponse:
|
||
try:
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
target_category = body.target_category
|
||
requested_name = SkillStorage.validate_skill_name(body.skill_name) if body.skill_name else None
|
||
if target_category == SkillCategory.PUBLIC and not is_admin:
|
||
raise HTTPException(status_code=403, detail="Only admin users can install skills to the public directory")
|
||
skill_file_path = await aresolve_thread_virtual_path(body.thread_id, body.path)
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
validation = storage.validate_skill_archive(skill_file_path, require_skill_suffix=False)
|
||
preview_name = requested_name or str(validation.get("skill_name") or "")
|
||
can_install_as_is, _can_overwrite, owned_by_me, owned_by_other, builtin_conflict, _owner_user_id, _owner_username, message = await _preview_install_target(
|
||
request=request,
|
||
skill_store=skill_store,
|
||
skill_name=preview_name,
|
||
target_category=target_category,
|
||
conflict_kind=(validation.get("conflict_kind") if requested_name is None or requested_name == validation.get("skill_name") else None),
|
||
)
|
||
if not can_install_as_is:
|
||
if requested_name and requested_name != validation.get("skill_name") and not owned_by_other and not builtin_conflict:
|
||
pass
|
||
else:
|
||
raise HTTPException(status_code=409, detail=message)
|
||
result = await _install_archive_with_optional_name(
|
||
storage=storage,
|
||
archive_path=skill_file_path,
|
||
requested_name=requested_name,
|
||
target_category=target_category,
|
||
require_skill_suffix=False,
|
||
)
|
||
try:
|
||
from deerflow.skills.usage import mark_uploaded
|
||
mark_uploaded(result.skill_name)
|
||
except Exception:
|
||
pass
|
||
if target_category == SkillCategory.PUBLIC:
|
||
await skill_store.ensure_legacy(
|
||
{"name": result.skill_name, "owner_user_id": None, "published": True}
|
||
)
|
||
else:
|
||
await skill_store.ensure_legacy(
|
||
{"name": result.skill_name, "owner_user_id": user_id, "published": False}
|
||
)
|
||
await refresh_skills_system_prompt_cache_async()
|
||
return result
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except SkillAlreadyExistsError as e:
|
||
raise HTTPException(status_code=409, detail=str(e))
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error(f"Failed to install skill: {e}", exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to install skill: {str(e)}")
|
||
|
||
|
||
def _rewrite_skill_frontmatter_name(content: str, new_name: str) -> str:
|
||
"""Return *content* with the YAML-frontmatter ``name:`` field set to *new_name*.
|
||
|
||
A duplicated skill must declare the new name in its SKILL.md frontmatter or
|
||
``validate_skill_markdown_content`` will reject it (frontmatter name must
|
||
match the directory name). We rewrite only the first ``name:`` line inside
|
||
the leading ``---`` block, leaving every other field byte-for-byte intact.
|
||
"""
|
||
match = re.match(r"^---\n(.*?)\n---", content, re.DOTALL)
|
||
if not match:
|
||
raise ValueError("SKILL.md is missing the leading YAML frontmatter block")
|
||
fm = match.group(1)
|
||
new_fm, replaced = re.subn(r"(?m)^(name:\s*).*$", rf"\g<1>{new_name}", fm, count=1)
|
||
if replaced == 0:
|
||
new_fm = f"name: {new_name}\n{fm}"
|
||
return content[: match.start(1)] + new_fm + content[match.end(1) :]
|
||
|
||
|
||
@router.post(
|
||
"/skills/custom/{source_name}/duplicate",
|
||
response_model=SkillInstallResponse,
|
||
summary="Duplicate Skill",
|
||
description=(
|
||
"Copy an existing skill (built-in/public OR custom) into a brand-new custom "
|
||
"skill named ``new_name``, owned by the caller and unpublished. All support "
|
||
"files are copied verbatim; only the SKILL.md frontmatter name is rewritten. "
|
||
"The source already passed the install-time security scan, so no re-scan is run."
|
||
),
|
||
)
|
||
async def duplicate_skill(
|
||
source_name: str,
|
||
body: SkillDuplicateRequest,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillInstallResponse:
|
||
import shutil
|
||
|
||
try:
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
|
||
source_name = SkillStorage.validate_skill_name(source_name.replace("\r\n", "").replace("\n", ""))
|
||
new_name = SkillStorage.validate_skill_name(body.new_name.replace("\r\n", "").replace("\n", ""))
|
||
|
||
# Resolve the source directory (custom takes precedence over public, so a
|
||
# user-shadowed copy is the one duplicated when both exist).
|
||
if storage.custom_skill_exists(source_name):
|
||
source_dir = storage.get_custom_skill_dir(source_name)
|
||
elif storage.public_skill_exists(source_name):
|
||
source_dir = storage.get_skills_root_path() / SkillCategory.PUBLIC.value / source_name
|
||
else:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{source_name}' not found")
|
||
|
||
if new_name == source_name or storage.custom_skill_exists(new_name) or storage.public_skill_exists(new_name):
|
||
raise HTTPException(status_code=409, detail=f"A skill named '{new_name}' already exists")
|
||
if not is_admin:
|
||
existing = await skill_store.get_any(new_name)
|
||
if existing and existing.get("owner_user_id") not in (None, user_id):
|
||
raise HTTPException(status_code=409, detail=f"Skill '{new_name}' already exists and is owned by another user")
|
||
|
||
dest_dir = storage.get_custom_skill_dir(new_name)
|
||
# copytree refuses to overwrite, which is the behaviour we want — the
|
||
# existence checks above already guarantee a clean target.
|
||
shutil.copytree(source_dir, dest_dir)
|
||
try:
|
||
source_content = (dest_dir / SKILL_MD_FILE).read_text(encoding="utf-8")
|
||
new_content = _rewrite_skill_frontmatter_name(source_content, new_name)
|
||
SkillStorage.validate_skill_markdown_content(new_name, new_content)
|
||
storage.write_custom_skill(new_name, SKILL_MD_FILE, new_content)
|
||
except Exception:
|
||
# Roll back the half-created directory so a failed copy never leaves
|
||
# an unparseable skill on disk.
|
||
shutil.rmtree(dest_dir, ignore_errors=True)
|
||
raise
|
||
|
||
await skill_store.ensure_legacy({"name": new_name, "owner_user_id": user_id, "published": False})
|
||
try:
|
||
storage.append_history(
|
||
new_name,
|
||
{
|
||
"action": "duplicate",
|
||
"author": "human",
|
||
"thread_id": None,
|
||
"file_path": SKILL_MD_FILE,
|
||
"prev_content": None,
|
||
"new_content": new_content,
|
||
"source_skill": source_name,
|
||
},
|
||
)
|
||
except Exception:
|
||
pass
|
||
try:
|
||
from deerflow.skills.usage import mark_uploaded
|
||
mark_uploaded(new_name)
|
||
except Exception:
|
||
pass
|
||
await refresh_skills_system_prompt_cache_async()
|
||
return SkillInstallResponse(
|
||
success=True,
|
||
skill_name=new_name,
|
||
message=f"Skill '{source_name}' duplicated as '{new_name}'",
|
||
overwritten=False,
|
||
category=SkillCategory.CUSTOM.value,
|
||
)
|
||
except HTTPException:
|
||
raise
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except Exception as e:
|
||
logger.error("Failed to duplicate skill %s: %s", source_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to duplicate skill: {str(e)}")
|
||
|
||
|
||
@router.post(
|
||
"/skills/reconcile-custom",
|
||
response_model=SkillReconcileResponse,
|
||
summary="Adopt Orphan Custom Skills",
|
||
description=(
|
||
"Scan the on-disk ``skills/custom/`` directory and register any skill that has a "
|
||
"valid SKILL.md but no DB ownership row. Adopted skills are recorded as owned by "
|
||
"the calling user with ``published=false``. Skips skills that already have a row "
|
||
"(regardless of owner). Useful after manually unzipping skills onto the filesystem."
|
||
),
|
||
)
|
||
async def reconcile_custom_skills(
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillReconcileResponse:
|
||
try:
|
||
user_id = _current_user_id(request)
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
custom_skills = [s for s in storage.load_skills(enabled_only=False) if s.category == SkillCategory.CUSTOM]
|
||
|
||
adopted: list[str] = []
|
||
already: list[str] = []
|
||
for skill in custom_skills:
|
||
existing = await skill_store.get_any(skill.name)
|
||
if existing is not None:
|
||
already.append(skill.name)
|
||
continue
|
||
await skill_store.ensure_legacy(
|
||
{"name": skill.name, "owner_user_id": user_id, "published": False}
|
||
)
|
||
try:
|
||
from deerflow.skills.usage import mark_uploaded
|
||
mark_uploaded(skill.name)
|
||
except Exception:
|
||
pass
|
||
adopted.append(skill.name)
|
||
|
||
if adopted:
|
||
await refresh_skills_system_prompt_cache_async()
|
||
return SkillReconcileResponse(adopted=adopted, already_registered=already)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to reconcile custom skills: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to reconcile custom skills: {e}")
|
||
|
||
|
||
@router.get("/skills/custom", response_model=SkillsListResponse, summary="List Custom Skills")
|
||
async def list_custom_skills(
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillsListResponse:
|
||
try:
|
||
user_id = _current_user_id(request)
|
||
custom = [s for s in get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False) if s.category == SkillCategory.CUSTOM]
|
||
ownership = await _build_ownership_index(skill_store, user_id)
|
||
display_adapters = await _list_skill_display_adapters([skill.name for skill in custom])
|
||
responses = []
|
||
for skill in custom:
|
||
if skill.name not in ownership:
|
||
continue
|
||
adapter = display_adapters.get(skill.name)
|
||
responses.append(
|
||
_skill_to_response(
|
||
skill,
|
||
ownership[skill.name],
|
||
display_override=_display_override_from_adapter(adapter),
|
||
display_sample=adapter.get("sample") if adapter else None,
|
||
)
|
||
)
|
||
responses = await _attach_owner_usernames(responses)
|
||
return SkillsListResponse(skills=responses)
|
||
except Exception as e:
|
||
logger.error("Failed to list custom skills: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to list custom skills: {str(e)}")
|
||
|
||
|
||
@router.get("/skills/custom/{skill_name}", response_model=CustomSkillContentResponse, summary="Get Custom Skill Content")
|
||
async def get_custom_skill(skill_name: str, config: AppConfig = Depends(get_config)) -> CustomSkillContentResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name and s.category == SkillCategory.CUSTOM), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Custom skill '{skill_name}' not found")
|
||
adapter = await _get_skill_display_adapter(skill_name)
|
||
return CustomSkillContentResponse(
|
||
**_skill_to_response(
|
||
skill,
|
||
display_override=_display_override_from_adapter(adapter),
|
||
display_sample=adapter.get("sample") if adapter else None,
|
||
).model_dump(),
|
||
content=get_or_new_skill_storage(app_config=config).read_custom_skill(skill_name),
|
||
)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to get custom skill %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to get custom skill: {str(e)}")
|
||
|
||
|
||
@router.put("/skills/custom/{skill_name}", response_model=CustomSkillContentResponse, summary="Edit Custom Skill")
|
||
async def update_custom_skill(skill_name: str, request: CustomSkillUpdateRequest, config: AppConfig = Depends(get_config)) -> CustomSkillContentResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
storage.ensure_custom_skill_is_editable(skill_name)
|
||
storage.validate_skill_markdown_content(skill_name, request.content)
|
||
scan = await scan_skill_content(request.content, executable=False, location=f"{skill_name}/{SKILL_MD_FILE}", app_config=config)
|
||
if scan.decision == "block":
|
||
raise HTTPException(status_code=400, detail=f"Security scan blocked the edit: {scan.reason}")
|
||
prev_content = storage.read_custom_skill(skill_name)
|
||
storage.write_custom_skill(skill_name, SKILL_MD_FILE, request.content)
|
||
storage.append_history(
|
||
skill_name,
|
||
{
|
||
"action": "human_edit",
|
||
"author": "human",
|
||
"thread_id": None,
|
||
"file_path": SKILL_MD_FILE,
|
||
"prev_content": prev_content,
|
||
"new_content": request.content,
|
||
"scanner": {"decision": scan.decision, "reason": scan.reason},
|
||
},
|
||
)
|
||
try:
|
||
from deerflow.skills.usage import bump_patch
|
||
bump_patch(skill_name)
|
||
except Exception:
|
||
pass
|
||
await refresh_skills_system_prompt_cache_async()
|
||
return await get_custom_skill(skill_name, config)
|
||
except HTTPException:
|
||
raise
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except Exception as e:
|
||
logger.error("Failed to update custom skill %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update custom skill: {str(e)}")
|
||
|
||
|
||
@router.delete("/skills/custom/{skill_name}", summary="Delete Custom Skill")
|
||
async def delete_custom_skill(
|
||
skill_name: str,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
agent_store=Depends(get_agent_store),
|
||
) -> dict[str, bool]:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
record, mode = await _get_editable_skill_or_raise(skill_store, skill_name, user_id, is_admin=is_admin)
|
||
operation = "admin_takedown" if mode == "admin" and record.get("owner_user_id") not in (None, user_id) else "delete"
|
||
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
storage.delete_custom_skill(
|
||
skill_name,
|
||
history_meta={
|
||
"action": "human_delete",
|
||
"author": "human",
|
||
"thread_id": None,
|
||
"file_path": SKILL_MD_FILE,
|
||
"prev_content": None,
|
||
"new_content": None,
|
||
"scanner": {"decision": "allow", "reason": "Deletion requested."},
|
||
},
|
||
)
|
||
try:
|
||
from deerflow.skills.usage import forget
|
||
forget(skill_name)
|
||
except Exception:
|
||
pass
|
||
|
||
# Persist ownership record removal so the DB stays consistent.
|
||
if mode == "admin":
|
||
await skill_store.delete_admin(skill_name)
|
||
else:
|
||
await skill_store.delete(skill_name, user_id)
|
||
|
||
try:
|
||
await get_tag_store(request).unassign_all("skill", skill_name)
|
||
except Exception: # noqa: BLE001 — tag cleanup must not block deletion
|
||
pass
|
||
|
||
# Strip references from every agent that pulled this skill in,
|
||
# and notify their owners (no-op for self).
|
||
await cascade_skill_unpublished_or_deleted(
|
||
skill_name,
|
||
actor_user_id=user_id,
|
||
operation=operation,
|
||
skill_owner_user_id=record.get("owner_user_id"),
|
||
agent_store=agent_store,
|
||
)
|
||
|
||
await refresh_skills_system_prompt_cache_async()
|
||
return {"success": True}
|
||
except HTTPException:
|
||
raise
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except Exception as e:
|
||
logger.error("Failed to delete custom skill %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to delete custom skill: {str(e)}")
|
||
|
||
|
||
def _remove_skill_state_config(skill_name: str) -> None:
|
||
"""Remove a deleted skill's enabled/display override from extensions config."""
|
||
extensions_config = get_extensions_config()
|
||
if skill_name not in extensions_config.skills:
|
||
return
|
||
|
||
config_path = ExtensionsConfig.resolve_config_path()
|
||
if config_path is None:
|
||
config_path = Path.cwd().parent / "extensions_config.json"
|
||
extensions_config.skills.pop(skill_name, None)
|
||
config_data = _dump_extensions_config(extensions_config)
|
||
|
||
tmp_path: Path | None = None
|
||
try:
|
||
with tempfile.NamedTemporaryFile(
|
||
"w",
|
||
encoding="utf-8",
|
||
delete=False,
|
||
dir=str(config_path.parent),
|
||
) as tmp_file:
|
||
json.dump(config_data, tmp_file, indent=2)
|
||
tmp_path = Path(tmp_file.name)
|
||
os.replace(tmp_path, config_path)
|
||
finally:
|
||
if tmp_path is not None and tmp_path.exists():
|
||
tmp_path.unlink()
|
||
reload_extensions_config()
|
||
|
||
|
||
@router.delete(
|
||
"/skills/{skill_name}",
|
||
summary="Delete Built-in Skill",
|
||
description=(
|
||
"Administrators may permanently delete a built-in skill. The skill is "
|
||
"removed from enabled configuration and detached from agents that used it."
|
||
),
|
||
)
|
||
async def delete_builtin_skill(
|
||
skill_name: str,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
agent_store=Depends(get_agent_store),
|
||
) -> dict[str, bool]:
|
||
try:
|
||
if not _is_admin_user(request):
|
||
raise HTTPException(status_code=403, detail="Only administrators can delete built-in skills")
|
||
|
||
skill_name = _normalize_skill_name_param(skill_name)
|
||
skill = _load_skill_or_404(skill_name, config)
|
||
if skill.category != SkillCategory.PUBLIC:
|
||
raise HTTPException(status_code=400, detail="Use the custom skill deletion endpoint for a custom skill")
|
||
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
public_root = (storage.get_skills_root_path() / SkillCategory.PUBLIC.value).resolve()
|
||
skill_dir = skill.skill_file.parent.resolve()
|
||
if skill_dir == public_root or public_root not in skill_dir.parents:
|
||
raise ValueError(f"Refusing to delete built-in skill outside the public skills directory: {skill.name}")
|
||
|
||
_remove_skill_state_config(skill.name)
|
||
shutil.rmtree(skill_dir)
|
||
|
||
try:
|
||
from deerflow.skills.usage import forget
|
||
|
||
forget(skill.name)
|
||
except Exception:
|
||
pass
|
||
|
||
await skill_store.delete_admin(skill.name)
|
||
try:
|
||
await get_tag_store(request).unassign_all("skill", skill.name)
|
||
except Exception: # noqa: BLE001 -- tag cleanup must not block deletion
|
||
pass
|
||
|
||
await cascade_skill_unpublished_or_deleted(
|
||
skill.name,
|
||
actor_user_id=_current_user_id(request),
|
||
operation="delete",
|
||
skill_owner_user_id=None,
|
||
agent_store=agent_store,
|
||
)
|
||
await refresh_skills_system_prompt_cache_async()
|
||
logger.info("Administrator deleted built-in skill %s", skill.name)
|
||
return {"success": True}
|
||
except HTTPException:
|
||
raise
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except Exception as e:
|
||
logger.error("Failed to delete built-in skill %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to delete built-in skill: {str(e)}")
|
||
|
||
|
||
class SkillPublishRequest(BaseModel):
|
||
"""Request body for toggling a custom skill's published flag."""
|
||
|
||
published: bool = Field(..., description="Whether to publish (True) or unpublish (False)")
|
||
square_id: str | None = Field(default=None, description="Publish square (发布广场) to file the skill under when publishing. Omit / empty = the default square.")
|
||
|
||
|
||
class SkillPublishResponse(BaseModel):
|
||
name: str
|
||
owner_user_id: str | None
|
||
published: bool
|
||
square_id: str = ""
|
||
|
||
|
||
@router.put(
|
||
"/skills/custom/{skill_name}/published",
|
||
response_model=SkillPublishResponse,
|
||
summary="Toggle Custom Skill Published State",
|
||
)
|
||
async def update_custom_skill_published(
|
||
skill_name: str,
|
||
body: SkillPublishRequest,
|
||
request: Request,
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
agent_store=Depends(get_agent_store),
|
||
) -> SkillPublishResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
record, mode = await _get_editable_skill_or_raise(skill_store, skill_name, user_id, is_admin=is_admin)
|
||
|
||
was_published = bool(record.get("published", False))
|
||
# When publishing, also file the skill into the chosen square (normalized
|
||
# to a valid configured square id; default when unset). Only one square
|
||
# per skill — a single column, so it can never be "published twice".
|
||
update_fields: dict[str, Any] = {"published": body.published}
|
||
if body.published:
|
||
from deerflow.config.system_settings import load_system_settings, normalize_square_id
|
||
|
||
update_fields["square_id"] = normalize_square_id(body.square_id, load_system_settings().publish_squares)
|
||
if mode == "admin":
|
||
updated = await skill_store.update_admin(skill_name, update_fields)
|
||
else:
|
||
updated = await skill_store.update(skill_name, user_id, update_fields)
|
||
if updated is None:
|
||
raise HTTPException(status_code=500, detail="Failed to persist published state")
|
||
|
||
# If we are *removing* visibility (was published, now not), strip from
|
||
# other users' agents and notify them. admin takedown is signaled by
|
||
# the actor not being the owner.
|
||
if was_published and not body.published:
|
||
skill_owner = record.get("owner_user_id")
|
||
operation = "admin_takedown" if mode == "admin" and skill_owner not in (None, user_id) else "unpublish"
|
||
await cascade_skill_unpublished_or_deleted(
|
||
skill_name,
|
||
actor_user_id=user_id,
|
||
operation=operation,
|
||
skill_owner_user_id=skill_owner,
|
||
agent_store=agent_store,
|
||
)
|
||
await refresh_skills_system_prompt_cache_async()
|
||
|
||
return SkillPublishResponse(
|
||
name=updated["name"],
|
||
owner_user_id=updated.get("owner_user_id"),
|
||
published=bool(updated.get("published", False)),
|
||
square_id=str(updated.get("square_id") or ""),
|
||
)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to update published state for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update published state: {str(e)}")
|
||
|
||
|
||
_ALLOWED_USER_STATES = {"active", "archived"}
|
||
|
||
|
||
class SkillStateRequest(BaseModel):
|
||
"""Request body for switching a custom skill's lifecycle state."""
|
||
|
||
state: str = Field(..., description="New state: 'active' (restore) or 'archived' (soft archive)")
|
||
|
||
|
||
class SkillPinRequest(BaseModel):
|
||
"""Request body for pinning / unpinning a custom skill."""
|
||
|
||
pinned: bool = Field(..., description="Whether to pin (true) so the curator never auto-archives this skill")
|
||
|
||
|
||
def _reload_skill_response(skill_name: str, ownership: dict[str, Any], config: AppConfig) -> SkillResponse:
|
||
"""Reload a single skill from disk and return its API response form."""
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found on disk")
|
||
return _skill_to_response(skill, ownership)
|
||
|
||
|
||
@router.put(
|
||
"/skills/custom/{skill_name}/state",
|
||
response_model=SkillResponse,
|
||
summary="Update Custom Skill Lifecycle State",
|
||
description="Restore (active) or soft-archive a custom skill. Only the owner or admins may invoke.",
|
||
)
|
||
async def update_custom_skill_state(
|
||
skill_name: str,
|
||
body: SkillStateRequest,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
if body.state not in _ALLOWED_USER_STATES:
|
||
raise HTTPException(
|
||
status_code=400,
|
||
detail=f"state must be one of {sorted(_ALLOWED_USER_STATES)}",
|
||
)
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
record, _mode = await _get_editable_skill_or_raise(skill_store, skill_name, user_id, is_admin=is_admin)
|
||
|
||
from deerflow.skills.usage import set_state as _usage_set_state
|
||
|
||
_usage_set_state(skill_name, body.state)
|
||
await refresh_skills_system_prompt_cache_async()
|
||
return _reload_skill_response(skill_name, record, config)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to update state for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update state: {e}")
|
||
|
||
|
||
@router.put(
|
||
"/skills/custom/{skill_name}/pinned",
|
||
response_model=SkillResponse,
|
||
summary="Pin or Unpin a Custom Skill",
|
||
description="Pin a skill so the curator never auto-archives it. Only the owner or admins may invoke.",
|
||
)
|
||
async def update_custom_skill_pinned(
|
||
skill_name: str,
|
||
body: SkillPinRequest,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
record, _mode = await _get_editable_skill_or_raise(skill_store, skill_name, user_id, is_admin=is_admin)
|
||
|
||
from deerflow.skills.usage import set_pinned as _usage_set_pinned
|
||
|
||
_usage_set_pinned(skill_name, body.pinned)
|
||
return _reload_skill_response(skill_name, record, config)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to update pinned for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update pinned: {e}")
|
||
|
||
|
||
class SkillFeaturedRequest(BaseModel):
|
||
"""Request body for pinning a custom skill to the top of the list."""
|
||
|
||
featured: bool = Field(..., description="Whether to pin the skill to the top of the list (sort only)")
|
||
|
||
|
||
@router.put(
|
||
"/skills/custom/{skill_name}/featured",
|
||
response_model=SkillResponse,
|
||
summary="Pin or Unpin a Custom Skill to Top",
|
||
description="Pin a skill to the top of the list (sort only; orthogonal to pinned/anti-archive). Only the owner or admins may invoke.",
|
||
)
|
||
async def update_custom_skill_featured(
|
||
skill_name: str,
|
||
body: SkillFeaturedRequest,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
record, _mode = await _get_editable_skill_or_raise(skill_store, skill_name, user_id, is_admin=is_admin)
|
||
|
||
from deerflow.skills.usage import set_featured as _usage_set_featured
|
||
|
||
_usage_set_featured(skill_name, body.featured)
|
||
return _reload_skill_response(skill_name, record, config)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to update featured for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update featured: {e}")
|
||
|
||
|
||
@router.get("/skills/custom/{skill_name}/history", response_model=CustomSkillHistoryResponse, summary="Get Custom Skill History")
|
||
async def get_custom_skill_history(skill_name: str, config: AppConfig = Depends(get_config)) -> CustomSkillHistoryResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
if not storage.custom_skill_exists(skill_name) and not storage.get_skill_history_file(skill_name).exists():
|
||
raise HTTPException(status_code=404, detail=f"Custom skill '{skill_name}' not found")
|
||
return CustomSkillHistoryResponse(history=storage.read_history(skill_name))
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to read history for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to read history: {str(e)}")
|
||
|
||
|
||
@router.post("/skills/custom/{skill_name}/rollback", response_model=CustomSkillContentResponse, summary="Rollback Custom Skill")
|
||
async def rollback_custom_skill(skill_name: str, request: SkillRollbackRequest, config: AppConfig = Depends(get_config)) -> CustomSkillContentResponse:
|
||
try:
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
if not storage.custom_skill_exists(skill_name) and not storage.get_skill_history_file(skill_name).exists():
|
||
raise HTTPException(status_code=404, detail=f"Custom skill '{skill_name}' not found")
|
||
history = storage.read_history(skill_name)
|
||
if not history:
|
||
raise HTTPException(status_code=400, detail=f"Custom skill '{skill_name}' has no history")
|
||
record = history[request.history_index]
|
||
target_content = record.get("prev_content")
|
||
if target_content is None:
|
||
raise HTTPException(status_code=400, detail="Selected history entry has no previous content to roll back to")
|
||
storage.validate_skill_markdown_content(skill_name, target_content)
|
||
scan = await scan_skill_content(target_content, executable=False, location=f"{skill_name}/{SKILL_MD_FILE}", app_config=config)
|
||
skill_file = storage.get_custom_skill_file(skill_name)
|
||
current_content = skill_file.read_text(encoding="utf-8") if skill_file.exists() else None
|
||
history_entry = {
|
||
"action": "rollback",
|
||
"author": "human",
|
||
"thread_id": None,
|
||
"file_path": SKILL_MD_FILE,
|
||
"prev_content": current_content,
|
||
"new_content": target_content,
|
||
"rollback_from_ts": record.get("ts"),
|
||
"scanner": {"decision": scan.decision, "reason": scan.reason},
|
||
}
|
||
if scan.decision == "block":
|
||
storage.append_history(skill_name, history_entry)
|
||
raise HTTPException(status_code=400, detail=f"Rollback blocked by security scanner: {scan.reason}")
|
||
storage.write_custom_skill(skill_name, SKILL_MD_FILE, target_content)
|
||
storage.append_history(skill_name, history_entry)
|
||
await refresh_skills_system_prompt_cache_async()
|
||
return await get_custom_skill(skill_name, config)
|
||
except HTTPException:
|
||
raise
|
||
except IndexError:
|
||
raise HTTPException(status_code=400, detail="history_index is out of range")
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except Exception as e:
|
||
logger.error("Failed to roll back custom skill %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to roll back custom skill: {str(e)}")
|
||
|
||
|
||
@router.get(
|
||
"/skills/{skill_name}",
|
||
response_model=SkillResponse,
|
||
summary="Get Skill Details",
|
||
description="Retrieve detailed information about a specific skill by its name.",
|
||
)
|
||
async def get_skill(
|
||
skill_name: str,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
|
||
# Attach ownership (owner/published/detail) when the caller may see it.
|
||
ownership = await skill_store.get_visible(skill_name, _current_user_id(request))
|
||
adapter = await _get_skill_display_adapter(skill_name)
|
||
return _skill_to_response(
|
||
skill,
|
||
ownership,
|
||
display_override=_display_override_from_adapter(adapter),
|
||
display_sample=adapter.get("sample") if adapter else None,
|
||
)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error(f"Failed to get skill {skill_name}: {e}", exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to get skill: {str(e)}")
|
||
|
||
|
||
class SkillContentResponse(BaseModel):
|
||
"""SKILL.md 原文内容,适用于公共和自定义技能。"""
|
||
|
||
name: str
|
||
category: SkillCategory
|
||
content: str
|
||
|
||
|
||
@router.get(
|
||
"/skills/{skill_name}/content",
|
||
response_model=SkillContentResponse,
|
||
summary="Get Skill Markdown Content",
|
||
description="Retrieve the raw SKILL.md content of a skill. Readable by any user the skill is visible to (built-in, published, or owned).",
|
||
)
|
||
async def get_skill_content(
|
||
skill_name: str,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillContentResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
user_id = _current_user_id(request)
|
||
# SKILL.md describes how to use a skill — visible to anyone who can see
|
||
# the skill itself (built-in, published, or owned). Admins always pass.
|
||
if not _is_admin_user(request):
|
||
record = await skill_store.get_any(skill_name)
|
||
if skill.category == SkillCategory.CUSTOM:
|
||
if record is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
if not record.get("published") and record.get("owner_user_id") != user_id:
|
||
raise HTTPException(status_code=403, detail="You can only view SKILL.md for skills visible to you")
|
||
content = skill.skill_file.read_text(encoding="utf-8")
|
||
return SkillContentResponse(name=skill.name, category=skill.category, content=content)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to get skill content %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to get skill content: {str(e)}")
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# Skill files browser (tree + per-file read / write)
|
||
# ---------------------------------------------------------------------------
|
||
|
||
|
||
_SKILL_FILE_TEXT_EXTS = {
|
||
".md", ".markdown", ".txt", ".rst",
|
||
".py", ".pyi",
|
||
".json", ".jsonl",
|
||
".yaml", ".yml",
|
||
".toml", ".ini", ".cfg", ".env",
|
||
".sh", ".bash", ".zsh", ".ps1",
|
||
".js", ".jsx", ".ts", ".tsx", ".mjs", ".cjs",
|
||
".html", ".htm", ".css", ".scss", ".sass", ".xml", ".svg",
|
||
".csv", ".tsv",
|
||
".sql",
|
||
}
|
||
_SKILL_FILE_TREE_SKIP_NAMES = {
|
||
"__pycache__", "node_modules", ".git", ".DS_Store", ".idea", ".vscode",
|
||
}
|
||
_SKILL_FILE_MAX_NODES = 2000
|
||
_SKILL_FILE_MAX_DEPTH = 8
|
||
_SKILL_FILE_READ_BYTES = 2 * 1024 * 1024 # 2 MiB cap on JSON content payload
|
||
_SKILL_FILE_WRITE_BYTES = 1 * 1024 * 1024 # 1 MiB cap on edits
|
||
|
||
|
||
class SkillFileNode(BaseModel):
|
||
"""One entry in a skill's file tree.
|
||
|
||
``children`` is only set for directories. ``path`` is POSIX-style and is
|
||
always relative to the skill's own directory (never absolute).
|
||
"""
|
||
|
||
name: str = Field(..., description="Base name (no slashes)")
|
||
type: str = Field(..., description="'file' or 'dir'")
|
||
path: str = Field(..., description="Path relative to the skill directory, POSIX-separated")
|
||
size: int | None = Field(default=None, description="File size in bytes (files only)")
|
||
mtime: float | None = Field(default=None, description="Modification time as POSIX seconds (files only)")
|
||
children: list["SkillFileNode"] | None = Field(default=None, description="Child entries (directories only)")
|
||
|
||
|
||
SkillFileNode.model_rebuild()
|
||
|
||
|
||
class SkillFilesResponse(BaseModel):
|
||
"""Tree of files under a skill directory."""
|
||
|
||
name: str
|
||
category: SkillCategory
|
||
editable: bool = Field(..., description="True when the caller may PUT new content into any file under this tree")
|
||
root: SkillFileNode
|
||
|
||
|
||
class SkillFileContentResponse(BaseModel):
|
||
"""JSON payload for reading a single file inside a skill."""
|
||
|
||
name: str
|
||
category: SkillCategory
|
||
path: str
|
||
mime: str
|
||
size: int
|
||
mtime: float
|
||
is_text: bool = Field(..., description="True when ``content`` carries the decoded text; otherwise download via the binary endpoint")
|
||
editable: bool = Field(..., description="True when the caller may PUT new content for this file")
|
||
content: str | None = Field(default=None, description="Decoded text content when ``is_text`` is True")
|
||
|
||
|
||
class SkillFileSaveRequest(BaseModel):
|
||
"""Request body for writing a text file inside a skill."""
|
||
|
||
path: str = Field(..., description="File path relative to the skill directory, POSIX-separated")
|
||
content: str = Field(..., description="Replacement text content")
|
||
|
||
|
||
class SkillFileSaveResponse(BaseModel):
|
||
name: str
|
||
path: str
|
||
size: int
|
||
mtime: float
|
||
|
||
|
||
def _normalize_skill_name_param(skill_name: str) -> str:
|
||
"""Strip CR/LF the same way other skill routes do before lookup."""
|
||
return skill_name.replace("\r\n", "").replace("\n", "")
|
||
|
||
|
||
def _load_skill_or_404(skill_name: str, config: AppConfig) -> Skill:
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
return skill
|
||
|
||
|
||
async def _ensure_files_view_allowed(
|
||
skill: Skill,
|
||
request: Request,
|
||
skill_store: SkillStore,
|
||
) -> dict[str, Any] | None:
|
||
"""Mirror _can_view_skill_source: admins or the publisher (custom owner)."""
|
||
record = await skill_store.get_any(skill.name)
|
||
owner_user_id = record.get("owner_user_id") if record else None
|
||
user_id = _current_user_id(request)
|
||
if not _can_view_skill_source(skill.category, owner_user_id, user_id, is_admin=_is_admin_user(request)):
|
||
raise HTTPException(status_code=403, detail="Only the skill publisher or an admin can browse skill files")
|
||
return record
|
||
|
||
|
||
def _caller_can_edit_skill_files(
|
||
skill: Skill,
|
||
record: dict[str, Any] | None,
|
||
request: Request,
|
||
) -> bool:
|
||
"""Whether the caller can edit or delete a file within *skill*.
|
||
|
||
Built-in skills are deployment-owned, so only administrators may change
|
||
their files. Custom skills remain editable by their owner or by an
|
||
administrator.
|
||
"""
|
||
if _is_admin_user(request):
|
||
return True
|
||
if skill.category != SkillCategory.CUSTOM:
|
||
return False
|
||
if record is None:
|
||
return False
|
||
return record.get("owner_user_id") == _current_user_id(request)
|
||
|
||
|
||
def _resolve_skill_relative_path(skill: Skill, relative_path: str) -> Path:
|
||
"""Validate *relative_path* against the skill directory and return the absolute Path.
|
||
|
||
Rejects empty, absolute, parent-traversing, or out-of-tree paths so a
|
||
malicious caller can never escape the skill's own directory.
|
||
"""
|
||
if not relative_path or relative_path.strip() == "":
|
||
raise HTTPException(status_code=400, detail="path must not be empty")
|
||
candidate = Path(relative_path)
|
||
if candidate.is_absolute() or any(part in {"..", ""} for part in candidate.parts):
|
||
raise HTTPException(status_code=400, detail="path must be a relative path without '..' segments")
|
||
skill_dir = skill.skill_file.parent.resolve()
|
||
target = (skill_dir / candidate).resolve()
|
||
try:
|
||
target.relative_to(skill_dir)
|
||
except ValueError:
|
||
raise HTTPException(status_code=400, detail="path must stay inside the skill directory")
|
||
return target
|
||
|
||
|
||
def _file_mime(path: Path) -> str:
|
||
"""Best-effort MIME guess; defaults to ``application/octet-stream``."""
|
||
guessed, _ = mimetypes.guess_type(path.name)
|
||
if guessed:
|
||
return guessed
|
||
if path.suffix.lower() in _SKILL_FILE_TEXT_EXTS:
|
||
return "text/plain"
|
||
return "application/octet-stream"
|
||
|
||
|
||
def _is_text_file(path: Path) -> bool:
|
||
"""Treat known text extensions as text; everything else as binary."""
|
||
return path.suffix.lower() in _SKILL_FILE_TEXT_EXTS
|
||
|
||
|
||
def _write_skill_text_file_atomically(target: Path, content: str) -> None:
|
||
"""Replace an already-validated skill file without exposing partial text."""
|
||
tmp_path: Path | None = None
|
||
try:
|
||
with tempfile.NamedTemporaryFile(
|
||
"w",
|
||
encoding="utf-8",
|
||
delete=False,
|
||
dir=str(target.parent),
|
||
) as tmp_file:
|
||
tmp_file.write(content)
|
||
tmp_path = Path(tmp_file.name)
|
||
os.replace(tmp_path, target)
|
||
finally:
|
||
if tmp_path is not None and tmp_path.exists():
|
||
tmp_path.unlink()
|
||
|
||
|
||
def _build_skill_file_tree(skill_dir: Path) -> SkillFileNode:
|
||
"""Walk *skill_dir* and build a SkillFileNode tree, depth/count capped.
|
||
|
||
Hidden files (``.foo``) and a small block-list of noise directories
|
||
(``__pycache__`` etc.) are skipped so the UI does not drown in housekeeping
|
||
artefacts.
|
||
"""
|
||
counter = {"n": 0}
|
||
|
||
def build(current_dir: Path, depth: int) -> SkillFileNode:
|
||
rel = current_dir.relative_to(skill_dir).as_posix()
|
||
node = SkillFileNode(
|
||
name=current_dir.name if rel != "." else skill_dir.name,
|
||
type="dir",
|
||
path="" if rel == "." else rel,
|
||
children=[],
|
||
)
|
||
if depth >= _SKILL_FILE_MAX_DEPTH:
|
||
return node
|
||
try:
|
||
entries = sorted(current_dir.iterdir(), key=lambda p: (not p.is_dir(), p.name.lower()))
|
||
except OSError:
|
||
return node
|
||
for entry in entries:
|
||
if counter["n"] >= _SKILL_FILE_MAX_NODES:
|
||
break
|
||
if entry.name.startswith(".") or entry.name in _SKILL_FILE_TREE_SKIP_NAMES:
|
||
continue
|
||
counter["n"] += 1
|
||
if entry.is_dir():
|
||
node.children.append(build(entry, depth + 1))
|
||
continue
|
||
try:
|
||
stat = entry.stat()
|
||
size = stat.st_size
|
||
mtime = stat.st_mtime
|
||
except OSError:
|
||
size = 0
|
||
mtime = 0.0
|
||
node.children.append(
|
||
SkillFileNode(
|
||
name=entry.name,
|
||
type="file",
|
||
path=entry.relative_to(skill_dir).as_posix(),
|
||
size=size,
|
||
mtime=mtime,
|
||
)
|
||
)
|
||
return node
|
||
|
||
return build(skill_dir, 0)
|
||
|
||
|
||
@router.get(
|
||
"/skills/{skill_name}/files",
|
||
response_model=SkillFilesResponse,
|
||
summary="List Skill Files (Tree)",
|
||
description=(
|
||
"Return the full file tree under a skill directory. Only an admin or the "
|
||
"skill's publisher (owner of a custom skill) may browse files; other users "
|
||
"get 403."
|
||
),
|
||
)
|
||
async def list_skill_files(
|
||
skill_name: str,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillFilesResponse:
|
||
try:
|
||
skill_name = _normalize_skill_name_param(skill_name)
|
||
skill = _load_skill_or_404(skill_name, config)
|
||
record = await _ensure_files_view_allowed(skill, request, skill_store)
|
||
skill_dir = skill.skill_file.parent
|
||
root = _build_skill_file_tree(skill_dir)
|
||
return SkillFilesResponse(
|
||
name=skill.name,
|
||
category=skill.category,
|
||
editable=_caller_can_edit_skill_files(skill, record, request),
|
||
root=root,
|
||
)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to list skill files for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to list skill files: {e}")
|
||
|
||
|
||
@router.get(
|
||
"/skills/{skill_name}/files/content",
|
||
response_model=SkillFileContentResponse,
|
||
summary="Read Skill File (Text)",
|
||
description=(
|
||
"Return JSON metadata for a single file under a skill. Text files include "
|
||
"their decoded ``content``; binary files omit ``content`` and should be "
|
||
"fetched from the ``files/download`` endpoint for inline display or "
|
||
"download. Same view authorisation as the file tree endpoint."
|
||
),
|
||
)
|
||
async def get_skill_file_content(
|
||
skill_name: str,
|
||
request: Request,
|
||
path: str = Query(..., description="File path relative to the skill directory"),
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillFileContentResponse:
|
||
try:
|
||
skill_name = _normalize_skill_name_param(skill_name)
|
||
skill = _load_skill_or_404(skill_name, config)
|
||
record = await _ensure_files_view_allowed(skill, request, skill_store)
|
||
target = _resolve_skill_relative_path(skill, path)
|
||
if not target.exists() or not target.is_file():
|
||
raise HTTPException(status_code=404, detail=f"File '{path}' not found in skill '{skill_name}'")
|
||
stat = target.stat()
|
||
mime = _file_mime(target)
|
||
is_text = _is_text_file(target)
|
||
content: str | None = None
|
||
if is_text:
|
||
if stat.st_size > _SKILL_FILE_READ_BYTES:
|
||
raise HTTPException(status_code=413, detail=f"File exceeds the {_SKILL_FILE_READ_BYTES // 1024} KiB preview limit")
|
||
try:
|
||
content = target.read_text(encoding="utf-8")
|
||
except UnicodeDecodeError:
|
||
# File has a text-y extension but is not UTF-8 — treat as binary.
|
||
is_text = False
|
||
return SkillFileContentResponse(
|
||
name=skill.name,
|
||
category=skill.category,
|
||
path=target.relative_to(skill.skill_file.parent.resolve()).as_posix(),
|
||
mime=mime,
|
||
size=stat.st_size,
|
||
mtime=stat.st_mtime,
|
||
is_text=is_text,
|
||
editable=is_text and _caller_can_edit_skill_files(skill, record, request),
|
||
content=content,
|
||
)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to read skill file %s/%s: %s", skill_name, path, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to read skill file: {e}")
|
||
|
||
|
||
@router.get(
|
||
"/skills/{skill_name}/files/download",
|
||
summary="Download Skill File (Raw Bytes)",
|
||
description=(
|
||
"Stream the raw bytes of a single file under a skill, suitable for "
|
||
"embedding in an ``<img>`` tag or for direct download. Same view "
|
||
"authorisation as the file tree endpoint."
|
||
),
|
||
)
|
||
async def download_skill_file(
|
||
skill_name: str,
|
||
request: Request,
|
||
path: str = Query(..., description="File path relative to the skill directory"),
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
):
|
||
skill_name = _normalize_skill_name_param(skill_name)
|
||
skill = _load_skill_or_404(skill_name, config)
|
||
await _ensure_files_view_allowed(skill, request, skill_store)
|
||
target = _resolve_skill_relative_path(skill, path)
|
||
if not target.exists() or not target.is_file():
|
||
raise HTTPException(status_code=404, detail=f"File '{path}' not found in skill '{skill_name}'")
|
||
return FileResponse(
|
||
path=target,
|
||
media_type=_file_mime(target),
|
||
filename=target.name,
|
||
)
|
||
|
||
|
||
@router.put(
|
||
"/skills/{skill_name}/files/content",
|
||
response_model=SkillFileSaveResponse,
|
||
summary="Write Skill File (Text)",
|
||
description=(
|
||
"Replace the contents of a text file under a skill. Administrators may "
|
||
"write built-in and custom skills; custom-skill owners may write their "
|
||
"own skills. Edits to ``SKILL.md`` re-validate its YAML frontmatter and "
|
||
"refresh the system-prompt cache."
|
||
),
|
||
)
|
||
async def save_skill_file(
|
||
skill_name: str,
|
||
body: SkillFileSaveRequest,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillFileSaveResponse:
|
||
try:
|
||
skill_name = _normalize_skill_name_param(skill_name)
|
||
skill = _load_skill_or_404(skill_name, config)
|
||
record = await skill_store.get_any(skill.name)
|
||
if not _caller_can_edit_skill_files(skill, record, request):
|
||
raise HTTPException(status_code=403, detail="Only an administrator or the custom skill owner can edit files in this skill")
|
||
|
||
target = _resolve_skill_relative_path(skill, body.path)
|
||
if not target.exists():
|
||
raise HTTPException(status_code=404, detail=f"File '{body.path}' not found in skill '{skill_name}'")
|
||
if target.is_dir():
|
||
raise HTTPException(status_code=400, detail="path points to a directory")
|
||
if not _is_text_file(target):
|
||
raise HTTPException(status_code=400, detail=f"File type '{target.suffix or target.name}' is not editable")
|
||
|
||
if len(body.content.encode("utf-8")) > _SKILL_FILE_WRITE_BYTES:
|
||
raise HTTPException(status_code=413, detail=f"Content exceeds the {_SKILL_FILE_WRITE_BYTES // 1024} KiB edit limit")
|
||
|
||
relative_posix = target.relative_to(skill.skill_file.parent.resolve()).as_posix()
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
|
||
is_skill_md = relative_posix.casefold() == SKILL_MD_FILE.casefold()
|
||
prev_content: str | None = None
|
||
if target.exists():
|
||
try:
|
||
prev_content = target.read_text(encoding="utf-8")
|
||
except UnicodeDecodeError:
|
||
prev_content = None
|
||
|
||
if is_skill_md:
|
||
# SKILL.md edits must keep their frontmatter valid (name must match).
|
||
storage.validate_skill_markdown_content(skill.name, body.content)
|
||
|
||
_write_skill_text_file_atomically(target, body.content)
|
||
|
||
try:
|
||
if skill.category == SkillCategory.CUSTOM:
|
||
storage.append_history(
|
||
skill.name,
|
||
{
|
||
"action": "human_edit_file",
|
||
"author": "human",
|
||
"thread_id": None,
|
||
"file_path": relative_posix,
|
||
"prev_content": prev_content,
|
||
"new_content": body.content,
|
||
},
|
||
)
|
||
except Exception as history_err: # noqa: BLE001 — history is best-effort
|
||
logger.warning("Failed to append edit history for %s/%s: %s", skill.name, relative_posix, history_err)
|
||
|
||
if is_skill_md:
|
||
try:
|
||
from deerflow.skills.usage import bump_patch
|
||
bump_patch(skill.name)
|
||
except Exception:
|
||
pass
|
||
await refresh_skills_system_prompt_cache_async()
|
||
|
||
stat = target.stat()
|
||
return SkillFileSaveResponse(
|
||
name=skill.name,
|
||
path=relative_posix,
|
||
size=stat.st_size,
|
||
mtime=stat.st_mtime,
|
||
)
|
||
except HTTPException:
|
||
raise
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except Exception as e:
|
||
logger.error("Failed to save skill file %s/%s: %s", skill_name, body.path, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to save skill file: {e}")
|
||
|
||
|
||
@router.delete(
|
||
"/skills/{skill_name}/files/content",
|
||
status_code=204,
|
||
summary="Delete Skill File",
|
||
description=(
|
||
"Delete one non-SKILL.md file under a skill. Administrators may delete "
|
||
"files from built-in and custom skills; custom-skill owners may delete "
|
||
"files from their own skills. The root ``SKILL.md`` manifest cannot be "
|
||
"deleted because it defines the skill."
|
||
),
|
||
)
|
||
async def delete_skill_file(
|
||
skill_name: str,
|
||
request: Request,
|
||
path: str = Query(..., description="File path relative to the skill directory"),
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> None:
|
||
try:
|
||
skill_name = _normalize_skill_name_param(skill_name)
|
||
skill = _load_skill_or_404(skill_name, config)
|
||
record = await skill_store.get_any(skill.name)
|
||
if not _caller_can_edit_skill_files(skill, record, request):
|
||
raise HTTPException(status_code=403, detail="Only an administrator or the custom skill owner can delete files in this skill")
|
||
|
||
target = _resolve_skill_relative_path(skill, path)
|
||
if not target.exists() or not target.is_file():
|
||
raise HTTPException(status_code=404, detail=f"File '{path}' not found in skill '{skill_name}'")
|
||
|
||
relative_posix = target.relative_to(skill.skill_file.parent.resolve()).as_posix()
|
||
if relative_posix.casefold() == SKILL_MD_FILE.casefold():
|
||
raise HTTPException(status_code=400, detail="SKILL.md defines the skill and cannot be deleted")
|
||
|
||
prev_content: str | None = None
|
||
if _is_text_file(target):
|
||
try:
|
||
prev_content = target.read_text(encoding="utf-8")
|
||
except UnicodeDecodeError:
|
||
pass
|
||
target.unlink()
|
||
|
||
if skill.category == SkillCategory.CUSTOM:
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
try:
|
||
storage.append_history(
|
||
skill.name,
|
||
{
|
||
"action": "human_delete_file",
|
||
"author": "human",
|
||
"thread_id": None,
|
||
"file_path": relative_posix,
|
||
"prev_content": prev_content,
|
||
},
|
||
)
|
||
except Exception as history_err: # noqa: BLE001
|
||
logger.warning("Failed to append deletion history for %s/%s: %s", skill.name, relative_posix, history_err)
|
||
else:
|
||
logger.info("Administrator deleted built-in skill file %s/%s", skill.name, relative_posix)
|
||
except HTTPException:
|
||
raise
|
||
except ValueError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
except FileNotFoundError as e:
|
||
raise HTTPException(status_code=404, detail=str(e))
|
||
except Exception as e:
|
||
logger.error("Failed to delete skill file %s/%s: %s", skill_name, path, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to delete skill file: {e}")
|
||
|
||
|
||
@router.put(
|
||
"/skills/{skill_name}",
|
||
response_model=SkillResponse,
|
||
summary="Update Skill",
|
||
description="Update a skill's enabled status by modifying the extensions_config.json file.",
|
||
)
|
||
async def update_skill(skill_name: str, request: SkillUpdateRequest, config: AppConfig = Depends(get_config)) -> SkillResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
|
||
config_path = ExtensionsConfig.resolve_config_path()
|
||
if config_path is None:
|
||
config_path = Path.cwd().parent / "extensions_config.json"
|
||
logger.info(f"No existing extensions config found. Creating new config at: {config_path}")
|
||
|
||
extensions_config = get_extensions_config()
|
||
current_skill_config = extensions_config.skills.get(skill_name)
|
||
extensions_config.skills[skill_name] = SkillStateConfig(
|
||
enabled=request.enabled,
|
||
display=current_skill_config.display if current_skill_config else None,
|
||
)
|
||
|
||
config_data = _dump_extensions_config(extensions_config)
|
||
|
||
with open(config_path, "w", encoding="utf-8") as f:
|
||
json.dump(config_data, f, indent=2)
|
||
|
||
logger.info(f"Skills configuration updated and saved to: {config_path}")
|
||
reload_extensions_config()
|
||
await refresh_skills_system_prompt_cache_async()
|
||
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
updated_skill = next((s for s in skills if s.name == skill_name), None)
|
||
|
||
if updated_skill is None:
|
||
raise HTTPException(status_code=500, detail=f"Failed to reload skill '{skill_name}' after update")
|
||
|
||
logger.info(f"Skill '{skill_name}' enabled status updated to {request.enabled}")
|
||
return _skill_to_response(updated_skill)
|
||
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error(f"Failed to update skill {skill_name}: {e}", exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update skill: {str(e)}")
|
||
|
||
|
||
@router.put(
|
||
"/skills/{skill_name}/always-on",
|
||
response_model=SkillResponse,
|
||
summary="Update Skill 常驻(免压缩) Flag",
|
||
description="Mark a skill as 常驻(免压缩) so it is injected full-text into the system prompt even when skill compression is enabled.",
|
||
)
|
||
async def update_skill_always_on(
|
||
skill_name: str,
|
||
request: SkillAlwaysOnUpdateRequest,
|
||
http_request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillResponse:
|
||
"""Persist the 常驻(免压缩) flag in the **DB** (skills table), not in
|
||
extensions_config.json — the source-tree config file is read-only / unsafe to
|
||
rewrite in some deployments, and built-in skills carry no file entry anyway.
|
||
Mirrors the detail/name_zh endpoints: custom skill → owner or admin; built-in
|
||
/ public skill → admin only (materializes a row on demand)."""
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
user_id = _current_user_id(http_request)
|
||
is_admin = _is_admin_user(http_request)
|
||
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
|
||
if skill.category == SkillCategory.CUSTOM:
|
||
_record, mode = await _get_editable_skill_or_raise(skill_store, skill_name, user_id, is_admin=is_admin)
|
||
if mode == "admin":
|
||
updated = await skill_store.update_admin(skill_name, {"always_on": request.always_on})
|
||
else:
|
||
updated = await skill_store.update(skill_name, user_id, {"always_on": request.always_on})
|
||
else:
|
||
# Built-in / public skill — admin only; no DB row by default.
|
||
if not is_admin:
|
||
raise HTTPException(status_code=403, detail="Only admin users can change a built-in skill's 常驻 flag")
|
||
if await skill_store.get_any(skill_name) is None:
|
||
await skill_store.ensure_legacy({"name": skill_name, "owner_user_id": None, "published": True})
|
||
updated = await skill_store.update_admin(skill_name, {"always_on": request.always_on})
|
||
|
||
if updated is None:
|
||
raise HTTPException(status_code=500, detail="Failed to persist skill 常驻 flag")
|
||
|
||
# always_on lives in the system prompt's skill injection — refresh the cache.
|
||
await refresh_skills_system_prompt_cache_async()
|
||
|
||
logger.info("Skill '%s' always_on updated to %s", skill_name, request.always_on)
|
||
return _skill_to_response(skill, updated)
|
||
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error(f"Failed to update skill always_on {skill_name}: {e}", exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update skill: {str(e)}")
|
||
|
||
|
||
@router.put(
|
||
"/skills/{skill_name}/display",
|
||
response_model=SkillResponse,
|
||
summary="Update Skill Display Adapter",
|
||
description="Update a skill's frontend display adapter in extensions_config.json.",
|
||
)
|
||
async def update_skill_display(
|
||
skill_name: str,
|
||
request: SkillDisplayUpdateRequest,
|
||
config: AppConfig = Depends(get_config),
|
||
) -> SkillResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
|
||
config_path = ExtensionsConfig.resolve_config_path()
|
||
if config_path is None:
|
||
config_path = Path.cwd().parent / "extensions_config.json"
|
||
logger.info(f"No existing extensions config found. Creating new config at: {config_path}")
|
||
|
||
extensions_config = get_extensions_config()
|
||
current_skill_config = extensions_config.skills.get(skill_name)
|
||
extensions_config.skills[skill_name] = SkillStateConfig(
|
||
enabled=current_skill_config.enabled if current_skill_config else skill.enabled,
|
||
display=request.display,
|
||
)
|
||
|
||
display_payload = request.display.model_dump(by_alias=True) if request.display is not None else None
|
||
sample_payload = request.sample if "sample" in request.model_fields_set else ...
|
||
await _upsert_skill_display_adapter(
|
||
skill_name,
|
||
display=display_payload,
|
||
sample=sample_payload,
|
||
)
|
||
|
||
with open(config_path, "w", encoding="utf-8") as f:
|
||
json.dump(_dump_extensions_config(extensions_config), f, indent=2, ensure_ascii=False)
|
||
|
||
reload_extensions_config()
|
||
await refresh_skills_system_prompt_cache_async()
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
updated_skill = next((s for s in skills if s.name == skill_name), None)
|
||
if updated_skill is None:
|
||
raise HTTPException(status_code=500, detail=f"Failed to reload skill '{skill_name}' after display update")
|
||
adapter = await _get_skill_display_adapter(skill_name)
|
||
return _skill_to_response(
|
||
updated_skill,
|
||
display_override=_display_override_from_adapter(adapter),
|
||
display_sample=adapter.get("sample") if adapter else None,
|
||
)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to update display adapter for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update skill display adapter: {str(e)}")
|
||
|
||
|
||
@router.get(
|
||
"/skills/{skill_name}/display-sample",
|
||
response_model=PersistedSkillDisplaySampleResponse,
|
||
summary="Get Persisted Skill Display Sample",
|
||
description="Return only the sample stored in the database. Does not generate a fallback sample.",
|
||
)
|
||
async def get_persisted_skill_display_sample(
|
||
skill_name: str,
|
||
config: AppConfig = Depends(get_config),
|
||
) -> PersistedSkillDisplaySampleResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
adapter = await _get_skill_display_adapter(skill_name)
|
||
sample = adapter.get("sample") if adapter else None
|
||
return PersistedSkillDisplaySampleResponse(
|
||
skill_name=skill_name,
|
||
sample=sample if isinstance(sample, dict) else None,
|
||
)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to get persisted display sample for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to get persisted display sample: {str(e)}")
|
||
|
||
|
||
@router.post(
|
||
"/skills/{skill_name}/display-sample",
|
||
response_model=SkillDisplaySampleResponse,
|
||
summary="Get Skill Display Sample",
|
||
description="Generate or return a structured sample response used by the frontend to infer display mapping.",
|
||
)
|
||
async def get_skill_display_sample(
|
||
skill_name: str,
|
||
config: AppConfig = Depends(get_config),
|
||
) -> SkillDisplaySampleResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
sample = await _generate_skill_display_sample(skill, config)
|
||
await _upsert_skill_display_adapter(skill_name, sample=sample)
|
||
return SkillDisplaySampleResponse(skill_name=skill_name, sample=sample)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to get display sample for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to get display sample: {str(e)}")
|
||
|
||
|
||
class SkillDetailRequest(BaseModel):
|
||
"""Request body for setting a skill's human-authored detail text."""
|
||
|
||
detail: str = Field(default="", description="详情说明文本,支持 Markdown;传空字符串表示清空")
|
||
|
||
|
||
class SkillDetailResponse(BaseModel):
|
||
"""Response model for the skill-detail update endpoint."""
|
||
|
||
name: str = Field(..., description="Skill name")
|
||
detail: str | None = Field(default=None, description="Persisted detail text (null when cleared)")
|
||
|
||
|
||
@router.put(
|
||
"/skills/{skill_name}/detail",
|
||
response_model=SkillDetailResponse,
|
||
summary="Update Skill Detail",
|
||
description="设置技能的详情说明。自定义技能由其所有者或管理员编辑;内置/公开技能仅管理员可编辑。",
|
||
)
|
||
async def update_skill_detail(
|
||
skill_name: str,
|
||
body: SkillDetailRequest,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillDetailResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
|
||
# Empty / whitespace-only input clears the annotation.
|
||
detail_value: str | None = body.detail.strip() or None
|
||
|
||
if skill.category == SkillCategory.CUSTOM:
|
||
# Owner edits their own skill; admins may edit anyone's.
|
||
_record, mode = await _get_editable_skill_or_raise(skill_store, skill_name, user_id, is_admin=is_admin)
|
||
if mode == "admin":
|
||
updated = await skill_store.update_admin(skill_name, {"detail": detail_value})
|
||
else:
|
||
updated = await skill_store.update(skill_name, user_id, {"detail": detail_value})
|
||
else:
|
||
# Built-in / public skill — only admins may annotate it. These carry
|
||
# no DB row by default, so materialize one before writing detail.
|
||
if not is_admin:
|
||
raise HTTPException(status_code=403, detail="Only admin users can edit built-in skill details")
|
||
if await skill_store.get_any(skill_name) is None:
|
||
await skill_store.ensure_legacy({"name": skill_name, "owner_user_id": None, "published": True})
|
||
updated = await skill_store.update_admin(skill_name, {"detail": detail_value})
|
||
|
||
if updated is None:
|
||
raise HTTPException(status_code=500, detail="Failed to persist skill detail")
|
||
return SkillDetailResponse(name=skill_name, detail=updated.get("detail"))
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to update detail for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update skill detail: {str(e)}")
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# 技能中文名称(用户可配置 + 大模型自动生成 + 一键补充)
|
||
# ---------------------------------------------------------------------------
|
||
|
||
_SKILL_NAME_ZH_SYSTEM_PROMPT = (
|
||
"你是一名中文产品命名助手。根据给定技能的英文名与功能描述,生成一个简洁、自然、"
|
||
"易懂的中文名称。要求:2 到 10 个汉字;准确概括技能用途;不要包含标点、引号、英文、"
|
||
"解释或多余文字;只输出中文名称本身。"
|
||
)
|
||
|
||
|
||
class SkillNameZhModelError(Exception):
|
||
"""The LLM call backing skill-name generation failed (a *model-side* problem).
|
||
|
||
Distinguishes a genuine model failure (timeout / API error / unbuildable
|
||
model) — which must be surfaced to the user as "大模型报错导致无法补全" — from
|
||
the model merely returning an unusable name (callers get ``None`` instead).
|
||
"""
|
||
|
||
|
||
def _clean_generated_name_zh(raw: str) -> str | None:
|
||
"""Normalise an LLM-produced Chinese name: first line, strip quotes/punct, cap length."""
|
||
if not raw:
|
||
return None
|
||
text = str(raw).strip()
|
||
# Take the first non-empty line only.
|
||
for line in text.splitlines():
|
||
line = line.strip()
|
||
if line:
|
||
text = line
|
||
break
|
||
# Strip wrapping quotes / common trailing punctuation.
|
||
text = text.strip().strip("\"'“”‘’`「」《》【】()()[]").strip()
|
||
text = text.rstrip("。.,,、!!??;;::")
|
||
if not text:
|
||
return None
|
||
# Keep it short — guard against the model returning a sentence.
|
||
if len(text) > 20:
|
||
text = text[:20]
|
||
return text or None
|
||
|
||
|
||
def _resolve_skill_name_zh_model_name(app_config: AppConfig, model_name: str | None) -> str | None:
|
||
"""Pick the model name for skill-name generation.
|
||
|
||
An explicit ``model_name`` (user-selected in the UI) wins; otherwise fall
|
||
back to the curator/security-scan model setting (itself falling back to the
|
||
main model when unset).
|
||
"""
|
||
if model_name and model_name.strip():
|
||
return model_name.strip()
|
||
if app_config.skill_evolution and app_config.skill_evolution.curator:
|
||
return app_config.skill_evolution.curator.model_name
|
||
return None
|
||
|
||
|
||
async def _generate_skill_name_zh(skill: Skill, app_config: AppConfig, model: Any = None, *, model_name: str | None = None) -> str | None:
|
||
"""Generate a Chinese display name for ``skill`` via the LLM.
|
||
|
||
An explicit ``model_name`` lets the caller pick the model; otherwise it
|
||
reuses the curator/security-scan model setting (falls back to the main
|
||
model). Returns ``None`` only when the model replies but the output cannot be
|
||
normalised into a short Chinese name. Raises :class:`SkillNameZhModelError`
|
||
when the model itself fails (cannot be built / times out / API error), so the
|
||
caller can tell the user it was a *model* problem.
|
||
"""
|
||
user_prompt = (
|
||
f"技能英文名:{skill.name}\n"
|
||
f"功能描述:{skill.description or '(无描述)'}\n\n"
|
||
"请为该技能起一个中文名称。"
|
||
)
|
||
if model is None:
|
||
try:
|
||
from deerflow.models import create_chat_model
|
||
|
||
resolved = _resolve_skill_name_zh_model_name(app_config, model_name)
|
||
model = create_chat_model(name=resolved, thinking_enabled=False, app_config=app_config)
|
||
except Exception as e:
|
||
logger.warning("Skill name_zh model build failed for %s", skill.name, exc_info=True)
|
||
raise SkillNameZhModelError(str(e) or "无法初始化所选大模型") from e
|
||
try:
|
||
response = await asyncio.wait_for(
|
||
model.ainvoke(
|
||
[
|
||
{"role": "system", "content": _SKILL_NAME_ZH_SYSTEM_PROMPT},
|
||
{"role": "user", "content": user_prompt},
|
||
],
|
||
config={"run_name": "skill_name_zh"},
|
||
),
|
||
timeout=60,
|
||
)
|
||
except asyncio.TimeoutError as e:
|
||
logger.warning("Skill name_zh LLM call timed out for %s", skill.name)
|
||
raise SkillNameZhModelError("大模型响应超时(60 秒)") from e
|
||
except Exception as e:
|
||
logger.warning("Skill name_zh LLM call failed for %s", skill.name, exc_info=True)
|
||
raise SkillNameZhModelError(str(e) or "大模型调用失败") from e
|
||
return _clean_generated_name_zh(str(getattr(response, "content", "") or ""))
|
||
|
||
|
||
async def _persist_skill_name_zh(
|
||
skill: Skill,
|
||
name_zh: str | None,
|
||
skill_store: SkillStore,
|
||
user_id: str,
|
||
*,
|
||
is_admin: bool,
|
||
) -> str | None:
|
||
"""Persist ``name_zh`` into the ``skills`` DB table (never the on-disk SKILL.md).
|
||
|
||
Authorization mirrors the ``detail`` annotation: a custom skill is editable by
|
||
its owner or an admin; a built-in / public skill only by an admin (its DB row
|
||
is materialized on demand). Returns the stored value.
|
||
"""
|
||
if skill.category == SkillCategory.CUSTOM:
|
||
_record, mode = await _get_editable_skill_or_raise(skill_store, skill.name, user_id, is_admin=is_admin)
|
||
if mode == "admin":
|
||
updated = await skill_store.update_admin(skill.name, {"name_zh": name_zh})
|
||
else:
|
||
updated = await skill_store.update(skill.name, user_id, {"name_zh": name_zh})
|
||
else:
|
||
if not is_admin:
|
||
raise HTTPException(status_code=403, detail="Only admin users can edit built-in skill Chinese names")
|
||
if await skill_store.get_any(skill.name) is None:
|
||
await skill_store.ensure_legacy({"name": skill.name, "owner_user_id": None, "published": True})
|
||
updated = await skill_store.update_admin(skill.name, {"name_zh": name_zh})
|
||
|
||
if updated is None:
|
||
raise HTTPException(status_code=500, detail="Failed to persist skill Chinese name")
|
||
return updated.get("name_zh")
|
||
|
||
|
||
class SkillNameZhRequest(BaseModel):
|
||
"""Request body for setting a skill's Chinese display name."""
|
||
|
||
name_zh: str = Field(default="", description="中文名称;传空字符串表示清空")
|
||
|
||
|
||
class SkillNameZhResponse(BaseModel):
|
||
"""Response model for a single skill's Chinese-name update."""
|
||
|
||
name: str = Field(..., description="Skill name")
|
||
name_zh: str | None = Field(default=None, description="中文名称(为空表示已清空)")
|
||
|
||
|
||
class SkillNameZhBulkResponse(BaseModel):
|
||
"""Response model for the one-click bulk fill of missing Chinese names."""
|
||
|
||
generated: list[SkillNameZhResponse] = Field(default_factory=list, description="本次成功生成中文名的技能")
|
||
failed: list[str] = Field(default_factory=list, description="生成失败的技能英文名")
|
||
skipped: int = Field(default=0, description="已有中文名而被跳过的技能数")
|
||
model_error: str | None = Field(default=None, description="若有技能因大模型报错而失败,这里给出代表性的报错信息")
|
||
|
||
|
||
class GenerateSkillNameZhRequest(BaseModel):
|
||
"""Optional body for skill-name generation: pick the model to use."""
|
||
|
||
model_name: str | None = Field(default=None, description="可选:指定用于生成的大模型名称;留空使用默认配置")
|
||
|
||
|
||
def _can_edit_skill_name_zh(skill: Skill, record: dict[str, Any] | None, user_id: str, *, is_admin: bool) -> bool:
|
||
"""Whether ``user_id`` may set this skill's Chinese name (mirrors detail auth)."""
|
||
if is_admin:
|
||
return True
|
||
if skill.category != SkillCategory.CUSTOM:
|
||
return False
|
||
return record is not None and record.get("owner_user_id") == user_id
|
||
|
||
|
||
@router.put(
|
||
"/skills/{skill_name}/name-zh",
|
||
response_model=SkillNameZhResponse,
|
||
summary="Set Skill Chinese Name",
|
||
description="设置技能的中文名称(存数据库,不改源文件)。自定义技能由其所有者或管理员编辑;内置/公开技能仅管理员可编辑。传空字符串表示清空,回退到使用英文名。",
|
||
)
|
||
async def set_skill_name_zh(
|
||
skill_name: str,
|
||
body: SkillNameZhRequest,
|
||
request: Request,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillNameZhResponse:
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
name_zh = body.name_zh.strip() or None
|
||
stored = await _persist_skill_name_zh(skill, name_zh, skill_store, user_id, is_admin=is_admin)
|
||
return SkillNameZhResponse(name=skill_name, name_zh=stored)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to set name_zh for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to set skill Chinese name: {str(e)}")
|
||
|
||
|
||
@router.post(
|
||
"/skills/{skill_name}/name-zh/generate",
|
||
response_model=SkillNameZhResponse,
|
||
summary="Generate Skill Chinese Name",
|
||
description="用大模型为该技能生成一个中文名称并保存到数据库。可在请求体里通过 model_name 指定使用的大模型。",
|
||
)
|
||
async def generate_skill_name_zh(
|
||
skill_name: str,
|
||
request: Request,
|
||
body: GenerateSkillNameZhRequest | None = None,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillNameZhResponse:
|
||
try:
|
||
body = body or GenerateSkillNameZhRequest()
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
try:
|
||
name_zh = await _generate_skill_name_zh(skill, config, model_name=body.model_name)
|
||
except SkillNameZhModelError as e:
|
||
raise HTTPException(status_code=502, detail=f"大模型报错,无法补全技能中文名:{e}")
|
||
if not name_zh:
|
||
raise HTTPException(status_code=502, detail="大模型未能生成有效的中文名称,请稍后重试或手动填写")
|
||
stored = await _persist_skill_name_zh(skill, name_zh, skill_store, user_id, is_admin=is_admin)
|
||
return SkillNameZhResponse(name=skill_name, name_zh=stored)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to generate name_zh for %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to generate skill Chinese name: {str(e)}")
|
||
|
||
|
||
@router.post(
|
||
"/skills/name-zh/generate-missing",
|
||
response_model=SkillNameZhBulkResponse,
|
||
summary="Bulk Generate Missing Skill Chinese Names",
|
||
description="管理员一键为库中所有尚未配置中文名称的技能生成中文名(存数据库)。可在请求体里通过 model_name 指定使用的大模型。",
|
||
)
|
||
async def generate_missing_skill_names_zh(
|
||
request: Request,
|
||
body: GenerateSkillNameZhRequest | None = None,
|
||
config: AppConfig = Depends(get_config),
|
||
skill_store: SkillStore = Depends(get_skill_store),
|
||
) -> SkillNameZhBulkResponse:
|
||
try:
|
||
from deerflow.models import create_chat_model
|
||
|
||
body = body or GenerateSkillNameZhRequest()
|
||
user_id = _current_user_id(request)
|
||
is_admin = _is_admin_user(request)
|
||
# 一键补充技能中文名仅对管理员开放。
|
||
if not is_admin:
|
||
raise HTTPException(status_code=403, detail="仅管理员可一键补充技能中文名")
|
||
skills = get_or_new_skill_storage(app_config=config).load_skills(enabled_only=False)
|
||
|
||
# Admin sees every row, so existing name_zh values come from the full set.
|
||
records = {r["name"]: r for r in await skill_store.list_all()}
|
||
|
||
pending: list[Skill] = []
|
||
skipped = 0
|
||
for skill in skills:
|
||
record = records.get(skill.name)
|
||
if record and record.get("name_zh"):
|
||
skipped += 1
|
||
continue
|
||
pending.append(skill)
|
||
|
||
generated: list[SkillNameZhResponse] = []
|
||
failed: list[str] = []
|
||
model_error: str | None = None
|
||
if pending:
|
||
# Build the model once and reuse it for every skill. A build failure
|
||
# is a model problem — fail fast with a clear message rather than
|
||
# silently producing zero results.
|
||
try:
|
||
resolved_name = _resolve_skill_name_zh_model_name(config, body.model_name)
|
||
model = create_chat_model(name=resolved_name, thinking_enabled=False, app_config=config)
|
||
except Exception as e:
|
||
logger.warning("Failed to build model for bulk name_zh generation", exc_info=True)
|
||
raise HTTPException(status_code=502, detail=f"大模型报错,无法补全技能中文名:{str(e) or '无法初始化所选大模型'}")
|
||
|
||
# Generate concurrently with a small bound to avoid overloading the model.
|
||
semaphore = asyncio.Semaphore(4)
|
||
|
||
async def _one(skill: Skill) -> tuple[Skill, str | None, str | None]:
|
||
async with semaphore:
|
||
try:
|
||
return skill, await _generate_skill_name_zh(skill, config, model=model), None
|
||
except SkillNameZhModelError as e:
|
||
return skill, None, str(e)
|
||
|
||
results = await asyncio.gather(*[_one(s) for s in pending])
|
||
for skill, name_zh, err in results:
|
||
if err and model_error is None:
|
||
model_error = err
|
||
if not name_zh:
|
||
failed.append(skill.name)
|
||
continue
|
||
try:
|
||
stored = await _persist_skill_name_zh(skill, name_zh, skill_store, user_id, is_admin=is_admin)
|
||
generated.append(SkillNameZhResponse(name=skill.name, name_zh=stored))
|
||
except HTTPException:
|
||
# Permission / persistence issue for this one skill — skip it.
|
||
failed.append(skill.name)
|
||
|
||
return SkillNameZhBulkResponse(generated=generated, failed=failed, skipped=skipped, model_error=model_error)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to bulk-generate skill Chinese names: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to bulk-generate skill Chinese names: {str(e)}")
|
||
|
||
|
||
class SkillPinRequest(BaseModel):
|
||
pinned: bool = Field(..., description="Whether to pin or unpin the skill")
|
||
|
||
|
||
@router.post("/skills/custom/{skill_name}/pin", response_model=SkillResponse, summary="Pin or Unpin Skill")
|
||
async def pin_skill(skill_name: str, request: SkillPinRequest, config: AppConfig = Depends(get_config)) -> SkillResponse:
|
||
"""Pin a skill to prevent it from being automatically archived by the curator."""
|
||
try:
|
||
skill_name = skill_name.replace("\r\n", "").replace("\n", "")
|
||
storage = get_or_new_skill_storage(app_config=config)
|
||
if not storage.custom_skill_exists(skill_name):
|
||
raise HTTPException(status_code=404, detail=f"Custom skill '{skill_name}' not found")
|
||
from deerflow.skills.usage import set_pinned
|
||
set_pinned(skill_name, request.pinned)
|
||
skills = storage.load_skills(enabled_only=False)
|
||
skill = next((s for s in skills if s.name == skill_name), None)
|
||
if skill is None:
|
||
raise HTTPException(status_code=404, detail=f"Skill '{skill_name}' not found")
|
||
return _skill_to_response(skill)
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to pin skill %s: %s", skill_name, e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to pin skill: {str(e)}")
|
||
|
||
|
||
@router.get("/curator/status", summary="Curator Status")
|
||
async def curator_status(request: Request) -> dict:
|
||
"""Return the current curator state and per-lifecycle skill counts for the
|
||
authenticated user."""
|
||
try:
|
||
from deerflow.skills.curator import get_curator_status
|
||
return get_curator_status(_current_user_id(request))
|
||
except Exception as e:
|
||
logger.error("Failed to get curator status: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to get curator status: {str(e)}")
|
||
|
||
|
||
@router.post("/curator/run", summary="Trigger Curator Run")
|
||
async def run_curator(request: Request) -> dict:
|
||
"""Manually trigger a curator review pass for the authenticated user
|
||
(bypasses the interval check, runs in a background thread).
|
||
|
||
Pre-checks the master switch / curator switch / run lock so the response
|
||
accurately reflects whether a run actually started — instead of always
|
||
reporting success. Returns ``{"triggered": bool, "reason": str | None}``.
|
||
"""
|
||
try:
|
||
import threading
|
||
|
||
from deerflow.skills.curator import _curator_thread_target, manual_trigger_blocked_reason
|
||
|
||
user_id = _current_user_id(request)
|
||
reason = manual_trigger_blocked_reason(user_id)
|
||
if reason is not None:
|
||
return {"triggered": False, "reason": reason}
|
||
thread = threading.Thread(
|
||
target=_curator_thread_target,
|
||
args=(user_id,),
|
||
daemon=True,
|
||
name=f"skill-curator-manual-{user_id}",
|
||
)
|
||
thread.start()
|
||
return {"triggered": True, "reason": None}
|
||
except Exception as e:
|
||
logger.error("Failed to trigger curator: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to trigger curator: {str(e)}")
|
||
|
||
|
||
class CuratorSubConfig(BaseModel):
|
||
"""Nested curator configuration mirroring CuratorConfig."""
|
||
|
||
# 与 CuratorConfig.enabled 保持一致:默认关闭,需显式 opt-in。
|
||
enabled: bool = Field(default=False, description="Whether the curator runs automatically.")
|
||
interval_hours: int = Field(default=168, ge=1, description="Hours between curator runs.")
|
||
stale_after_days: int = Field(default=30, ge=1, description="Days of inactivity before a skill is marked stale.")
|
||
archive_after_days: int = Field(default=90, ge=1, description="Days of inactivity before a stale skill is archived.")
|
||
review_timeout_seconds: int = Field(default=600, ge=1, description="Overall timeout (seconds) for the LLM review pass.")
|
||
model_name: str | None = Field(default=None, description="Model name for the curator LLM review pass; null = primary model.")
|
||
|
||
|
||
class CuratorConfigResponse(BaseModel):
|
||
"""Response model exposing the effective skill_evolution config."""
|
||
|
||
enabled: bool
|
||
curator: CuratorSubConfig
|
||
|
||
|
||
class CuratorConfigRequest(BaseModel):
|
||
"""Payload for updating user-facing curator settings in skill_evolution.json.
|
||
|
||
All other settings (enabled, timeouts) are global-only and controlled via config.yaml.
|
||
"""
|
||
|
||
interval_hours: int = Field(..., ge=1, description="Hours between curator runs.")
|
||
stale_after_days: int = Field(..., ge=1, description="Days of inactivity before a skill is marked stale.")
|
||
archive_after_days: int = Field(..., ge=1, description="Days of inactivity before a stale skill is archived.")
|
||
model_name: str | None = Field(default=None, description="Model for the curator LLM review; null = primary model.")
|
||
|
||
|
||
class CuratorLogRecord(BaseModel):
|
||
ts: str
|
||
level: str
|
||
message: str
|
||
|
||
|
||
class CuratorLogsResponse(BaseModel):
|
||
records: list[CuratorLogRecord]
|
||
run_count: int
|
||
last_run_at: str | None
|
||
|
||
|
||
def _build_curator_config_response(evo) -> CuratorConfigResponse:
|
||
"""Pure helper so both the GET route and the post-PUT refresh share one shape."""
|
||
return CuratorConfigResponse(
|
||
enabled=evo.enabled,
|
||
curator=CuratorSubConfig(
|
||
enabled=evo.curator.enabled,
|
||
interval_hours=evo.curator.interval_hours,
|
||
stale_after_days=evo.curator.stale_after_days,
|
||
archive_after_days=evo.curator.archive_after_days,
|
||
review_timeout_seconds=evo.curator.review_timeout_seconds,
|
||
model_name=evo.curator.model_name,
|
||
),
|
||
)
|
||
|
||
|
||
@router.get("/curator/config", response_model=CuratorConfigResponse, summary="Get Skill Evolution Config")
|
||
async def get_curator_config(request: Request) -> CuratorConfigResponse:
|
||
"""Return the current curator config for the authenticated user (yaml defaults + per-user overrides)."""
|
||
try:
|
||
from deerflow.config.skill_evolution_runtime import load_skill_evolution_config_for_user
|
||
|
||
evo = load_skill_evolution_config_for_user(_current_user_id(request))
|
||
return _build_curator_config_response(evo)
|
||
except Exception as e:
|
||
logger.error("Failed to read curator config: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to read curator config: {str(e)}")
|
||
|
||
|
||
@router.put("/curator/config", response_model=CuratorConfigResponse, summary="Update Skill Evolution Config")
|
||
async def update_curator_config(body: CuratorConfigRequest, request: Request) -> CuratorConfigResponse:
|
||
"""Persist the three user-facing curator thresholds to the per-user ``skill_evolution.json``.
|
||
|
||
Other settings (enabled, model, timeouts) are global and controlled via config.yaml.
|
||
"""
|
||
try:
|
||
from deerflow.config.skill_evolution_runtime import load_skill_evolution_config_for_user, save_user_overrides
|
||
|
||
user_id = _current_user_id(request)
|
||
save_user_overrides(user_id, {
|
||
"curator": {
|
||
"interval_hours": body.interval_hours,
|
||
"stale_after_days": body.stale_after_days,
|
||
"archive_after_days": body.archive_after_days,
|
||
"model_name": body.model_name,
|
||
},
|
||
})
|
||
return _build_curator_config_response(load_skill_evolution_config_for_user(user_id))
|
||
except HTTPException:
|
||
raise
|
||
except Exception as e:
|
||
logger.error("Failed to update curator config: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to update curator config: {str(e)}")
|
||
|
||
|
||
@router.get("/curator/logs", response_model=CuratorLogsResponse, summary="Recent Curator Logs")
|
||
async def get_curator_logs(request: Request) -> CuratorLogsResponse:
|
||
"""Return in-memory log records from the most recent curator run for the
|
||
authenticated user.
|
||
|
||
The buffer is process-local and cleared when curator starts a new run or
|
||
when the backend restarts. Records are ordered oldest → newest.
|
||
"""
|
||
try:
|
||
from deerflow.skills import curator_log
|
||
from deerflow.skills.curator import get_curator_status
|
||
|
||
user_id = _current_user_id(request)
|
||
records = curator_log.get_recent(user_id)
|
||
status = get_curator_status(user_id)
|
||
return CuratorLogsResponse(
|
||
records=[CuratorLogRecord(**r) for r in records],
|
||
run_count=int(status.get("run_count", 0)),
|
||
last_run_at=status.get("last_run_at"),
|
||
)
|
||
except Exception as e:
|
||
logger.error("Failed to read curator logs: %s", e, exc_info=True)
|
||
raise HTTPException(status_code=500, detail=f"Failed to read curator logs: {str(e)}")
|