/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/schemas/flow.py
260 строк
8 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
"""Pydantic schemas для flows.""" from __future__ import annotations from datetime import datetime from typing import Any, Literal from uuid import UUID from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator from src.primitives import ( Flow, FlowEdge, FlowNode, FlowValidator, is_safe_expression, ) from src.primitives import ( FlowEdgeType as FlowEdgeTypeEnum, ) from src.primitives import ( FlowNodeType as FlowNodeTypeEnum, ) # ============================================================================ # Types # ============================================================================ FlowNodeType = Literal[ "start", "end", "agent", "llm", "tool", "condition", "parallel", "loop", "subflow", "transform", "wait", ] FlowEdgeType = Literal["normal", "conditional", "error", "default"] FlowStatus = Literal["draft", "active", "paused", "archived"] # ============================================================================ # Flow Graph (nodes / edges) # ============================================================================ class FlowNodeSchema(BaseModel): """Узел графа flow.""" id: str = Field(..., min_length=1) type: FlowNodeType label: str = "" data: dict[str, Any] = Field(default_factory=dict) position: dict[str, float] = Field(default_factory=dict) config: dict[str, Any] = Field(default_factory=dict) class FlowEdgeSchema(BaseModel): """Ребро графа flow.""" id: str = Field(..., min_length=1) source: str target: str label: str = "" edge_type: FlowEdgeType = "normal" condition: str | None = None # безопасное выражение (safe_eval) priority: int = 0 @field_validator("condition") @classmethod def validate_condition(cls, v: str | None) -> str | None: """ Проверить безопасность выражения условия. Использует ``is_safe_expression`` (AST whitelist) — отклоняет import, lambda, dunder-атрибуты, произвольные вызовы. """ if v is None or not v.strip(): return v if not is_safe_expression(v): raise ValueError( f"Небезопасное или невалидное условие: {v!r}. " "Допустимы только сравнения (==, !=, <, >, in), " "логические операторы (and, or, not) и переменные." ) return v @model_validator(mode="after") def validate_conditional_edge(self) -> FlowEdgeSchema: """Conditional edge обязан иметь выражение условия.""" if self.edge_type == "conditional" and not (self.condition and self.condition.strip()): raise ValueError(f"Edge '{self.id}': conditional edge должен иметь выражение condition") return self # ============================================================================ # Flow # ============================================================================ class FlowCreate(BaseModel): """Запрос на создание flow.""" name: str = Field(..., min_length=1, max_length=200) description: str | None = Field(default=None, max_length=2000) nodes: list[FlowNodeSchema] = Field(default_factory=list) edges: list[FlowEdgeSchema] = Field(default_factory=list) status: FlowStatus = "draft" tags: list[str] = Field(default_factory=list) metadata: dict[str, Any] = Field(default_factory=dict) def _to_flow(self) -> Flow: """Сконвертировать в primitives.Flow для валидации графа.""" nodes = [ FlowNode( id=n.id, type=FlowNodeTypeEnum(n.type), label=n.label, data=n.data, position=n.position, config=n.config, ) for n in self.nodes ] edges = [ FlowEdge( id=e.id, source=e.source, target=e.target, label=e.label, edge_type=FlowEdgeTypeEnum(e.edge_type), condition=e.condition, priority=e.priority, ) for e in self.edges ] return Flow(id="validation", name=self.name, nodes=nodes, edges=edges) @model_validator(mode="after") def validate_graph(self) -> FlowCreate: """ Валидировать граф flow (если задан nodes/edges). Проверяет через FlowValidator: - ровно один START node - минимум один END node - нет недостижимых узлов - рёбра ссылаются на существующие узлы - условия безопасны (is_safe_expression) """ if not self.nodes: return self # runner-based flow (без графа) errors = FlowValidator.validate(self._to_flow()) if errors: raise ValueError(f"Невалидный граф flow: {'; '.join(errors)}") return self class FlowUpdate(BaseModel): """Запрос на обновление flow (все поля опциональны).""" name: str | None = Field(default=None, min_length=1, max_length=200) description: str | None = Field(default=None, max_length=2000) nodes: list[FlowNodeSchema] | None = None edges: list[FlowEdgeSchema] | None = None status: FlowStatus | None = None tags: list[str] | None = None metadata: dict[str, Any] | None = None @model_validator(mode="after") def validate_graph(self) -> FlowUpdate: """ Валидировать граф, если переданы И nodes, И edges (полное обновление). Частичные обновления (только name/status) не валидируют граф. """ if self.nodes is None or self.edges is None: return self nodes = [ FlowNode( id=n.id, type=FlowNodeTypeEnum(n.type), label=n.label, data=n.data, position=n.position, config=n.config, ) for n in self.nodes ] edges = [ FlowEdge( id=e.id, source=e.source, target=e.target, label=e.label, edge_type=FlowEdgeTypeEnum(e.edge_type), condition=e.condition, priority=e.priority, ) for e in self.edges ] flow = Flow(id="validation", name=self.name or "validation", nodes=nodes, edges=edges) errors = FlowValidator.validate(flow) if errors: raise ValueError(f"Невалидный граф flow: {'; '.join(errors)}") return self class FlowResponse(BaseModel): """Ответ с информацией о flow.""" model_config = ConfigDict(from_attributes=True, populate_by_name=True) id: UUID workspace_id: str name: str description: str | None nodes: list[dict[str, Any]] = Field(default_factory=list) edges: list[dict[str, Any]] = Field(default_factory=list) status: str tags: list[str] = Field(default_factory=list) # В модели атрибут metadata_ (колонка "metadata") metadata: dict[str, Any] = Field( default_factory=dict, validation_alias="metadata_", serialization_alias="metadata", ) created_at: datetime updated_at: datetime # ============================================================================ # Flow Run # ============================================================================ class FlowRunRequest(BaseModel): """Запрос на запуск flow.""" input: dict[str, Any] = Field(default_factory=dict) stream: bool = Field(default=False) class FlowRunResponse(BaseModel): """Ответ с результатом запуска flow (non-streaming).""" flow_id: UUID status: str output: dict[str, Any] = Field(default_factory=dict) steps: list[dict[str, Any]] = Field(default_factory=list) tokens_used: int = 0 duration_ms: float = 0.0 error: str | None = None