/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/api/tasks.py
402 строки
13 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""API endpoints для управления пользовательскими задачами. Использует TaskService (CRUD + статусы + runs) и schemas/task.py. Эндпоинты: - GET /api/v1/tasks — список задач - POST /api/v1/tasks — создать задачу - GET /api/v1/tasks/stats — статистика - GET /api/v1/tasks/overdue — просроченные задачи - GET /api/v1/tasks/{id} — получить задачу - PATCH /api/v1/tasks/{id} — обновить задачу - DELETE /api/v1/tasks/{id} — удалить задачу - POST /api/v1/tasks/{id}/start — в работу - POST /api/v1/tasks/{id}/complete — завершить - POST /api/v1/tasks/{id}/cancel — отменить - GET /api/v1/tasks/{id}/runs — история выполнений - POST /api/v1/tasks/{id}/runs — добавить run """ from __future__ import annotations import uuid from datetime import UTC, datetime from typing import Any import structlog from fastapi import APIRouter, Depends, HTTPException, Query, status from pydantic import BaseModel, Field from src.api.dependencies import get_task_repo, get_task_service, get_user_ctx from src.db.models import Task as TaskModel from src.db.repositories import TaskRepository from src.middleware.auth import UserContext from src.schemas.task import ( TaskCreate, TaskResponse, TaskRunResponse, TaskUpdate, ) from src.services import TaskService logger = structlog.get_logger() router = APIRouter(prefix="/api/v1/tasks", tags=["tasks"]) # ============================================================================ # REQUEST/RESPONSE MODELS (специфичные для API) # ============================================================================ class TaskStatsResponse(BaseModel): """Статистика по задачам.""" total_tasks: int by_status: dict[str, int] by_priority: dict[str, int] overdue_count: int completed_today: int class TaskRunCreate(BaseModel): """Запрос на добавление записи о выполнении.""" status: str = Field(default="running", description="Статус run") output: str = Field(default="", description="Результат выполнения") tokens_used: int | None = Field(default=None, ge=0) error_message: str | None = None # ============================================================================ # HELPER FUNCTIONS # ============================================================================ def _utc_now() -> datetime: """Timezone-aware UTC now.""" return datetime.now(UTC) def _build_task_response(task: TaskModel) -> TaskResponse: """Собрать TaskResponse из DB-модели (через from_attributes).""" return TaskResponse.model_validate(task) # ============================================================================ # CRUD ENDPOINTS # ============================================================================ @router.get("", response_model=list[TaskResponse]) async def list_tasks( status_filter: str | None = Query(None, alias="status", description="Фильтр по статусу"), priority: str | None = Query(None, description="Фильтр по приоритету"), assigned_agent_id: str | None = Query(None, description="Фильтр по агенту"), limit: int = Query(100, ge=1, le=500), offset: int = Query(0, ge=0), user_ctx: UserContext = Depends(get_user_ctx), task_repo: TaskRepository = Depends(get_task_repo), ) -> list[TaskResponse]: """ Список задач с фильтрами и пагинацией. Note: используется task_repo напрямую для поддержки фильтра assigned_agent_id (TaskService.list его не поддерживает). """ tasks = await task_repo.list( status=status_filter, priority=priority, assigned_agent_id=assigned_agent_id, limit=limit, offset=offset, ) logger.info( "tasks.listed", count=len(tasks), status=status_filter, priority=priority, assigned_agent_id=assigned_agent_id, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return [_build_task_response(task) for task in tasks] @router.post("", response_model=TaskResponse, status_code=status.HTTP_201_CREATED) async def create_task( data: TaskCreate, user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> TaskResponse: """Создать новую задачу (статус по умолчанию — todo).""" task = await task_service.create(data) logger.info( "task.created", task_id=str(task.id), title=task.title, priority=task.priority, user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return _build_task_response(task) @router.get("/stats", response_model=TaskStatsResponse) async def get_tasks_stats( user_ctx: UserContext = Depends(get_user_ctx), task_repo: TaskRepository = Depends(get_task_repo), ) -> TaskStatsResponse: """ Статистика по задачам. TODO: оптимизировать через SQL-агрегации в repository. """ all_tasks = await task_repo.list(limit=10_000) by_status: dict[str, int] = {} by_priority: dict[str, int] = {} completed_today = 0 now = _utc_now() today_start = now.replace(hour=0, minute=0, second=0, microsecond=0) for task in all_tasks: by_status[task.status] = by_status.get(task.status, 0) + 1 by_priority[task.priority] = by_priority.get(task.priority, 0) + 1 if task.completed_at and task.completed_at >= today_start: completed_today += 1 overdue_tasks = await task_repo.get_overdue() logger.info( "tasks.stats", total=len(all_tasks), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return TaskStatsResponse( total_tasks=len(all_tasks), by_status=by_status, by_priority=by_priority, overdue_count=len(overdue_tasks), completed_today=completed_today, ) @router.get("/overdue", response_model=list[TaskResponse]) async def list_overdue_tasks( user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> list[TaskResponse]: """Просроченные задачи (deadline прошёл, статус не done/cancelled).""" tasks = await task_service.get_overdue() logger.info( "tasks.overdue.listed", count=len(tasks), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) return [_build_task_response(task) for task in tasks] @router.get("/{task_id}", response_model=TaskResponse) async def get_task( task_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> TaskResponse: """Получить задачу по ID.""" task = await task_service.get(task_id) if task is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Task {task_id} not found", ) return _build_task_response(task) @router.patch("/{task_id}", response_model=TaskResponse) async def update_task( task_id: uuid.UUID, data: TaskUpdate, user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> TaskResponse: """ Обновить задачу (частичное обновление). Автоматически проставляет completed_at при переводе в статус done. """ # Автоматически проставляем completed_at при статусе done if data.status == "done" and data.completed_at is None: data.completed_at = _utc_now() task = await task_service.update(task_id, data) if task is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Task {task_id} not found", ) logger.info( "task.updated", task_id=str(task_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_task_response(task) @router.delete("/{task_id}", status_code=status.HTTP_204_NO_CONTENT) async def delete_task( task_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> None: """Удалить задачу.""" deleted = await task_service.delete(task_id) if not deleted: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Task {task_id} not found", ) logger.info( "task.deleted", task_id=str(task_id), user_id=str(user_ctx.user_id), workspace_id=user_ctx.workspace_id, ) # ============================================================================ # STATUS TRANSITION ENDPOINTS # ============================================================================ @router.post("/{task_id}/start", response_model=TaskResponse) async def start_task( task_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> TaskResponse: """Перевести задачу в статус in_progress.""" task = await task_service.mark_in_progress(task_id) if task is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Task {task_id} not found", ) logger.info("task.started", task_id=str(task_id), user_id=str(user_ctx.user_id)) return _build_task_response(task) @router.post("/{task_id}/complete", response_model=TaskResponse) async def complete_task( task_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> TaskResponse: """Завершить задачу (статус done + completed_at).""" task = await task_service.mark_done(task_id) if task is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Task {task_id} not found", ) logger.info("task.completed", task_id=str(task_id), user_id=str(user_ctx.user_id)) return _build_task_response(task) @router.post("/{task_id}/cancel", response_model=TaskResponse) async def cancel_task( task_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> TaskResponse: """Отменить задачу (статус cancelled).""" task = await task_service.cancel(task_id) if task is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Task {task_id} not found", ) logger.info("task.cancelled", task_id=str(task_id), user_id=str(user_ctx.user_id)) return _build_task_response(task) # ============================================================================ # TASK RUNS ENDPOINTS # ============================================================================ @router.get("/{task_id}/runs", response_model=list[TaskRunResponse]) async def list_task_runs( task_id: uuid.UUID, user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> list[TaskRunResponse]: """История выполнений задачи (из JSONB Task.runs).""" task = await task_service.get(task_id) if task is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Task {task_id} not found", ) runs = await task_service.get_runs(task_id) return [TaskRunResponse(**run) for run in runs] @router.post("/{task_id}/runs", status_code=status.HTTP_201_CREATED) async def add_task_run( task_id: uuid.UUID, data: TaskRunCreate, user_ctx: UserContext = Depends(get_user_ctx), task_service: TaskService = Depends(get_task_service), ) -> dict[str, Any]: """ Добавить запись о выполнении задачи. Использует атомарный JSONB concat (защита от race condition). """ task = await task_service.add_run( task_id=task_id, status=data.status, output=data.output, tokens_used=data.tokens_used, error_message=data.error_message, ) if task is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=f"Task {task_id} not found", ) logger.info( "task.run.added", task_id=str(task_id), run_status=data.status, user_id=str(user_ctx.user_id), ) return { "status": "added", "task_id": str(task_id), "runs_count": len(task.runs or []), }