/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/primitives/flow.py
806 строк
28 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""Flow — композиция агентов для выполнения задачи. Использует: - ``src.primitives._safe_eval``: безопасные conditional edges (замена eval) - ``src.primitives._time``: timezone-aware UTC и monotonic для duration """ from __future__ import annotations from collections.abc import AsyncIterator, Callable from dataclasses import dataclass, field from datetime import datetime from enum import StrEnum from typing import Any from pydantic import Field from src.primitives._safe_eval import is_safe_expression, safe_eval_bool from src.primitives._time import elapsed_ms, monotonic, to_iso, utc_now from src.primitives.context import RunContext # ============================================================================ # Flow Node Types # ============================================================================ class FlowNodeType(StrEnum): """Типы узлов в графе flow.""" START = "start" END = "end" AGENT = "agent" LLM = "llm" TOOL = "tool" CONDITION = "condition" PARALLEL = "parallel" LOOP = "loop" SUBFLOW = "subflow" TRANSFORM = "transform" WAIT = "wait" # ============================================================================ # Flow Edge Types # ============================================================================ class FlowEdgeType(StrEnum): """Типы ребер в графе flow.""" NORMAL = "normal" CONDITIONAL = "conditional" ERROR = "error" DEFAULT = "default" # ============================================================================ # Flow Status # ============================================================================ class FlowStatus(StrEnum): """Статус выполнения flow.""" PENDING = "pending" RUNNING = "running" PAUSED = "paused" COMPLETED = "completed" FAILED = "failed" CANCELLED = "cancelled" # ============================================================================ # Flow Node # ============================================================================ @dataclass class FlowNode: """ Узел в графе flow. Представляет шаг выполнения: агент, LLM вызов, tool, условие и т.д. """ id: str type: FlowNodeType label: str = "" data: dict[str, Any] = field(default_factory=dict) position: dict[str, float] = field(default_factory=lambda: {"x": 0.0, "y": 0.0}) # Дополнительные поля config: dict[str, Any] = field(default_factory=dict) timeout_seconds: int | None = None retry_count: int = 0 retry_delay_seconds: int = 1 def get_config(self, key: str, default: Any = None) -> Any: """Получить значение из config.""" return self.config.get(key, default) def set_config(self, key: str, value: Any) -> None: """Установить значение в config.""" self.config[key] = value def is_executable(self) -> bool: """Проверить, является ли узел исполняемым.""" return self.type in [ FlowNodeType.AGENT, FlowNodeType.LLM, FlowNodeType.TOOL, FlowNodeType.SUBFLOW, FlowNodeType.TRANSFORM, ] def is_control_flow(self) -> bool: """Проверить, является ли узел управляющим.""" return self.type in [ FlowNodeType.START, FlowNodeType.END, FlowNodeType.CONDITION, FlowNodeType.PARALLEL, FlowNodeType.LOOP, FlowNodeType.WAIT, ] def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "type": self.type.value, "label": self.label, "data": self.data, "position": self.position, "config": self.config, "timeout_seconds": self.timeout_seconds, "retry_count": self.retry_count, "retry_delay_seconds": self.retry_delay_seconds, } @classmethod def from_dict(cls, data: dict[str, Any]) -> FlowNode: """Создать из словаря.""" return cls( id=data["id"], type=FlowNodeType(data["type"]), label=data.get("label", ""), data=data.get("data", {}), position=data.get("position", {"x": 0.0, "y": 0.0}), config=data.get("config", {}), timeout_seconds=data.get("timeout_seconds"), retry_count=data.get("retry_count", 0), retry_delay_seconds=data.get("retry_delay_seconds", 1), ) # ============================================================================ # Flow Edge # ============================================================================ @dataclass class FlowEdge: """ Ребро в графе flow. Соединяет два узла и может содержать условия перехода. Условия вычисляются через безопасный ``safe_eval_bool`` (без eval). """ id: str source: str target: str label: str = "" data: dict[str, Any] = field(default_factory=dict) # Дополнительные поля edge_type: FlowEdgeType = FlowEdgeType.NORMAL condition: str | None = None # Безопасное выражение для conditional edges priority: int = 0 # для выбора при нескольких исходящих ребрах def is_conditional(self) -> bool: """Проверить, является ли ребро условным.""" return self.edge_type == FlowEdgeType.CONDITIONAL and self.condition is not None def evaluate_condition(self, context: dict[str, Any]) -> bool: """ Безопасно вычислить условие перехода. Использует ``safe_eval_bool`` (AST whitelist) — произвольный код выполнить невозможно. При ошибке возвращает False (fail-safe). Args: context: контекст выполнения с переменными Returns: True если условие выполнено. """ if not self.is_conditional(): return True if self.condition is None: return False return safe_eval_bool(self.condition, context) def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "source": self.source, "target": self.target, "label": self.label, "data": self.data, "edge_type": self.edge_type.value, "condition": self.condition, "priority": self.priority, } @classmethod def from_dict(cls, data: dict[str, Any]) -> FlowEdge: """Создать из словаря.""" return cls( id=data["id"], source=data["source"], target=data["target"], label=data.get("label", ""), data=data.get("data", {}), edge_type=FlowEdgeType(data.get("edge_type", "normal")), condition=data.get("condition"), priority=data.get("priority", 0), ) # ============================================================================ # Flow Step # ============================================================================ @dataclass class FlowStep: """ Шаг выполнения flow. Содержит информацию о выполненном узле. """ node_id: str node_type: FlowNodeType started_at: datetime = field(default_factory=utc_now) completed_at: datetime | None = None status: FlowStatus = FlowStatus.RUNNING output: dict[str, Any] = field(default_factory=dict) error: str | None = None duration_ms: int | None = None tokens_used: int = 0 llm_calls: int = 0 tool_calls: int = 0 _start_time: float = field(default_factory=monotonic, repr=False) def complete(self, output: dict[str, Any] | None = None) -> None: """Завершить шаг успешно.""" self.status = FlowStatus.COMPLETED self.completed_at = utc_now() self.duration_ms = int(elapsed_ms(self._start_time)) if output: self.output = output def fail(self, error: str) -> None: """Завершить шаг с ошибкой.""" self.status = FlowStatus.FAILED self.completed_at = utc_now() self.duration_ms = int(elapsed_ms(self._start_time)) self.error = error def cancel(self) -> None: """Отменить шаг.""" self.status = FlowStatus.CANCELLED self.completed_at = utc_now() self.duration_ms = int(elapsed_ms(self._start_time)) def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "node_id": self.node_id, "node_type": self.node_type.value, "started_at": to_iso(self.started_at), "completed_at": to_iso(self.completed_at), "status": self.status.value, "output": self.output, "error": self.error, "duration_ms": self.duration_ms, "tokens_used": self.tokens_used, "llm_calls": self.llm_calls, "tool_calls": self.tool_calls, } # ============================================================================ # Flow Result # ============================================================================ @dataclass class FlowResult: """ Результат выполнения flow. Содержит финальный output и метрики. """ flow_id: str status: FlowStatus output: dict[str, Any] = field(default_factory=dict) error: str | None = None steps: list[FlowStep] = field(default_factory=list) # Метрики total_duration_ms: int = 0 total_tokens_used: int = 0 total_llm_calls: int = 0 total_tool_calls: int = 0 started_at: datetime = field(default_factory=utc_now) completed_at: datetime | None = None def add_step(self, step: FlowStep) -> None: """Добавить шаг (агрегирует метрики, включая duration).""" self.steps.append(step) self.total_tokens_used += step.tokens_used self.total_llm_calls += step.llm_calls self.total_tool_calls += step.tool_calls self.total_duration_ms += step.duration_ms or 0 def get_successful_steps(self) -> list[FlowStep]: """Получить успешные шаги.""" return [s for s in self.steps if s.status == FlowStatus.COMPLETED] def get_failed_steps(self) -> list[FlowStep]: """Получить неудачные шаги.""" return [s for s in self.steps if s.status == FlowStatus.FAILED] def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "flow_id": self.flow_id, "status": self.status.value, "output": self.output, "error": self.error, "steps": [step.to_dict() for step in self.steps], "total_duration_ms": self.total_duration_ms, "total_tokens_used": self.total_tokens_used, "total_llm_calls": self.total_llm_calls, "total_tool_calls": self.total_tool_calls, "started_at": to_iso(self.started_at), "completed_at": to_iso(self.completed_at), } # ============================================================================ # Flow Context # ============================================================================ class FlowContext(RunContext): """ Контекст выполнения flow. Расширяет RunContext специфичными для flow полями: flow_id, flow_name, current_node_id, variables, steps. Note: это Pydantic-модель (наследует RunContext) — используется Pydantic ``Field``, а не dataclass ``field``. """ flow_id: str | None = None flow_name: str | None = None current_node_id: str | None = None variables: dict[str, Any] = Field(default_factory=dict) steps: list[FlowStep] = Field(default_factory=list) def set_variable(self, name: str, value: Any) -> None: """Установить переменную flow.""" self.variables[name] = value def get_variable(self, name: str, default: Any = None) -> Any: """Получить переменную flow.""" return self.variables.get(name, default) def has_variable(self, name: str) -> bool: """Проверить наличие переменной.""" return name in self.variables def delete_variable(self, name: str) -> None: """Удалить переменную.""" if name in self.variables: del self.variables[name] def add_step(self, step: FlowStep) -> None: """Добавить шаг.""" self.steps.append(step) def get_current_step(self) -> FlowStep | None: """Получить текущий шаг.""" return self.steps[-1] if self.steps else None def set_current_node(self, node_id: str) -> None: """Установить текущий узел.""" self.current_node_id = node_id def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" data = super().to_dict() data.update( { "flow_id": self.flow_id, "flow_name": self.flow_name, "current_node_id": self.current_node_id, "variables": self.variables, "steps": [step.to_dict() for step in self.steps], } ) return data # ============================================================================ # Flow Metadata # ============================================================================ @dataclass class FlowMetadata: """Метаданные flow.""" version: str = "1.0.0" author: str | None = None tags: list[str] = field(default_factory=list) description: str = "" # Статистика usage_count: int = 0 average_duration_ms: int = 0 success_rate: float = 1.0 # Конфигурация max_concurrent_runs: int = 1 timeout_seconds: int = 300 retry_on_failure: bool = False max_retries: int = 3 def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "version": self.version, "author": self.author, "tags": self.tags, "description": self.description, "usage_count": self.usage_count, "average_duration_ms": self.average_duration_ms, "success_rate": self.success_rate, "max_concurrent_runs": self.max_concurrent_runs, "timeout_seconds": self.timeout_seconds, "retry_on_failure": self.retry_on_failure, "max_retries": self.max_retries, } # ============================================================================ # Flow Validator # ============================================================================ class FlowValidator: """Валидатор flow графа.""" @staticmethod def validate(flow: Flow) -> list[str]: """ Валидировать flow граф. Returns: Список ошибок (пустой если валиден). """ errors: list[str] = [] # Проверка наличия start узла start_nodes = [n for n in flow.nodes if n.type == FlowNodeType.START] if len(start_nodes) == 0: errors.append("Flow must have at least one START node") elif len(start_nodes) > 1: errors.append("Flow must have exactly one START node") # Проверка наличия end узла end_nodes = [n for n in flow.nodes if n.type == FlowNodeType.END] if len(end_nodes) == 0: errors.append("Flow must have at least one END node") # Проверка связности графа if flow.nodes: reachable = FlowValidator._get_reachable_nodes( flow, start_nodes[0].id if start_nodes else None ) all_node_ids = {n.id for n in flow.nodes} unreachable = all_node_ids - reachable if unreachable: errors.append(f"Unreachable nodes: {', '.join(unreachable)}") # Проверка ребер node_ids = {n.id for n in flow.nodes} for edge in flow.edges: if edge.source not in node_ids: errors.append(f"Edge {edge.id} has invalid source: {edge.source}") if edge.target not in node_ids: errors.append(f"Edge {edge.id} has invalid target: {edge.target}") # ✅ Проверка безопасности условий (защита от опасных выражений) if edge.condition and not is_safe_expression(edge.condition): errors.append(f"Edge {edge.id} has unsafe condition: {edge.condition!r}") return errors @staticmethod def _get_reachable_nodes(flow: Flow, start_id: str | None) -> set[str]: """Получить все достижимые узлы из start (BFS).""" if not start_id: return set() reachable: set[str] = set() queue = [start_id] while queue: node_id = queue.pop(0) if node_id in reachable: continue reachable.add(node_id) for edge in flow.get_outgoing_edges(node_id): if edge.target not in reachable: queue.append(edge.target) return reachable @staticmethod def _has_cycle(flow: Flow) -> bool: """Проверить наличие циклов в графе (DFS).""" visited: set[str] = set() rec_stack: set[str] = set() def dfs(node_id: str) -> bool: visited.add(node_id) rec_stack.add(node_id) for edge in flow.get_outgoing_edges(node_id): if edge.target not in visited: if dfs(edge.target): return True elif edge.target in rec_stack: return True rec_stack.remove(node_id) return False for node in flow.nodes: if node.id not in visited: if dfs(node.id): return True return False # ============================================================================ # Flow Runner Type # ============================================================================ # Тип функции-исполнителя flow FlowRunner = Callable[[dict[str, Any]], AsyncIterator[dict[str, Any]]] # ============================================================================ # Flow # ============================================================================ @dataclass class Flow: """ Flow — композиция агентов для выполнения задачи. Может быть задан двумя способами: 1. Через ``runner`` — async generator, который yield'ит события 2. Через ``nodes``/``edges`` — граф для визуального редактора (Flow Composer) """ id: str name: str description: str = "" # Визуальный граф (для Flow Composer в UI) nodes: list[FlowNode] = field(default_factory=list) edges: list[FlowEdge] = field(default_factory=list) # Список имён агентов (для совместимости со старым кодом) agents: list[str] = field(default_factory=list) # Исполнитель flow (async generator событий) runner: FlowRunner | None = None # Метаданные metadata: FlowMetadata = field(default_factory=FlowMetadata) # ======================================================================== # Node Operations # ======================================================================== def get_node(self, node_id: str) -> FlowNode | None: """Получить узел по ID.""" return next((n for n in self.nodes if n.id == node_id), None) def add_node(self, node: FlowNode) -> None: """Добавить узел.""" self.nodes.append(node) def remove_node(self, node_id: str) -> bool: """Удалить узел и все связанные ребра.""" node = self.get_node(node_id) if not node: return False self.nodes.remove(node) self.edges = [e for e in self.edges if node_id not in (e.source, e.target)] return True def get_start_node(self) -> FlowNode | None: """Получить start узел.""" return next((n for n in self.nodes if n.type == FlowNodeType.START), None) def get_end_nodes(self) -> list[FlowNode]: """Получить все end узлы.""" return [n for n in self.nodes if n.type == FlowNodeType.END] # ======================================================================== # Edge Operations # ======================================================================== def get_edge(self, edge_id: str) -> FlowEdge | None: """Получить ребро по ID.""" return next((e for e in self.edges if e.id == edge_id), None) def get_outgoing_edges(self, node_id: str) -> list[FlowEdge]: """Получить исходящие ребра.""" return [e for e in self.edges if e.source == node_id] def get_incoming_edges(self, node_id: str) -> list[FlowEdge]: """Получить входящие ребра.""" return [e for e in self.edges if e.target == node_id] def add_edge(self, edge: FlowEdge) -> None: """Добавить ребро.""" self.edges.append(edge) def remove_edge(self, edge_id: str) -> bool: """Удалить ребро.""" edge = self.get_edge(edge_id) if not edge: return False self.edges.remove(edge) return True def get_next_nodes(self, node_id: str, context: dict[str, Any] | None = None) -> list[FlowNode]: """ Получить следующие узлы для перехода. Учитывает conditional edges и их условия (безопасное вычисление). """ edges = self.get_outgoing_edges(node_id) if not edges: return [] # Сортировка по priority edges.sort(key=lambda e: e.priority, reverse=True) next_nodes: list[FlowNode] = [] for edge in edges: if edge.is_conditional(): if context and edge.evaluate_condition(context): node = self.get_node(edge.target) if node: next_nodes.append(node) else: node = self.get_node(edge.target) if node: next_nodes.append(node) return next_nodes # ======================================================================== # Validation # ======================================================================== def validate(self) -> list[str]: """Валидировать flow граф.""" return FlowValidator.validate(self) def is_valid(self) -> bool: """Проверить, валиден ли flow.""" return len(self.validate()) == 0 # ======================================================================== # Runner # ======================================================================== def has_runner(self) -> bool: """Проверить, есть ли runner у flow.""" return self.runner is not None # ======================================================================== # Metadata # ======================================================================== def increment_usage(self) -> None: """Инкрементировать счётчик использования.""" self.metadata.usage_count += 1 def update_average_duration(self, duration_ms: int) -> None: """Обновить среднюю длительность (экспоненциальное сглаживание).""" if self.metadata.usage_count == 0: self.metadata.average_duration_ms = duration_ms else: alpha = 0.1 self.metadata.average_duration_ms = int( alpha * duration_ms + (1 - alpha) * self.metadata.average_duration_ms ) def update_success_rate(self, success: bool) -> None: """Обновить процент успешных выполнений (экспоненциальное сглаживание).""" if self.metadata.usage_count == 0: self.metadata.success_rate = 1.0 if success else 0.0 else: alpha = 0.1 self.metadata.success_rate = ( alpha * (1.0 if success else 0.0) + (1 - alpha) * self.metadata.success_rate ) # ======================================================================== # Serialization # ======================================================================== def to_dict(self) -> dict[str, Any]: """Преобразовать в словарь.""" return { "id": self.id, "name": self.name, "description": self.description, "nodes": [node.to_dict() for node in self.nodes], "edges": [edge.to_dict() for edge in self.edges], "agents": self.agents, "metadata": self.metadata.to_dict(), } @classmethod def from_dict(cls, data: dict[str, Any]) -> Flow: """Создать из словаря.""" metadata_data = data.get("metadata", {}) metadata = FlowMetadata( version=metadata_data.get("version", "1.0.0"), author=metadata_data.get("author"), tags=metadata_data.get("tags", []), description=metadata_data.get("description", ""), usage_count=metadata_data.get("usage_count", 0), average_duration_ms=metadata_data.get("average_duration_ms", 0), success_rate=metadata_data.get("success_rate", 1.0), max_concurrent_runs=metadata_data.get("max_concurrent_runs", 1), timeout_seconds=metadata_data.get("timeout_seconds", 300), retry_on_failure=metadata_data.get("retry_on_failure", False), max_retries=metadata_data.get("max_retries", 3), ) return cls( id=data["id"], name=data["name"], description=data.get("description", ""), nodes=[FlowNode.from_dict(n) for n in data.get("nodes", [])], edges=[FlowEdge.from_dict(e) for e in data.get("edges", [])], agents=data.get("agents", []), metadata=metadata, ) # ============================================================================ # Exports # ============================================================================ __all__ = [ "Flow", "FlowContext", "FlowEdge", "FlowEdgeType", "FlowMetadata", "FlowNode", "FlowNodeType", "FlowResult", "FlowRunner", "FlowStatus", "FlowStep", "FlowValidator", ]