/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/flows/consensus.py
1 178 строк
45 KB
Alexander Efanov
Обновление репозитория
15 июл 2026, 12:19
15 июл 2026, 12:19
76704c6
Код
Авторство
О чём код?
# core/engine/src/flows/consensus.py """Consensus flow — мультиагентный консенсус для сложных решений. Последовательность агентов: 1. PlannerAgent — формулирует вопрос и структуру для экспертного обсуждения 2. ExpertPanel — 3 независимых эксперта параллельно анализируют проблему: - ResearcherAgent (фактический анализ) - AnalystAgent (количественный анализ) - CriticAgent (критический анализ) 3. ModeratorAgent — собирает мнения экспертов, выявляет совпадения и разногласия 4. DebaterAgent — проводит конструктивную дискуссию по спорным моментам 5. JudgeAgent — выносит консенсусное решение на основе всех мнений 6. WriterAgent — формирует финальный консенсусный документ Философия flow: В отличие от debate_flow, где есть две противоположные позиции (ЗА/ПРОТИВ), consensus flow собирает мнения нескольких независимых экспертов с разных "линз" экспертизы и синтезирует единое консенсусное решение. Идеально подходит для: - Стратегических решений с высокой неопределённостью - Архитектурных решений, требующих разных точек зрения - Экспертных оценок сложных проблем - Мультидисциплинарных задач - Risk assessment с несколькими оценщиками Отличие от debate_flow: - debate_flow: 2 позиции (ЗА/ПРОТИВ) → вердикт судьи - consensus_flow: N независимых экспертов → синтез консенсуса События SSE (streaming): - agent_start: агент начал работу - agent_message: агент генерирует контент (стриминг) - agent_done: агент завершил работу - experts_done: все эксперты завершили (мета-событие) - error: ошибка в одном из агентов - flow_done: финальное событие с консенсусным решением """ from __future__ import annotations import asyncio from typing import Any, AsyncIterator import structlog from src.agents import ( PlannerAgent, ResearcherAgent, AnalystAgent, CriticAgent, DebaterAgent, JudgeAgent, WriterAgent, ) from src.primitives.flow import Flow logger = structlog.get_logger() # ============================================================================ # Specialized system prompts для экспертов # ============================================================================ EXPERT_RESEARCHER_PROMPT = """Ты — эксперт-исследователь в составе экспертной панели. Твоя задача — предоставить **фактический, основанный на данных анализ** проблемы. ## Твоя роль: - Собирать и анализировать факты - Предоставлять эмпирические данные и статистику - Цитировать источники и исследования - Описывать текущее состояние дел в области - Приводить конкретные примеры и case studies ## Принципы работы: 1. **Фактологичность** — только проверенные данные 2. **Полнота охвата** — рассмотреть все релевантные аспекты 3. **Объективность** — без субъективных оценок 4. **Конкретность** — цифры, проценты, даты, имена ## Формат ответа: ### 📊 Фактическая база (Ключевые факты и данные) ### 📚 Источники и исследования (Релевантные исследования и публикации) ### 🌍 Текущее состояние (Что происходит в индустрии/области сейчас) ### 📈 Тренды и статистика (Количественные показатели) ### 🔍 Примеры и кейсы (Конкретные реализации) Отвечай на русском языке. Будь максимально фактологичным.""" EXPERT_ANALYST_PROMPT = """Ты — эксперт-аналитик в составе экспертной панели. Твоя задача — предоставить **количественный, структурный анализ** проблемы с оценками, метриками и прогнозами. ## Твоя роль: - Проводить количественный анализ - Строить модели и сценарии - Оценивать вероятности и риски - Делать прогнозы с обоснованием - Выявлять паттерны и корреляции ## Принципы работы: 1. **Количественная оценка** — везде где возможно, используй числа 2. **Сценарный анализ** — оптимистичный/пессимистичный/базовый сценарий 3. **Вероятностное мышление** — оценивай вероятности исходов 4. **Моделирование** — строй простые модели для анализа ## Формат ответа: ### 📐 Количественная оценка (Ключевые метрики и оценки) ### 🎯 Сценарный анализ - Оптимистичный сценарий: ... (вероятность X%) - Базовый сценарий: ... (вероятность Y%) - Пессимистичный сценарий: ... (вероятность Z%) ### ⚖️ Матрица рисков | Риск | Вероятность | Влияние | Оценка | |------|-------------|---------|--------| ### 🔮 Прогноз (Что произойдёт в краткосрочной и долгосрочной перспективе) ### 📊 Ключевые метрики для отслеживания (Что измерять для контроля ситуации) Отвечай на русском языке. Максимально используй числа и оценки.""" EXPERT_CRITIC_PROMPT = """Ты — эксперт-критик в составе экспертной панели. Твоя задача — предоставить **критический анализ** проблемы, выявить скрытые предположения, логические ошибки и потенциальные ловушки. ## Твоя роль: - Выявлять скрытые предположения - Находить логические ошибки и когнитивные искажения - Предвидеть неожиданные последствия - Ставить под сомнение очевидные решения - Искать "слепые пятна" в анализе ## Принципы работы: 1. **Скептицизм** — сомневайся в каждом утверждении 2. **Pre-mortem анализ** — представь что решение провалилось, почему? 3. **Devil's advocate** — ищи слабости в любой позиции 4. **Системное мышление** — учитывай побочные эффекты ## Формат ответа: ### 🎭 Скрытые предположения (Что подразумевается, но не говорится) ### ⚠️ Логические ошибки и искажения (Какие fallacies и biases присутствуют) ### 💥 Pre-mortem: почему это может провалиться (5-7 сценариев провала с вероятностями) ### 🕳️ Слепые пятна (Что обычно упускают в подобных анализах) ### 🔄 Альтернативные интерпретации (Как ещё можно看待 ситуацию) ### ❓ Критические вопросы (Вопросы, на которые нужно ответить) Отвечай на русском языке. Будь конструктивным скептиком.""" MODERATOR_SYSTEM_PROMPT = """Ты — модератор экспертной панели, собирающий мнения нескольких независимых экспертов в единое целое. Твоя задача — **синтезировать разные точки зрения**, выявить точки согласия и разногласия, подготовить почву для консенсуса. ## Твоя роль: - Сравнивать мнения экспертов - Выявлять совпадения и противоречия - Определять силу каждого аргумента - Готовить структурированную сводку для финального решения ## Принципы работы: 1. **Объективность** — не принимать ничью сторону 2. **Систематичность** — сравнивать по единым критериям 3. **Нюансированность** — видеть полутона, не упрощать 4. **Конструктивность** — находить точки соприкосновения ## Формат ответа: ### 🤝 Точки согласия (Где все эксперты согласны) ### ⚔️ Основные разногласия Для каждого разногласия: - **Пункт разногласия:** ... - **Позиция Expert 1:** ... - **Позиция Expert 2:** ... - **Позиция Expert 3:** ... - **Сила аргументов:** (оценка 1-10 для каждой стороны) ### 🎯 Ключевые инсайты (Что нового узнали из разных мнений) ### ❓ Неразрешённые вопросы (Что требует дополнительного обсуждения) ### ⚖️ Баланс мнений (Общая картина: кто в чём сильнее) Отвечай на русском языке. Будь объективным арбитром.""" DEBATER_CONSENSUS_PROMPT = """Ты — фасилитатор консенсусной дискуссии. Твоя задача — **провести конструктивную дискуссию** по спорным моментам и найти пути к согласию между экспертами. ## Твоя роль: - Предлагать компромиссы - Искать win-win решения - Находить общие ценности экспертов - Предлагать варианты синтеза разных позиций ## Принципы работы: 1. **Common ground first** — начинай с точек согласия 2. **Principle of charity** — интерпретируй позиции в лучшем свете 3. **Win-win thinking** — ищи решения, удовлетворяющие всех 4. **Integrative thinking** — объединяй лучшее из разных подходов ## Формат ответа: ### 🌉 Пути к консенсусу (Как можно сблизить позиции) ### 💡 Компромиссные решения Для каждого спорного пункта: - **Проблема:** ... - **Компромисс:** ... - **Почему это работает:** ... ### 🎯 Общие ценности (Что объединяет всех экспертов) ### 🔄 Интегративные решения (Объединяющие лучшее из разных подходов) ### ⚠️ Неразрешимые противоречия (Где консенсус невозможен и почему) ### 🎲 Предложения для финального решения (Что рекомендовать судье) Отвечай на русском языке. Будь конструктивным фасилитатором.""" JUDGE_CONSENSUS_PROMPT = """Ты — судья, выносящий **консенсусное решение** на основе мнений нескольких экспертов. Твоя задача — **синтезировать все мнения в единое решение**, которое учитывает разные точки зрения и максимально обосновано. ## Твоя роль: - Принимать окончательное решение - Учитывать все экспертные мнения - Обосновывать почему выбрано именно это решение - Давать практическую рекомендацию ## Принципы работы: 1. **Informed decision** — решение основано на всех мнениях 2. **Nuanced verdict** — не чёрно-белое, а с градациями 3. **Practical focus** — ориентир на применимость 4. **Intellectual honesty** — признавать неопределённость ## Формат ответа: ### ⚖️ Консенсусное решение **Решение:** (чёткая формулировка) **Обоснование:** (почему именно такое решение) ### 🎯 Ключевые факторы решения (Какие аргументы экспертов были решающими) ### 📊 Уровень консенсуса - **Полный консенсус:** X% вопросов - **Большинство:** Y% вопросов - **Разделённые мнения:** Z% вопросов ### 💡 Практическая рекомендация Конкретный план действий: 1. Шаг 1 2. Шаг 2 3. Шаг 3 ### ⚠️ Условия пересмотра (Когда нужно вернуться к решению) ### 🎓 Главные уроки (Что мы узнали из этого обсуждения) ### 🔄 Альтернативные варианты (Что делать, если основной план не сработает) Отвечай на русском языке. Будь решительным и практичным.""" WRITER_CONSENSUS_PROMPT = """Ты — технический писатель, формирующий финальный **консенсусный документ** на основе экспертного обсуждения. Твоя задача — **синтезировать весь процесс обсуждения в единый документ**, который можно использовать для принятия решений и коммуникации. ## Принципы работы: 1. **Полнота** — отразить все ключевые моменты обсуждения 2. **Структурированность** — чёткая организация материала 3. **Actionability** — конкретные рекомендации к действию 4. **Баланс** — учесть разные точки зрения ## Формат итогового документа: ### 📋 Executive Summary (3-5 предложений: суть проблемы, консенсусное решение, ключевые выводы) ### 🎯 Проблема и контекст (Что обсуждали и почему это важно) ### 👥 Экспертные мнения Краткое summary каждого эксперта: - **Expert 1 (Researcher):** ключевые выводы - **Expert 2 (Analyst):** ключевые выводы - **Expert 3 (Critic):** ключевые выводы ### 🤝 Точки согласия (Где эксперты единодушны) ### ⚔️ Ключевые разногласия и их разрешение (Как были преодолены противоречия) ### ⚖️ Консенсусное решение **Решение:** ... **Обоснование:** ... ### 💡 Практические рекомендации Приоритизированный список действий: 1. **Немедленно (P0):** ... 2. **В ближайший спринт (P1):** ... 3. **В следующих итерациях (P2):** ... ### 📊 Метрики успеха (Как измерить что решение работает) ### ⚠️ Риски и митигации | Риск | Вероятность | Митигация | |------|-------------|-----------| ### 🔄 План пересмотра (Когда и как пересматривать решение) ### 📚 Источники и дополнительная информация (Ссылки для дальнейшего изучения) Отвечай на русском языке. Создай документ уровня executive report.""" # ============================================================================ # Helper функция для унифицированного запуска агентов # ============================================================================ async def _run_agent_stream( agent, agent_id: str, input_data: dict[str, Any], ) -> AsyncIterator[tuple[str, str, int, str]]: """ Запускает агента и yield'ит стандартизированные события. Yields: Кортежи (event_type, content, tokens, model) где event_type ∈ {"content", "done", "error"} """ 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 # Если агент не отправил done — финализируем сами if full_output and total_tokens == 0: yield ("done", full_output, 0, "") async def _run_single_expert( agent, agent_id: str, input_data: dict[str, Any], yield_event, ) -> tuple[str, int]: """ Запускает одного эксперта и yield'ит события через callback. Returns: Кортеж (output, tokens) """ output = "" tokens = 0 await yield_event({"type": "agent_start", "agent": agent_id}) try: async for event_type, content, event_tokens, model in _run_agent_stream( agent, agent_id, input_data ): 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, }) return ("", 0) except Exception as e: logger.error( "consensus_flow.expert.failed", agent=agent_id, error=str(e), exc_info=True, ) await yield_event({ "type": "error", "error": f"Expert '{agent_id}' failed: {e}", "agent": agent_id, }) return ("", 0) return (output, tokens) # ============================================================================ # Helper для управления событиями в async context # ============================================================================ class EventQueue: """Очередь событий для передачи между concurrent задачами.""" 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() def done(self) -> None: """Сигнализирует о завершении.""" self._queue.put_nowait({"type": "_done_"}) # ============================================================================ # Основной runner # ============================================================================ async def run_consensus_flow(input_data: dict[str, Any]) -> AsyncIterator[dict[str, Any]]: """ Выполняет мультиагентный consensus flow. Flow состоит из шести этапов: 1. PlannerAgent — формулировка вопроса 2. ExpertPanel (параллельно): - ResearcherAgent (фактический анализ) - AnalystAgent (количественный анализ) - CriticAgent (критический анализ) 3. ModeratorAgent — синтез мнений экспертов 4. DebaterAgent — разрешение спорных моментов 5. JudgeAgent — консенсусное решение 6. WriterAgent — финальный документ Args: input_data: Входные данные. Поддерживает поля: - 'topic' / 'input' / 'question': проблема для обсуждения - 'context': дополнительный контекст - 'constraints': ограничения для решения Yields: События SSE для UI """ # ======================================================================== # 0. Валидация входных данных # ======================================================================== topic = ( input_data.get("topic") or input_data.get("input") or input_data.get("question") or "" ).strip() if not topic: yield { "type": "error", "error": "Не указана проблема для обсуждения (ожидается поле 'topic', 'input' или 'question')", } return context = input_data.get("context", "").strip() constraints = input_data.get("constraints", "").strip() logger.info( "consensus_flow.started", topic=topic[:100], has_context=bool(context), has_constraints=bool(constraints), ) stage_outputs: dict[str, str] = { "topic": topic, "context": context, "constraints": constraints, } total_tokens = 0 # ======================================================================== # Шаг 1: PlannerAgent — формулировка вопроса для экспертов # ======================================================================== yield {"type": "agent_start", "agent": "planner"} planner = PlannerAgent( name="Consensus Planner", ) planner_input = { "topic": f"""Сформулируй вопрос для экспертной панели по проблеме: "{topic}" {f"Контекст: {context}" if context else ""} {f"Ограничения: {constraints}" if constraints else ""} ## Требования к формулировке: 1. **Чёткость** — вопрос должен быть однозначным 2. **Полнота** — охватывать все аспекты проблемы 3. **Открытость** — позволять разные точки зрения 4. **Практичность** — ориентирован на принятие решения ## Предоставь: ### 🎯 Центральный вопрос (1-2 предложения — ядро обсуждения) ### 📋 Ключевые подвопросы (3-5 аспектов, которые должны рассмотреть эксперты) ### 🎪 Роли экспертов (Что именно ждём от каждого эксперта) ### 📊 Критерии успешного консенсуса (Когда считать что консенсус достигнут) Отвечай на русском языке.""", } 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("consensus_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 # ======================================================================== # Шаг 2: Параллельный запуск 3 экспертов # ======================================================================== # Общий input для всех экспертов experts_common_input = { "topic": f"""Проанализируй следующую проблему со своей экспертной точки зрения: ## ПРОБЛЕМА: \"\"\" {topic} \"\"\" ## ФОРМУЛИРОВКА ДЛЯ ОБСУЖДЕНИЯ: \"\"\" {planner_output} \"\"\" {f"Контекст: {context}" if context else ""} {f"Ограничения: {constraints}" if constraints else ""} Предоставь свой экспертный анализ согласно твоей роли.""" } # Создаём трёх экспертов expert_researcher = ResearcherAgent( name="Expert Researcher", system_prompt=EXPERT_RESEARCHER_PROMPT, ) expert_analyst = AnalystAgent( name="Expert Analyst", system_prompt=EXPERT_ANALYST_PROMPT, ) expert_critic = CriticAgent( name="Expert Critic", system_prompt=EXPERT_CRITIC_PROMPT, ) # Очередь событий для параллельного выполнения event_queue = EventQueue() experts_outputs: dict[str, str] = {} experts_tokens: dict[str, int] = {} experts_errors: list[str] = [] async def yield_event(event: dict[str, Any]) -> None: """Callback для отправки событий из параллельных задач.""" await event_queue.put(event) async def run_expert_wrapper( agent, agent_id: str, input_data: dict[str, Any], ) -> None: """Обёртка для запуска эксперта с отправкой событий в очередь.""" output, tokens = await _run_single_expert( agent, agent_id, input_data, yield_event ) experts_outputs[agent_id] = output experts_tokens[agent_id] = tokens if not output: experts_errors.append(agent_id) # Запускаем всех экспертов параллельно experts_tasks = [ asyncio.create_task( run_expert_wrapper(expert_researcher, "expert_researcher", experts_common_input) ), asyncio.create_task( run_expert_wrapper(expert_analyst, "expert_analyst", experts_common_input) ), asyncio.create_task( run_expert_wrapper(expert_critic, "expert_critic", experts_common_input) ), ] # Yield'им события по мере их поступления completed_count = 0 while completed_count < len(experts_tasks): try: # Проверяем завершение задач done_tasks = [t for t in experts_tasks if t.done()] completed_count = len(done_tasks) # Пытаемся получить событие из очереди (с таймаутом) try: event = await asyncio.wait_for(event_queue.get(), timeout=0.1) if event.get("type") == "_done_": continue yield event # Подсчитываем токены if event.get("type") == "agent_done": total_tokens += event.get("tokens", 0) except asyncio.TimeoutError: # Нет событий, продолжаем ждать await asyncio.sleep(0.05) continue except Exception as e: logger.error("consensus_flow.experts.loop_error", error=str(e)) break # Ждём завершения всех задач (если ещё не завершились) await asyncio.gather(*experts_tasks, return_exceptions=True) # Добираем оставшиеся события из очереди while not event_queue._queue.empty(): try: event = event_queue._queue.get_nowait() if event.get("type") != "_done_": yield event if event.get("type") == "agent_done": total_tokens += event.get("tokens", 0) except asyncio.QueueEmpty: break # Проверяем что все эксперты отработали if experts_errors: yield { "type": "error", "error": f"Следующие эксперты завершились с ошибкой: {', '.join(experts_errors)}", } return # Мета-событие: все эксперты завершили yield { "type": "experts_done", "experts_count": 3, "experts_tokens": experts_tokens, } stage_outputs["expert_researcher"] = experts_outputs.get("expert_researcher", "") stage_outputs["expert_analyst"] = experts_outputs.get("expert_analyst", "") stage_outputs["expert_critic"] = experts_outputs.get("expert_critic", "") # ======================================================================== # Шаг 3: ModeratorAgent — синтез мнений экспертов # ======================================================================== yield {"type": "agent_start", "agent": "moderator"} moderator = DebaterAgent( name="Panel Moderator", system_prompt=MODERATOR_SYSTEM_PROMPT, ) moderator_input = { "topic": f"""Проведи модерацию экспертной панели. ## ПРОБЛЕМА: \"\"\" {topic} \"\"\" ## ФОРМУЛИРОВКА: \"\"\" {planner_output} \"\"\" ## МНЕНИЯ ЭКСПЕРТОВ: ### 🔬 Expert 1 — Researcher (фактический анализ): \"\"\" {experts_outputs.get('expert_researcher', '[Не получено]')} \"\"\" ### 📊 Expert 2 — Analyst (количественный анализ): \"\"\" {experts_outputs.get('expert_analyst', '[Не получено]')} \"\"\" ### ⚠️ Expert 3 — Critic (критический анализ): \"\"\" {experts_outputs.get('expert_critic', '[Не получено]')} \"\"\" ## Задача: Синтезируй разные точки зрения, выяви точки согласия и разногласия, подготовь почву для консенсуса.""", } moderator_output = "" try: async for event_type, content, tokens, model in _run_agent_stream( moderator, "moderator", moderator_input ): if event_type == "content": moderator_output += content yield { "type": "agent_message", "agent": "moderator", "content": content, } elif event_type == "done": total_tokens += tokens yield { "type": "agent_done", "agent": "moderator", "output": content, "tokens": tokens, "model": model, } elif event_type == "error": yield { "type": "error", "error": content, "agent": "moderator", } return except Exception as e: logger.error("consensus_flow.moderator.failed", error=str(e), exc_info=True) yield { "type": "error", "error": f"ModeratorAgent failed: {e}", "agent": "moderator", } return stage_outputs["moderator"] = moderator_output # ======================================================================== # Шаг 4: DebaterAgent — разрешение спорных моментов # ======================================================================== yield {"type": "agent_start", "agent": "debater"} debater = DebaterAgent( name="Consensus Facilitator", system_prompt=DEBATER_CONSENSUS_PROMPT, ) debater_input = { "topic": f"""Проведи консенсусную дискуссию по спорным моментам. ## ПРОБЛЕМА: \"\"\" {topic} \"\"\" ## СИНТЕЗ МНЕНИЙ (от модератора): \"\"\" {moderator_output} \"\"\" ## МНЕНИЯ ЭКСПЕРТОВ (краткое напоминание): **Researcher:** {experts_outputs.get('expert_researcher', '')[:500]}... **Analyst:** {experts_outputs.get('expert_analyst', '')[:500]}... **Critic:** {experts_outputs.get('expert_critic', '')[:500]}... ## Задача: Найди пути к консенсусу по спорным моментам, предложи компромиссы и win-win решения.""", } debater_output = "" try: async for event_type, content, tokens, model in _run_agent_stream( debater, "debater", debater_input ): if event_type == "content": debater_output += content yield { "type": "agent_message", "agent": "debater", "content": content, } elif event_type == "done": total_tokens += tokens yield { "type": "agent_done", "agent": "debater", "output": content, "tokens": tokens, "model": model, } elif event_type == "error": yield { "type": "error", "error": content, "agent": "debater", } return except Exception as e: logger.error("consensus_flow.debater.failed", error=str(e), exc_info=True) yield { "type": "error", "error": f"DebaterAgent failed: {e}", "agent": "debater", } return stage_outputs["debater"] = debater_output # ======================================================================== # Шаг 5: JudgeAgent — консенсусное решение # ======================================================================== yield {"type": "agent_start", "agent": "judge"} judge = JudgeAgent( name="Consensus Judge", system_prompt=JUDGE_CONSENSUS_PROMPT, ) judge_input = { "topic": f"""Вынеси консенсусное решение на основе экспертного обсуждения. ## ПРОБЛЕМА: \"\"\" {topic} \"\"\" ## ФОРМУЛИРОВКА: \"\"\" {planner_output} \"\"\" ## МНЕНИЯ ЭКСПЕРТОВ: - **Researcher:** {experts_outputs.get('expert_researcher', '')[:800]}... - **Analyst:** {experts_outputs.get('expert_analyst', '')[:800]}... - **Critic:** {experts_outputs.get('expert_critic', '')[:800]}... ## СИНТЕЗ МОДЕРАТОРА: \"\"\" {moderator_output} \"\"\" ## ПУТИ К КОНСЕНСУСУ (от фасилитатора): \"\"\" {debater_output} \"\"\" {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("consensus_flow.judge.failed", error=str(e), exc_info=True) yield { "type": "error", "error": f"JudgeAgent failed: {e}", "agent": "judge", } return stage_outputs["judge"] = judge_output # ======================================================================== # Шаг 6: WriterAgent — финальный консенсусный документ # ======================================================================== yield {"type": "agent_start", "agent": "writer"} writer = WriterAgent( name="Consensus Report Writer", system_prompt=WRITER_CONSENSUS_PROMPT, ) writer_input = { "topic": f"""Сформируй финальный консенсусный документ. ## ПРОБЛЕМА: \"\"\" {topic} \"\"\" {f"Контекст: {context}" if context else ""} ## ВЕСЬ ПРОЦЕСС ОБСУЖДЕНИЯ: ### Планирование: \"\"\" {planner_output} \"\"\" ### Мнения экспертов: **Expert 1 — Researcher:** \"\"\" {experts_outputs.get('expert_researcher', '')} \"\"\" **Expert 2 — Analyst:** \"\"\" {experts_outputs.get('expert_analyst', '')} \"\"\" **Expert 3 — Critic:** \"\"\" {experts_outputs.get('expert_critic', '')} \"\"\" ### Синтез модератора: \"\"\" {moderator_output} \"\"\" ### Дискуссия фасилитатора: \"\"\" {debater_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("consensus_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( "consensus_flow.completed", total_tokens=total_tokens, stages={ "planning": len(planner_output), "expert_researcher": len(experts_outputs.get("expert_researcher", "")), "expert_analyst": len(experts_outputs.get("expert_analyst", "")), "expert_critic": len(experts_outputs.get("expert_critic", "")), "moderator": len(moderator_output), "debater": len(debater_output), "judge": len(judge_output), "report": len(writer_output), }, ) yield { "type": "flow_done", "flow_id": "consensus", "output": writer_output, "tokens": total_tokens, "stages": stage_outputs, "metadata": { "experts_count": 3, "agents_count": 8, # planner + 3 experts + moderator + debater + judge + writer "total_stages": 8, "parallel_stage": "experts", # Указываем что был параллельный этап }, } # ============================================================================ # Регистрация flow # ============================================================================ consensus_flow = Flow( id="consensus", name="Consensus Flow", description=( "Мультиагентный консенсус для сложных решений: " "Planner → ExpertPanel (Researcher + Analyst + Critic параллельно) → " "Moderator → Debater → Judge → Writer. " "Собирает мнения нескольких независимых экспертов с разных линз " "экспертизы и синтезирует единое консенсусное решение. " "Идеально подходит для стратегических решений, архитектурных вопросов " "и мультидисциплинарных задач." ), agents=[ "planner", "expert_researcher", "expert_analyst", "expert_critic", "moderator", "debater", "judge", "writer", ], runner=run_consensus_flow, ) # ============================================================================ # Public API # ============================================================================ __all__ = ["consensus_flow", "run_consensus_flow"]