/
Watashicuvu
/
agentic-tools
Обзор
Документация
Войти
/
Watashicuvu
/
agentic-tools
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/cli_agent/interactive_executor.py
1 618 строк
67 KB
Your Name
added come roles
28 май 2026, 15:16
28 май 2026, 15:16
745363b
Код
Авторство
О чём код?
"""Interactive Executor — обёртка над WorkerExecutor с интерактивными точками одобрения. Позволяет пользователю: 1. Видеть план выполнения (strategist → executor → reviewer) 2. Одобрять/отклонять каждый шаг 3. Возвращать на доработку с feedback 4. Наблюдать прогресс в реальном времени Workflow: 1. RoleRouter определяет план → показывает пользователю 2. Пользователь одобряет (/approve) или отклоняет (/reject) 3. InteractiveExecutor выполняет задачи по одной 4. После каждой роли → пауза и ожидание одобрения 5. Reviewer blocker → возврат на executor с feedback 6. Финальный результат → возврат в чат Usage: from src.cli_agent.interactive_executor import InteractiveExecutor executor = InteractiveExecutor( smart_client=client, project_root="/path/to/project", mcp_url="http://localhost:8080" ) # Запуск плана с интерактивными паузами result = await executor.execute_with_approval( plan=role_plan, on_step_complete=show_step_result ) """ import asyncio import hashlib import logging from dataclasses import dataclass, field from datetime import datetime, timezone from enum import Enum from typing import Optional, Dict, Any, List, Callable, Awaitable, Tuple from rich.console import Console from rich.panel import Panel from rich.markdown import Markdown from src.services.async_smart_client import AsyncSmartOpenAI from src.swarm.classifier import AgentRole from src.swarm.planner import TaskPlan, PlanTask, PlanStatus, DAGStructure from src.swarm.worker_executor import WorkerExecutor, ExecutionResult as WorkerExecutionResult from src.cli_agent.role_router import RolePlan, ExecutionMode from src.cli_agent.llm_streaming_client import StreamEvent, ToolCallError from src.cli_agent.llm_reviewer import ReviewViolation logger = logging.getLogger(__name__) console = Console() async def _with_progress(coro, message: str = "Получаю данные..."): """Выполняет coroutine с показом прогресс-спиннера. Args: coro: Already-awaited coroutine (передавать без await) message: Текст для отображения во время выполнения Returns: Результат coro Usage: result = await _with_progress(some_async_func(), message="Working...") """ import itertools import sys spinner = itertools.cycle(["⠋", "⠙", "⠹", "⠸", "⠼", "⠴", "⠦", "⠧", "⠇", "⠏"]) # Оборачиваем coroutine в task task = asyncio.ensure_future(coro) if not asyncio.isfuture(coro) else coro first = True while not task.done(): symbol = next(spinner) # Используем sys.stdout.write для корректного \r if first: sys.stdout.write(f"\r{symbol} {message}") sys.stdout.flush() first = False else: sys.stdout.write(f"\r{symbol} {message}") sys.stdout.flush() await asyncio.sleep(0.3) # Очищаем строку спиннера и печатаем newline sys.stdout.write("\r" + " " * 60 + "\r") sys.stdout.flush() return await task # ============================================================================ # Models # ============================================================================ class StepStatus(str, Enum): """Статус шага выполнения.""" PENDING = "pending" IN_PROGRESS = "in_progress" COMPLETED = "completed" FAILED = "failed" BLOCKED = "blocked" # Reviewer blocker CANCELLED = "cancelled" WAITING_APPROVAL = "waiting_approval" @dataclass class ExecutionStep: """Один шаг выполнения (одна роль).""" step_id: str role: AgentRole status: StepStatus = StepStatus.PENDING result: Optional[Dict[str, Any]] = None error: Optional[str] = None started_at: Optional[str] = None completed_at: Optional[str] = None feedback: Optional[str] = None # Feedback от reviewer def to_summary(self) -> str: """Краткое описание шага.""" status_emoji = { StepStatus.PENDING: "⏳", StepStatus.IN_PROGRESS: "🔄", StepStatus.COMPLETED: "✅", StepStatus.FAILED: "❌", StepStatus.BLOCKED: "🚫", StepStatus.CANCELLED: "⛔", StepStatus.WAITING_APPROVAL: "⏸️", } return f"{status_emoji.get(self.status, '❓')} {self.role.value}: {self.status.value}" @dataclass class ExecutionResult: """Результат интерактивного выполнения.""" plan_id: str status: str # "completed", "failed", "cancelled" steps: List[ExecutionStep] = field(default_factory=list) total_duration_ms: float = 0.0 error: Optional[str] = None reviewer_verdict: Optional[str] = None # "APPROVED" / "BLOCKER" @property def success(self) -> bool: return self.status == "completed" and self.reviewer_verdict == "APPROVED" def to_summary(self) -> str: """Итоговое резюме.""" emoji = "✅" if self.success else "❌" lines = [ f"{emoji} **Результат выполнения:**", f"", f" **Статус:** {self.status}", f" **Reviewer:** {self.reviewer_verdict or 'N/A'}", f" **Шаги:**", ] for step in self.steps: lines.append(f" {step.to_summary()}") if step.error: lines.append(f" Ошибка: {step.error}") if step.feedback: lines.append(f" Feedback: {step.feedback}") if self.error: lines.append(f"") lines.append(f" **Критическая ошибка:** {self.error}") return "\n".join(lines) # ============================================================================ # Interactive Executor # ============================================================================ class InteractiveExecutor: """Интерактивный исполнитель с точками одобрения. Для каждой роли в плане: 1. Показывает пользователю что будет делаться 2. Ждёт /approve или /reject 3. Выполняет роль (через WorkerExecutor или напрямую) 4. Показывает результат 5. Если reviewer blocker → возврат на executor с feedback Args: smart_client: Клиент для LLM project_root: Корень проекта mcp_url: URL MCP сервера (опционально) model: Модель LLM """ def __init__( self, smart_client: AsyncSmartOpenAI, project_root: str, mcp_url: Optional[str] = None, model: Optional[str] = None, ): self.smart_client = smart_client self.project_root = project_root self.mcp_url = mcp_url self.model = model # Конфигурация retry self.max_retries: int = 3 # Максимум попыток на шаг # Состояние self._current_step: Optional[ExecutionStep] = None self._approval_event = asyncio.Event() self._approval_action: str = "approve" # "approve" или "reject" self._reject_reason: Optional[str] = None # Retry tracking self._retry_count: int = 0 # Сколько раз уже retried текущий шаг self._feedback_history: List[str] = [] # История всех feedback сообщений # Callbacks self._on_event: Optional[Callable[[Dict[str, Any]], Awaitable[None]]] = None # Результаты self._execution_result: Optional[ExecutionResult] = None def set_event_handler(self, handler: Callable[[Dict[str, Any]], Awaitable[None]]): """Устанавливает обработчик событий.""" self._on_event = handler async def execute_with_approval( self, role_plan: RolePlan, on_step_complete: Optional[Callable[[ExecutionStep], Awaitable[None]]] = None, ) -> ExecutionResult: """Выполняет план с интерактивными одобрениями и retry логикой. Feedback Loop: 1. Показывает шаг → ждёт /approve или /reject 2. /reject → сохраняет feedback → retry (max 3) 3. retry: enriched query с feedback history → executor retry 4. Max retries исчерпан → manual intervention Args: role_plan: План от RoleRouter on_step_complete: Callback после каждого шага Returns: ExecutionResult с итогами """ if not role_plan.roles: return ExecutionResult( plan_id="empty", status="completed", steps=[], reviewer_verdict="N/A" ) # Создаём шаги steps = [ ExecutionStep( step_id=f"step_{i+1}", role=role, ) for i, role in enumerate(role_plan.roles) ] execution_result = ExecutionResult( plan_id=f"plan_{datetime.now(timezone.utc).strftime('%Y%m%d_%H%M%S')}", status="in_progress", steps=steps, ) console.print(Panel( f"🚀 **Начало выполнения**\n\n" f"План ID: {execution_result.plan_id}\n" f"Роли: {' → '.join(s.role.value for s in steps)}\n" f"Запрос: {role_plan.query}\n" f"Max retries на шаг: {self.max_retries}", title="Interactive Executor", border_style="blue" )) # Пошаговое выполнение с retry i = 0 while i < len(steps): step = steps[i] logger.info( f"[EXECUTOR] Цикл i={i}, шаг={step.step_id}, роль={step.role.value}, " f"статус={step.status.value}, retry_count={self._retry_count}" ) # Проверяем, есть ли retry для этого шага attempt_number = self._retry_count + 1 is_retry = self._retry_count > 0 if is_retry: console.print(Panel( f"🔄 **Retry попытка {attempt_number}/{self.max_retries}**\n\n" f"Шаг: {step.role.value}\n" f"Предыдущие feedback:\n" + "\n".join(f" {idx+1}. {fb}" for idx, fb in enumerate(self._feedback_history)), title="Retry Mode", border_style="yellow" )) # Формируем enriched query с feedback enriched_query = self._enrich_query_with_feedback( original_query=role_plan.query, step=step, previous_results={s.step_id: s.result for s in steps[:i] if s.result}, ) # Показываем план и ждём одобрения logger.info( f"[EXECUTOR] Устанавливаю _current_step={step.step_id} ({step.role.value}), " f"вызываю _wait_for_approval(step_number={i+1}/{len(steps)})" ) self._current_step = step approval = await self._wait_for_approval( step=step, step_number=i + 1, total_steps=len(steps), is_retry=is_retry, attempt_number=attempt_number, ) logger.info( f"[EXECUTOR] _wait_for_approval вернул approval={approval}, " f"_current_step={self._current_step.step_id if self._current_step else None}" ) if approval == "reject": # Проверяем лимит retry if self._retry_count >= self.max_retries: # Max retries исчерпан → manual intervention execution_result.status = "failed" step.status = StepStatus.BLOCKED step.error = ( f"Max retries ({self.max_retries}) exceeded. " f"Требуется ручное вмешательство.\n" f"Последний feedback: {self._reject_reason}" ) console.print(Panel( f"🚫 **Max retries exceeded**\n\n" f"Шаг: {step.role.value}\n" f"Попыток: {self._retry_count}/{self.max_retries}\n" f"Последний feedback:\n" + "\n".join(f" {idx+1}. {fb}" for idx, fb in enumerate(self._feedback_history)) + f"\n\n⚠️ **Требуется ручное вмешательство**\n" f"Исправьте вручную и продолжите с /approve", title="Manual Intervention Required", border_style="red" )) self._current_step = None break # Есть попытки retry → сохраняем feedback и повторяем шаг self._retry_count += 1 feedback = self._reject_reason or "Не указано" self._feedback_history.append(feedback) step.status = StepStatus.FAILED step.error = f"Rejected (attempt {attempt_number}): {feedback}" step.completed_at = datetime.now(timezone.utc).isoformat() console.print(Panel( f"↩️ **Возврат на шаг {i+1} с feedback**\n\n" f"Feedback: {feedback}\n" f"Попытка {attempt_number + 1}/{self.max_retries}\n" f"Следующая попытка учтёт этот feedback", title="Retry Scheduled", border_style="yellow" )) # НЕ увеличиваем i — повторяем тот же шаг continue # Успешное одобрение — сбрасываем retry счётчик для следующего шага self._retry_count = 0 self._feedback_history.clear() # Выполняем шаг step.status = StepStatus.IN_PROGRESS step.started_at = datetime.now(timezone.utc).isoformat() console.print(f"\n[cyan]🔄 Выполняю шаг {i+1}/{len(steps)}: {step.role.value}...[/cyan]") try: result = await self._execute_role( role=step.role, query=enriched_query, step=step, previous_results={s.step_id: s.result for s in steps[:i] if s.result}, ) step.result = result step.status = StepStatus.COMPLETED step.completed_at = datetime.now(timezone.utc).isoformat() console.print(f"[green]✅ Шаг {i+1} завершён: {step.role.value}[/green]") # SPECIAL HANDLING: Reviewer BLOCKER if step.role == AgentRole.REVIEWER and step.result: reviewer_verdict = step.result.get("verdict", "APPROVED") if reviewer_verdict == "BLOCKER": console.print(Panel( f"🚫 **Reviewer BLOCKER**\n\n" f"Обнаружены критические нарушения:\n" + "\n".join(f" {idx+1}. {v}" for idx, v in enumerate(step.result.get("violations", [])[:5])) + (f"\n ... и ещё {len(step.result.get('violations', [])) - 5}" if len(step.result.get("violations", [])) > 5 else ""), title="Review BLOCKER", border_style="red" )) # Предлагаем варианты action = await self._handle_reviewer_blocker(step) if action == "fix": # Auto-fix через LintAutoFixPipeline console.print("[cyan]🔧 Запускаю auto-fix...[/cyan]") fix_result = await self._auto_fix_reviewer_blocker(step) if fix_result and fix_result.success: console.print("[green]✅ Auto-fix успешен, повторная проверка...[/green]") # Re-review re_review_result = await self._execute_reviewer( query=role_plan.query, previous_results={s.step_id: s.result for s in steps[:i] if s.result}, ) step.result = re_review_result if re_review_result.get("verdict") != "BLOCKER": console.print("[green]✅ Re-review прошёл![/green]") else: console.print("[yellow]⚠ Re-review всё ещё BLOCKER[/yellow]") else: console.print("[yellow]⚠ Auto-fix не успешен[/yellow]") elif action == "reject": # Вернуть executor с feedback feedback = "Reviewer BLOCKER: " + "; ".join( step.result.get("violations", [])[:3] ) self._retry_count = 0 # Сбросить для executor retry self._feedback_history.append(feedback) # Находим последний executor шаг last_executor_idx = None for idx in range(i - 1, -1, -1): if steps[idx].role == AgentRole.EXECUTOR: last_executor_idx = idx break if last_executor_idx is not None: console.print(Panel( f"↩️ **Возврат на Executor с feedback**\n\n" f"Feedback: {feedback[:100]}...", title="Executor Retry", border_style="yellow" )) # Переходим на executor шаг i = last_executor_idx continue else: console.print("[yellow]⚠ Нет executor шага для возврата[/yellow]") elif action == "skip": console.print("[yellow]⚠️ Пропускаем BLOCKER (принято с рисками)[/yellow]") step.result["verdict"] = "APPROVED_WITH_RISKS" # Продолжаем elif action == "abort": execution_result.status = "failed" step.error = "Aborted due to BLOCKER" console.print("[red]🛑 Выполнение прервано[/red]") break # Callback if on_step_complete: await on_step_complete(step) # Переходим к следующему шагу logger.info(f"[EXECUTOR] Шаг {step.step_id} завершён успешно, очищаю _current_step, i={i}→{i+1}") self._current_step = None i += 1 except Exception as e: logger.error(f"[InteractiveExecutor] Step {step.role.value} failed: {e}", exc_info=True) step.status = StepStatus.FAILED step.error = str(e) step.completed_at = datetime.now(timezone.utc).isoformat() logger.info(f"[EXECUTOR] Ошибка шага {step.step_id}, очищаю _current_step") self._current_step = None execution_result.status = "failed" execution_result.error = str(e) console.print(f"[red]❌ Шаг {i+1} провален: {e}[/red]") break # Итог if execution_result.status == "in_progress": execution_result.status = "completed" execution_result.reviewer_verdict = "APPROVED" self._execution_result = execution_result return execution_result async def _wait_for_approval( self, step: ExecutionStep, step_number: int, total_steps: int, is_retry: bool = False, attempt_number: int = 1, ) -> str: """Ждёт одобрения пользователя на шаг. Args: step: Текущий шаг step_number: Номер шага (1-based) total_steps: Всего шагов is_retry: Это retry попытка? attempt_number: Номер попытки Returns: "approve" или "reject" """ step.status = StepStatus.WAITING_APPROVAL logger.info( f"[WAIT_APPROЛ] Вход: step={step.step_id}, роль={step.role.value}, " f"step_number={step_number}/{total_steps}, is_retry={is_retry}" ) # Показываем информацию о шаге role_description = self._get_role_description(step.role) retry_info = "" if is_retry: retry_info = ( f"\n🔄 **Retry попытка {attempt_number}**\n\n" f"⚠️ Предыдущие попытки были отклонены.\n" f"Убедитесь, что учтён весь feedback." ) panel_content = ( f"**Шаг {step_number}/{total_steps}: {step.role.value.upper()}**{retry_info}\n\n" f"{role_description}\n\n" f"⏳ **Введите:**\n" f" `/approve` — начать выполнение\n" f" `/reject <причина>` — вернуть на доработку с feedback" ) console.print(Panel(panel_content, title="Ожидание одобрения", border_style="yellow")) # Ждём события одобрения self._approval_event.clear() self._approval_action = "pending" self._reject_reason = None logger.info( f"[WAIT_APPROЛ] Ожидание _approval_event (timeout=300s), " f"_approval_action={self._approval_action}" ) # Читаем ввод пользователя параллельно с ожиданием event async def _read_input(): """Читает ввод из консоли и обрабатывает команду.""" try: # Пробуем использовать интерактивный ввод с prompt_toolkit try: from src.cli_agent.interactive_input import interactive_approval_prompt cmd, reason = await interactive_approval_prompt( message="Command", timeout=300, ) except ImportError: # Fallback на простой ввод from src.cli_agent.interactive_input import simple_approval_input cmd, reason = await simple_approval_input(console, timeout=300) logger.info(f"[WAIT_APPROЛ] Прочитан ввод: cmd={cmd!r}, reason={reason!r}") if cmd == "approve": logger.info(f"[WAIT_APPROЛ] Одобрение") self._approval_action = "approve" self._approval_event.set() return "approve" if cmd == "reject": reason_text = reason or "Не указано" logger.info(f"[WAIT_APPROЛ] Отклонение, reason={reason_text!r}") self._approval_action = "reject" self._reject_reason = reason_text self._approval_event.set() return "reject" if cmd == "unknown": console.print("[yellow]⚠ Введите `/approve` или `/reject <причина>`[/yellow]") except EOFError: return "reject" except Exception as e: logger.warning(f"[WAIT_APPROЛ] Ошибка ввода: {e}, fallback на простой ввод") # Fallback на простой ввод try: def _blocking_read(): return console.input("\n[yellow]Command:[/yellow] ").strip() cmd = await asyncio.to_thread(_blocking_read) if cmd == '/approve': self._approval_action = "approve" self._approval_event.set() return "approve" elif cmd.startswith('/reject'): reason = cmd[8:].strip() or "Не указано" self._approval_action = "reject" self._reject_reason = reason self._approval_event.set() return "reject" elif cmd.lower() in ['exit', 'quit']: self._approval_action = "reject" self._approval_event.set() return "reject" else: console.print("[yellow]⚠ Введите `/approve` или `/reject <причина>`[/yellow]") except EOFError: return "reject" # Запускаем две задачи параллельно: external signal + user input external_wait = asyncio.ensure_future(self._approval_event.wait()) input_task = asyncio.ensure_future(_read_input()) try: done, pending_tasks = await asyncio.wait( [external_wait, input_task], timeout=300, return_when=asyncio.FIRST_COMPLETED, ) # Отменяем оставшиеся задачи for task in pending_tasks: task.cancel() try: await task except (asyncio.CancelledError, Exception): pass except asyncio.TimeoutError: console.print("[yellow]⏱️ Таймаут ожидания одобрения (5 минут)[/yellow]") logger.warning(f"[WAIT_APPROЛ] Таймаут для step={step.step_id}") return "reject" logger.info( f"[WAIT_APPROЛ] Получен сигнал! _approval_action={self._approval_action}, " f"_reject_reason={self._reject_reason}" ) return self._approval_action def _enrich_query_with_feedback( self, original_query: str, step: ExecutionStep, previous_results: Dict[str, Any], ) -> str: """Формирует enriched query с feedback history. Для executor retry: - Добавляет весь feedback history - Добавляет предыдущие результаты для контекста - Формирует инструкцию для LLM учесть feedback Returns: Enriched query string """ if not self._feedback_history: # Нет feedback — возвращаем оригинальный query return original_query # Формируем feedback блок feedback_section = "\n\n".join([ f"**FEEDBACK #{idx+1}:** {fb}" for idx, fb in enumerate(self._feedback_history) ]) # Предыдущие результаты (если есть) prev_context = "" if previous_results: prev_lines = ["\n**PREVIOUS RESULTS:**"] for step_id, result in previous_results.items(): if isinstance(result, dict) and result.get("status") == "completed": prev_lines.append(f"- {step_id}: {result.get('message', 'OK')}") prev_context = "\n".join(prev_lines) # Формируем итоговый query enriched_query = ( f"{original_query}\n\n" f"{'='*60}\n" f"⚠️ **IMPORTANT: RETRY MODE** ⚠️\n" f"{'='*60}\n\n" f"Это повторная попытка выполнения. " f"Предыдущие попытки были отклонены по следующим причинам:\n\n" f"{feedback_section}\n\n" f"{'='*60}\n" f"ТРЕБОВАНИЯ К ИСПРАВЛЕНИЮ:\n" f"{'='*60}\n\n" f"1. ВНИМАТЕЛЬНО изучите весь feedback выше\n" f"2. Убедитесь, что ВСЕ замечания учтены в новой версии\n" f"3. НЕ повторяйте те же ошибки\n" f"4. Протестируйте изменения перед завершением\n" f"{prev_context}" ) logger.info(f"[Feedback] Enriched query with {len(self._feedback_history)} feedback(s)") return enriched_query def approve_step(self): """Одобрить текущий шаг (вызывается из чата).""" logger.info( f"[APPROVE_STEP] Вызван approve_step(), _current_step={self._current_step.step_id if self._current_step else None}" ) self._approval_action = "approve" self._approval_event.set() console.print("[green]✓ Шаг одобрен[/green]") def reject_step(self, reason: str = ""): """Отклонить текущий шаг (вызывается из чата).""" logger.info( f"[REJECT_STEP] Вызван reject_step(), reason={reason!r}, _current_step={self._current_step.step_id if self._current_step else None}" ) self._approval_action = "reject" self._reject_reason = reason self._approval_event.set() console.print(f"[yellow]✗ Шаг отклонён: {reason}[/yellow]") def get_feedback_summary(self) -> str: """Возвращает сводку feedback для отображения пользователю.""" if not self._feedback_history: return "Нет feedback" lines = [ f"📋 **Feedback History** ({len(self._feedback_history)} entries):", f"", f"Retry attempts: {self._retry_count}/{self.max_retries}", f"", ] for idx, fb in enumerate(self._feedback_history): lines.append(f" {idx+1}. {fb}") return "\n".join(lines) async def _handle_reviewer_blocker( self, step: ExecutionStep, ) -> str: """Показывает варианты действий при BLOCKER и ждёт выбора. Returns: "fix", "reject", "skip", или "abort" """ violations = step.result.get("violations", []) files_reviewed = step.result.get("files_reviewed", []) panel_content = ( f"🚫 **Reviewer обнаружил BLOCKER**\n\n" f"📁 Files: {', '.join(files_reviewed) if files_reviewed else 'N/A'}\n" f"⚠️ Violations: {len(violations)}\n\n" f"**Выберите действие:**\n\n" f" 🔧 `/fix` — Auto-fix через LintAutoFixPipeline\n" f" (попытка исправить автоматически)\n\n" f" ↩️ `/reject` — Вернуть Executor с feedback\n" f" (ручное исправление с учётом нарушений)\n\n" f" ⚠️ `/skip` — Принять с рисками\n" f" (продолжить, зафиксировав нарушения)\n\n" f" 🛑 `/abort` — Прервать выполнение\n" f" (требуется ручное вмешательство)" ) console.print(Panel(panel_content, title="BLOCKER Actions", border_style="red")) # Ждём выбор пользователя self._approval_event.clear() self._approval_action = "pending" try: await asyncio.wait_for(self._approval_event.wait(), timeout=300) except asyncio.TimeoutError: console.print("[yellow]⏱️ Таймаут — выбираю /abort[/yellow]") return "abort" # Парсим выбор action = self._approval_action.lower().strip() if action in ("fix", "reject", "skip", "abort"): return action console.print(f"[yellow]⚠ Неизвестное действие '{action}', использую /abort[/yellow]") return "abort" def choose_blocker_action(self, action: str): """Вызывается из чата для выбора действия при BLOCKER.""" valid_actions = ("fix", "reject", "skip", "abort") if action in valid_actions: self._approval_action = action self._approval_event.set() console.print(f"[green]✓ Выбрано действие: /{action}[/green]") else: console.print(f"[red]✗ Неизвестное действие: {action}. Выберите: {', '.join(valid_actions)}[/red]") async def _auto_fix_reviewer_blocker( self, step: ExecutionStep, ) -> Optional[Any]: """Запускает LintAutoFixPipeline для исправления BLOCKER violations.""" try: from src.cli_agent.lint_auto_fix import LintAutoFixPipeline files_reviewed = step.result.get("files_reviewed", []) violations = step.result.get("violations", []) if not files_reviewed: console.print("[yellow]⚠ No files to fix[/yellow]") return None # Создаём ReviewViolation объекты из строк review_violations = [] for v_str in violations: review_violations.append(ReviewViolation( severity="critical", category="llm_review", message=v_str, file_path=files_reviewed[0] if files_reviewed else "unknown", )) # Запускаем auto-fix для каждого файла for file_path in files_reviewed: console.print(f"[cyan]🔧 Fixing: {file_path}[/cyan]") pipeline = LintAutoFixPipeline( project_root=self.project_root, smart_client=self.smart_client, model=self.model, max_retries=2, ) fix_result = await pipeline.fix_from_reviewer_violations( file_path=file_path, violations=review_violations, max_retries=2, ) if fix_result.success: console.print(f"[green]✅ {file_path} fixed![/green]") else: console.print(f"[yellow]⚠ {file_path}: {fix_result.issues_remaining} issues remaining[/yellow]") return fix_result except Exception as e: logger.error(f"Auto-fix failed: {e}", exc_info=True) console.print(f"[red]✗ Auto-fix error: {e}[/red]") return None async def _execute_role( self, role: AgentRole, query: str, step: ExecutionStep, previous_results: Dict[str, Any], ) -> Dict[str, Any]: """Выполняет одну роль. Это заглушка — реальная интеграция с WorkerExecutor будет позже. Сейчас используем прямые вызовы LLM + инструменты. """ logger.info( f"[EXECUTE_ROLE] Вход: role={role.value}, step={step.step_id}, " f"query_len={len(query)}, prev_results={len(previous_results)}" ) if role == AgentRole.STRATEGIST: return await self._execute_strategist(query, previous_results) elif role == AgentRole.EXECUTOR: return await self._execute_executor(query, previous_results, step) elif role == AgentRole.REVIEWER: return await self._execute_reviewer(query, previous_results) elif role == AgentRole.EXPLORER: return await self._execute_explorer(query, previous_results) elif role == AgentRole.EXPLAINER: return await self._execute_explainer(query, previous_results) elif role == AgentRole.ARCH_ANALYZER: return await self._execute_arch_analyzer(query, previous_results) else: raise ValueError(f"Unknown role: {role}") async def _execute_strategist( self, query: str, previous_results: Dict[str, Any], ) -> Dict[str, Any]: """Strategist: анализ архитектуры и планирование.""" from src.swarm.classifier import classify_request, TaskContext from src.swarm.planner import TaskPlanner # 1. Классификация запроса (LLM) classification = await _with_progress( classify_request( query=query, smart_client=self.smart_client, ), message="Strategist: классифицирую запрос..." ) console.print(f"[cyan]📊 Классификация: primary_role={classification.primary_role}[/cyan]") # 2. Анализ target files (если есть контекст) if classification.task_context and classification.task_context.target_files: files = [fc.path for fc in classification.task_context.target_files] console.print(f" Target files: {', '.join(files)}") # 3. Создание плана planner = TaskPlanner() try: plan = await _with_progress( asyncio.to_thread(planner.generate_plan, classification), message="Strategist: создаю план выполнения..." ) console.print(f"[green]✓ План создан: {len(plan.tasks)} задач[/green]") return { "classification": classification.to_dict(), "plan": { "task_count": len(plan.tasks), "tasks": [t.model_dump() for t in plan.tasks], "edges": plan.dag.edges if plan.dag else [], }, "status": "completed", } except Exception as e: logger.error(f"[Strategist] Plan generation failed: {e}") return { "classification": classification.to_dict(), "error": str(e), "status": "failed", } async def _execute_executor( self, query: str, previous_results: Dict[str, Any], step: ExecutionStep, ) -> Dict[str, Any]: """Executor: генерация кода + инструменты через WorkerExecutor. Интеграция: 1. Извлекает plan от Strategist из previous_results 2. Создаёт TaskPlan с executor tasks из strategist plan 3. Вызывает WorkerExecutor.execute() 4. LLM вызывает tools через ToolOrchestrator (MCP + local) 5. Возвращает результат с list of changes """ console.print(f"[yellow]🔧 Executor: генерация кода для '{query}'[/yellow]") try: # Извлекаем strategist plan из previous_results strategist_plan = self._extract_strategist_plan(previous_results) # Создаём TaskPlan с executor tasks из strategist plan plan = self._create_executor_plan_from_strategist( query=query, strategist_plan=strategist_plan, previous_results=previous_results, ) # Callback для streaming async def on_event(event: StreamEvent): if event.type == "text" and event.content: console.print(event.content, end="") elif event.type == "tool_call_start": console.print(f"\n [cyan]🔧 Tool: {event.tool_name}[/cyan]") elif event.type == "tool_call_end": console.print(f" [green]✓ Tool {event.tool_name} completed[/green]") # Callback при завершении задачи async def on_task_complete(task_result): console.print(f"\n[green]✅ Task completed: {task_result.task_id}[/green]") # Создаём и выполняем WorkerExecutor executor = WorkerExecutor( plan=plan, smart_client=self.smart_client, model=self.model, project_root=self.project_root, mcp_url=self.mcp_url, on_event=on_event, on_task_complete=on_task_complete, auto_confirm_unsafe=True, max_tool_call_rounds=15, max_task_retries=2, ) console.print("[cyan]🔧 Executor: инициализация инструментов...[/cyan]") await _with_progress( executor.initialize(), message="Executor: инициализирую инструменты..." ) console.print("[cyan]🔧 Executor: выполнение задач...[/cyan]") result, guardrails = await _with_progress( executor.execute(), message="Executor: генерирую код и вызываю инструменты..." ) # Конвертируем результат if result.success: task_result = result.task_results[0] if result.task_results else None # Извлекаем информацию об изменениях files_changed = [] tool_calls_info = [] if task_result and task_result.tool_calls: for tc in task_result.tool_calls: tool_calls_info.append({ "tool": tc.get("function", {}).get("name", "unknown"), "arguments": tc.get("function", {}).get("arguments", {}), }) # Ищем file.edit вызовы if tc.get("function", {}).get("name") in ("file.edit", "edit_file"): args = tc.get("function", {}).get("arguments", {}) if isinstance(args, str): import json try: args = json.loads(args) except: pass if isinstance(args, dict) and "file_path" in args: files_changed.append(args["file_path"]) return { "status": "completed", "files_changed": files_changed, "tool_calls": tool_calls_info, "message": f"Executor completed: {query}", "worker_result": result.to_dict(), "guardrails": guardrails, } else: return { "status": "failed", "error": result.error, "message": f"Executor failed: {result.error}", "worker_result": result.to_dict(), } except ToolCallError as e: logger.error(f"Tool call error: {e}") return { "status": "failed", "error": str(e), "message": f"Tool call error: {e}", } except Exception as e: logger.error(f"Executor execution failed: {e}", exc_info=True) return { "status": "failed", "error": str(e), "message": f"Executor failed: {e}", } def _extract_strategist_plan( self, previous_results: Dict[str, Any], ) -> Optional[Any]: """Извлекает plan от Strategist из previous_results. Ищет результат strategist шага и извлекает TaskPlan. Returns: TaskPlan или None если strategist не выполнялся """ for step_id, result in previous_results.items(): if not isinstance(result, dict): continue # Проверяем есть ли plan в результате if "plan" in result: plan_data = result["plan"] # Если это dict — пробуем восстановить TaskPlan if isinstance(plan_data, dict): try: from src.swarm.planner import TaskPlan return TaskPlan(**plan_data) except Exception: logger.warning(f"Failed to parse TaskPlan from {step_id}") return plan_data # Возвращаем как dict return plan_data return None def _create_executor_plan_from_strategist( self, query: str, strategist_plan: Optional[Any], previous_results: Dict[str, Any], ) -> TaskPlan: """Создаёт TaskPlan используя результаты Strategist. Извлекает: - target_files из strategist plan - target_zone из strategist plan - task_context (zones, constraints, existing_implementations) - enriched description из plan Returns: TaskPlan для WorkerExecutor """ # Если нет strategist plan — fallback на старый метод if strategist_plan is None: logger.info("[Executor] No strategist plan found, using fallback") return self._create_executor_plan_fallback(query, previous_results) # Извлекаем executor задачи из strategist plan executor_tasks = [] if hasattr(strategist_plan, 'tasks'): executor_tasks = [ t for t in strategist_plan.tasks if t.role == "executor" ] elif isinstance(strategist_plan, dict) and 'tasks' in strategist_plan: executor_tasks = [ t for t in strategist_plan['tasks'] if t.get('role') == "executor" ] # Если нет executor задач — создаём одну из original query if not executor_tasks: logger.info("[Executor] No executor tasks in strategist plan, creating from query") return self._create_executor_plan_fallback(query, previous_results) # Извлекаем task_context для обогащения task_context = None if hasattr(strategist_plan, 'task_context'): task_context = strategist_plan.task_context elif isinstance(strategist_plan, dict) and 'task_context' in strategist_plan: task_context = strategist_plan.get('task_context') # Создаём план с executor задачами plan_id = f"interactive_executor_{datetime.now(timezone.utc).strftime('%Y%m%d_%H%M%S')}" plan_tasks = [] for idx, exec_task in enumerate(executor_tasks): # Извлекаем поля if hasattr(exec_task, 'model_dump'): task_dict = exec_task.model_dump() elif isinstance(exec_task, dict): task_dict = exec_task else: continue # Обогащаем description контекстом от strategist enriched_description = self._enrich_task_description( base_description=task_dict.get('description', query), task_context=task_context, ) plan_tasks.append(PlanTask( task_id=f"exec_{task_dict.get('task_id', f'task_{idx}')}", role="executor", description=enriched_description, target_files=task_dict.get('target_files', []), dependencies=task_dict.get('dependencies', []), estimated_complexity=task_dict.get('estimated_complexity', 'medium'), target_zone=task_dict.get('target_zone'), )) # Строим DAG dag_nodes = [t.task_id for t in plan_tasks] dag_edges = [] for t in plan_tasks: for dep in t.dependencies: dag_edges.append((dep, t.task_id)) plan = TaskPlan( plan_id=plan_id, original_query=query, created_at=datetime.now(timezone.utc).isoformat(), status=PlanStatus.APPROVED, tasks=plan_tasks, dag=DAGStructure(nodes=dag_nodes, edges=dag_edges), task_context=task_context, ) logger.info( f"[Executor] Created plan from strategist: " f"{len(plan_tasks)} tasks, " f"target_files: {sum(len(t.target_files) for t in plan_tasks)}" ) return plan def _enrich_task_description( self, base_description: str, task_context: Optional[Any], ) -> str: """Обогащает описание задачи контекстом от Strategist.""" if task_context is None: return base_description # Извлекаем контекст context_parts = [base_description] # Target files if hasattr(task_context, 'target_files') and task_context.target_files: files = [ f.path for f in task_context.target_files if hasattr(f, 'path') ] if files: context_parts.append( f"\n\n📁 **Target files (from Strategist):**\n" + "\n".join(f" - {f}" for f in files) ) # Architectural constraints if hasattr(task_context, 'architectural_constraints') and task_context.architectural_constraints: constraints = task_context.architectural_constraints if isinstance(constraints, list): context_parts.append( f"\n\n⚠️ **Architectural constraints:**\n" + "\n".join(f" - {c}" for c in constraints) ) # Existing implementations if hasattr(task_context, 'existing_implementations') and task_context.existing_implementations: impls = task_context.existing_implementations if isinstance(impls, list) and len(impls) > 0: context_parts.append( f"\n\n💡 **Existing implementations (reuse if possible):**\n" + "\n".join( f" - {i.file_path}: {i.function_name} ({i.description})" for i in impls if hasattr(i, 'file_path') ) ) return "\n".join(context_parts) def _create_executor_plan_fallback( self, query: str, previous_results: Dict[str, Any], ) -> TaskPlan: """Fallback: создаёт TaskPlan без strategist plan.""" # Формируем enriched query с предыдущими результатами enriched_query = query if previous_results: strat_result = previous_results.get("step_1") if strat_result and "plan" in strat_result: enriched_query += ( f"\n\nStrategist plan:\n{strat_result['plan']}" ) # Создаём план с одной задачей task_id = f"executor_{hashlib.md5(query.encode()).hexdigest()[:8]}" plan = TaskPlan( plan_id=f"interactive_{task_id}", original_query=enriched_query, created_at=datetime.now(timezone.utc).isoformat(), status=PlanStatus.APPROVED, tasks=[ PlanTask( task_id=task_id, role="executor", description=query, target_files=[], # Определяется в процессе dependencies=[], estimated_complexity="medium", ) ], dag=DAGStructure( nodes=[task_id], edges=[], # Нет зависимостей — одна задача ), ) return plan async def _execute_reviewer( self, query: str, previous_results: Dict[str, Any], ) -> Dict[str, Any]: """Reviewer: проверка изменений через LLMReviewer. Интеграция: 1. Извлекаем files_changed из previous executor results 2. Для каждого файла вызываем LLMReviewer.review_file() 3. Агрегируем результаты 4. Возвращаем вердикт: APPROVED или BLOCKER """ console.print(f"[magenta]🔍 Reviewer: проверка изменений[/magenta]") try: # Импортируем LLMReviewer from src.cli_agent.llm_reviewer import LLMReviewer # Извлекаем файлы из предыдущих результатов files_to_review = [] for step_id, result in previous_results.items(): if isinstance(result, dict) and "files_changed" in result: files_to_review.extend(result["files_changed"]) # Если нет files_changed — пытаемся определить из query if not files_to_review: console.print(f"[yellow]⚠ No files changed detected, using query analysis[/yellow]") # Попробуем извлечь путь к файлу из query import re file_match = re.search(r'[\w/\.]+\.py', query) if file_match: files_to_review.append(file_match.group()) if not files_to_review: return { "status": "completed", "verdict": "APPROVED", "violations": [], "message": "No files to review", "files_reviewed": [], } # Создаём Reviewer reviewer = LLMReviewer( smart_client=self.smart_client, project_root=self.project_root, model=self.model, include_impact=True, include_duplication=True, include_tests=True, ) # Ревьюим каждый файл all_results = [] has_blocker = False all_violations = [] for file_path in files_to_review: console.print(f"[magenta] 📄 Reviewing: {file_path}[/magenta]") try: review_result = await _with_progress( reviewer.review_file(file_path), message=f"Reviewer: проверяю {file_path}..." ) all_results.append({ "file": file_path, "verdict": review_result.verdict, "violations_count": len(review_result.violations), "summary": review_result.summary, }) all_violations.extend([ v.to_markdown() for v in review_result.violations ]) if review_result.is_blocker: has_blocker = True except Exception as e: logger.warning(f"Review failed for {file_path}: {e}") all_results.append({ "file": file_path, "verdict": "ERROR", "error": str(e), }) # Формируем итог verdict = "BLOCKER" if has_blocker else "APPROVED" return { "status": "completed", "verdict": verdict, "violations": all_violations, "files_reviewed": files_to_review, "review_details": all_results, "message": f"Reviewer verdict: {verdict}", } except Exception as e: logger.error(f"Reviewer execution failed: {e}", exc_info=True) return { "status": "failed", "verdict": "ERROR", "violations": [], "error": str(e), "message": f"Reviewer failed: {e}", } async def _execute_explorer( self, query: str, previous_results: Dict[str, Any], ) -> Dict[str, Any]: """Explorer: обзор кодовой базы и поиск существующих реализаций. Использует: - MCP tools (search.code_implementations, search.code_pattern) если доступны - Fallback на grep_search если MCP недоступен Возвращает: Markdown report с architectural insights и duplication risks """ from src.cli_agent.explorer_agent import ExplorerAgent console.print(f"[cyan]🔍 Explorer: поиск по запросу '{query}'[/cyan]") console.print(f" MCP URL: {self.mcp_url or 'недоступен'}") try: explorer = ExplorerAgent( smart_client=self.smart_client, project_root=self.project_root, mcp_url=self.mcp_url, model=self.model, ) result = await _with_progress( explorer.explore(query), message="Explorer: ищу по кодовой базе..." ) # Показываем отчёт console.print(Panel( result.report, title="Explorer Report", border_style="blue" )) return { "status": "completed" if result.success else "partial", "query": query, "report": result.report, "files_explored": result.files_explored, "implementations_found": result.implementations_found, "tools_used": result.search_tools_used, "errors": result.errors, "message": result.to_summary(), } except Exception as e: logger.error(f"Explorer execution failed: {e}", exc_info=True) return { "status": "failed", "error": str(e), "message": f"Explorer failed: {e}", } async def _execute_explainer( self, query: str, previous_results: Dict[str, Any], ) -> Dict[str, Any]: """Explainer: объяснение кода на естественном языке.""" from src.cli_agent.code_explainer import CodeExplainer console.print(f"[cyan]💡 Explainer: объяснение '{query}'[/cyan]") try: explainer = CodeExplainer( smart_client=self.smart_client, project_root=self.project_root, mcp_url=self.mcp_url, model=self.model, ) result = await _with_progress( explainer.explain(query), message="Explainer: анализирую код..." ) console.print(Panel( result.explanation, title="Code Explanation", border_style="green" )) return { "status": "completed" if result.success else "partial", "query": query, "explanation": result.explanation, "target_file": result.target_file, "target_function": result.target_function, "similar_implementations": result.similar_implementations, "impact_files": result.impact_files, "message": result.to_summary(), } except Exception as e: logger.error(f"Explainer execution failed: {e}", exc_info=True) return { "status": "failed", "error": str(e), "message": f"Explainer failed: {e}", } async def _execute_arch_analyzer( self, query: str, previous_results: Dict[str, Any], ) -> Dict[str, Any]: """Arch Analyzer: анализ архитектуры.""" from src.cli_agent.architecture_analyzer import ArchitectureAnalyzer console.print(f"[cyan]🏛️ Arch Analyzer: анализ '{query}'[/cyan]") try: analyzer = ArchitectureAnalyzer( smart_client=self.smart_client, project_root=self.project_root, mcp_url=self.mcp_url, model=self.model, ) # Определяем scope: весь проект или конкретный модуль if query.strip() in (".", "project", ""): result = await analyzer.analyze_project() else: result = await analyzer.analyze_module(query) console.print(Panel( result.report, title="Architecture Analysis", border_style="blue" )) return { "status": "completed" if result.success else "partial", "scope": result.scope, "report": result.report, "zones_found": result.zones_found, "dependencies_count": result.dependencies_count, "violations_count": result.violations_count, "message": result.to_summary(), } except Exception as e: logger.error(f"Arch Analyzer execution failed: {e}", exc_info=True) return { "status": "failed", "error": str(e), "message": f"Arch Analyzer failed: {e}", } def _get_role_description(self, role: AgentRole) -> str: """Возвращает описание роли для пользователя.""" descriptions = { AgentRole.STRATEGIST: ( "🧠 **Strategist** проведёт анализ:\n" " • Поиск похожих реализаций\n" " • Анализ impact radius\n" " • Определение architectural constraints\n" " • Создание детального плана" ), AgentRole.EXECUTOR: ( "🔧 **Executor** выполнит работу:\n" " • Генерация кода\n" " • Использование инструментов (search, edit, test)\n" " • Сохранение файлов" ), AgentRole.REVIEWER: ( "🔍 **Reviewer** проверит качество:\n" " • Статический анализ (mypy, ruff)\n" " • Architectural boundary checks\n" " • LLM-based review\n" " • Вердикт: APPROVED или BLOCKER" ), AgentRole.EXPLORER: ( "🔎 **Explorer** исследует кодовую базу:\n" " • Семантический поиск существующих реализаций\n" " • Pattern search (grep)\n" " • Анализ архитектурных зон\n" " • Отчёт о duplication risks" ), AgentRole.EXPLAINER: ( "💡 **Explainer** объяснит код:\n" " • Поиск целевой функции/модуля\n" " • Чтение и анализ кода\n" " • Объяснение на естественном языке\n" " • Примеры использования" ), AgentRole.ARCH_ANALYZER: ( "🏛️ **Arch Analyzer** проанализирует архитектуру:\n" " • Карта проекта по зонам\n" " • Граф зависимостей (DAG)\n" " • Проверка границ модулей\n" " • Рекомендации по улучшению" ), } return descriptions.get(role, f"Неизвестная роль: {role.value}")