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 ```` 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)}")