/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/primitives/task.py
452 строки
15 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""Задачи (Tasks) — единицы выполнения в Engine и пользовательские задачи для БД. Использует: - ``src.primitives._time``: timezone-aware UTC, monotonic для duration - ``StrEnum`` (Python 3.11+) для статусов/приоритетов/типов Два вида задач: - ``Task`` — runtime (внутри Engine, для отслеживания шагов flow/агента) - ``UserTask`` — пользовательская (для БД и UI, с приоритетом/дедлайном/историей) """ from __future__ import annotations from dataclasses import dataclass, field from datetime import datetime from enum import StrEnum from typing import Any from uuid import uuid4 from src.primitives._time import ( elapsed_ms, monotonic, parse_datetime, to_iso, utc_now, ) from src.primitives._time import ( is_overdue as check_overdue, ) from src.primitives.state import TaskStatus # ============================================================================ # Runtime Task (для внутреннего использования в Engine) # ============================================================================ @dataclass class Task: """ Задача для выполнения (runtime, не для БД). Используется внутри Engine для отслеживания выполнения отдельных шагов flow/агента. """ id: str = field(default_factory=lambda: str(uuid4())) title: str = "" flow_id: str = "" flow_name: str = "" status: TaskStatus = TaskStatus.PENDING input_data: dict[str, Any] = field(default_factory=dict) output: dict[str, Any] = field(default_factory=dict) error: str | None = None # Метрики (monotonic — точное duration, не зависит от системных часов) _start_time: float = field(default_factory=monotonic, repr=False) _end_time: float | None = field(default=None, repr=False) @property def duration_ms(self) -> float: """Длительность выполнения в миллисекундах.""" if self._end_time is None: return elapsed_ms(self._start_time) return (self._end_time - self._start_time) * 1000 def start(self) -> None: """Начать выполнение.""" self.status = TaskStatus.RUNNING self._start_time = monotonic() def complete(self, output: dict[str, Any] | None = None) -> None: """Завершить успешно.""" self.status = TaskStatus.COMPLETED self._end_time = monotonic() if output is not None: self.output = output def fail(self, error: str) -> None: """Завершить с ошибкой.""" self.status = TaskStatus.FAILED self._end_time = monotonic() self.error = error def cancel(self) -> None: """Отменить.""" self.status = TaskStatus.CANCELLED self._end_time = monotonic() def is_running(self) -> bool: """Проверить, выполняется ли задача.""" return self.status == TaskStatus.RUNNING def is_completed(self) -> bool: """Проверить, завершена ли задача успешно.""" return self.status == TaskStatus.COMPLETED def is_failed(self) -> bool: """Проверить, завершилась ли задача с ошибкой.""" return self.status == TaskStatus.FAILED def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "title": self.title, "flow_id": self.flow_id, "flow_name": self.flow_name, "status": self.status.value, "input_data": self.input_data, "output": self.output, "error": self.error, "duration_ms": self.duration_ms, } # ============================================================================ # UserTask Enums (для пользовательских задач в БД) # ============================================================================ class TaskPriority(StrEnum): """Приоритет задачи.""" LOW = "low" MEDIUM = "medium" HIGH = "high" URGENT = "urgent" class UserTaskStatus(StrEnum): """Статус пользовательской задачи.""" TODO = "todo" IN_PROGRESS = "in_progress" DONE = "done" CANCELLED = "cancelled" class TaskType(StrEnum): """Тип задачи.""" SIMPLE = "simple" AGENT = "agent" RECURRING = "recurring" class TaskRunStatus(StrEnum): """Статус выполнения задачи (run).""" RUNNING = "running" COMPLETED = "completed" FAILED = "failed" # ============================================================================ # TaskRun (запись о выполнении) # ============================================================================ @dataclass class TaskRun: """ Запись о выполнении задачи. Содержит информацию о конкретном запуске задачи: статус, output, ошибки, метрики, timestamps. """ id: str = field(default_factory=lambda: str(uuid4())) task_id: str = "" status: TaskRunStatus = TaskRunStatus.COMPLETED output: str = "" started_at: datetime = field(default_factory=utc_now) completed_at: datetime | None = None duration_ms: int | None = None tokens_used: int | None = None error_message: str | None = None model_used: str | None = None def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "task_id": self.task_id, "status": self.status.value, "output": self.output, "started_at": to_iso(self.started_at), "completed_at": to_iso(self.completed_at), "duration_ms": self.duration_ms, "tokens_used": self.tokens_used, "error_message": self.error_message, "model_used": self.model_used, } @classmethod def from_dict(cls, data: dict[str, Any]) -> TaskRun: """Создать TaskRun из словаря.""" return cls( id=data.get("id", str(uuid4())), task_id=data.get("task_id", ""), status=TaskRunStatus(data.get("status", "completed")), output=data.get("output", ""), started_at=parse_datetime(data.get("started_at")) or utc_now(), completed_at=parse_datetime(data.get("completed_at")), duration_ms=data.get("duration_ms"), tokens_used=data.get("tokens_used"), error_message=data.get("error_message"), model_used=data.get("model_used"), ) # ============================================================================ # UserTask (для пользовательских задач в БД) # ============================================================================ @dataclass class UserTask: """ Пользовательская задача (для БД и UI). Отличается от runtime Task: - Имеет приоритет, дедлайн, теги - Может быть назначена агенту - Может использовать skill или workflow - Имеет историю выполнения (runs) - Может быть повторяющейся (cron) Вычисления (is_overdue, get_latest_run) — чистые функции во время render. """ id: str = field(default_factory=lambda: str(uuid4())) workspace_id: str = "default" user_id: str = "" # Основные поля title: str = "" description: str = "" status: UserTaskStatus = UserTaskStatus.TODO priority: TaskPriority = TaskPriority.MEDIUM type: TaskType = TaskType.SIMPLE # Опциональные поля deadline: datetime | None = None assigned_agent_id: str | None = None skill_id: str | None = None workflow_id: str | None = None chat_id: str | None = None # Повторяющиеся задачи recurrence_pattern: str | None = None # cron expression # Метаданные tags: list[str] = field(default_factory=list) estimated_duration_minutes: int | None = None actual_duration_minutes: int | None = None # История выполнения runs: list[TaskRun] = field(default_factory=list) # Timestamps created_at: datetime = field(default_factory=utc_now) updated_at: datetime = field(default_factory=utc_now) completed_at: datetime | None = None def add_run(self, run: TaskRun) -> None: """Добавить run в историю.""" self.runs.append(run) def get_latest_run(self) -> TaskRun | None: """Получить последний run.""" return self.runs[-1] if self.runs else None def mark_as_done(self) -> None: """Отметить как выполненную.""" self.status = UserTaskStatus.DONE self.completed_at = utc_now() self.updated_at = utc_now() def mark_as_in_progress(self) -> None: """Отметить как выполняемую.""" self.status = UserTaskStatus.IN_PROGRESS self.updated_at = utc_now() def cancel(self) -> None: """Отменить задачу.""" self.status = UserTaskStatus.CANCELLED self.updated_at = utc_now() def is_overdue(self) -> bool: """ Проверить, просрочена ли задача. Чистая функция — вычисляется во время render. Задача просрочена, если дедлайн прошёл И она не завершена. """ return check_overdue(self.deadline, is_done=self.status == UserTaskStatus.DONE) def has_agent(self) -> bool: """Проверить, назначен ли агент.""" return self.assigned_agent_id is not None def has_skill(self) -> bool: """Проверить, привязан ли skill.""" return self.skill_id is not None def has_workflow(self) -> bool: """Проверить, привязан ли workflow.""" return self.workflow_id is not None def is_recurring(self) -> bool: """Проверить, является ли задача повторяющейся.""" return self.type == TaskType.RECURRING and self.recurrence_pattern is not None def add_tag(self, tag: str) -> None: """Добавить тег.""" if tag not in self.tags: self.tags.append(tag) def remove_tag(self, tag: str) -> None: """Удалить тег.""" if tag in self.tags: self.tags.remove(tag) def has_tag(self, tag: str) -> bool: """Проверить наличие тега.""" return tag in self.tags def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "workspace_id": self.workspace_id, "user_id": self.user_id, "title": self.title, "description": self.description, "status": self.status.value, "priority": self.priority.value, "type": self.type.value, "deadline": to_iso(self.deadline), "assigned_agent_id": self.assigned_agent_id, "skill_id": self.skill_id, "workflow_id": self.workflow_id, "chat_id": self.chat_id, "recurrence_pattern": self.recurrence_pattern, "tags": self.tags, "estimated_duration_minutes": self.estimated_duration_minutes, "actual_duration_minutes": self.actual_duration_minutes, "runs": [run.to_dict() for run in self.runs], "created_at": to_iso(self.created_at), "updated_at": to_iso(self.updated_at), "completed_at": to_iso(self.completed_at), } @classmethod def from_dict(cls, data: dict[str, Any]) -> UserTask: """Создать UserTask из словаря.""" return cls( id=data.get("id", str(uuid4())), workspace_id=data.get("workspace_id", "default"), user_id=data.get("user_id", ""), title=data.get("title", ""), description=data.get("description", ""), status=UserTaskStatus(data.get("status", "todo")), priority=TaskPriority(data.get("priority", "medium")), type=TaskType(data.get("type", "simple")), deadline=parse_datetime(data.get("deadline")), assigned_agent_id=data.get("assigned_agent_id"), skill_id=data.get("skill_id"), workflow_id=data.get("workflow_id"), chat_id=data.get("chat_id"), recurrence_pattern=data.get("recurrence_pattern"), tags=data.get("tags", []), estimated_duration_minutes=data.get("estimated_duration_minutes"), actual_duration_minutes=data.get("actual_duration_minutes"), runs=[TaskRun.from_dict(run) for run in data.get("runs", [])], created_at=parse_datetime(data.get("created_at")) or utc_now(), updated_at=parse_datetime(data.get("updated_at")) or utc_now(), completed_at=parse_datetime(data.get("completed_at")), ) # ============================================================================ # Helper Functions # ============================================================================ def create_task( title: str, flow_id: str = "", flow_name: str = "", input_data: dict[str, Any] | None = None, ) -> Task: """Создать runtime задачу.""" return Task( title=title, flow_id=flow_id, flow_name=flow_name, input_data=input_data or {}, ) def create_user_task( title: str, workspace_id: str = "default", user_id: str = "", description: str = "", priority: TaskPriority = TaskPriority.MEDIUM, **kwargs: Any, ) -> UserTask: """Создать пользовательскую задачу.""" return UserTask( workspace_id=workspace_id, user_id=user_id, title=title, description=description, priority=priority, **kwargs, ) def create_task_run( task_id: str, status: TaskRunStatus | str = TaskRunStatus.COMPLETED, output: str = "", **kwargs: Any, ) -> TaskRun: """Создать запись о выполнении задачи.""" return TaskRun( task_id=task_id, status=TaskRunStatus(status), output=output, **kwargs, ) # ============================================================================ # Exports # ============================================================================ __all__ = [ "Task", "TaskPriority", "TaskRun", "TaskRunStatus", "TaskType", "UserTask", "UserTaskStatus", "create_task", "create_task_run", "create_user_task", ]