/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/engine/src/flows/research.py
491 строка
18 KB
Alexander Efanov
Обновление репозитория
15 июл 2026, 12:19
15 июл 2026, 12:19
76704c6
Код
Авторство
О чём код?
# core/engine/src/flows/research.py """Research flow — мультиагентное исследование темы. Последовательность агентов: 1. PlannerAgent — создаёт план исследования с декомпозицией задачи 2. ResearcherAgent — собирает информацию по каждому пункту плана 3. AnalystAgent — анализирует данные, выявляет паттерны и инсайты 4. WriterAgent — структурирует результаты в итоговый отчёт 5. ReviewerAgent — проверяет качество и полноту результата Каждый агент получает на вход результат предыдущего, создавая глубокий, всесторонний анализ с верификацией качества. """ from __future__ import annotations from typing import Any, AsyncIterator import structlog from src.agents import PlannerAgent, ResearcherAgent, AnalystAgent, WriterAgent, ReviewerAgent from src.primitives.flow import Flow logger = structlog.get_logger() async def run_research_flow(input_data: dict[str, Any]) -> AsyncIterator[dict[str, Any]]: """ Выполняет мультиагентный research flow. Flow состоит из пяти последовательных шагов: 1. PlannerAgent — декомпозиция задачи и создание плана исследования 2. ResearcherAgent — сбор информации по плану 3. AnalystAgent — анализ данных и выявление паттернов 4. WriterAgent — формирование структурированного отчёта 5. ReviewerAgent — проверка качества и финальные рекомендации Args: input_data: Входные данные. Поддерживает поля: - 'topic' или 'input' или 'query': тема для исследования - 'depth' (опционально): глубина анализа (default: 'comprehensive') - 'focus_areas' (опционально): конкретные аспекты для фокуса Yields: События выполнения flow для SSE streaming: - agent_start: агент начал работу - agent_message: агент генерирует контент (стриминг) - agent_done: агент завершил работу - error: ошибка в одном из агентов - flow_done: финальное событие с итоговым результатом """ # Извлекаем тему исследования topic = ( input_data.get("topic") or input_data.get("input") or input_data.get("query") or "" ) if not topic.strip(): yield { "type": "error", "error": "Не указана тема для исследования (ожидается поле 'topic', 'input' или 'query')", } return depth = input_data.get("depth", "comprehensive") focus_areas = input_data.get("focus_areas", "") logger.info( "research_flow.started", topic=topic[:100], depth=depth, has_focus=bool(focus_areas), ) # Контейнер для накопления результатов каждого этапа stage_outputs: dict[str, str] = { "topic": topic, "depth": depth, "focus_areas": focus_areas, } total_tokens = 0 # ======================================================================== # Шаг 1: PlannerAgent — создание плана исследования # ======================================================================== yield {"type": "agent_start", "agent": "planner"} planner_input = { "topic": f"""Создай детальный план исследования по теме: "{topic}" Глубина анализа: {depth} {f"Фокус на аспектах: {focus_areas}" if focus_areas else ""} План должен включать: 1. Ключевые вопросы для исследования (5-7 вопросов) 2. Источники информации, которые нужно изучить 3. Методологию анализа данных 4. Критерии оценки качества результатов 5. Структуру итогового отчёта Создай практический, выполнимый план с чёткими шагами. """, } planner = PlannerAgent() planner_output = "" try: async for chunk in planner.run_stream(planner_input): chunk_type = chunk.get("type") if chunk_type == "content": content = chunk.get("content", "") planner_output += content yield { "type": "agent_message", "agent": "planner", "content": content, } elif chunk_type == "done": tokens = chunk.get("tokens_total", 0) total_tokens += tokens yield { "type": "agent_done", "agent": "planner", "output": planner_output, "tokens": tokens, "model": chunk.get("model", ""), } elif chunk_type == "error": yield { "type": "error", "error": chunk.get("error", "Unknown error in planner"), "agent": "planner", } return except Exception as e: logger.error("research_flow.planner.failed", error=str(e)) yield { "type": "error", "error": f"PlannerAgent failed: {e}", "agent": "planner", } return stage_outputs["plan"] = planner_output # ======================================================================== # Шаг 2: ResearcherAgent — сбор информации # ======================================================================== yield {"type": "agent_start", "agent": "researcher"} researcher_input = { "topic": f"""Проведи исследование по теме: "{topic}" ПЛАН ИССЛЕДОВАНИЯ: \"\"\" {planner_output} \"\"\" {f"Особое внимание удели: {focus_areas}" if focus_areas else ""} Собери информацию по каждому пункту плана: - Ключевые факты и данные - Экспертные мнения и источники - Статистика и примеры - Альтернативные точки зрения - Актуальные тренды и developments Предоставь структурированный результат с цитированием источников. """, } researcher = ResearcherAgent() researcher_output = "" try: async for chunk in researcher.run_stream(researcher_input): chunk_type = chunk.get("type") if chunk_type == "content": content = chunk.get("content", "") researcher_output += content yield { "type": "agent_message", "agent": "researcher", "content": content, } elif chunk_type == "done": tokens = chunk.get("tokens_total", 0) total_tokens += tokens yield { "type": "agent_done", "agent": "researcher", "output": researcher_output, "tokens": tokens, "model": chunk.get("model", ""), } elif chunk_type == "error": yield { "type": "error", "error": chunk.get("error", "Unknown error in researcher"), "agent": "researcher", } return except Exception as e: logger.error("research_flow.researcher.failed", error=str(e)) yield { "type": "error", "error": f"ResearcherAgent failed: {e}", "agent": "researcher", } return stage_outputs["research"] = researcher_output # ======================================================================== # Шаг 3: AnalystAgent — анализ данных # ======================================================================== yield {"type": "agent_start", "agent": "analyst"} analyst_input = { "topic": f"""Проанализируй собранные данные по теме: "{topic}" СОБРАННАЯ ИНФОРМАЦИЯ: \"\"\" {researcher_output} \"\"\" Проведи глубокий анализ: 1. Выяви ключевые паттерны и тренды 2. Определи причинно-следственные связи 3. Найди противоречия и пробелы в данных 4. Сформулируй основные инсайты (5-7 штук) 5. Оцени надёжность источников 6. Определи implications для практики Предоставь аналитический отчёт с чёткими выводами. """, } analyst = AnalystAgent() analyst_output = "" try: async for chunk in analyst.run_stream(analyst_input): chunk_type = chunk.get("type") if chunk_type == "content": content = chunk.get("content", "") analyst_output += content yield { "type": "agent_message", "agent": "analyst", "content": content, } elif chunk_type == "done": tokens = chunk.get("tokens_total", 0) total_tokens += tokens yield { "type": "agent_done", "agent": "analyst", "output": analyst_output, "tokens": tokens, "model": chunk.get("model", ""), } elif chunk_type == "error": yield { "type": "error", "error": chunk.get("error", "Unknown error in analyst"), "agent": "analyst", } return except Exception as e: logger.error("research_flow.analyst.failed", error=str(e)) yield { "type": "error", "error": f"AnalystAgent failed: {e}", "agent": "analyst", } return stage_outputs["analysis"] = analyst_output # ======================================================================== # Шаг 4: WriterAgent — формирование отчёта # ======================================================================== yield {"type": "agent_start", "agent": "writer"} writer_input = { "topic": f"""Сформируй итоговый исследовательский отчёт по теме: "{topic}" ПЛАН ИССЛЕДОВАНИЯ: \"\"\" {planner_output} \"\"\" СОБРАННАЯ ИНФОРМАЦИЯ: \"\"\" {researcher_output} \"\"\" АНАЛИЗ ДАННЫХ: \"\"\" {analyst_output} \"\"\" Создай профессиональный отчёт в markdown формате со следующими разделами: ## 📋 Executive Summary (Краткая выжимка основных выводов, 3-5 предложений) ## 🎯 Ключевые находки (Топ-5 наиболее важных открытий с объяснением значимости) ## 📊 Детальный анализ (Структурированный анализ по каждому пункту плана) ## 💡 Инсайты и паттерны (Выявленные закономерности и неочевидные связи) ## ⚠️ Ограничения и пробелы (Что не удалось исследовать, где нужна дополнительная информация) ## 🚀 Практические рекомендации (Конкретные действия на основе исследования) ## 📚 Источники (Список использованных источников с оценкой надёжности) Объедини результаты всех этапов, убери дублирование, создай цельный документ. """, } writer = WriterAgent() writer_output = "" try: async for chunk in writer.run_stream(writer_input): chunk_type = chunk.get("type") if chunk_type == "content": content = chunk.get("content", "") writer_output += content yield { "type": "agent_message", "agent": "writer", "content": content, } elif chunk_type == "done": tokens = chunk.get("tokens_total", 0) total_tokens += tokens yield { "type": "agent_done", "agent": "writer", "output": writer_output, "tokens": tokens, "model": chunk.get("model", ""), } elif chunk_type == "error": yield { "type": "error", "error": chunk.get("error", "Unknown error in writer"), "agent": "writer", } return except Exception as e: logger.error("research_flow.writer.failed", error=str(e)) yield { "type": "error", "error": f"WriterAgent failed: {e}", "agent": "writer", } return stage_outputs["report"] = writer_output # ======================================================================== # Шаг 5: ReviewerAgent — проверка качества # ======================================================================== yield {"type": "agent_start", "agent": "reviewer"} reviewer_input = { "topic": f"""Проведи финальную проверку качества исследовательского отчёта по теме: "{topic}" ИТОГОВЫЙ ОТЧЁТ: \"\"\" {writer_output} \"\"\" Проверь по критериям: 1. **Полнота**: покрыты ли все пункты плана? 2. **Точность**: соответствуют ли выводы данным? 3. **Логика**: нет ли противоречий? 4. **Практичность**: можно ли применить рекомендации? 5. **Источники**: достаточно ли доказательств? Предоставь: - Оценку качества по шкале 1-10 - Конкретные замечания (если есть) - Рекомендации по улучшению (опционально) - Финальный вердикт: готов к использованию или требует доработки """, } reviewer = ReviewerAgent() reviewer_output = "" try: async for chunk in reviewer.run_stream(reviewer_input): chunk_type = chunk.get("type") if chunk_type == "content": content = chunk.get("content", "") reviewer_output += content yield { "type": "agent_message", "agent": "reviewer", "content": content, } elif chunk_type == "done": tokens = chunk.get("tokens_total", 0) total_tokens += tokens yield { "type": "agent_done", "agent": "reviewer", "output": reviewer_output, "tokens": tokens, "model": chunk.get("model", ""), } elif chunk_type == "error": yield { "type": "error", "error": chunk.get("error", "Unknown error in reviewer"), "agent": "reviewer", } return except Exception as e: logger.error("research_flow.reviewer.failed", error=str(e)) yield { "type": "error", "error": f"ReviewerAgent failed: {e}", "agent": "reviewer", } return stage_outputs["review"] = reviewer_output # ======================================================================== # Финальное событие flow # ======================================================================== logger.info( "research_flow.completed", total_tokens=total_tokens, plan_length=len(planner_output), research_length=len(researcher_output), analysis_length=len(analyst_output), report_length=len(writer_output), review_length=len(reviewer_output), ) yield { "type": "flow_done", "flow_id": "research", "output": writer_output, "tokens": total_tokens, "stages": { "plan": planner_output, "research": researcher_output, "analysis": analyst_output, "report": writer_output, "review": reviewer_output, }, } # ============================================================================ # Регистрация flow # ============================================================================ research_flow = Flow( id="research", name="Research Flow", description=( "Мультиагентное исследование: Planner → Researcher → Analyst → Writer → Reviewer. " "Создаёт план, собирает информацию, анализирует данные, формирует отчёт и проверяет качество." ), agents=["planner", "researcher", "analyst", "writer", "reviewer"], runner=run_research_flow, )