/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/flows/registry.py
260 строк
8 KB
Alexander Efanov
Обновление репозитория
15 июл 2026, 12:19
15 июл 2026, 12:19
76704c6
Код
Авторство
О чём код?
# core/engine/src/flows/registry.py """Реестр сценариев (Flows) для FlowStack Engine. Этот модуль предоставляет централизованный реестр всех доступных сценариев, которые могут быть запущены через API `/api/v1/runs`. Каждый flow представляет собой мультиагентный workflow с определённой последовательностью агентов и специфичной логикой выполнения. Доступные flows: - research: Полное исследование (Planner → Researcher → Analyst → Writer → Reviewer) - review: Критический обзор (Critic → Reviewer → Writer) - debate: Дебаты и принятие решения (Planner → Debater PRO + CON → Analyst → Judge) - code_review: Ревью кода (Researcher → Analyst → Reviewer → Critic → Writer) """ from __future__ import annotations from typing import Any, AsyncIterator, Callable import structlog from src.primitives.flow import Flow logger = structlog.get_logger() # ============================================================================ # Типы # ============================================================================ FlowRunner = Callable[[dict[str, Any]], AsyncIterator[dict[str, Any]]] # ============================================================================ # Импорт flows # ============================================================================ # Импортируем все flow модули # Каждый модуль экспортирует экземпляр Flow с атрибутом runner try: from src.flows.research import research_flow except ImportError as e: logger.warning("flows.registry.research_import_failed", error=str(e)) research_flow = None try: from src.flows.review_flow import review_flow except ImportError as e: logger.warning("flows.registry.review_import_failed", error=str(e)) review_flow = None try: from src.flows.debate_flow import debate_flow except ImportError as e: logger.warning("flows.registry.debate_import_failed", error=str(e)) debate_flow = None try: from src.flows.code_review_flow import code_review_flow except ImportError as e: logger.warning("flows.registry.code_review_import_failed", error=str(e)) code_review_flow = None # ============================================================================ # Реестр # ============================================================================ # Централизованный реестр всех доступных flows # Ключ: flow_id (используется в API) # Значение: Flow экземпляр FLOW_REGISTRY: dict[str, Flow] = {} # Регистрируем flows (только те, что успешно импортированы) if research_flow is not None: FLOW_REGISTRY["research"] = research_flow logger.info("flows.registry.registered", flow_id="research", agents=research_flow.agents) if review_flow is not None: FLOW_REGISTRY["review"] = review_flow logger.info("flows.registry.registered", flow_id="review", agents=review_flow.agents) if debate_flow is not None: FLOW_REGISTRY["debate"] = debate_flow logger.info("flows.registry.registered", flow_id="debate", agents=debate_flow.agents) if code_review_flow is not None: FLOW_REGISTRY["code_review"] = code_review_flow logger.info("flows.registry.registered", flow_id="code_review", agents=code_review_flow.agents) # ============================================================================ # Public API # ============================================================================ def get_flow(flow_id: str) -> Flow | None: """ Получить flow по ID. Args: flow_id: Идентификатор сценария (например, "research", "review") Returns: Flow экземпляр или None, если flow не найден Example: >>> flow = get_flow("research") >>> if flow: ... async for event in flow.runner(input_data): ... print(event) """ flow = FLOW_REGISTRY.get(flow_id) if flow is None: logger.warning("flows.registry.not_found", flow_id=flow_id, available=list(FLOW_REGISTRY.keys())) return flow def list_flows() -> list[Flow]: """ Получить список всех доступных flows. Returns: Список Flow экземпляров, отсортированный по id Example: >>> flows = list_flows() >>> for flow in flows: ... print(f"{flow.id}: {flow.name}") """ return sorted(FLOW_REGISTRY.values(), key=lambda f: f.id) def get_flow_ids() -> list[str]: """ Получить список всех доступных flow ID. Returns: Отсортированный список строк — ID сценариев Example: >>> ids = get_flow_ids() >>> print(ids) # ['code_review', 'debate', 'research', 'review'] """ return sorted(FLOW_REGISTRY.keys()) def flow_exists(flow_id: str) -> bool: """ Проверить существование flow по ID. Args: flow_id: Идентификатор сценария Returns: True если flow существует, False иначе """ return flow_id in FLOW_REGISTRY def get_flows_summary() -> dict[str, dict[str, Any]]: """ Получить краткую сводку всех flows для API/UI. Returns: Словарь {flow_id: {name, description, agents, agents_count}} Example: >>> summary = get_flows_summary() >>> print(summary["research"]) { "name": "Research Flow", "description": "Мультиагентное исследование...", "agents": ["planner", "researcher", "analyst", "writer", "reviewer"], "agents_count": 5 } """ summary = {} for flow_id, flow in FLOW_REGISTRY.items(): summary[flow_id] = { "name": flow.name, "description": flow.description, "agents": flow.agents, "agents_count": len(flow.agents), } return summary def register_flow(flow: Flow) -> None: """ Зарегистрировать новый flow в реестре (для динамического добавления). Args: flow: Flow экземпляр для регистрации Raises: ValueError: Если flow с таким ID уже существует Example: >>> custom_flow = Flow(id="custom", name="Custom", ...) >>> register_flow(custom_flow) """ if flow.id in FLOW_REGISTRY: raise ValueError(f"Flow with id '{flow.id}' already exists in registry") FLOW_REGISTRY[flow.id] = flow logger.info("flows.registry.dynamic_registered", flow_id=flow.id, agents=flow.agents) def unregister_flow(flow_id: str) -> bool: """ Удалить flow из реестра (для динамического удаления). Args: flow_id: ID сценария для удаления Returns: True если flow был удалён, False если не найден """ if flow_id in FLOW_REGISTRY: del FLOW_REGISTRY[flow_id] logger.info("flows.registry.unregistered", flow_id=flow_id) return True logger.warning("flows.registry.unregister_not_found", flow_id=flow_id) return False # ============================================================================ # Инициализация # ============================================================================ # Логируем итоговое состояние реестра при импорте модуля logger.info( "flows.registry.initialized", total_flows=len(FLOW_REGISTRY), flow_ids=list(FLOW_REGISTRY.keys()), total_agents=sum(len(f.agents) for f in FLOW_REGISTRY.values()), ) # ============================================================================ # Exports # ============================================================================ __all__ = [ # Public API "get_flow", "list_flows", "get_flow_ids", "flow_exists", "get_flows_summary", "register_flow", "unregister_flow", # Реестр (для прямого доступа) "FLOW_REGISTRY", # Типы "FlowRunner", ]