/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/api/skills.py
380 строк
13 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""API endpoints для управления скиллами (переиспользуемыми промпт-шаблонами). Использует SkillService (CRUD + запуск: render prompt → LLM) и schemas/skill.py. Эндпоинты: - GET /api/v1/skills — список скиллов - POST /api/v1/skills — создать скилл - GET /api/v1/skills/stats — статистика - GET /api/v1/skills/public — публичные скиллы - GET /api/v1/skills/{id} — получить скилл - PATCH /api/v1/skills/{id} — обновить скилл - DELETE /api/v1/skills/{id} — удалить скилл - POST /api/v1/skills/{id}/use — инкрементировать usage - POST /api/v1/skills/{id}/run — запустить скилл (sync/stream) """ from __future__ import annotations import json import uuid from collections.abc import AsyncIterator from typing import Any import structlog from fastapi import APIRouter, Depends, HTTPException, Query, status from pydantic import BaseModel from sse_starlette.sse import EventSourceResponse from src.api.dependencies import get_skill_repo, get_skill_service, get_user_ctx from src.db.models import Skill as SkillModel from src.db.repositories import SkillRepository from src.middleware.auth import UserContext from src.schemas.skill import ( SkillCreate, SkillResponse, SkillRunRequest, SkillRunResponse, SkillUpdate, ) from src.services import SkillService logger = structlog.get_logger() router = APIRouter(prefix="/api/v1/skills", tags=["skills"]) # ============================================================================ # RESPONSE MODELS (специфичные для API) # ============================================================================ class SkillStatsResponse(BaseModel): """Агрегированная статистика по скиллам.""" total_skills: int public_skills: int private_skills: int archived_skills: int total_usage: int categories: dict[str, int] # ============================================================================ # HELPER FUNCTIONS # ============================================================================ def _build_skill_response(skill: SkillModel) -> SkillResponse: """Собрать SkillResponse из DB-модели (через from_attributes).""" return SkillResponse.model_validate(skill) # ============================================================================ # CRUD ENDPOINTS # ============================================================================ @router.get("", response_model=list[SkillResponse]) async def list_skills( category: str | None = Query(None, description="Фильтр по категории"), is_public: bool | None = Query(None, description="Фильтр по публичности"), is_archived: bool = Query(False, description="Включать архивные"), limit: int = Query(100, ge=1, le=500), offset: int = Query(0, ge=0), user_ctx: UserContext = Depends(get_user_ctx), skill_repo: SkillRepository = Depends(get_skill_repo), ) -> list[SkillResponse]: """Список скиллов с фильтрами и пагинацией.""" skills = await skill_repo.list( category=category, is_public=is_public, is_archived=is_archived, limit=limit, offset=offset, ) logger.info( "skills.listed", count=len(skills), category=category, is_public=is_public, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return [_build_skill_response(skill) for skill in skills] @router.post("", response_model=SkillResponse, status_code=status.HTTP_201_CREATED) async def create_skill( data: SkillCreate, user_ctx: UserContext = Depends(get_user_ctx), skill_service: SkillService = Depends(get_skill_service), ) -> SkillResponse: """Создать новый скилл.""" skill = await skill_service.create(data) logger.info( "skill.created", skill_id=str(skill.id), name=skill.name, category=skill.category, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return _build_skill_response(skill) @router.get("/stats", response_model=SkillStatsResponse) async def get_skills_stats( user_ctx: UserContext = Depends(get_user_ctx), skill_repo: SkillRepository = Depends(get_skill_repo), ) -> SkillStatsResponse: """ Агрегированная статистика по скиллам. TODO: оптимизировать через SQL-агрегации в repository. """ all_skills = await skill_repo.list(limit=10_000) total = len(all_skills) public = sum(1 for s in all_skills if s.is_public and not s.is_archived) private = sum(1 for s in all_skills if not s.is_public and not s.is_archived) archived = sum(1 for s in all_skills if s.is_archived) total_usage = sum(s.usage_count for s in all_skills) categories: dict[str, int] = {} for s in all_skills: if not s.is_archived: categories[s.category] = categories.get(s.category, 0) + 1 logger.info( "skills.stats", total=total, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return SkillStatsResponse( total_skills=total, public_skills=public, private_skills=private, archived_skills=archived, total_usage=total_usage, categories=categories, ) @router.get("/public", response_model=list[SkillResponse]) async def list_public_skills( category: str | None = Query(None, description="Фильтр по категории"), limit: int = Query(100, ge=1, le=500), offset: int = Query(0, ge=0), skill_repo: SkillRepository = Depends(get_skill_repo), ) -> list[SkillResponse]: """Публичные скиллы из всех workspace (не требует изоляции).""" skills = await skill_repo.get_public_skills(limit=limit, offset=offset) if category: skills = [s for s in skills if s.category == category] logger.info("skills.public.listed", count=len(skills), category=category) return [_build_skill_response(skill) for skill in skills] @router.get("/{skill_id}", response_model=SkillResponse) async def get_skill( skill_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), skill_service: SkillService = Depends(get_skill_service), ) -> SkillResponse: """Получить скилл по ID.""" skill = await skill_service.get(skill_id) if skill is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Skill {skill_id} not found", ) return _build_skill_response(skill) @router.patch("/{skill_id}", response_model=SkillResponse) async def update_skill( skill_id: uuid.UUID, data: SkillUpdate, user_ctx: UserContext = Depends(get_user_ctx), skill_service: SkillService = Depends(get_skill_service), ) -> SkillResponse: """Обновить скилл (частичное обновление).""" skill = await skill_service.update(skill_id, data) if skill is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Skill {skill_id} not found", ) logger.info( "skill.updated", skill_id=str(skill_id), updated_fields=list(data.model_dump(exclude_unset=True).keys()), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return _build_skill_response(skill) @router.delete("/{skill_id}", status_code=status.HTTP_204_NO_CONTENT) async def delete_skill( skill_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), skill_service: SkillService = Depends(get_skill_service), ) -> None: """Удалить скилл.""" deleted = await skill_service.delete(skill_id) if not deleted: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Skill {skill_id} not found", ) logger.info( "skill.deleted", skill_id=str(skill_id), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) # ============================================================================ # USAGE ENDPOINT # ============================================================================ @router.post("/{skill_id}/use") async def increment_skill_usage( skill_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), skill_repo: SkillRepository = Depends(get_skill_repo), ) -> dict[str, Any]: """Инкрементировать счётчик использования скилла.""" skill = await skill_repo.get(skill_id) if skill is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Skill {skill_id} not found", ) await skill_repo.increment_usage(skill_id) logger.info( "skill.usage.incremented", skill_id=str(skill_id), new_usage_count=skill.usage_count + 1, user_id=str(user_ctx.user_id), ) return { "status": "incremented", "skill_id": str(skill_id), "new_usage_count": skill.usage_count + 1, } # ============================================================================ # RUN ENDPOINT (render prompt → LLM) # ============================================================================ # ✅ response_model=None — функция возвращает Union (SkillRunResponse | EventSourceResponse), # FastAPI не может вывести единую Pydantic-схему из EventSourceResponse @router.post("/{skill_id}/run", response_model=None) async def run_skill( skill_id: uuid.UUID, data: SkillRunRequest, user_ctx: UserContext = Depends(get_user_ctx), skill_service: SkillService = Depends(get_skill_service), ) -> SkillRunResponse | EventSourceResponse: """ Запустить скилл: валидация параметров → render prompt → LLM. Два режима: - **stream=true**: SSE поток (reasoning, content, done, error) - **stream=false**: JSON с полным ответом """ skill = await skill_service.get(skill_id) if skill is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Skill {skill_id} not found", ) logger.info( "skill.run.requested", skill_id=str(skill_id), stream=data.stream, parameters_keys=list(data.parameters.keys()), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) if data.stream: return EventSourceResponse( _stream_skill_run(skill_id, data.parameters, data.model, skill_service) ) # Non-streaming mode try: result = await skill_service.run(skill_id, data.parameters, data.model) except ValueError as e: raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(e)) return SkillRunResponse( skill_id=skill_id, output=result["output"], rendered_prompt=result.get("rendered_prompt"), tokens_used=result.get("tokens_used", 0), ) async def _stream_skill_run( skill_id: uuid.UUID, parameters: dict[str, Any], model: str | None, skill_service: SkillService, ) -> AsyncIterator[dict[str, str]]: """Генератор SSE событий для streaming запуска скилла.""" try: async for chunk in skill_service.run_stream(skill_id, parameters, model): event_type = chunk.get("type", "content") yield { "event": event_type, "data": json.dumps(chunk, ensure_ascii=False, default=str), } except ValueError as e: yield { "event": "error", "data": json.dumps({"type": "error", "error": str(e)}, ensure_ascii=False), } except Exception as e: logger.exception( "skill.run.stream.error", skill_id=str(skill_id), error=str(e), error_type=type(e).__name__, ) yield { "event": "error", "data": json.dumps( { "type": "error", "error": str(e), "error_type": type(e).__name__, }, ensure_ascii=False, ), }