266 lines
10 KiB
Python
266 lines
10 KiB
Python
# orchestrator.py
|
||
from typing import Dict, Any, Literal, TypedDict, List, Optional
|
||
from langgraph.graph import StateGraph, END
|
||
from langgraph.checkpoint.memory import MemorySaver
|
||
import uuid
|
||
from datetime import datetime
|
||
from models import Task, TaskStatus, CodeSolution, ReviewResult, IterationRecord
|
||
from agent.developer import DeveloperAgent
|
||
from agent.reviewer import ReviewerAgent
|
||
from agent.planner import PlannerAgent
|
||
from agent.executor import ExecutorAgent
|
||
from agent.plan_reviewer import PlanReviewerAgent
|
||
from tool.filesystem import read_file, write_file, delete_lines, replace_text_in_file, search_in_files
|
||
|
||
class AgentState(TypedDict):
|
||
task: Task
|
||
plan: List[Dict]
|
||
execution_results: List[Dict]
|
||
review: Dict
|
||
iteration: int
|
||
max_iterations: int
|
||
|
||
class Orchestrator:
|
||
"""Оркестратор - управляет процессом разработки"""
|
||
|
||
def __init__(self):
|
||
self.developer = DeveloperAgent()
|
||
self.reviewer = ReviewerAgent()
|
||
self.planner = PlannerAgent()
|
||
self.plan_reviewer = PlanReviewerAgent()
|
||
|
||
# Инструменты для ExecutorAgent
|
||
self.tools = [read_file, write_file, delete_lines, replace_text_in_file, search_in_files]
|
||
self.executor = ExecutorAgent(self.tools)
|
||
|
||
self.max_retries = 3
|
||
|
||
# Создаем граф состояний
|
||
self.workflow = self._build_workflow()
|
||
self.checkpointer = MemorySaver()
|
||
|
||
def _add_history(self, task: Task, agent: str, action: str,
|
||
input_summary: str = None, output_summary: str = None,
|
||
details: Dict = None, error: str = None):
|
||
"""Добавляет запись в историю задачи"""
|
||
record = IterationRecord(
|
||
iteration=task.iteration,
|
||
timestamp=datetime.now().isoformat(),
|
||
agent=agent,
|
||
action=action,
|
||
input_summary=input_summary,
|
||
output_summary=output_summary,
|
||
details=details or {},
|
||
error=error
|
||
)
|
||
task.history.append(record)
|
||
|
||
def _build_workflow(self):
|
||
"""Строит граф процесса разработки"""
|
||
|
||
# Определяем состояния
|
||
workflow = StateGraph(AgentState)
|
||
|
||
# Добавляем узлы
|
||
workflow.add_node("planner", self._plan_node)
|
||
workflow.add_node("executor", self._exec_node)
|
||
workflow.add_node("reviewer", self._review_node)
|
||
|
||
# Определяем переходы
|
||
workflow.set_entry_point("planner")
|
||
workflow.add_edge("planner", "executor")
|
||
workflow.add_edge("executor", "reviewer")
|
||
|
||
workflow.add_conditional_edges("reviewer", self._after_review, {
|
||
"approved": END,
|
||
"rework": "executor", # повторить выполнение с тем же планом? или перепланировать?
|
||
"replan": "planner",
|
||
"reject": END
|
||
})
|
||
|
||
return workflow.compile()
|
||
|
||
def _plan_node(self, state: AgentState):
|
||
task = state["task"]
|
||
print(f"\n🔧 [Iteration {task.iteration}] Планировщик начал работу...")
|
||
|
||
# Создаем запись в истории
|
||
task.history.append(IterationRecord(
|
||
iteration=task.iteration,
|
||
timestamp=datetime.now().isoformat(),
|
||
agent="developer",
|
||
action="start",
|
||
input_summary=f"Requirement: {task.requirement[:100]}...",
|
||
details={}
|
||
))
|
||
|
||
# Описываем доступные инструменты для планировщика
|
||
tools_desc = [
|
||
{"name": "read_file", "description": "Прочитать файл", "args": ["path", "start_line", "end_line"]},
|
||
{"name": "write_file", "description": "Записать файл", "args": ["path", "content"]},
|
||
{"name": "delete_lines", "description": "Удалить строки", "args": ["path", "start_line", "end_line"]},
|
||
{"name": "replace_text_in_file", "description": "Заменить текст", "args": ["path", "old", "new"]},
|
||
{"name": "search_in_files", "description": "Поиск в файлах", "args": ["directory", "pattern", "file_pattern"]}
|
||
]
|
||
|
||
# Создаем план
|
||
plan = self.planner.create_plan(task.requirement, tools_desc)
|
||
state["plan"] = plan
|
||
|
||
print(f"📋 План создан: {len(plan)} шагов")
|
||
return state
|
||
|
||
def _exec_node(self, state: AgentState):
|
||
"""Узел исполнителя"""
|
||
task = state["task"]
|
||
print(f"\n🔨 [Iteration {task.iteration}] Исполнитель начал работу...")
|
||
|
||
# Выполняем план через ExecutorAgent
|
||
results = self.executor.execute_plan(state["plan"])
|
||
state["execution_results"] = results
|
||
|
||
# Записываем результат выполнения
|
||
success_count = sum(1 for r in results if r.get("success", False))
|
||
print(f"✅ Выполнено шагов: {success_count}/{len(results)}")
|
||
|
||
return state
|
||
|
||
def _review_node(self, state: AgentState):
|
||
task = state["task"]
|
||
print(f"\n🔍 [Iteration {task.iteration}] План-ревьювер начал работу...")
|
||
|
||
# Добавляем запись о начале ревью
|
||
task.history.append(IterationRecord(
|
||
iteration=task.iteration,
|
||
timestamp=datetime.now().isoformat(),
|
||
agent="reviewer",
|
||
action="start",
|
||
input_summary=f"Reviewing plan with {len(state['plan'])} steps",
|
||
details={}
|
||
))
|
||
|
||
# Проводим ревью плана через PlanReviewerAgent
|
||
review = self.plan_reviewer.review(task.requirement, state["plan"], state["execution_results"])
|
||
state["review"] = review
|
||
|
||
# Обновляем статус задачи
|
||
if review.get("status") == "approved":
|
||
task.status = TaskStatus.APPROVED
|
||
elif review.get("status") == "changes_requested":
|
||
task.status = TaskStatus.REWORK
|
||
else:
|
||
task.status = TaskStatus.FAILED
|
||
|
||
# Увеличиваем счётчик итераций
|
||
state["iteration"] = state.get("iteration", 0) + 1
|
||
|
||
# Добавляем запись в историю
|
||
task.history.append(IterationRecord(
|
||
iteration=task.iteration,
|
||
timestamp=datetime.now().isoformat(),
|
||
agent="reviewer",
|
||
action="complete",
|
||
output_summary=f"Verdict: {review.get('status')}",
|
||
details={
|
||
"verdict": review.get("status"),
|
||
"feedback": review.get("feedback", "")[:200]
|
||
}
|
||
))
|
||
|
||
print(f"📊 Вердикт: {review.get('status')}")
|
||
if review.get("feedback"):
|
||
print(f" Feedback: {review['feedback'][:100]}...")
|
||
|
||
return state
|
||
|
||
def _after_review(self, state: AgentState):
|
||
task = state["task"]
|
||
review = state["review"]
|
||
|
||
# Определяем статус
|
||
status = review.get("status") if isinstance(review, dict) else None
|
||
|
||
if status == "approved":
|
||
return "approved"
|
||
elif status == "changes_requested":
|
||
if state.get("iteration", 0) >= state.get("max_iterations", 3):
|
||
task.status = TaskStatus.FAILED
|
||
return "reject"
|
||
task.status = TaskStatus.REWORK
|
||
return "rework"
|
||
elif status == "rejected":
|
||
if state.get("iteration", 0) >= state.get("max_iterations", 3):
|
||
task.status = TaskStatus.FAILED
|
||
return "reject"
|
||
task.status = TaskStatus.REWORK
|
||
return "replan"
|
||
else:
|
||
task.status = TaskStatus.FAILED
|
||
return "reject"
|
||
|
||
def _save_solution(self, solution: CodeSolution):
|
||
"""Сохраняет решение в файлы"""
|
||
import os
|
||
output_dir = "generated_code"
|
||
os.makedirs(output_dir, exist_ok=True)
|
||
|
||
for filename, content in solution.files.items():
|
||
filepath = os.path.join(output_dir, filename)
|
||
with open(filepath, "w", encoding="utf-8") as f:
|
||
f.write(content)
|
||
print(f" 💾 Сохранен файл: {filepath}")
|
||
|
||
for testname, test_content in solution.tests.items():
|
||
testpath = os.path.join(output_dir, testname)
|
||
with open(testpath, "w", encoding="utf-8") as f:
|
||
f.write(test_content)
|
||
print(f" 💾 Сохранен тест: {testpath}")
|
||
|
||
def run(self, task_text: str, max_iterations: int = 3) -> Dict:
|
||
"""Запускает процесс разработки"""
|
||
|
||
task_id = str(uuid.uuid4())
|
||
|
||
# Создаем объект Task
|
||
task = Task(
|
||
id=task_id,
|
||
requirement=task_text,
|
||
status=TaskStatus.PENDING,
|
||
iteration=0,
|
||
solution=None,
|
||
last_review=None,
|
||
error_log=[],
|
||
history=[]
|
||
)
|
||
|
||
initial_state = {
|
||
"task": task,
|
||
"plan": [],
|
||
"execution_results": [],
|
||
"review": {},
|
||
"iteration": 0,
|
||
"max_iterations": max_iterations
|
||
}
|
||
|
||
print(f"\n🚀 Запуск мультиагентной системы")
|
||
print(f"📝 Задача: {task_text[:100]}...")
|
||
print(f"🆔 ID: {task_id}")
|
||
|
||
# Запускаем граф
|
||
final_state = self.workflow.invoke(initial_state)
|
||
|
||
# Обновляем итоговый статус в задаче
|
||
final_task = final_state["task"]
|
||
|
||
print("\n" + "="*50)
|
||
print("📋 ИТОГИ РАБОТЫ")
|
||
print("="*50)
|
||
print(f"Статус: {final_task.status.value}")
|
||
print(f"Итераций: {final_task.iteration}")
|
||
|
||
if final_task.error_log:
|
||
print(f"\n❌ Ошибки: {len(final_task.error_log)}")
|
||
for error in final_task.error_log:
|
||
print(f" - {error}")
|
||
|
||
return {"task": final_task} |