/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/services/task_service.py
121 строка
4 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""Бизнес-логика задач (CRUD + статусы + runs).""" from __future__ import annotations import uuid from typing import Any import structlog from sqlalchemy.ext.asyncio import AsyncSession from src.db.repositories import TaskRepository from src.primitives import utc_now from src.schemas import TaskCreate, TaskUpdate logger = structlog.get_logger() class TaskService: """Сервис для работы с задачами.""" def __init__( self, session: AsyncSession, user_id: uuid.UUID, workspace_id: str, ) -> None: self.session = session self.user_id = user_id self.workspace_id = workspace_id self.task_repo = TaskRepository(session, user_id, workspace_id) # ======================================================================== # CRUD # ======================================================================== async def create(self, data: TaskCreate) -> Any: """Создать задачу.""" task = await self.task_repo.create( user_id=self.user_id, workspace_id=self.workspace_id, title=data.title, description=data.description, priority=data.priority, deadline=data.deadline, assigned_agent_id=str(data.assigned_agent_id) if data.assigned_agent_id else None, skill_id=data.skill_id, workflow_id=data.workflow_id, chat_id=data.chat_id, recurrence_pattern=data.recurrence_pattern, tags=data.tags, estimated_duration_minutes=data.estimated_duration_minutes, ) logger.info("task.created", task_id=str(task.id), title=data.title) return task async def get(self, task_id: uuid.UUID) -> Any | None: return await self.task_repo.get(task_id) async def list( self, status: str | None = None, priority: str | None = None, limit: int = 50, offset: int = 0, ) -> list[Any]: return await self.task_repo.list( status=status, priority=priority, limit=limit, offset=offset ) async def update(self, task_id: uuid.UUID, data: TaskUpdate) -> Any | None: updates = data.model_dump(exclude_unset=True) if not updates: return await self.task_repo.get(task_id) return await self.task_repo.update(task_id, **updates) async def delete(self, task_id: uuid.UUID) -> bool: return await self.task_repo.delete(task_id) # ======================================================================== # Status transitions # ======================================================================== async def mark_in_progress(self, task_id: uuid.UUID) -> Any | None: return await self.task_repo.update(task_id, status="in_progress") async def mark_done(self, task_id: uuid.UUID) -> Any | None: return await self.task_repo.update(task_id, status="done", completed_at=utc_now()) async def cancel(self, task_id: uuid.UUID) -> Any | None: return await self.task_repo.update(task_id, status="cancelled") # ======================================================================== # Runs (JSONB в Task.runs) # ======================================================================== async def add_run( self, task_id: uuid.UUID, status: str, output: str = "", tokens_used: int | None = None, error_message: str | None = None, ) -> Any | None: """Добавить run в историю (JSONB).""" run = { "status": status, "output": output, "tokens_used": tokens_used, "error_message": error_message, "created_at": utc_now().isoformat(), } return await self.task_repo.add_run(task_id, run) async def get_runs(self, task_id: uuid.UUID) -> list[dict[str, Any]]: """Получить историю выполнений (из Task.runs JSONB).""" task = await self.task_repo.get(task_id) return list(task.runs) if task and task.runs else [] async def get_overdue(self) -> list[Any]: """Получить просроченные задачи.""" return await self.task_repo.get_overdue()