/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/flows/custom.py
886 строк
30 KB
Alexander Efanov
Обновление репозитория
15 июл 2026, 12:19
15 июл 2026, 12:19
76704c6
Код
Авторство
О чём код?
# core/engine/src/flows/custom.py """Custom Flow Builder — динамическое создание мультиагентных flows. Предоставляет инструменты для создания пользовательских flows без написания кода — через Builder pattern, YAML/JSON конфигурации или API. Основные возможности: - Builder pattern для пошагового создания flow - Функции-хелперы для быстрого создания - Загрузка flow из YAML/JSON конфигураций - Валидация конфигурации - Динамический runner для выполнения цепочки агентов - Экспорт/импорт flows - Регистрация в глобальном registry Примеры использования: 1. Builder pattern: from src.flows.custom import CustomFlowBuilder flow = ( CustomFlowBuilder() .with_id("my_research") .with_name("Custom Research") .with_description("Мой кастомный research flow") .add_agent("planner", system_prompt="Ты — эксперт по планированию") .add_agent("researcher") .add_agent("writer") .chain() .build() ) 2. Быстрое создание: from src.flows.custom import create_custom_flow flow = create_custom_flow( id="quick_flow", name="Quick Flow", agents=["planner", "researcher", "writer"], prompts={"planner": "Custom prompt for planner"} ) 3. Из YAML конфигурации: from src.flows.custom import load_flow_from_yaml flow = load_flow_from_yaml("my_flow.yaml") 4. Динамический runner: async for event in run_custom_flow(agents_config, input_data): print(event) """ from __future__ import annotations import json from pathlib import Path from typing import Any, AsyncIterator import structlog import yaml from src.agents import AGENT_REGISTRY, create_agent from src.flows.registry import register_flow, unregister_flow from src.primitives.flow import Flow logger = structlog.get_logger() # ============================================================================ # Exceptions # ============================================================================ class CustomFlowError(Exception): """Базовое исключение для custom flow ошибок.""" pass class FlowValidationError(CustomFlowError): """Ошибка валидации конфигурации flow.""" pass class AgentNotFoundError(CustomFlowError): """Агент не найден в registry.""" pass class FlowAlreadyExistsError(CustomFlowError): """Flow с таким ID уже существует.""" pass # ============================================================================ # Types # ============================================================================ class AgentConfig: """Конфигурация агента в custom flow.""" def __init__( self, agent_type: str, name: str | None = None, system_prompt: str | None = None, input_template: str | None = None, ): """ Args: agent_type: Тип агента (ключ из AGENT_REGISTRY) name: Кастомное имя агента system_prompt: Кастомный system prompt input_template: Шаблон для формирования input (может использовать {original_input}, {prev_output}, {prev_outputs.XXX}) """ self.agent_type = agent_type self.name = name self.system_prompt = system_prompt self.input_template = input_template def to_dict(self) -> dict[str, Any]: """Сериализация в dict.""" result: dict[str, Any] = {"agent_type": self.agent_type} if self.name: result["name"] = self.name if self.system_prompt: result["system_prompt"] = self.system_prompt if self.input_template: result["input_template"] = self.input_template return result @classmethod def from_dict(cls, data: dict[str, Any]) -> "AgentConfig": """Десериализация из dict.""" return cls( agent_type=data["agent_type"], name=data.get("name"), system_prompt=data.get("system_prompt"), input_template=data.get("input_template"), ) # ============================================================================ # Dynamic runner для custom flows # ============================================================================ async def run_custom_flow( agents_config: list[AgentConfig], input_data: dict[str, Any], flow_id: str = "custom", ) -> AsyncIterator[dict[str, Any]]: """ Динамический runner для custom flow. Выполняет агентов последовательно, передавая output предыдущего агента следующему. Args: agents_config: Список конфигураций агентов input_data: Входные данные flow_id: ID flow (для логов и событий) Yields: События SSE: - agent_start: агент начал работу - agent_message: стриминг контента - agent_done: агент завершил - error: ошибка - flow_done: финальное событие """ logger.info( "custom_flow.started", flow_id=flow_id, agents_count=len(agents_config), ) # Контейнер для outputs каждого агента prev_outputs: dict[str, str] = {} total_tokens = 0 # Исходный input original_input = ( input_data.get("topic") or input_data.get("input") or input_data.get("content") or "" ) if not original_input.strip(): yield { "type": "error", "error": "Не указан input (ожидается поле 'topic', 'input' или 'content')", } return for i, agent_config in enumerate(agents_config): agent_id = agent_config.name or f"{agent_config.agent_type}_{i}" yield {"type": "agent_start", "agent": agent_id} # Валидация типа агента if agent_config.agent_type not in AGENT_REGISTRY: logger.error( "custom_flow.agent.not_found", agent_type=agent_config.agent_type, ) yield { "type": "error", "error": f"Агент '{agent_config.agent_type}' не найден в registry", "agent": agent_id, } return # Создаём агента с кастомными параметрами try: agent = create_agent( agent_type=agent_config.agent_type, name=agent_config.name, system_prompt=agent_config.system_prompt, ) except Exception as e: logger.error( "custom_flow.agent.create_failed", agent_type=agent_config.agent_type, error=str(e), ) yield { "type": "error", "error": f"Не удалось создать агента '{agent_config.agent_type}': {e}", "agent": agent_id, } return # Формируем input для агента if agent_config.input_template: # Используем кастомный шаблон prev_output = list(prev_outputs.values())[-1] if prev_outputs else "" try: agent_input_text = agent_config.input_template.format( original_input=original_input, prev_output=prev_output, prev_outputs=prev_outputs, ) except (KeyError, IndexError) as e: logger.warning( "custom_flow.template.error", template=agent_config.input_template, error=str(e), ) # Fallback на простой input agent_input_text = f"{original_input}\n\n{prev_output}" elif i == 0: # Первый агент получает оригинальный input agent_input_text = original_input else: # Последующие получают оригинальный input + output предыдущего prev_output = list(prev_outputs.values())[-1] agent_input_text = f"""Исходная задача: {original_input} Результат предыдущего этапа: \"\"\" {prev_output} \"\"\" Продолжи работу над задачей, учитывая результат предыдущего этапа.""" # Запускаем агента со стримингом agent_output = "" try: async for chunk in agent.run_stream({"topic": agent_input_text}): chunk_type = chunk.get("type") if chunk_type == "content": content = chunk.get("content", "") agent_output += content yield { "type": "agent_message", "agent": agent_id, "content": content, } elif chunk_type == "done": tokens = chunk.get("tokens_total", 0) total_tokens += tokens yield { "type": "agent_done", "agent": agent_id, "output": agent_output, "tokens": tokens, "model": chunk.get("model", ""), } elif chunk_type == "error": yield { "type": "error", "error": chunk.get("error", f"Unknown error in {agent_id}"), "agent": agent_id, } return except Exception as e: logger.error( "custom_flow.agent.execution_failed", agent=agent_id, error=str(e), exc_info=True, ) yield { "type": "error", "error": f"Ошибка выполнения агента '{agent_id}': {e}", "agent": agent_id, } return # Сохраняем output для следующего агента prev_outputs[agent_id] = agent_output # Финальное событие final_output = list(prev_outputs.values())[-1] if prev_outputs else "" logger.info( "custom_flow.completed", flow_id=flow_id, total_tokens=total_tokens, stages_count=len(prev_outputs), ) yield { "type": "flow_done", "flow_id": flow_id, "output": final_output, "tokens": total_tokens, "stages": prev_outputs, "metadata": { "agents_count": len(agents_config), "total_stages": len(prev_outputs), "is_custom": True, }, } # ============================================================================ # Builder pattern # ============================================================================ class CustomFlowBuilder: """ Builder для пошагового создания custom flow. Пример: builder = CustomFlowBuilder() flow = ( builder .with_id("my_flow") .with_name("My Custom Flow") .with_description("Описание") .add_agent("planner", system_prompt="Custom prompt") .add_agent("researcher") .add_agent("writer") .chain() .build() ) """ def __init__(self): self._id: str | None = None self._name: str | None = None self._description: str | None = None self._agents: list[AgentConfig] = [] self._metadata: dict[str, Any] = {} def with_id(self, flow_id: str) -> "CustomFlowBuilder": """Установить ID flow.""" self._id = flow_id return self def with_name(self, name: str) -> "CustomFlowBuilder": """Установить имя flow.""" self._name = name return self def with_description(self, description: str) -> "CustomFlowBuilder": """Установить описание flow.""" self._description = description return self def with_metadata(self, key: str, value: Any) -> "CustomFlowBuilder": """Добавить метаданные.""" self._metadata[key] = value return self def add_agent( self, agent_type: str, name: str | None = None, system_prompt: str | None = None, input_template: str | None = None, ) -> "CustomFlowBuilder": """Добавить агента в flow.""" self._agents.append(AgentConfig( agent_type=agent_type, name=name, system_prompt=system_prompt, input_template=input_template, )) return self def chain(self) -> "CustomFlowBuilder": """ Пометить агентов как последовательную цепочку. В текущей реализации все агенты выполняются последовательно по умолчанию, этот метод существует для будущего расширения (графы, ветвления). """ # В текущей версии run_custom_flow всегда выполняет агентов последовательно # Этот метод зарезервирован для будущих расширений return self def validate(self) -> list[str]: """ Валидировать конфигурацию flow. Returns: Список ошибок валидации (пустой если всё ок) """ errors = [] if not self._id: errors.append("ID flow не установлен") if not self._name: errors.append("Name flow не установлен") if not self._agents: errors.append("Не добавлено ни одного агента") # Проверяем что все агенты существуют for agent_config in self._agents: if agent_config.agent_type not in AGENT_REGISTRY: errors.append( f"Агент '{agent_config.agent_type}' не найден в registry. " f"Доступные: {list(AGENT_REGISTRY.keys())}" ) return errors def build(self, auto_register: bool = False) -> Flow: """ Построить Flow. Args: auto_register: Автоматически зарегистрировать в registry Returns: Flow instance Raises: FlowValidationError: Если конфигурация невалидна """ errors = self.validate() if errors: raise FlowValidationError( "Ошибки валидации flow:\n" + "\n".join(f" - {e}" for e in errors) ) # Type narrowing: после validate() эти поля гарантированно не None # (validate() проверяет self._id и self._name и бросает ошибку если они пустые) assert self._id is not None, "ID должен быть установлен (проверено в validate)" assert self._name is not None, "Name должен быть установлен (проверено в validate)" assert self._agents, "Список агентов не должен быть пустым (проверено в validate)" # Сохраняем в локальные переменные для ясности flow_id: str = self._id flow_name: str = self._name # Формируем описание агентов для description agents_list = [a.agent_type for a in self._agents] agents_str = " → ".join(agents_list) description = self._description or f"Custom flow: {agents_str}" # Создаём runner замыканием agents_config = self._agents.copy() async def custom_runner(input_data: dict[str, Any]) -> AsyncIterator[dict[str, Any]]: async for event in run_custom_flow(agents_config, input_data, flow_id): yield event flow = Flow( id=flow_id, name=flow_name, description=description, agents=agents_list, runner=custom_runner, ) # Сохраняем metadata в атрибутах flow flow.agents_config = agents_config # type: ignore flow.metadata = self._metadata # type: ignore flow.is_custom = True # type: ignore if auto_register: register_flow(flow) logger.info("custom_flow.registered", flow_id=self._id) return flow def reset(self) -> "CustomFlowBuilder": """Сбросить builder для повторного использования.""" self._id = None self._name = None self._description = None self._agents = [] self._metadata = {} return self # ============================================================================ # Helper functions # ============================================================================ def create_custom_flow( id: str, name: str, agents: list[str], description: str | None = None, prompts: dict[str, str] | None = None, names: dict[str, str] | None = None, auto_register: bool = False, ) -> Flow: """ Быстрое создание custom flow из списка агентов. Args: id: ID flow name: Имя flow agents: Список типов агентов (последовательно) description: Описание flow prompts: Словарь {agent_type: system_prompt} для кастомизации names: Словарь {agent_type: custom_name} для именования auto_register: Автоматически зарегистрировать в registry Returns: Flow instance Example: flow = create_custom_flow( id="my_flow", name="My Flow", agents=["planner", "researcher", "writer"], prompts={"planner": "Custom planner prompt"}, ) """ prompts = prompts or {} names = names or {} builder = CustomFlowBuilder() builder.with_id(id).with_name(name) if description: builder.with_description(description) for agent_type in agents: builder.add_agent( agent_type=agent_type, name=names.get(agent_type), system_prompt=prompts.get(agent_type), ) return builder.build(auto_register=auto_register) def load_flow_from_yaml(path: str | Path, auto_register: bool = False) -> Flow: """ Загрузить custom flow из YAML файла. Формат YAML: id: my_custom_flow name: My Custom Flow description: Описание flow agents: - type: planner name: Custom Planner system_prompt: "Ты — эксперт по планированию" input_template: "Задача: {original_input}" - type: researcher - type: writer metadata: author: user@example.com version: "1.0" Args: path: Путь к YAML файлу auto_register: Автоматически зарегистрировать Returns: Flow instance Raises: FileNotFoundError: Файл не найден FlowValidationError: Ошибка валидации """ path = Path(path) if not path.exists(): raise FileNotFoundError(f"Файл не найден: {path}") with open(path, "r", encoding="utf-8") as f: data = yaml.safe_load(f) return _load_flow_from_dict(data, auto_register=auto_register) def load_flow_from_json(path: str | Path, auto_register: bool = False) -> Flow: """ Загрузить custom flow из JSON файла. Формат аналогичен YAML. Args: path: Путь к JSON файлу auto_register: Автоматически зарегистрировать Returns: Flow instance """ path = Path(path) if not path.exists(): raise FileNotFoundError(f"Файл не найден: {path}") with open(path, "r", encoding="utf-8") as f: data = json.load(f) return _load_flow_from_dict(data, auto_register=auto_register) def _load_flow_from_dict(data: dict[str, Any], auto_register: bool = False) -> Flow: """Внутренняя функция загрузки flow из dict.""" builder = CustomFlowBuilder() # Обязательные поля if "id" not in data: raise FlowValidationError("Отсутствует обязательное поле 'id'") if "name" not in data: raise FlowValidationError("Отсутствует обязательное поле 'name'") if "agents" not in data: raise FlowValidationError("Отсутствует обязательное поле 'agents'") builder.with_id(data["id"]).with_name(data["name"]) if "description" in data: builder.with_description(data["description"]) # Добавляем metadata if "metadata" in data and isinstance(data["metadata"], dict): for key, value in data["metadata"].items(): builder.with_metadata(key, value) # Добавляем агентов for agent_data in data["agents"]: if isinstance(agent_data, str): # Простая форма: ["planner", "researcher", "writer"] builder.add_agent(agent_type=agent_data) elif isinstance(agent_data, dict): # Полная форма с кастомизацией if "type" not in agent_data: raise FlowValidationError( f"В конфигурации агента отсутствует поле 'type': {agent_data}" ) builder.add_agent( agent_type=agent_data["type"], name=agent_data.get("name"), system_prompt=agent_data.get("system_prompt"), input_template=agent_data.get("input_template"), ) else: raise FlowValidationError( f"Неверный формат агента: {agent_data}. " "Ожидается строка или dict" ) return builder.build(auto_register=auto_register) def load_flow_from_dict(data: dict[str, Any], auto_register: bool = False) -> Flow: """ Загрузить custom flow из dict (публичная версия). Args: data: Конфигурация flow в виде dict auto_register: Автоматически зарегистрировать Returns: Flow instance """ return _load_flow_from_dict(data, auto_register=auto_register) def export_flow_to_yaml(flow: Flow, path: str | Path) -> None: """ Экспортировать custom flow в YAML файл. Args: flow: Flow для экспорта path: Путь для сохранения """ path = Path(path) # Получаем agents_config если это custom flow agents_config = getattr(flow, "agents_config", None) metadata = getattr(flow, "metadata", {}) if agents_config is None: # Это не custom flow, экспортируем только базовую информацию agents_data = [{"type": a} for a in flow.agents] else: agents_data = [ac.to_dict() for ac in agents_config] data = { "id": flow.id, "name": flow.name, "description": flow.description, "agents": agents_data, } if metadata: data["metadata"] = metadata with open(path, "w", encoding="utf-8") as f: yaml.dump(data, f, allow_unicode=True, sort_keys=False, default_flow_style=False) logger.info("custom_flow.exported.yaml", flow_id=flow.id, path=str(path)) def export_flow_to_json(flow: Flow, path: str | Path) -> None: """ Экспортировать custom flow в JSON файл. Args: flow: Flow для экспорта path: Путь для сохранения """ path = Path(path) agents_config = getattr(flow, "agents_config", None) metadata = getattr(flow, "metadata", {}) if agents_config is None: agents_data = [{"type": a} for a in flow.agents] else: agents_data = [ac.to_dict() for ac in agents_config] data = { "id": flow.id, "name": flow.name, "description": flow.description, "agents": agents_data, } if metadata: data["metadata"] = metadata with open(path, "w", encoding="utf-8") as f: json.dump(data, f, ensure_ascii=False, indent=2) logger.info("custom_flow.exported.json", flow_id=flow.id, path=str(path)) # ============================================================================ # Управление custom flows # ============================================================================ def register_custom_flow(flow: Flow) -> None: """ Зарегистрировать custom flow в registry. Args: flow: Flow для регистрации Raises: FlowAlreadyExistsError: Если flow с таким ID уже существует """ try: register_flow(flow) logger.info("custom_flow.registered", flow_id=flow.id) except ValueError as e: raise FlowAlreadyExistsError(str(e)) from e def unregister_custom_flow(flow_id: str) -> bool: """ Удалить custom flow из registry. Args: flow_id: ID flow Returns: True если flow был удалён """ result = unregister_flow(flow_id) if result: logger.info("custom_flow.unregistered", flow_id=flow_id) return result def list_custom_flows() -> list[Flow]: """ Получить список всех custom flows. Returns: Список custom flows (с атрибутом is_custom=True) """ from src.flows.registry import list_flows return [f for f in list_flows() if getattr(f, "is_custom", False)] # ============================================================================ # CLI для управления custom flows # ============================================================================ def _print_custom_flow_info(flow: Flow) -> None: """Красиво вывести информацию о custom flow.""" from src.flows.debug import Colors, color, header, subheader print(header(f"Custom Flow: {flow.name}")) print(subheader("Basic Info")) print(f" {color('ID:', Colors.BOLD)} {flow.id}") print(f" {color('Name:', Colors.BOLD)} {flow.name}") print(f" {color('Description:', Colors.BOLD)} {flow.description}") print(f" {color('Is Custom:', Colors.BOLD)} Yes") print(subheader("Agents")) for i, agent_type in enumerate(flow.agents, 1): print(f" {i}. {color(agent_type, Colors.BRIGHT_CYAN)}") agents_config = getattr(flow, "agents_config", None) if agents_config: print(subheader("Agent Configurations")) for ac in agents_config: print(f" {color(ac.agent_type, Colors.BRIGHT_CYAN)}:") if ac.name: print(f" name: {ac.name}") if ac.system_prompt: prompt_preview = ac.system_prompt[:100] + "..." if len(ac.system_prompt) > 100 else ac.system_prompt print(f" system_prompt: {prompt_preview}") if ac.input_template: print(f" input_template: {ac.input_template}") # ============================================================================ # Public API # ============================================================================ __all__ = [ # Exceptions "CustomFlowError", "FlowValidationError", "AgentNotFoundError", "FlowAlreadyExistsError", # Types "AgentConfig", # Builder "CustomFlowBuilder", # Quick creation "create_custom_flow", # Loaders "load_flow_from_yaml", "load_flow_from_json", "load_flow_from_dict", # Exporters "export_flow_to_yaml", "export_flow_to_json", # Management "register_custom_flow", "unregister_custom_flow", "list_custom_flows", # Runner "run_custom_flow", ]