/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/primitives/context.py
493 строки
17 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""Контекст выполнения — пробрасывается через весь pipeline. Использует: - ``src.primitives._time``: timezone-aware UTC (``utc_now``) и monotonic-время для точного измерения duration (не зависит от системных часов). - Pydantic v2 (``ConfigDict``). Контексты: - ``RunContext`` — базовый (идентификация, счётчики, budget, тайминги) - ``SkillContext`` — выполнение скилла - ``TaskContext`` — выполнение задачи - ``MCPContext`` — MCP вызовы - ``AgentContext`` — выполнение агента """ from __future__ import annotations from collections.abc import Callable from datetime import datetime from enum import Enum from typing import Any from uuid import uuid4 from pydantic import BaseModel, ConfigDict, Field from src.primitives._time import monotonic, to_iso, utc_now # ============================================================================ # Cancellation Token # ============================================================================ class CancellationToken: """ Токен для отмены длительных операций. Используется для graceful остановки выполнения flow/агента/skill. Передаётся через весь pipeline и проверяется на каждой итерации. Это НЕ Pydantic модель — обычный класс с приватными полями. """ def __init__(self) -> None: self._cancelled: bool = False self._callbacks: list[Callable[[], None]] = [] def cancel(self) -> None: """Отменить операцию (вызывает все зарегистрированные callback).""" self._cancelled = True for callback in self._callbacks: try: callback() except Exception: # Callback не должен ломать отмену — игнорируем ошибки pass def is_cancelled(self) -> bool: """Проверить, отменена ли операция.""" return self._cancelled def check(self) -> None: """ Проверить отмену и выбросить исключение, если отменено. Raises: CancelledError: если операция отменена """ if self._cancelled: raise CancelledError("Operation was cancelled") def on_cancel(self, callback: Callable[[], None]) -> None: """Зарегистрировать callback для вызова при отмене.""" self._callbacks.append(callback) class CancelledError(Exception): """Исключение для отменённых операций.""" # ============================================================================ # Run Status # ============================================================================ class RunStatus(str, Enum): """Статус выполнения.""" PENDING = "pending" RUNNING = "running" PAUSED = "paused" COMPLETED = "completed" FAILED = "failed" CANCELLED = "cancelled" # ============================================================================ # Budget Limits # ============================================================================ class BudgetLimits(BaseModel): """ Лимиты бюджета для выполнения. Используется для предотвращения бесконечных циклов и контроля расходов. """ max_tokens: int = Field(default=100_000, description="Максимум токенов") max_llm_calls: int = Field(default=50, description="Максимум LLM вызовов") max_tool_calls: int = Field(default=100, description="Максимум tool вызовов") max_duration_seconds: int = Field( default=300, description="Максимальная длительность в секундах", ) def check( self, tokens_used: int, llm_calls: int, tool_calls: int, duration_seconds: float, ) -> str | None: """ Проверить, не превышены ли лимиты. Returns: None если всё ок, строка с описанием нарушения если превышено. """ if tokens_used >= self.max_tokens: return f"Token limit exceeded: {tokens_used}/{self.max_tokens}" if llm_calls >= self.max_llm_calls: return f"LLM call limit exceeded: {llm_calls}/{self.max_llm_calls}" if tool_calls >= self.max_tool_calls: return f"Tool call limit exceeded: {tool_calls}/{self.max_tool_calls}" if duration_seconds >= self.max_duration_seconds: return f"Duration limit exceeded: {duration_seconds:.1f}/{self.max_duration_seconds}s" return None # ============================================================================ # Base Run Context # ============================================================================ class RunContext(BaseModel): """ Базовый контекст выполнения flow/агента/skill. Пробрасывается через весь pipeline и содержит: - Идентификацию запуска - Счётчики ресурсов - Тайминги - Budget limits - Metadata - Parent context (для вложенных вызовов) Note: ``start_time`` использует ``time.monotonic()`` (точное duration, не зависит от системных часов) и исключён из сериализации. """ model_config = ConfigDict(arbitrary_types_allowed=True) # Идентификация run_id: str = Field(default_factory=lambda: str(uuid4())) workspace_id: str = "default" session_id: str | None = None user_id: str | None = None # Metadata metadata: dict[str, Any] = Field(default_factory=dict) # Счётчики ресурсов tokens_used: int = 0 llm_calls: int = 0 tool_calls: int = 0 # Тайминги started_at: datetime = Field(default_factory=utc_now) # start_time — monotonic, для точного измерения duration (exclude из JSON) start_time: float = Field(default_factory=monotonic, exclude=True) # Budget budget_limits: BudgetLimits = Field(default_factory=BudgetLimits) # Статус status: RunStatus = RunStatus.PENDING # Parent context (для вложенных вызовов) parent_run_id: str | None = None # ======================================================================== # Resource counters # ======================================================================== def add_tokens(self, count: int) -> None: """Добавить использованные токены.""" self.tokens_used += count def add_llm_call(self) -> None: """Добавить LLM вызов.""" self.llm_calls += 1 def add_tool_call(self) -> None: """Добавить tool вызов.""" self.tool_calls += 1 # ======================================================================== # Timing # ======================================================================== def get_duration_seconds(self) -> float: """Получить длительность выполнения в секундах (monotonic).""" return monotonic() - self.start_time def get_duration_ms(self) -> int: """Получить длительность выполнения в миллисекундах.""" return int(self.get_duration_seconds() * 1000) # ======================================================================== # Budget # ======================================================================== def check_budget(self) -> str | None: """ Проверить budget limits. Returns: None если всё ок, строка с описанием нарушения если превышено. """ return self.budget_limits.check( tokens_used=self.tokens_used, llm_calls=self.llm_calls, tool_calls=self.tool_calls, duration_seconds=self.get_duration_seconds(), ) # ======================================================================== # Status & Metadata # ======================================================================== def set_status(self, status: RunStatus) -> None: """Установить статус выполнения.""" self.status = status def get_metadata(self, key: str, default: Any = None) -> Any: """Получить значение из metadata.""" return self.metadata.get(key, default) def set_metadata(self, key: str, value: Any) -> None: """Установить значение в metadata.""" self.metadata[key] = value # ======================================================================== # Child context # ======================================================================== def create_child_context(self) -> RunContext: """ Создать дочерний контекст для вложенных вызовов. Budget limits и metadata копируются (не разделяются с parent). """ return RunContext( workspace_id=self.workspace_id, session_id=self.session_id, user_id=self.user_id, parent_run_id=self.run_id, budget_limits=self.budget_limits.model_copy(), metadata=self.metadata.copy(), ) # ======================================================================== # Serialization # ======================================================================== def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "run_id": self.run_id, "workspace_id": self.workspace_id, "session_id": self.session_id, "user_id": self.user_id, "tokens_used": self.tokens_used, "llm_calls": self.llm_calls, "tool_calls": self.tool_calls, "started_at": to_iso(self.started_at), "duration_seconds": self.get_duration_seconds(), "duration_ms": self.get_duration_ms(), "status": self.status.value, "parent_run_id": self.parent_run_id, "metadata": self.metadata, } # ============================================================================ # Skill Context # ============================================================================ class SkillContext(RunContext): """ Контекст для выполнения скилла. Расширяет RunContext специфичными для скилла полями: skill_id, skill_name, parameters, rendered_prompt, model, agent_id. """ skill_id: str | None = None skill_name: str | None = None parameters: dict[str, Any] = Field(default_factory=dict) rendered_prompt: str | None = None model: str = "deepseek-ai/DeepSeek-V4-Pro" agent_id: str | None = None def set_rendered_prompt(self, prompt: str) -> None: """Установить отрендеренный промпт.""" self.rendered_prompt = prompt def get_parameter(self, name: str, default: Any = None) -> Any: """Получить параметр по имени.""" return self.parameters.get(name, default) def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" data = super().to_dict() data.update( { "skill_id": self.skill_id, "skill_name": self.skill_name, "parameters": self.parameters, "rendered_prompt": self.rendered_prompt, "model": self.model, "agent_id": self.agent_id, } ) return data # ============================================================================ # Task Context # ============================================================================ class TaskContext(RunContext): """ Контекст для выполнения задачи. Расширяет RunContext специфичными для задачи полями: task_id, task_type, assigned_agent_id, skill_id, workflow_id. """ task_id: str | None = None task_type: str = "simple" # simple, agent, recurring assigned_agent_id: str | None = None skill_id: str | None = None workflow_id: str | None = None def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" data = super().to_dict() data.update( { "task_id": self.task_id, "task_type": self.task_type, "assigned_agent_id": self.assigned_agent_id, "skill_id": self.skill_id, "workflow_id": self.workflow_id, } ) return data # ============================================================================ # MCP Context # ============================================================================ class MCPContext(RunContext): """ Контекст для MCP вызовов. Расширяет RunContext специфичными для MCP полями: server_id, server_name, tool_name, parameters, timeout_seconds. """ server_id: str | None = None server_name: str | None = None tool_name: str | None = None parameters: dict[str, Any] = Field(default_factory=dict) timeout_seconds: int = 30 def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" data = super().to_dict() data.update( { "server_id": self.server_id, "server_name": self.server_name, "tool_name": self.tool_name, "parameters": self.parameters, "timeout_seconds": self.timeout_seconds, } ) return data # ============================================================================ # Agent Context # ============================================================================ class AgentContext(RunContext): """ Контекст для выполнения агента. Расширяет RunContext специфичными для агента полями: agent_id, agent_name, system_prompt, tools, memory. """ agent_id: str | None = None agent_name: str | None = None system_prompt: str | None = None tools: list[str] = Field(default_factory=list) memory: dict[str, Any] = Field(default_factory=dict) # ======================================================================== # Tools # ======================================================================== def add_tool(self, tool_name: str) -> None: """Добавить tool в список доступных.""" if tool_name not in self.tools: self.tools.append(tool_name) def remove_tool(self, tool_name: str) -> None: """Удалить tool из списка доступных.""" if tool_name in self.tools: self.tools.remove(tool_name) def has_tool(self, tool_name: str) -> bool: """Проверить, доступен ли tool.""" return tool_name in self.tools # ======================================================================== # Memory # ======================================================================== def get_memory(self, key: str, default: Any = None) -> Any: """Получить значение из памяти.""" return self.memory.get(key, default) def set_memory(self, key: str, value: Any) -> None: """Установить значение в память.""" self.memory[key] = value def clear_memory(self) -> None: """Очистить память.""" self.memory.clear() # ======================================================================== # Serialization # ======================================================================== def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" data = super().to_dict() data.update( { "agent_id": self.agent_id, "agent_name": self.agent_name, "system_prompt": self.system_prompt, "tools": self.tools, "memory": self.memory, } ) return data # ============================================================================ # Exports # ============================================================================ __all__ = [ "AgentContext", "BudgetLimits", "CancellationToken", "CancelledError", "MCPContext", "RunContext", "RunStatus", "SkillContext", "TaskContext", ]