Files
MultiAgent/orchestrator.py

348 lines
14 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# orchestrator.py
from typing import Dict, Any, Literal, TypedDict, List
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 import DeveloperAgent, ReviewerAgent
from config import MAX_REVIEW_RETRIES
class AgentState(TypedDict):
task: str
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.max_retries = MAX_REVIEW_RETRIES
# Создаем граф состояний
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):
print(f"\n🔧 [Iteration {state.get('iteration', 0)}] Планировщик начал работу...")
# Сначала получим описания инструментов для планировщика
tools_desc = [{"name": t.name, "description": t.description, "args": t.args} for t in self.tools]
plan = self.planner.create_plan(state["task"], tools_desc)
state["plan"] = plan
return state
def _exec_node(self, state: AgentState):
"""Узел разработчика"""
print(f"\n🔧 [Iteration {state.get('iteration', 0)}] Исполнитель начал работу...")
# Логируем начало работы разработчика
results = self.executor.execute_plan(state["plan"])
state["execution_results"] = results
return state
def _review_node(self, state: AgentState):
print(f"\n🔧 [Iteration {state.get('iteration', 0)}] Ревьювер начал работу...")
review = self.reviewer.review(state["task"], state["plan"], state["execution_results"])
state["review"] = review
state["iteration"] = state.get("iteration", 0) + 1
return state
def _after_review(self, state: AgentState):
status = state["review"].get("status")
if status == "approved":
return "approved"
elif status == "changes_requested":
if state["iteration"] >= state.get("max_iterations", 3):
return "reject"
# пробуем перевыполнить те же шаги (может, ошибка временная)
return "rework"
elif status == "rejected":
if state["iteration"] >= state.get("max_iterations", 3):
return "reject"
# нужен новый план
return "replan"
else:
return "reject"
def _developer_node(self, state: Dict) -> Dict:
"""Узел разработчика"""
task = state["task"]
print(f"\n🔧 [Iteration {task.iteration}] Разработчик начал работу...")
# Логируем начало работы разработчика
self._add_history(
task, "developer", "start",
input_summary=f"Requirement: {task.requirement[:100]}...",
details={"iteration": task.iteration, "previous_comments_count": len(task.last_review.comments) if task.last_review else 0}
)
# Получаем замечания с предыдущего ревью
previous_comments = []
if task.last_review and task.last_review.comments:
previous_comments = [
{"file": c.file, "line": c.line, "text": c.text}
for c in task.last_review.comments
]
try:
# Генерируем код
solution = self.developer.develop(task, previous_comments)
task.solution = solution
task.status = TaskStatus.REVIEW
# Логируем успех
self._add_history(
task, "developer", "complete",
output_summary=f"Generated {len(solution.files)} file(s), {len(solution.tests)} test(s)",
details={"files": list(solution.files.keys()), "tests": list(solution.tests.keys())}
)
print(f"✅ Разработчик создал решение: {len(solution.files)} файлов")
for filename in solution.files:
print(f" - {filename}")
except Exception as e:
task.status = TaskStatus.FAILED
task.error_log.append(f"Developer error: {str(e)}")
# Логируем ошибку
self._add_history(
task, "developer", "error",
error=error_msg,
details={"exception": str(e)}
)
print(f"❌ Ошибка разработчика: {e}")
return {"task": task}
def _reviewer_node(self, state: Dict) -> Dict:
"""Узел ревьювера"""
task = state["task"]
print(f"\n🔍 [Iteration {task.iteration}] Ревьювер проверяет код...")
# Логируем начало ревью
self._add_history(
task, "reviewer", "start",
input_summary=f"Reviewing {len(task.solution.files)} file(s)" if task.solution else "No solution",
details={"iteration": task.iteration}
)
if not task.solution:
task.status = TaskStatus.FAILED
error_msg = "No solution to review"
task.error_log.append(error_msg)
self._add_history(task, "reviewer", "error", error=error_msg)
return {"task": task}
try:
# Проводим ревью
review = self.reviewer.review(task, task.solution)
task.last_review = review
task.iteration += 1
# Логируем результат ревью
self._add_history(
task, "reviewer", "complete",
output_summary=f"Verdict: {review.status}, Comments: {len(review.comments)}",
details={
"verdict": review.status,
"comments_count": len(review.comments),
"summary": review.summary
}
)
print(f"📊 Вердикт: {review.status}")
if review.comments:
print(f" Найдено замечаний: {len(review.comments)}")
for comment in review.comments[:3]: # Показываем первые 3
print(f" - [{comment.severity}] {comment.text[:100]}")
# Обновляем статус
if review.status == "approved":
task.status = TaskStatus.APPROVED
print("✅ Код одобрен!")
elif review.status == "rejected":
task.status = TaskStatus.REJECTED
print("❌ Код отклонен архитектурно")
else:
task.status = TaskStatus.REWORK
print("🔄 Требуются доработки")
except Exception as e:
task.status = TaskStatus.FAILED
error_msg = f"Reviewer error: {str(e)}"
task.error_log.append(error_msg)
self._add_history(task, "reviewer", "error", error=error_msg, details={"exception": str(e)})
return {"task": task}
def _decide_next_step(self, state: Dict) -> Literal["approved", "rework", "rejected", "failed"]:
"""Решает, что делать дальше"""
task = state["task"]
if task.status == TaskStatus.APPROVED:
return "approved"
elif task.status == TaskStatus.REJECTED:
return "rejected"
elif task.status == TaskStatus.FAILED:
return "failed"
elif task.iteration >= self.max_retries:
print(f"\n⚠️ Достигнут лимит итераций ({self.max_retries}). Останавливаемся.")
task.status = TaskStatus.FAILED
return "failed"
else:
return "rework"
def _finalize_node(self, state: Dict) -> Dict:
"""Финальный узел - подведение итогов"""
task = state["task"]
print("\n" + "="*50)
print("📋 ИТОГИ РАБОТЫ")
print("="*50)
print(f"Задача: {task.requirement[:100]}...")
print(f"Статус: {task.status}")
print(f"Итераций: {task.iteration}")
if task.status == TaskStatus.APPROVED and task.solution:
print(f"\n✅ Решение принято!")
print(f" Файлов: {len(task.solution.files)}")
print(f" Тестов: {len(task.solution.tests)}")
# Сохраняем решение в файлы
self._save_solution(task.solution)
elif task.last_review:
print(f"\n📝 Финальный вердикт: {task.last_review.summary}")
if task.error_log:
print(f"\n❌ Ошибки: {len(task.error_log)}")
for error in task.error_log:
print(f" - {error}")
return {"task": task, "finalized": True}
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: str, max_iterations:int = None) -> Dict:
"""Запускает процесс разработки"""
task_id = str(uuid.uuid4())
initial_state = {
"task": task,
"plan": [],
"execution_results": [],
"review": {},
"iteration": 0,
"max_iterations": max_iterations
}
print(f"\n🚀 Запуск мультиагентной системы")
print(f"📝 Задача: {task}")
print(f"🆔 ID: {task_id}")
# Запускаем граф
final_state = self.workflow.invoke(initial_state)
return final_state
# main.py
def main():
"""Пример использования"""
# Создаем оркестратора
orchestrator = Orchestrator()
# Пример задачи
requirement = """
Реализуй класс BankAccount с методами:
- deposit(amount): пополнение счета
- withdraw(amount): снятие средств (нельзя снять больше, чем есть)
- get_balance(): получение баланса
- add_interest(rate): добавление процентов (rate в процентах)
Требования:
- Нельзя создать счет с отрицательным балансом
- Все операции должны быть потокобезопасными (используй threading.Lock)
- Напиши юнит-тесты (pytest) для всех методов, включая граничные случаи
- Добавь docstring для всех методов
"""
# Запускаем процесс
result = orchestrator.run(requirement)
# Анализируем результат
task = result["task"]
if task.status == TaskStatus.APPROVED:
print("\n🎉 Успех! Код принят и сохранен в папке 'generated_code'")
else:
print(f"\n💔 Неудача. Статус: {task.status}")
if task.last_review:
print(f"Причина: {task.last_review.summary}")
if __name__ == "__main__":
main()