/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/flows/compare.py
1 066 строк
43 KB
Alexander Efanov
Обновление репозитория
15 июл 2026, 12:19
15 июл 2026, 12:19
76704c6
Код
Авторство
О чём код?
# core/engine/src/flows/compare.py """Compare flow — мультиагентное сравнение вариантов для принятия решений. Последовательность агентов: 1. PlannerAgent — определяет критерии сравнения и структуру анализа 2. Options Research — параллельное исследование каждого варианта: для каждой опции запускается отдельный ResearcherAgent 3. AnalystAgent — проводит сравнительный анализ по критериям 4. CriticAgent — выявляет скрытые trade-offs и подводные камни 5. JudgeAgent — выносит рекомендацию с обоснованием 6. WriterAgent — формирует финальный сравнительный отчёт Философия flow: В отличие от debate_flow (ЗА/ПРОТИВ одной идеи) и consensus_flow (несколько экспертов по одной проблеме), compare_flow фокусируется на **сравнении N различных вариантов** решения одной задачи. Идеально подходит для: - Выбор технологий (PostgreSQL vs MongoDB vs Cassandra) - Выбор облачных провайдеров (AWS vs GCP vs Azure) - Сравнение архитектурных подходов (монолит vs микросервисы vs serverless) - Выбор библиотек и фреймворков - Сравнение бизнес-стратегий - Оценка vendor-ов для enterprise решений Отличие от других flows: - debate_flow: 2 позиции (ЗА/ПРОТИВ) → вердикт судьи - consensus_flow: N экспертов по 1 проблеме → консенсус - compare_flow: 1 задача, N вариантов → сравнительный анализ и выбор Особенности реализации: - Параллельное исследование всех опций (ускорение в N раз) - Автоматическое определение опций из input - Матрица сравнения с оценками по критериям - Trade-off анализ для каждой опции - Контекстно-зависимые рекомендации События SSE (streaming): - agent_start: агент начал работу - agent_message: агент генерирует контент (стриминг) - agent_done: агент завершил работу - options_researched — все опции исследованы (мета-событие) - error: ошибка в одном из агентов - flow_done: финальное событие с сравнительным отчётом """ from __future__ import annotations import asyncio import re from typing import Any, AsyncIterator import structlog from src.agents import ( PlannerAgent, ResearcherAgent, AnalystAgent, CriticAgent, JudgeAgent, WriterAgent, ) from src.primitives.flow import Flow logger = structlog.get_logger() # ============================================================================ # Specialized system prompts # ============================================================================ PLANNER_COMPARE_PROMPT = """Ты — эксперт по планированию сравнительных анализов. Твоя задача — разработать **методологию сравнения** нескольких вариантов решения одной задачи. ## Твоя роль: - Определить ключевые критерии сравнения - Разработать шкалу оценок для каждого критерия - Определить структуру сравнительного анализа - Учесть контекст и ограничения ## Принципы работы: 1. **Релевантность** — критерии должны быть важны для конкретной задачи 2. **Измеримость** — каждый критерий должен быть измерим/оценим 3. **Полнота** — охватить все важные аспекты (технические, бизнес, риски) 4. **Сбалансированность** — не отдавать предпочтение заранее ## Формат ответа: ### 🎯 Цель сравнения (1-2 предложения: зачем мы сравниваем и что хотим выбрать) ### 📋 Ключевые критерии сравнения (5-8 штук) Для каждого критерия: - **Название:** ... - **Вес:** (1-10, насколько важен этот критерий) - **Как оценивать:** (методология оценки 1-10) - **Почему важен:** (обоснование) ### 🎯 Критерии исключения (Условия, при которых опция сразу отклоняется) ### 📊 Структура анализа (Как будет организован сравнительный отчёт) ### ⚠️ Типичные ловушки (На что обратить внимание при сравнении) Отвечай на русском языке. Создай практичную методологию.""" RESEARCHER_OPTION_PROMPT_TEMPLATE = """Ты — исследователь, изучающий конкретный вариант решения в рамках сравнительного анализа. ## Исследуемый вариант: "{option_name}" ## Контекст задачи: \"\"\" {task_context} \"\"\" ## Критерии для анализа (от планировщика): \"\"\" {criteria} \"\"\" ## Твоя задача: Проведи глубокий анализ именно этого варианта по всем критериям. ## Формат ответа: ### 📋 Паспорт варианта - **Название:** {option_name} - **Категория:** (тип решения) - **Краткое описание:** (1-2 предложения) ### 💪 Сильные стороны (5-7 конкретных преимуществ с примерами) ### ⚠️ Слабые стороны и ограничения (5-7 конкретных недостатков с примерами) ### 📊 Оценка по критериям (1-10 для каждого) Для каждого критерия: - **Критерий:** оценка X/10 - **Обоснование:** почему именно такая оценка - **Доказательства:** факты, примеры, бенчмарки ### 💰 Стоимость внедрения - **Время:** оценка времени внедрения - **Ресурсы:** что потребуется - **Обучение:** нужно ли обучение команды ### 🎯 Идеальный use case (В каких ситуациях этот вариант лучший выбор) ### 🚫 Анти-patterns (В каких ситуациях этот вариант НЕ подходит) ### 📚 Источники и референсы (Ссылки на документацию, кейсы, бенчмарки) Отвечай на русском языке. Будь максимально конкретным и объективным.""" ANALYST_COMPARE_PROMPT = """Ты — аналитик, проводящий **сравнительный анализ** нескольких вариантов по единой методологии. Твоя задача — создать **объективную матрицу сравнения** с оценками и выявить паттерны, которые не очевидны при рассмотрении опций по отдельности. ## Твоя роль: - Построить сравнительную матрицу - Выявить корреляции и trade-offs - Определить явных лидеров и аутсайдеров по критериям - Найти неочевидные зависимости ## Принципы работы: 1. **Единая шкала** — все опции оцениваются одинаково 2. **Взвешивание** — учитывай важность каждого критерия 3. **Контекстность** — учитывай контекст задачи 4. **Объективность** — опирайся на факты, не мнения ## Формат ответа: ### 📊 Матрица сравнения | Критерий (вес) | Опция 1 | Опция 2 | Опция 3 | ... | |----------------|---------|---------|---------|-----| | Критерий 1 (10) | X/10 | X/10 | X/10 | ... | | Критерий 2 (8) | X/10 | X/10 | X/10 | ... | | ... | ... | ... | ... | ... | | **ИТОГО** | **Y** | **Y** | **Y** | ... | ### 🏆 Рейтинг опций (Суммарные баллы с учётом весов критериев) ### 📈 Паттерны и корреляции (Что интересного видно в сравнении) ### ⚖️ Ключевые trade-offs Для каждой пары опций: - **Опция A vs Опция B:** в чём основной trade-off - **Когда выбирать A:** ... - **Когда выбирать B:** ... ### 🎯 Лидеры по критериям (Какая опция лучшая для каждого критерия) ### 🔍 Неочевидные инсайты (Что стало понятно только в сравнении) ### 📊 Визуализация сильных сторон (Какая опция в чём доминирует) Отвечай на русском языке. Будь максимально объективным и структурированным.""" CRITIC_TRADEOFF_PROMPT = """Ты — критик, специализирующийся на выявлении **скрытых trade-offs и подводных камней** в сравнительных анализах. Твоя задача — найти то, что упущено в формальном сравнении: скрытые costs, долгосрочные последствия, vendor lock-in, etc. ## Твоя роль: - Выявлять скрытые costs и последствия - Находить vendor lock-in и зависимости - Предвидеть долгосрочные implications - Ставить под сомнение очевидные выводы ## Формат ответа: ### 💸 Скрытые costs каждого варианта Для каждой опции: - **Опция X:** - Скрытый cost 1: ... - Скрытый cost 2: ... ### 🔒 Vendor lock-in и зависимости (Насколько сложно будет переключиться на другой вариант в будущем) ### 🕰️ Долгосрочные последствия (1-3-5 лет) (Как каждый вариант будет вести себя со временем) ### 🎭 Marketing vs Reality (Где заявленные преимущества не соответствуют действительности) ### ⚠️ Типичные ошибки при выборе (Какие ошибки обычно совершают при выборе между этими опциями) ### 🔄 Сценарии миграции (Что потребуется для переключения между опциями) ### 🎯 Критические вопросы перед выбором (Что обязательно нужно прояснить до принятия решения) Отвечай на русском языке. Будь конструктивным скептиком.""" JUDGE_RECOMMENDATION_PROMPT = """Ты — судья, выносящий **финальную рекомендацию** по выбору из нескольких вариантов. Твоя задача — принять обоснованное решение и дать **практическую рекомендацию**, учитывая все проведённые анализы. ## Принципы работы: 1. **Контекстность** — рекомендация должна учитывать конкретную ситуацию 2. **Нюансированность** — не "один размер для всех", а "для X лучше Y" 3. **Практичность** — чёткие инструкции по внедрению 4. **Честность** — признавать неопределённость и условия ## Формат ответа: ### 🏆 Основная рекомендация **Выбор:** (название опции) **Обоснование:** (почему именно этот вариант) ### 🎯 Альтернативные сценарии - **Если [условие A]:** выбирай [опция X] - **Если [условие B]:** выбирай [опция Y] - **Если [условие C]:** выбирай [опция Z] ### 📊 Итоговый рейтинг с обоснованием 1. **Опция X** (балл) — почему на этом месте 2. **Опция Y** (балл) — почему на этом месте 3. **Опция Z** (балл) — почему на этом месте ### 🚀 План внедрения рекомендуемого варианта 1. **Шаг 1:** ... (срок, ресурсы) 2. **Шаг 2:** ... (срок, ресурсы) 3. **Шаг 3:** ... (срок, ресурсы) ### ⚠️ Критические риски и митигации | Риск | Вероятность | Митигация | |------|-------------|-----------| ### 🔄 Exit strategy (План отката/миграции если выбор окажется неверным) ### 📈 Метрики успеха (Как понять через 6 месяцев что выбор был правильный) ### ❓ Когда пересматривать решение (Условия для возврата к сравнению) Отвечай на русском языке. Будь решительным и практичным.""" WRITER_COMPARE_PROMPT = """Ты — технический писатель, формирующий **финальный сравнительный отчёт** для принятия решений. Твоя задача — синтезировать весь процесс сравнения в единый документ, который можно использовать для презентации стейкхолдерам и принятия решений. ## Формат итогового отчёта: ### 📋 Executive Summary (3-5 предложений: суть выбора, рекомендуемый вариант, ключевые причины) ### 🎯 Задача и контекст выбора (Что выбираем и зачем) ### 📊 Краткая матрица сравнения | Критерий | Опция 1 | Опция 2 | Опция 3 | |----------|---------|---------|---------| ### 🏆 Итоговый рейтинг 1. **Опция X** — краткое обоснование 2. **Опция Y** — краткое обоснование 3. **Опция Z** — краткое обоснование ### 💪 Детальный анализ каждого варианта #### Вариант 1: [Название] - **Описание:** ... - **Плюсы:** ... - **Минусы:** ... - **Идеален для:** ... #### Вариант 2: [Название] ... ### ⚖️ Ключевые trade-offs (Когда какой вариант выбирать) ### 🎯 Финальная рекомендация **Выбор:** ... **Обоснование:** ... **Условия применения:** ... ### 🚀 План внедрения 1. Шаг 1 2. Шаг 2 3. Шаг 3 ### ⚠️ Риски и митигации (Что может пойти не так и как это предотвратить) ### 📈 Метрики успеха (Как измерить правильность выбора) ### 🔄 Exit strategy (План отката если выбор окажется неверным) ### 📚 Дополнительная информация (Ссылки, документация, кейсы) Отвечай на русском языке. Создай документ уровня executive report.""" # ============================================================================ # Helpers # ============================================================================ def _extract_options(input_text: str, explicit_options: list[str] | None) -> list[str]: """ Извлекает список опций для сравнения из input. Поддерживает несколько форматов: 1. Явный список: ["PostgreSQL", "MongoDB", "Cassandra"] 2. Через "vs": "PostgreSQL vs MongoDB vs Cassandra" 3. Через запятые с маркерами: "Сравнить: PostgreSQL, MongoDB, Cassandra" 4. Маркированный список в тексте """ # Приоритет 1: явный список if explicit_options: return [opt.strip() for opt in explicit_options if opt.strip()] # Приоритет 2: разделитель "vs" if " vs " in input_text.lower() or " против " in input_text.lower(): separator = " vs " if " vs " in input_text.lower() else " против " parts = re.split(re.escape(separator), input_text, flags=re.IGNORECASE) if len(parts) >= 2: # Берём первые слова каждой части (без контекста) options = [] for part in parts: # Берём первую строку или предложение first_line = part.strip().split("\n")[0].split(".")[0].split(",")[0] options.append(first_line.strip()) return [opt for opt in options if opt] # Приоритет 3: "Сравнить:" или "Варианты:" for marker in ["Сравнить:", "Варианты:", "Options:", "Compare:", "Опции:"]: if marker.lower() in input_text.lower(): idx = input_text.lower().find(marker.lower()) after_marker = input_text[idx + len(marker):] # Берём первую строку после маркера first_line = after_marker.strip().split("\n")[0] # Разбиваем по запятым или "и" options = re.split(r'[,;]| и ', first_line, flags=re.IGNORECASE) options = [opt.strip() for opt in options if opt.strip()] if len(options) >= 2: return options # Приоритет 4: маркированный список (- или *) bullet_pattern = re.compile(r'^\s*[-*]\s+(.+)$', re.MULTILINE) bullets = bullet_pattern.findall(input_text) if len(bullets) >= 2: return [b.strip() for b in bullets[:10]] # Fallback: возвращаем пустой список, Planner сам определит return [] async def _run_agent_stream( agent, agent_id: str, input_data: dict[str, Any], ) -> AsyncIterator[tuple[str, str, int, str]]: """ Запускает агента и yield'ит стандартизированные события. """ full_output = "" total_tokens = 0 model_name = "" async for chunk in agent.run_stream(input_data): chunk_type = chunk.get("type") if chunk_type == "content": content = chunk.get("content", "") full_output += content yield ("content", content, 0, "") elif chunk_type == "done": total_tokens = chunk.get("tokens_total", 0) model_name = chunk.get("model", "") yield ("done", full_output, total_tokens, model_name) elif chunk_type == "error": error_msg = chunk.get("error", f"Unknown error in {agent_id}") yield ("error", error_msg, 0, "") return if full_output and total_tokens == 0: yield ("done", full_output, 0, "") class EventQueue: """Очередь событий для параллельных задач.""" def __init__(self): self._queue: asyncio.Queue[dict[str, Any]] = asyncio.Queue() async def put(self, event: dict[str, Any]) -> None: await self._queue.put(event) async def get(self) -> dict[str, Any]: return await self._queue.get() @property def empty(self) -> bool: return self._queue.empty() # ============================================================================ # Main runner # ============================================================================ async def run_compare_flow(input_data: dict[str, Any]) -> AsyncIterator[dict[str, Any]]: """ Выполняет мультиагентный compare flow. Flow состоит из шести этапов: 1. PlannerAgent — методология сравнения 2. Parallel Options Research — параллельное исследование каждой опции 3. AnalystAgent — сравнительная матрица 4. CriticAgent — trade-off анализ 5. JudgeAgent — финальная рекомендация 6. WriterAgent — сравнительный отчёт Args: input_data: Входные данные. Поддерживает поля: - 'input' / 'topic' / 'question': задача для сравнения - 'options': явный список опций (опционально) - 'criteria': пользовательские критерии (опционально) - 'context': контекст выбора (опционально) - 'constraints': ограничения (опционально) Yields: События SSE для UI """ # ======================================================================== # 0. Валидация и извлечение опций # ======================================================================== raw_input = ( input_data.get("input") or input_data.get("topic") or input_data.get("question") or "" ).strip() if not raw_input: yield { "type": "error", "error": "Не указана задача для сравнения (ожидается поле 'input', 'topic' или 'question')", } return explicit_options = input_data.get("options") if isinstance(explicit_options, str): # Поддержка строкового формата: "A, B, C" explicit_options = [o.strip() for o in explicit_options.split(",") if o.strip()] options = _extract_options(raw_input, explicit_options) criteria = input_data.get("criteria", "").strip() context = input_data.get("context", "").strip() constraints = input_data.get("constraints", "").strip() logger.info( "compare_flow.started", input_length=len(raw_input), options_count=len(options), options_preview=options[:3] if options else None, has_criteria=bool(criteria), has_context=bool(context), ) stage_outputs: dict[str, Any] = { "input": raw_input, "options": options, "criteria": criteria, "context": context, "constraints": constraints, } total_tokens = 0 # ======================================================================== # Шаг 1: PlannerAgent — методология сравнения # ======================================================================== yield {"type": "agent_start", "agent": "planner"} planner = PlannerAgent(name="Comparison Planner") planner_input = { "topic": f"""Разработай методологию сравнительного анализа. ## ЗАДАЧА: \"\"\" {raw_input} \"\"\" {f"## ВАРИАНТЫ ДЛЯ СРАВНЕНИЯ:\n" + chr(10).join(f"- {opt}" for opt in options) if options else "## ВАРИАНТЫ: определи самостоятельно из контекста"} {f"## ПОЛЬЗОВАТЕЛЬСКИЕ КРИТЕРИИ:\n{criteria}" if criteria else ""} {f"## КОНТЕКСТ:\n{context}" if context else ""} {f"## ОГРАНИЧЕНИЯ:\n{constraints}" if constraints else ""} Создай методологию сравнения, которая: 1. Определит релевантные критерии (если не указаны) 2. Задаст шкалу оценок 3. Учтёт контекст задачи 4. Выявит опции для сравнения (если не указаны явно) Отвечай на русском языке.""", } planner_output = "" try: async for event_type, content, tokens, model in _run_agent_stream( planner, "planner", planner_input ): if event_type == "content": planner_output += content yield {"type": "agent_message", "agent": "planner", "content": content} elif event_type == "done": total_tokens += tokens yield { "type": "agent_done", "agent": "planner", "output": content, "tokens": tokens, "model": model, } elif event_type == "error": yield {"type": "error", "error": content, "agent": "planner"} return except Exception as e: logger.error("compare_flow.planner.failed", error=str(e), exc_info=True) yield {"type": "error", "error": f"PlannerAgent failed: {e}", "agent": "planner"} return stage_outputs["planning"] = planner_output # Если опции не были определены явно — попробуем извлечь из planner output if not options: # Простая эвристика: ищем маркированный список в planner output bullet_pattern = re.compile(r'^\s*[-*]\s*(.+)$', re.MULTILINE) bullets = bullet_pattern.findall(planner_output) # Берём первые 5-7 пунктов как опции options = [b.strip() for b in bullets[:7] if len(b.strip()) > 3 and len(b.strip()) < 100] if not options: yield { "type": "error", "error": "Не удалось определить варианты для сравнения. Укажи их явно в поле 'options' или через 'vs'.", } return stage_outputs["options"] = options logger.info("compare_flow.options.determined", options=options, count=len(options)) # ======================================================================== # Шаг 2: Параллельное исследование всех опций # ======================================================================== event_queue = EventQueue() options_outputs: dict[str, str] = {} options_tokens: dict[str, int] = {} options_errors: list[str] = [] async def yield_event(event: dict[str, Any]) -> None: await event_queue.put(event) async def run_option_researcher(option_name: str, option_idx: int) -> None: """Запускает исследование одной опции.""" agent_id = f"option_{option_idx + 1}_{option_name[:20].replace(' ', '_')}" researcher = ResearcherAgent( name=f"Researcher: {option_name}", system_prompt=RESEARCHER_OPTION_PROMPT_TEMPLATE.format( option_name=option_name, task_context=raw_input + (f"\n\nКонтекст: {context}" if context else ""), criteria=planner_output, ), ) input_data_agent = { "topic": f"""Проведи глубокий анализ варианта: "{option_name}" Контекст задачи: {raw_input} {f"Дополнительный контекст: {context}" if context else ""} Используй методологию от планировщика для структурированного анализа. Будь максимально конкретным и объективным.""", } await yield_event({"type": "agent_start", "agent": agent_id}) output = "" tokens = 0 try: async for event_type, content, event_tokens, model in _run_agent_stream( researcher, agent_id, input_data_agent ): if event_type == "content": output += content await yield_event({ "type": "agent_message", "agent": agent_id, "content": content, }) elif event_type == "done": tokens = event_tokens await yield_event({ "type": "agent_done", "agent": agent_id, "output": content, "tokens": tokens, "model": model, }) elif event_type == "error": await yield_event({ "type": "error", "error": content, "agent": agent_id, }) options_errors.append(option_name) return except Exception as e: logger.error( "compare_flow.option_research.failed", option=option_name, error=str(e), exc_info=True, ) await yield_event({ "type": "error", "error": f"Research for '{option_name}' failed: {e}", "agent": agent_id, }) options_errors.append(option_name) return options_outputs[option_name] = output options_tokens[option_name] = tokens # Запускаем исследование всех опций параллельно option_tasks = [ asyncio.create_task(run_option_researcher(opt, idx)) for idx, opt in enumerate(options) ] # Yield'им события по мере поступления completed_count = 0 while completed_count < len(option_tasks): done_tasks = [t for t in option_tasks if t.done()] completed_count = len(done_tasks) try: event = await asyncio.wait_for(event_queue.get(), timeout=0.1) yield event if event.get("type") == "agent_done": total_tokens += event.get("tokens", 0) except asyncio.TimeoutError: await asyncio.sleep(0.05) continue # Ждём завершения всех задач await asyncio.gather(*option_tasks, return_exceptions=True) # Добираем оставшиеся события while not event_queue.empty: try: event = event_queue._queue.get_nowait() yield event if event.get("type") == "agent_done": total_tokens += event.get("tokens", 0) except asyncio.QueueEmpty: break if options_errors: yield { "type": "error", "error": f"Не удалось исследовать опции: {', '.join(options_errors)}", } return # Мета-событие: все опции исследованы yield { "type": "options_researched", "options_count": len(options), "options": list(options_outputs.keys()), "options_tokens": options_tokens, } stage_outputs["options_research"] = options_outputs # ======================================================================== # Шаг 3: AnalystAgent — сравнительная матрица # ======================================================================== yield {"type": "agent_start", "agent": "analyst"} # Формируем сводку исследований research_summary = "\n\n".join( f"### {opt}:\n\"\"\"\n{output}\n\"\"\"" for opt, output in options_outputs.items() ) analyst = AnalystAgent(name="Comparison Analyst") analyst_input = { "topic": f"""Проведи сравнительный анализ вариантов. ## ЗАДАЧА: \"\"\" {raw_input} \"\"\" ## МЕТОДОЛОГИЯ СРАВНЕНИЯ: \"\"\" {planner_output} \"\"\" ## РЕЗУЛЬТАТЫ ИССЛЕДОВАНИЯ КАЖДОГО ВАРИАНТА: {research_summary} {f"Контекст: {context}" if context else ""} Построй объективную сравнительную матрицу по всем критериям и выяви ключевые trade-offs между вариантами.""", } analyst_output = "" try: async for event_type, content, tokens, model in _run_agent_stream( analyst, "analyst", analyst_input ): if event_type == "content": analyst_output += content yield {"type": "agent_message", "agent": "analyst", "content": content} elif event_type == "done": total_tokens += tokens yield { "type": "agent_done", "agent": "analyst", "output": content, "tokens": tokens, "model": model, } elif event_type == "error": yield {"type": "error", "error": content, "agent": "analyst"} return except Exception as e: logger.error("compare_flow.analyst.failed", error=str(e), exc_info=True) yield {"type": "error", "error": f"AnalystAgent failed: {e}", "agent": "analyst"} return stage_outputs["analysis"] = analyst_output # ======================================================================== # Шаг 4: CriticAgent — trade-off анализ # ======================================================================== yield {"type": "agent_start", "agent": "critic"} critic = CriticAgent(name="Trade-off Critic") critic_input = { "topic": f"""Выяви скрытые trade-offs и подводные камни в сравнении. ## ЗАДАЧА: \"\"\" {raw_input} \"\"\" ## ВАРИАНТЫ: {chr(10).join(f"- {opt}" for opt in options)} ## СРАВНИТЕЛЬНЫЙ АНАЛИЗ: \"\"\" {analyst_output} \"\"\" {f"Ограничения: {constraints}" if constraints else ""} Найди то, что упущено в формальном сравнении: скрытые costs, vendor lock-in, долгосрочные последствия.""", } critic_output = "" try: async for event_type, content, tokens, model in _run_agent_stream( critic, "critic", critic_input ): if event_type == "content": critic_output += content yield {"type": "agent_message", "agent": "critic", "content": content} elif event_type == "done": total_tokens += tokens yield { "type": "agent_done", "agent": "critic", "output": content, "tokens": tokens, "model": model, } elif event_type == "error": yield {"type": "error", "error": content, "agent": "critic"} return except Exception as e: logger.error("compare_flow.critic.failed", error=str(e), exc_info=True) yield {"type": "error", "error": f"CriticAgent failed: {e}", "agent": "critic"} return stage_outputs["tradeoffs"] = critic_output # ======================================================================== # Шаг 5: JudgeAgent — финальная рекомендация # ======================================================================== yield {"type": "agent_start", "agent": "judge"} judge = JudgeAgent(name="Decision Judge") judge_input = { "topic": f"""Вынеси финальную рекомендацию по выбору варианта. ## ЗАДАЧА: \"\"\" {raw_input} \"\"\" ## ВАРИАНТЫ: {chr(10).join(f"- {opt}" for opt in options)} ## СРАВНИТЕЛЬНЫЙ АНАЛИЗ: \"\"\" {analyst_output} \"\"\" ## TRADE-OFF АНАЛИЗ: \"\"\" {critic_output} \"\"\" {f"Контекст: {context}" if context else ""} {f"Ограничения: {constraints}" if constraints else ""} Прими обоснованное решение с учётом всех анализов. Дай практичную рекомендацию с планом внедрения.""", } judge_output = "" try: async for event_type, content, tokens, model in _run_agent_stream( judge, "judge", judge_input ): if event_type == "content": judge_output += content yield {"type": "agent_message", "agent": "judge", "content": content} elif event_type == "done": total_tokens += tokens yield { "type": "agent_done", "agent": "judge", "output": content, "tokens": tokens, "model": model, } elif event_type == "error": yield {"type": "error", "error": content, "agent": "judge"} return except Exception as e: logger.error("compare_flow.judge.failed", error=str(e), exc_info=True) yield {"type": "error", "error": f"JudgeAgent failed: {e}", "agent": "judge"} return stage_outputs["recommendation"] = judge_output # ======================================================================== # Шаг 6: WriterAgent — финальный сравнительный отчёт # ======================================================================== yield {"type": "agent_start", "agent": "writer"} writer = WriterAgent(name="Comparison Report Writer") # Формируем краткую сводку исследований для финального отчёта options_summary = "\n\n".join( f"**{opt}:** {output[:500]}..." if len(output) > 500 else f"**{opt}:** {output}" for opt, output in options_outputs.items() ) writer_input = { "topic": f"""Сформируй финальный сравнительный отчёт. ## ЗАДАЧА: \"\"\" {raw_input} \"\"\" {f"Контекст: {context}" if context else ""} {f"Ограничения: {constraints}" if constraints else ""} ## ВАРИАНТЫ: {chr(10).join(f"- {opt}" for opt in options)} ## РЕЗУЛЬТАТЫ ИССЛЕДОВАНИЯ (краткая сводка): {options_summary} ## СРАВНИТЕЛЬНЫЙ АНАЛИЗ: \"\"\" {analyst_output} \"\"\" ## TRADE-OFF АНАЛИЗ: \"\"\" {critic_output} \"\"\" ## ФИНАЛЬНАЯ РЕКОМЕНДАЦИЯ: \"\"\" {judge_output} \"\"\" Создай профессиональный сравнительный отчёт уровня executive report, который можно использовать для презентации стейкхолдерам.""", } writer_output = "" try: async for event_type, content, tokens, model in _run_agent_stream( writer, "writer", writer_input ): if event_type == "content": writer_output += content yield {"type": "agent_message", "agent": "writer", "content": content} elif event_type == "done": total_tokens += tokens yield { "type": "agent_done", "agent": "writer", "output": content, "tokens": tokens, "model": model, } elif event_type == "error": yield {"type": "error", "error": content, "agent": "writer"} return except Exception as e: logger.error("compare_flow.writer.failed", error=str(e), exc_info=True) yield {"type": "error", "error": f"WriterAgent failed: {e}", "agent": "writer"} return stage_outputs["report"] = writer_output # ======================================================================== # Финальное событие flow # ======================================================================== logger.info( "compare_flow.completed", total_tokens=total_tokens, options_count=len(options), stages_count=len(stage_outputs), ) yield { "type": "flow_done", "flow_id": "compare", "output": writer_output, "tokens": total_tokens, "stages": stage_outputs, "metadata": { "options_count": len(options), "options": options, "agents_count": 5 + len(options), # planner + analysts + critics + judge + writer + N researchers "total_stages": 6, "parallel_stage": "options_research", }, } # ============================================================================ # Flow registration # ============================================================================ compare_flow = Flow( id="compare", name="Compare Flow", description=( "Мультиагентное сравнение вариантов для принятия решений: " "Planner → Parallel Options Research (N × Researcher) → " "Analyst → Critic → Judge → Writer. " "Проводит глубокий сравнительный анализ нескольких вариантов решения " "одной задачи с построением матрицы сравнения, trade-off анализом " "и финальной рекомендацией. Идеально подходит для выбора технологий, " "вендоров, архитектурных подходов и бизнес-стратегий." ), agents=[ "planner", "option_researcher", # Параллельно для каждой опции "analyst", "critic", "judge", "writer", ], runner=run_compare_flow, ) # ============================================================================ # Public API # ============================================================================ __all__ = ["compare_flow", "run_compare_flow"]