diff --git a/core/agent.py b/core/agent.py index 9f79f56..78e48de 100644 --- a/core/agent.py +++ b/core/agent.py @@ -2,9 +2,10 @@ import asyncio import logging import lmstudio as lms from typing import Any, Callable +from core.interfaces import Agent from .prompt import CAVEMAN_PROMPT -logger: logging.Logger = logging.getLogger("agent-caveman") +logger: logging.Logger = logging.getLogger("agent-base") class _ActResponseCapture: @@ -53,17 +54,17 @@ class _ActResponseCapture: return '\n'.join(self.responses) if self.responses else "No response captured." -class CavemanAgent: - """Caveman AI agent - minimal token usage variant.""" +class BaseAgent(Agent): + """Base AI agent implementing common LMStudio interaction patterns.""" def __init__(self, model_name: str) -> None: self.model_name: str = model_name self.model: Any | None = None - self.system_prompt: str = CAVEMAN_PROMPT + self.system_prompt: str = "" async def initialize(self) -> None: """Initialize the LM Studio model.""" - logger.info(f"Initializing CavemanAgent with model: {self.model_name}") + logger.info(f"Initializing agent with model: {self.model_name}") self.model = lms.llm(self.model_name) async def run(self, user_input: str) -> str: @@ -78,12 +79,12 @@ class CavemanAgent: ] try: - logger.info(f"Running CavemanAgent interactively (input length: {len(user_input)})") + logger.info(f"Running agent interactively (input length: {len(user_input)})") response = await self.model.respond(user_input, messages=messages) - logger.info(f"CavemanAgent responded successfully (response length: {len(response)})") + logger.info(f"Agent responded successfully (response length: {len(response)})") return response except Exception as e: - logger.error(f"CavemanAgent execution error: {e}") + logger.error(f"Agent execution error: {e}") return f"Error in agent execution: {str(e)}" async def run_with_tools(self, user_input: str, tools: list[Any]) -> str: @@ -94,14 +95,22 @@ class CavemanAgent: try: capture = _ActResponseCapture() - logger.info(f"Calling LMStudio act() on CavemanAgent with {len(tools)} tools...") + logger.info(f"Calling LMStudio act() on agent with {len(tools)} tools...") result: lms.ActResult = self.model.act(user_input, tools=tools, on_message=capture) - logger.info(f"act() on CavemanAgent returned: {result}") + logger.info(f"act() on agent returned: {result}") response: str = capture.full_response if not response or response == "No response captured.": logger.warning(f"Act completed with {result.rounds} rounds but no response was captured.") return f"Act completed with {result.rounds} rounds but no response captured." return response except Exception as e: - logger.error(f"CavemanAgent tool execution error: {e}") + logger.error(f"Agent tool execution error: {e}") return f"Error in agent tool execution: {str(e)}" + + +class CavemanAgent(BaseAgent): + """Caveman AI agent - minimal token usage variant.""" + + def __init__(self, model_name: str) -> None: + super().__init__(model_name) + self.system_prompt = CAVEMAN_PROMPT diff --git a/core/coding_agent.py b/core/coding_agent.py index cc40999..b7cd64d 100644 --- a/core/coding_agent.py +++ b/core/coding_agent.py @@ -1,108 +1,13 @@ -import asyncio import logging -import lmstudio as lms -from typing import Any, Callable -from .prompt import CAVEMAN_PROMPT +from core.agent import BaseAgent from .coding_prompt import CODING_AGENT_SYSTEM_PROMPT logger: logging.Logger = logging.getLogger("agent-coding") -class _ActResponseCapture: - """Captures the AI response from LMStudio act() callback.""" - - def __init__(self) -> None: - self.responses: list[str] = [] - - def __call__(self, message: Any) -> None: - content: str = "" - if hasattr(message, 'content'): - content = message.content - elif hasattr(message, 'text'): - content = message.text - elif hasattr(message, 'response'): - content = message.response - elif hasattr(message, 'message'): - content = message.message - else: - return - - if isinstance(content, list): - parts: list[str] = [] - for item in content: - if isinstance(item, dict): - text: str = item.get('text', '') - if isinstance(text, list): - parts.extend([str(t) for t in text]) - else: - parts.append(str(text)) - elif isinstance(item, str): - parts.append(item) - elif hasattr(item, 'text'): - parts.append(str(item.text)) - elif hasattr(item, 'content'): - parts.append(str(item.content)) - content = ''.join(parts) - elif not isinstance(content, str): - content = str(content) - - if content.strip(): - self.responses.append(content.strip()) - - @property - def full_response(self) -> str: - return '\n'.join(self.responses) if self.responses else "No response captured." - - -class CodingAgent: - """AI agent that interacts with LMStudio models and tools.""" +class CodingAgent(BaseAgent): + """AI agent that interacts with LMStudio models and tools for coding tasks.""" def __init__(self, model_name: str) -> None: - self.model_name: str = model_name - self.model: Any | None = None - self.system_prompt: str = CODING_AGENT_SYSTEM_PROMPT - - async def initialize(self) -> None: - """Initialize the LM Studio model.""" - logger.info(f"Initializing CodingAgent with model: {self.model_name}") - self.model = lms.llm(self.model_name) - - async def run(self, user_input: str) -> str: - """Run a single interaction with the agent.""" - if self.model is None: - await self.initialize() - assert self.model is not None - - messages: list[dict[str, str]] = [ - {"role": "system", "content": self.system_prompt}, - {"role": "user", "content": user_input}, - ] - - try: - logger.info(f"Running CodingAgent interactively (input length: {len(user_input)})") - response = await self.model.respond(user_input, messages=messages) - logger.info(f"CodingAgent responded successfully (response length: {len(response)})") - return response - except Exception as e: - logger.error(f"CodingAgent execution error: {e}") - return f"Error in agent execution: {str(e)}" - - async def run_with_tools(self, user_input: str, tools: list[Callable[..., Any]]) -> str: - """Run the agent with tool calling capability.""" - if self.model is None: - await self.initialize() - assert self.model is not None - - try: - capture = _ActResponseCapture() - logger.info(f"Calling LMStudio act() on CodingAgent with {len(tools)} tools...") - result: lms.ActResult = self.model.act(user_input, tools=tools, on_message=capture) - logger.info(f"act() on CodingAgent returned: {result}") - response: str = capture.full_response - if not response or response == "No response captured.": - logger.warning(f"Act completed with {result.rounds} rounds but no response was captured.") - return f"Act completed with {result.rounds} rounds but no response captured." - return response - except Exception as e: - logger.error(f"CodingAgent tool execution error: {e}") - return f"Error in agent tool execution: {str(e)}" + super().__init__(model_name) + self.system_prompt = CODING_AGENT_SYSTEM_PROMPT diff --git a/core/coordinator_agent.py b/core/coordinator_agent.py new file mode 100644 index 0000000..e50f1b1 --- /dev/null +++ b/core/coordinator_agent.py @@ -0,0 +1,47 @@ +import logging +from typing import Any, Callable +from core.agent import BaseAgent +from core.prompts import COORDINATOR_SYSTEM_PROMPT +from core.coordinator_tools import CoordinatorTools + +logger: logging.Logger = logging.getLogger("agent-coordinator") + + +class CoordinatorNoToolCalledError(Exception): + """Raised when the Coordinator Agent completes execution without calling any routing tool.""" + pass + + +class CoordinatorAgent(BaseAgent): + """AI agent that coordinates Gitea issues and decides the next action.""" + + def __init__(self, model_name: str) -> None: + super().__init__(model_name) + self.system_prompt = COORDINATOR_SYSTEM_PROMPT + + async def decide_action( + self, + mission: str, + planning_tools: list[Callable[..., Any]], + coord_tools: CoordinatorTools, + ) -> str: + """Run the Coordinator Agent and ensure a tool is called.""" + coord_tools_list: list[Callable[..., Any]] = [ + coord_tools.propose_plan, + coord_tools.start_implementation, + coord_tools.answer_question, + coord_tools.close_issue, + coord_tools.take_no_action, + ] + combined_tools: list[Callable[..., Any]] = planning_tools + coord_tools_list + + logger.info("Running CoordinatorAgent to decide action...") + response_text: str = await self.run_with_tools(mission, combined_tools) + + if not coord_tools.tool_called: + logger.warning("CoordinatorAgent did not call any tools!") + raise CoordinatorNoToolCalledError( + "CoordinatorAgent failed to call a routing tool during execution." + ) + + return response_text diff --git a/core/dispatcher.py b/core/dispatcher.py index 221e0a0..0c01ee4 100644 --- a/core/dispatcher.py +++ b/core/dispatcher.py @@ -1,630 +1,200 @@ -"""Dispatches work to a single CodingAgent, one repo at a time.""" +"""Dispatches work to a specialized task processor, one repo at a time.""" import logging import re -from typing import Any -from core.coding_agent import CodingAgent +import os +import subprocess +from abc import ABC, abstractmethod +from typing import Any, Callable + from core.queue import WorkItem +from core.coding_agent import CodingAgent +from core.planning_agent import PlanningAgent +from core.coordinator_agent import CoordinatorAgent, CoordinatorNoToolCalledError +from core.factory import AgentFactory + from gitea.tools.coding_tools import CodingTools from gitea.tools.research_tools import ResearchTools from gitea.tools.gitea_tools import GiteaTools from gitea.client import GiteaClient -from core.coding_prompt import CODING_AGENT_SYSTEM_PROMPT from core.coordinator_tools import CoordinatorTools -from core.prompts import COORDINATOR_SYSTEM_PROMPT -from gitea.config import AGENT_MODEL_ID from gitea.workspace import WorkspaceManager +from gitea.config import AGENT_MODEL_ID from gitea.models import CommentModel, PullRequestFileModel, PullRequestModel, IssueModel + logger: logging.Logger = logging.getLogger("agent-dispatcher") AGENT_USERNAMES: frozenset[str] = frozenset({"meeks-ai", "agent-bot"}) -CLOSE_KEYWORDS_PATTERN = re.compile( +CLOSE_KEYWORDS_PATTERN: re.Pattern[str] = re.compile( rf"\b(?:close|closes|closed|fix|fixes|fixed|resolve|resolves|resolved)\s+#(\d+)\b", re.IGNORECASE ) -class AgentDispatcher: - """Dispatches work to a single CodingAgent, one repo at a time.""" +def _find_pr_for_issue_helper(client: GiteaClient, repo_full_name: str, issue_number: int) -> PullRequestModel | None: + """Find an open pull request that addresses the given issue number.""" + owner, repo_name = repo_full_name.split("/") + try: + prs = client.list_repo_pull_requests(owner, repo_name) + for pr in prs: + ref = pr.head.get("ref", "") if pr.head else "" + if re.search(rf"(? list[int]: + """Extract referenced issue numbers from the PR body.""" + matches = CLOSE_KEYWORDS_PATTERN.findall(pr_body) + return list(set(int(m) for m in matches)) + + +def _is_awaiting_reply_helper(comments: list[CommentModel]) -> bool: + """Return True if the agent's most recent comment contains the + awaiting-reply marker AND no human has commented after it. + """ + if not comments: + return False + # Find the last agent comment index + last_agent_idx: int = -1 + for i, c in enumerate(comments): + if c.user and c.user.login in AGENT_USERNAMES: + last_agent_idx = i + if last_agent_idx == -1: + return False + last_agent_comment = comments[last_agent_idx] + body = last_agent_comment.body or "" + # The agent explicitly embeds this marker when it is waiting for input + if "" not in body: + return False + # Check if any human replied AFTER the last agent comment + for c in comments[last_agent_idx + 1:]: + if c.user and c.user.login not in AGENT_USERNAMES: + return False # Human replied — we can proceed + return True # Agent signalled wait, no human replied yet + + +class TaskProcessor(ABC): + """Abstract base class for processing Gitea tasks.""" def __init__( self, client: GiteaClient, tools: GiteaTools, - model_name: str = AGENT_MODEL_ID, - max_retries: int = 2, - ) -> None: - self._client = client - self._tools = tools - self._model_name = model_name - self._max_retries = max_retries - - def _find_pr_for_issue(self, repo_full_name: str, issue_number: int) -> PullRequestModel | None: - """Find an open pull request that addresses the given issue number.""" - owner, repo_name = repo_full_name.split("/") - try: - prs = self._client.list_repo_pull_requests(owner, repo_name) - for pr in prs: - ref = pr.head.get("ref", "") if pr.head else "" - if re.search(rf"(? list[int]: - """Extract referenced issue numbers from the PR body.""" - matches = CLOSE_KEYWORDS_PATTERN.findall(pr_body) - return list(set(int(m) for m in matches)) - - def _is_awaiting_reply(self, comments: list[CommentModel]) -> bool: - """Return True if the agent's most recent comment contains the - awaiting-reply marker AND no human has commented after it. - The agent itself embeds the marker when it needs human input. - """ - if not comments: - return False - # Find the last agent comment index - last_agent_idx: int = -1 - for i, c in enumerate(comments): - if c.user and c.user.login in AGENT_USERNAMES: - last_agent_idx = i - if last_agent_idx == -1: - return False - last_agent_comment = comments[last_agent_idx] - body = (last_agent_comment.body or "") - # The agent explicitly embeds this marker when it is waiting for input - if "" not in body: - return False - # Check if any human replied AFTER the last agent comment - for c in comments[last_agent_idx + 1:]: - if c.user and c.user.login not in AGENT_USERNAMES: - return False # Human replied — we can proceed - return True # Agent signalled wait, no human replied yet - - async def dispatch( - self, + model_name: str, repo: str, - work_items: list[WorkItem], - ) -> list[str]: - """Dispatch all work for a single repo to a fresh agent, then discard it.""" - workspace = WorkspaceManager() - repo_path = workspace.get_repo_path(repo) - coding_tools = CodingTools(str(repo_path)) - research_tools = ResearchTools() + item: WorkItem, + ai_username: str, + ) -> None: + self.client = client + self.tools = tools + self.model_name = model_name + self.repo = repo + self.item = item + self.ai_username = ai_username - planning_tools: list[Any] = [ - self._tools.get_issue, - self._tools.get_pull_request, - self._tools.list_issues, - self._tools.list_pull_requests, - self._tools.get_file_content, - self._tools.get_issue_comments, - self._tools.get_pull_request_comments, - self._tools.get_pull_request_diff, - self._tools.get_pull_request_patch, - coding_tools.list_files, - coding_tools.read_file, - coding_tools.grep_search, - coding_tools.get_working_directory, - coding_tools.run_command, - research_tools.web_search, - research_tools.fetch_url, + self.owner, self.repo_name = repo.split("/") + self.workspace = WorkspaceManager() + self.repo_path = self.workspace.get_repo_path(repo) + self.coding_tools = CodingTools(str(self.repo_path)) + self.research_tools = ResearchTools() + + self.planning_tools: list[Callable[..., Any]] = [ + self.tools.get_issue, + self.tools.get_pull_request, + self.tools.list_issues, + self.tools.list_pull_requests, + self.tools.get_file_content, + self.tools.get_issue_comments, + self.tools.get_pull_request_comments, + self.tools.get_pull_request_diff, + self.tools.get_pull_request_patch, + self.coding_tools.list_files, + self.coding_tools.read_file, + self.coding_tools.grep_search, + self.coding_tools.get_working_directory, + self.coding_tools.run_command, + self.research_tools.web_search, + self.research_tools.fetch_url, ] - coding_tools_list: list[Any] = [ - self._tools.get_issue, - self._tools.get_pull_request, - self._tools.list_issues, - self._tools.list_pull_requests, - self._tools.get_file_content, - self._tools.create_pull_request, - self._tools.update_pull_request, - self._tools.add_label_to_issue, - self._tools.add_label_to_pr, - self._tools.create_branch, - self._tools.commit_file, - self._tools.create_issue, - self._tools.add_comment_to_issue, - self._tools.close_issue, - self._tools.close_pull_request, - self._tools.get_issue_comments, - self._tools.get_pull_request_comments, - self._tools.add_comment, - self._tools.add_label, - self._tools.update_file, - self._tools.get_pull_request_diff, - self._tools.get_pull_request_patch, - self._tools.approve_pull_request, - self._tools.request_changes, - coding_tools.list_files, - coding_tools.read_file, - coding_tools.write_file, - coding_tools.edit_file, - coding_tools.run_command, - coding_tools.grep_search, - coding_tools.get_working_directory, - research_tools.web_search, - research_tools.fetch_url, + self.coding_tools_list: list[Callable[..., Any]] = [ + self.tools.get_issue, + self.tools.get_pull_request, + self.tools.list_issues, + self.tools.list_pull_requests, + self.tools.get_file_content, + self.tools.create_pull_request, + self.tools.update_pull_request, + self.tools.add_label_to_issue, + self.tools.add_label_to_pr, + self.tools.create_branch, + self.tools.commit_file, + self.tools.create_issue, + self.tools.add_comment_to_issue, + self.tools.close_issue, + self.tools.close_pull_request, + self.tools.get_issue_comments, + self.tools.get_pull_request_comments, + self.tools.add_comment, + self.tools.add_label, + self.tools.update_file, + self.tools.get_pull_request_diff, + self.tools.get_pull_request_patch, + self.tools.approve_pull_request, + self.tools.request_changes, + self.coding_tools.list_files, + self.coding_tools.read_file, + self.coding_tools.write_file, + self.coding_tools.edit_file, + self.coding_tools.run_command, + self.coding_tools.grep_search, + self.coding_tools.get_working_directory, + self.research_tools.web_search, + self.research_tools.fetch_url, ] - results: list[str] = [] + @abstractmethod + async def process(self, attempt_limit: int) -> str: + """Execute the task flow, including planning, coding, or coordination.""" + pass - import os - original_cwd = os.getcwd() - changed_dir = False - if os.path.isdir(str(repo_path)): - os.chdir(str(repo_path)) - changed_dir = True + +class PRTaskProcessor(TaskProcessor): + """Processes Gitea pull requests (code reviews and bug fixes).""" + + def _build_pr_mission(self, pr_info: PullRequestModel, is_own_pr: bool) -> str: + pr_number = self.item.task_number + pr_details = pr_info.model_dump_json(indent=2) + pr_diff = "" try: - # Get authenticated username for reviewer filter - ai_username = "meeks-ai" - try: - user = self._client.get_authenticated_user() - if user: - ai_username = user.login - except Exception: - pass - - for item in work_items: - owner, repo_name = repo.split("/") - - if item.task_type == "issue": - existing_pr = self._find_pr_for_issue(repo, item.task_number) - is_wip = False - has_request_changes = False - - if existing_pr: - if existing_pr.title: - title_upper = existing_pr.title.strip().upper() - is_wip = title_upper.startswith("WIP:") or "[WIP]" in title_upper - - try: - reviews = self._client.get_pr_reviews(owner, repo_name, existing_pr.number) - has_request_changes = any(r.get("state") == "REQUEST_CHANGES" for r in reviews) - except Exception as e: - logger.warning(f"Error checking reviews for PR #{existing_pr.number}: {e}") - - if not is_wip and not has_request_changes: - logger.info(f"Issue #{item.task_number} already has open PR #{existing_pr.number}. Skipping.") - results.append(f"SKIP: A pull request (PR #{existing_pr.number}) addressing issue #{item.task_number} already exists.") - continue - - # Check comments on issue and PR - issue_comments = [] - try: - issue_comments = self._client.get_issue_comments(owner, repo_name, item.task_number) - except Exception: - pass - - pr_comments = [] - if existing_pr: - try: - pr_comments = self._client.get_pull_request_comments(owner, repo_name, existing_pr.number) - except Exception: - pass - - if self._is_awaiting_reply(issue_comments) or self._is_awaiting_reply(pr_comments): - logger.info(f"Issue #{item.task_number}: awaiting human reply. Skipping.") - results.append(f"SKIP: Awaiting human reply on issue #{item.task_number} or PR.") - continue - - elif item.task_type == "pr": - # Check if the agent is requested/assigned as a reviewer - try: - pr_detail = self._client.get_pull_request(owner, repo_name, item.task_number) - except Exception as e: - logger.warning(f"Error fetching PR #{item.task_number} detail: {e}") - results.append(f"FAILED: Could not fetch details for PR #{item.task_number}.") - continue - - is_own_pr = (pr_detail.user and pr_detail.user.login == ai_username) - is_requested_reviewer = any(r.login == ai_username for r in pr_detail.requested_reviewers) - - if not is_own_pr and not is_requested_reviewer: - logger.info(f"PR #{item.task_number}: Agent is not a requested reviewer. Skipping.") - results.append(f"SKIP: Agent is not a requested reviewer on PR #{item.task_number}.") - continue - - # Check if we're waiting for a human reply before acting on a PR - pr_comments = [] - try: - pr_comments = self._client.get_pull_request_comments(owner, repo_name, item.task_number) - except Exception: - pass - if self._is_awaiting_reply(pr_comments): - logger.info(f"PR #{item.task_number}: agent asked a question and is awaiting a human reply. Skipping.") - results.append(f"SKIP: Awaiting human reply on PR #{item.task_number}.") - continue - - for attempt in range(1, self._max_retries + 1): - try: - if item.task_type == "pr": - # Process PR as before - base_mission = self._build_pr_mission(item) - if base_mission.startswith("SKIP:"): - logger.info(f"Skipping task #{item.task_number}: {base_mission}") - results.append(base_mission) - break - - logger.info(f"Starting Planning Phase for PR #{item.task_number} (attempt {attempt})") - planning_mission = ( - f"PHASE 1: PLANNING PHASE\n\n" - f"Your task is to research the problem, analyse the repository structure, and produce a detailed implementation plan.\n" - f"Original Mission details:\n{base_mission}\n\n" - f"CRITICAL RULES:\n" - f"1. You are ONLY generating an implementation plan. DO NOT write files, DO NOT edit files, DO NOT commit, DO NOT push, and DO NOT create branches or PRs.\n" - f"2. RESEARCH FIRST (mandatory before writing the plan):\n" - f" - Use web_search to find relevant documentation, known solutions, library APIs, error explanations, and best practices.\n" - f" - Use fetch_url to read specific documentation pages, changelogs, or Stack Overflow answers in full.\n" - f"3. EXPLORE the codebase using read_file, list_files, grep_search, or run_command.\n" - f"4. Output your final plan clearly.\n" - ) - planning_agent = CodingAgent(self._model_name) - plan = await planning_agent.run_with_tools(planning_mission, planning_tools) - - logger.info(f"Starting Coding Phase for PR #{item.task_number} (attempt {attempt})") - coding_mission = ( - f"PHASE 2: EXECUTION/CODING PHASE\n\n" - f"You must now implement the changes based on the following plan:\n" - f"--- PLAN ---\n{plan}\n--- PLAN END ---\n\n" - f"Original Mission details:\n{base_mission}\n\n" - f"Follow the workflow to implement changes, verify, and complete PR review/updates.\n" - ) - coding_agent = CodingAgent(self._model_name) - response = await coding_agent.run_with_tools(coding_mission, coding_tools_list) - results.append(response) - break - - else: - # Handle Issue Task (Redesigned Planning/Question Board Workflow) - issue_info = item.task_info - assert isinstance(issue_info, IssueModel) - title = issue_info.title - issue_body = issue_info.body or "No description provided." - issue_user = issue_info.user.login if issue_info.user else "unknown" - - # Format comments and reviews for the prompt - issue_comments_str = "\n".join([ - f"- @{c.user.login} ({c.created_at}): {c.body}" for c in issue_comments - ]) if issue_comments else "No comments yet." - - pr_info_str = "No existing PR." - if existing_pr: - pr_comments_str = "\n".join([ - f"- @{c.user.login} ({c.created_at}): {c.body}" for c in pr_comments - ]) if pr_comments else "No PR comments yet." - try: - reviews = self._client.get_pr_reviews(owner, repo_name, existing_pr.number) - reviews_str = "\n".join([ - f"- @{r.get('user', {}).get('login')} ({r.get('submitted_at')}): [{r.get('state')}] {r.get('body')}" - for r in reviews - ]) if reviews else "No reviews yet." - except Exception: - reviews_str = "No reviews available." - pr_info_str = ( - f"PR Number: #{existing_pr.number}\n" - f"PR Title: {existing_pr.title}\n" - f"PR Branch: {existing_pr.head.get('ref', 'unknown')}\n" - f"PR State: {existing_pr.state}\n" - f"PR Comments:\n{pr_comments_str}\n" - f"PR Reviews:\n{reviews_str}" - ) - - state_analysis_mission = ( - f"Analyzing issue #{item.task_number} in '{repo}'.\n\n" - f"Issue Title: {title}\n" - f"Description:\n{issue_body}\n\n" - f"Issue Comments:\n{issue_comments_str}\n\n" - f"Existing PR Details:\n{pr_info_str}\n" - ) - - # Run Coordinator Agent with CoordinatorTools - coord_tools = CoordinatorTools() - coord_tools_list = [ - coord_tools.propose_plan, - coord_tools.start_implementation, - coord_tools.answer_question, - coord_tools.close_issue, - coord_tools.take_no_action, - ] - combined_tools = planning_tools + coord_tools_list - - logger.info(f"Analyzing conversation state for issue #{item.task_number}...") - planning_agent = CodingAgent(self._model_name) - planning_agent.system_prompt = COORDINATOR_SYSTEM_PROMPT - - # Let the coordinator analyze and route - response_text = await planning_agent.run_with_tools(state_analysis_mission, combined_tools) - - # Fallback if no tool was called - if not coord_tools.tool_called: - logger.info(f"Coordinator agent did not call any tools. Falling back to JSON text parsing.") - import json - decision = {} - json_match = re.search(r"```json\s*(.*?)\s*```", response_text, re.DOTALL) - if json_match: - json_str = json_match.group(1).strip() - else: - json_str = response_text.strip() - try: - decision = json.loads(json_str) - except Exception as e: - try: - start_idx = json_str.find('{') - end_idx = json_str.rfind('}') - if start_idx != -1 and end_idx != -1: - decision = json.loads(json_str[start_idx:end_idx+1]) - except Exception: - pass - - if decision and "action" in decision: - coord_tools.action = decision["action"] - if coord_tools.action == "PROPOSE_PLAN": - coord_tools.arguments = { - "comment_body": decision.get("comment_body", "") or decision.get("reasoning", ""), - "issue_number": item.task_number - } - elif coord_tools.action == "ANSWER_QUESTION": - coord_tools.arguments = { - "comment_body": decision.get("comment_body", "") or decision.get("reasoning", ""), - "issue_number": item.task_number - } - elif coord_tools.action == "CLOSE_ISSUE": - coord_tools.arguments = { - "comment": decision.get("comment_body", "Closing the issue as resolved."), - "issue_number": item.task_number - } - elif coord_tools.action == "EXECUTE_PLAN": - coord_tools.arguments = { - "approved_plan": decision.get("approved_plan", ""), - "issue_number": item.task_number - } - - action = coord_tools.action - logger.info(f"Coordinator Decided Action: {action} (tool_called={coord_tools.tool_called})") - - if action == "PROPOSE_PLAN": - comment_body = coord_tools.arguments.get("comment_body", "") - if not comment_body: - plan = coord_tools.arguments.get("plan", "") - comment_body = ( - f"### Proposed Implementation Plan\n\n" - f"{plan}\n\n" - f"Is this plan ok for implementation or do you have any comments/changes?\n" - f"\n" - f"" - ) - self._client.add_comment(owner, repo_name, item.task_number, comment_body) - results.append(f"POSTED_COMMENT: PROPOSE_PLAN comment posted to issue #{item.task_number}.") - break - - elif action == "ANSWER_QUESTION": - comment_body = coord_tools.arguments.get("comment_body", "") - if not comment_body: - answer = coord_tools.arguments.get("answer", "") - comment_body = ( - f"{answer}\n\n" - f"Is this answer satisfactory?\n" - f"\n" - f"" - ) - self._client.add_comment(owner, repo_name, item.task_number, comment_body) - results.append(f"POSTED_COMMENT: ANSWER_QUESTION comment posted to issue #{item.task_number}.") - break - - elif action == "CLOSE_ISSUE": - comment = coord_tools.arguments.get("comment", "Closing the issue as resolved.") - self._client.add_comment(owner, repo_name, item.task_number, comment) - self._client.close_issue(owner, repo_name, item.task_number) - results.append(f"CLOSED_ISSUE: Issue #{item.task_number} closed.") - break - - elif action == "NO_ACTION": - results.append(f"NO_ACTION: No action taken on issue #{item.task_number}.") - break - - elif action == "EXECUTE_PLAN": - approved_plan = coord_tools.arguments.get("approved_plan", "") - pr_to_use = existing_pr - branch_name = "" - - if pr_to_use: - branch_name = pr_to_use.head.get("ref", "") - logger.info(f"Resuming work on existing PR #{pr_to_use.number} on branch '{branch_name}'") - else: - logger.info(f"Creating new WIP PR for issue #{item.task_number}") - clean_title = re.sub(r'[^a-zA-Z0-9\s-]', '', title).strip().lower() - title_words = clean_title.split()[:5] - desc_suffix = "-".join(title_words) - if not desc_suffix: - desc_suffix = "fix-issue" - branch_name = f"fix/issue-{item.task_number}-{desc_suffix}" - - try: - import subprocess - # Clean branch if exists, create fresh from master, and commit empty to push - subprocess.run(["git", "checkout", "master"], cwd=str(repo_path), check=True) - subprocess.run(["git", "pull", "origin", "master"], cwd=str(repo_path), check=True) - subprocess.run(["git", "branch", "-D", branch_name], cwd=str(repo_path), stderr=subprocess.DEVNULL) - subprocess.run(["git", "checkout", "-b", branch_name], cwd=str(repo_path), check=True) - subprocess.run(["git", "commit", "--allow-empty", "-m", f"WIP: start implementation for issue #{item.task_number}"], cwd=str(repo_path), check=True) - subprocess.run(["git", "push", "origin", branch_name], cwd=str(repo_path), check=True) - - # Create PR via Gitea client - pr_title = f"WIP: {title}" - pr_description = f"Work in progress for issue #{item.task_number}." - pr_to_use = self._client.create_pull_request( - owner, repo_name, head=branch_name, base="master", title=pr_title, description=pr_description - ) - - # Comment on the issue - pr_link = pr_to_use.html_url or f"{self._client.base_url}/{repo}/pulls/{pr_to_use.number}" - start_comment = f"Started work on PR #{pr_to_use.number} ({pr_link})." - self._client.add_comment(owner, repo_name, item.task_number, start_comment) - - logger.info(f"Successfully created WIP PR #{pr_to_use.number} on branch '{branch_name}'") - except Exception as e: - logger.error(f"Failed to create WIP PR for issue #{item.task_number}: {e}") - results.append(f"FAILED to create WIP PR: {e}") - break - - # Now run the Coding Phase on the PR branch - base_mission = self._build_issue_mission(item) - coding_mission = ( - f"PHASE 2: EXECUTION/CODING PHASE\n\n" - f"You are implementing changes for issue #{item.task_number} in repository '{repo}'.\n" - f"You are working on the existing Pull Request #{pr_to_use.number} on branch '{branch_name}'.\n\n" - f"--- APPROVED PLAN ---\n{approved_plan}\n--- APPROVED PLAN END ---\n\n" - f"Original Mission details:\n{base_mission}\n\n" - f"DIRECTIONS:\n" - f"1. Checkout the branch '{branch_name}' (it should already be checked out, or run `git checkout {branch_name}`).\n" - f"2. Implement the changes according to the APPROVED PLAN.\n" - f"3. Run verification/tests (check AGENTS.md for conventions).\n" - f"4. Commit and push your changes to origin on the branch '{branch_name}'.\n" - f"5. ONCE COMPLETED SUCCESSFULLY:\n" - f" - Call `update_pull_request(owner='{owner}', repo='{repo_name}', pull_number={pr_to_use.number}, title='{title}', body='')`.\n" - f" Note: The PR title must not contain 'WIP:'. The PR body must follow the mandatory PR template in CODING AGENT SYSTEM PROMPT.\n" - f" - Call `add_comment_to_issue(owner='{owner}', repo='{repo_name}', issue_number={item.task_number}, body='Work has been completed in PR #{pr_to_use.number}.')`.\n" - f"6. IF YOU ENCOUNTER A BLOCKER OR FAIL:\n" - f" - You are allowed (and encouraged) to comment on the WIP PR #{pr_to_use.number} (using `add_comment`) with any details, logs, or context to help the next agent resume the work.\n\n" - f"AVAILABLE RESEARCH TOOLS:\n" - f" - web_search(query, time_range, categories) — search the web via SearXNG/DuckDuckGo\n" - f" - fetch_url(url) — read any documentation page in full\n" - ) - logger.info(f"Starting Execution/Coding Phase for issue #{item.task_number} on branch '{branch_name}'") - coding_agent = CodingAgent(self._model_name) - response = await coding_agent.run_with_tools(coding_mission, coding_tools_list) - logger.info(f"Agent response for issue #{item.task_number}: {response}") - results.append(response) - break - - except Exception as e: - logger.error(f"Error processing issue/PR #{item.task_number} (attempt {attempt}/{self._max_retries}): {e}") - if attempt == self._max_retries: - results.append(f"FAILED after {self._max_retries} attempts: {str(e)}") - finally: - if changed_dir: - os.chdir(original_cwd) - - return results - - def _build_issue_mission(self, item: WorkItem) -> str: - issue_info = item.task_info - assert isinstance(issue_info, IssueModel) - - - repo_full_name: str = item.repo_full_name - issue_number: int = item.task_number - - issue_body: str = issue_info.body or "No description provided." - issue_labels: list[str] = [lbl.name for lbl in issue_info.labels] - issue_user: str = issue_info.user.login if issue_info.user else "unknown" - issue_created: str = issue_info.created_at or "unknown" - title: str = issue_info.title - - - owner: str = repo_full_name.split("/")[0] - repo_name: str = repo_full_name.split("/")[1] - - comments: list[CommentModel] = [] - try: - comments = self._client.get_issue_comments(owner, repo_name, issue_number) - except Exception: - pass - - labels_str: str = f"Labels: {', '.join(issue_labels)}" if issue_labels else "Labels: none" - comments_str: str = "\n".join([ - f"- @{c.user.login} ({c.created_at}): {c.body}" - for c in comments - ]) if comments else "No comments yet." - - clean_title = re.sub(r'[^a-zA-Z0-9\s-]', '', title).strip().lower() - title_words = clean_title.split()[:5] - desc_suffix = "-".join(title_words) - if not desc_suffix: - desc_suffix = "fix-issue" - branch_name: str = f"fix/issue-{issue_number}-{desc_suffix}" - - workspace = WorkspaceManager() - repo_path = workspace.get_repo_path(repo_full_name) - - return ( - f"Your mission is to resolve issue #{issue_number} in {repo_full_name}.\n\n" - f"Issue: {title}\n" - f"Author: @{issue_user} (created {issue_created})\n" - f"{labels_str}\n\n" - f"Description:\n{issue_body}\n\n" - f"Comments ({len(comments)}):\n{comments_str}\n\n" - f"Branch name: {branch_name}.\n\n" - "BEFORE WRITING ANY CODE:\n" - " - Search online for relevant documentation, known solutions, library APIs, and platform-specific behavior.\n" - " - If ANY part of the issue is unclear, ambiguous, or has multiple valid approaches:\n" - " → Post a comment on the issue using `add_comment_to_issue` with your specific question(s).\n" - " → List the approaches you are considering.\n" - " → End the comment with the marker: on its own line.\n" - " → STOP. Do NOT proceed until a human replies. The system will re-dispatch you once a human responds.\n" - " - Never assume or guess. Always prefer asking over guessing.\n\n" - "CRITICAL INSTRUCTIONS:\n" - f"1. The repo is already cloned locally at '{repo_path}'. DO NOT create a new repository.\n" - f" The repository directory is your current working directory (cwd). You can verify this via the `get_working_directory` tool.\n" - "2. Always start from master: `git checkout master && git pull origin master`\n" - "3. Branch from master: `git checkout -b /issue--`\n" - " Branch names MUST include a descriptive name (words/hyphens), not just the issue number.\n" - " Types: feat, fix, chore, docs, style, refactor, test, build, ci, perf\n" - "4. Use `edit_file`/`write_file` for code changes, then `git add` and `git commit` via `run_command`.\n" - "5. Push: `git push origin ` via `run_command`.\n" - "6. Create PR: Use the `create_pull_request` tool (do NOT use Gitea's `tea` CLI or GitHub's `gh` CLI in run_command, as they can hang/freeze interactively).\n" - " PR description MUST include the Gitea automation template with 'closes #'.\n" - "7. IMPORTANT: After successfully creating the pull request, you MUST comment on the issue (using the `add_comment_to_issue` tool) with the PR number, PR link, and summary.\n" - "8. Check AGENTS.md in repo root for project conventions and verification steps.\n" - "9. Grade severity: Critical/High = fix, Medium = review/fix, Low = skip.\n" - "An issue is DONE when the connected PR is merged (you cannot merge yourself).\n" - "DO NOT edit .git files unless explicitly resolving a git issue.\n" - "DO NOT work on non-meeks organization repos." - ) - - def _build_pr_mission(self, item: WorkItem) -> str: - pr_info = item.task_info - assert isinstance(pr_info, PullRequestModel) - - - repo_full_name: str = item.repo_full_name - pr_number: int = item.task_number - - owner: str = repo_full_name.split("/")[0] - repo_name: str = repo_full_name.split("/")[1] - - pr_model = self._client.get_pull_request(owner, repo_name, pr_number) - pr_details: str = pr_model.model_dump_json(indent=2) - pr_diff: str = "" - try: - pr_diff = self._client.get_pull_request_diff(owner, repo_name, pr_number) + pr_diff = self.client.get_pull_request_diff(self.owner, self.repo_name, pr_number) except Exception as e: logger.warning(f"Could not fetch PR diff for #{pr_number}: {e}") pr_diff = f"Error fetching diff: {e}" - pr_files: list[PullRequestFileModel] = [] try: - pr_files = self._client.get_pull_request_files(owner, repo_name, pr_number) + pr_files = self.client.get_pull_request_files(self.owner, self.repo_name, pr_number) except Exception: pass - files_summary: str = "\n".join([f"- {f.filename}" for f in pr_files]) if pr_files else "No files available." + files_summary = "\n".join([f"- {f.filename}" for f in pr_files]) if pr_files else "No files available." comments: list[CommentModel] = [] try: - comments = self._client.get_pull_request_comments(owner, repo_name, pr_number) + comments = self.client.get_pull_request_comments(self.owner, self.repo_name, pr_number) if not isinstance(comments, list): comments = [] except Exception: @@ -632,21 +202,12 @@ class AgentDispatcher: reviews: list[dict[str, Any]] = [] try: - reviews = self._client.get_pr_reviews(owner, repo_name, pr_number) + reviews = self.client.get_pr_reviews(self.owner, self.repo_name, pr_number) if not isinstance(reviews, list): reviews = [] except Exception: pass - ai_username = "meeks-ai" - try: - user = self._client.get_authenticated_user() - if user: - ai_username = user.login - except Exception: - pass - - # Create a combined, sorted timeline of timeline comments and reviews timeline: list[dict[str, Any]] = [] for c in comments: timeline.append({ @@ -654,7 +215,7 @@ class AgentDispatcher: "user": c.user.login, "type": "comment", "body": c.body, - "by_ai": "Reviewed by AI Agent" in c.body or c.user.login == ai_username + "by_ai": "Reviewed by AI Agent" in c.body or c.user.login == self.ai_username }) for r in reviews: r_user = (r.get("user") or {}).get("login", "unknown") @@ -665,7 +226,7 @@ class AgentDispatcher: "user": r_user, "type": "review", "body": f"[{r_state}] {r_body}", - "by_ai": r_user == ai_username + "by_ai": r_user == self.ai_username }) timeline.sort(key=lambda x: x["timestamp"]) @@ -674,33 +235,29 @@ class AgentDispatcher: if timeline: last_action_by_ai = timeline[-1]["by_ai"] - pr_author: str = pr_info.user.login if pr_info.user else "unknown" - is_own_pr = (pr_author == ai_username) - - # Skip if the latest action is already by AI (waiting for human turn) if last_action_by_ai: logger.info(f"PR #{pr_number} already addressed by AI. Skipping.") return f"SKIP: PR #{pr_number} has already been addressed by AI. No new action needed." - comments_str: str = "\n".join([ + comments_str = "\n".join([ f"- @{c.user.login} ({c.created_at}): {c.body}" for c in comments ]) if comments else "No comments yet." - reviews_str: str = "\n".join([ + reviews_str = "\n".join([ f"- @{(r.get('user') or {}).get('login')} ({r.get('submitted_at')}): [{r.get('state')}] {r.get('body')}" for r in reviews ]) if reviews else "No reviews yet." connected_issues_ctx = "" pr_body = pr_info.body or "" - linked_issues = self._find_issues_for_pr(pr_body) + linked_issues = _find_issues_for_pr_helper(pr_body) if linked_issues: issues_details = [] for issue_num in linked_issues: try: - issue = self._client.get_issue(owner, repo_name, issue_num) - issue_comments = self._client.get_issue_comments(owner, repo_name, issue_num) + issue = self.client.get_issue(self.owner, self.repo_name, issue_num) + issue_comments = self.client.get_issue_comments(self.owner, self.repo_name, issue_num) comments_list = "\n".join([ f" - @{c.user.login} ({c.created_at}): {c.body}" for c in issue_comments @@ -717,19 +274,16 @@ class AgentDispatcher: if issues_details: connected_issues_ctx = "\n---\n\n## 📋 CONNECTED ISSUE CONTEXT\n" + "\n\n".join(issues_details) - workspace = WorkspaceManager() - repo_path = workspace.get_repo_path(repo_full_name) - pr_head_branch: str = pr_info.head.get('ref', 'unknown') if pr_info.head else "unknown" - pr_base_branch: str = pr_info.base.get('ref', 'unknown') if pr_info.base else "unknown" - pr_state: str = pr_info.state - pr_created: str = pr_info.created_at or "unknown" + pr_head_branch = pr_info.head.get('ref', 'unknown') if pr_info.head else "unknown" + pr_base_branch = pr_info.base.get('ref', 'unknown') if pr_info.base else "unknown" + pr_state = pr_info.state + pr_created = pr_info.created_at or "unknown" - # Dynamically determine instructions based on ownership/comments is_fixing_pr = is_own_pr or any(r.get("state") == "REQUEST_CHANGES" for r in reviews) if is_fixing_pr: instructions = ( - f" Note: The repository is located locally at '{repo_path}'. This repository directory is your current working directory (cwd). You can verify this via the `get_working_directory` tool.\n" + f" Note: The repository is located locally at '{self.repo_path}'. This repository directory is your current working directory (cwd). You can verify this via the `get_working_directory` tool.\n" f" Your task is to FIX/UPDATE this PR by addressing comments/change requests.\n" f" DO NOT create a new branch or PR. Follow this exact workflow:\n" f" 1. Checkout the PR's head branch: `git checkout {pr_head_branch}`\n" @@ -740,7 +294,7 @@ class AgentDispatcher: ) else: instructions = ( - f" Note: The repository is located locally at '{repo_path}'. This repository directory is your current working directory (cwd). You can verify this via the `get_working_directory` tool.\n" + f" Note: The repository is located locally at '{self.repo_path}'. This repository directory is your current working directory (cwd). You can verify this via the `get_working_directory` tool.\n" " CRITICAL: You are ONLY reviewing this PR. DO NOT edit files, DO NOT make commits, DO NOT push branches, and DO NOT create any new PRs.\n" "1. Read the PR diff carefully.\n" "2. Analyze the changes for correctness, quality, and potential issues.\n" @@ -752,9 +306,9 @@ class AgentDispatcher: ) return ( - f"Your mission is to process PR #{pr_number} in {repo_full_name}.\n\n" + f"Your mission is to process PR #{pr_number} in {self.repo}.\n\n" f"PR: {pr_info.title}\n" - f"Author: @{pr_author}\n" + f"Author: @{pr_info.user.login if pr_info.user else 'unknown'}\n" f"Branch: {pr_head_branch} → {pr_base_branch}\n" f"State: {pr_state} (created {pr_created})\n\n" f"Description:\n{pr_info.body or 'No description'}\n\n" @@ -773,4 +327,437 @@ class AgentDispatcher: f"Instructions:\n{instructions}" ) + async def process(self, attempt_limit: int) -> str: + try: + pr_detail = self.client.get_pull_request(self.owner, self.repo_name, self.item.task_number) + except Exception as e: + logger.warning(f"Error fetching PR #{self.item.task_number} detail: {e}") + return f"FAILED: Could not fetch details for PR #{self.item.task_number}." + is_own_pr = (pr_detail.user and pr_detail.user.login == self.ai_username) + is_requested_reviewer = any(r.login == self.ai_username for r in pr_detail.requested_reviewers) + + if not is_own_pr and not is_requested_reviewer: + logger.info(f"PR #{self.item.task_number}: Agent is not a requested reviewer. Skipping.") + return f"SKIP: Agent is not a requested reviewer on PR #{self.item.task_number}." + + pr_comments = [] + try: + pr_comments = self.client.get_pull_request_comments(self.owner, self.repo_name, self.item.task_number) + except Exception: + pass + if _is_awaiting_reply_helper(pr_comments): + logger.info(f"PR #{self.item.task_number}: agent asked a question and is awaiting a human reply. Skipping.") + return f"SKIP: Awaiting human reply on PR #{self.item.task_number}." + + base_mission = self._build_pr_mission(pr_detail, is_own_pr) + if base_mission.startswith("SKIP:"): + logger.info(f"Skipping task #{self.item.task_number}: {base_mission}") + return base_mission + + for attempt in range(1, attempt_limit + 1): + try: + logger.info(f"Starting Planning Phase for PR #{self.item.task_number} (attempt {attempt})") + planning_mission = ( + f"PHASE 1: PLANNING PHASE\n\n" + f"Your task is to research the problem, analyse the repository structure, and produce a detailed implementation plan.\n" + f"Original Mission details:\n{base_mission}\n\n" + f"CRITICAL RULES:\n" + f"1. You are ONLY generating an implementation plan. DO NOT write files, DO NOT edit files, DO NOT commit, DO NOT push, and DO NOT create branches or PRs.\n" + f"2. RESEARCH FIRST (mandatory before writing the plan):\n" + f" - Use web_search to find relevant documentation, known solutions, library APIs, error explanations, and best practices.\n" + f" - Use fetch_url to read specific documentation pages, changelogs, or Stack Overflow answers in full.\n" + f"3. EXPLORE the codebase using read_file, list_files, grep_search, or run_command.\n" + f"4. Output your final plan clearly.\n" + ) + planning_agent = PlanningAgent(self.model_name) + plan = await planning_agent.run_with_tools(planning_mission, self.planning_tools) + + logger.info(f"Starting Coding Phase for PR #{self.item.task_number} (attempt {attempt})") + coding_mission = ( + f"PHASE 2: EXECUTION/CODING PHASE\n\n" + f"You must now implement the changes based on the following plan:\n" + f"--- PLAN ---\n{plan}\n--- PLAN END ---\n\n" + f"Original Mission details:\n{base_mission}\n\n" + f"Follow the workflow to implement changes, verify, and complete PR review/updates.\n" + ) + coding_agent = CodingAgent(self.model_name) + response = await coding_agent.run_with_tools(coding_mission, self.coding_tools_list) + return response + except Exception as e: + logger.error(f"Error processing PR #{self.item.task_number} (attempt {attempt}/{attempt_limit}): {e}") + if attempt == attempt_limit: + return f"FAILED after {attempt_limit} attempts: {str(e)}" + return f"FAILED: PR #{self.item.task_number} not processed." + + +class IssueTaskProcessor(TaskProcessor): + """Processes Gitea issues (acting as coordinator, then coding).""" + + def _build_issue_mission(self, issue_info: IssueModel, branch_name: str) -> str: + issue_number = self.item.task_number + issue_body = issue_info.body or "No description provided." + issue_labels = [lbl.name for lbl in issue_info.labels] + issue_user = issue_info.user.login if issue_info.user else "unknown" + issue_created = issue_info.created_at or "unknown" + title = issue_info.title + + comments: list[CommentModel] = [] + try: + comments = self.client.get_issue_comments(self.owner, self.repo_name, issue_number) + except Exception: + pass + + labels_str = f"Labels: {', '.join(issue_labels)}" if issue_labels else "Labels: none" + comments_str = "\n".join([ + f"- @{c.user.login} ({c.created_at}): {c.body}" + for c in comments + ]) if comments else "No comments yet." + + return ( + f"Your mission is to resolve issue #{issue_number} in {self.repo}.\n\n" + f"Issue: {title}\n" + f"Author: @{issue_user} (created {issue_created})\n" + f"{labels_str}\n\n" + f"Description:\n{issue_body}\n\n" + f"Comments ({len(comments)}):\n{comments_str}\n\n" + f"Branch name: {branch_name}.\n\n" + "BEFORE WRITING ANY CODE:\n" + " - Search online for relevant documentation, known solutions, library APIs, and platform-specific behavior.\n" + " - If ANY part of the issue is unclear, ambiguous, or has multiple valid approaches:\n" + " → Post a comment on the issue using `add_comment_to_issue` with your specific question(s).\n" + " → List the approaches you are considering.\n" + " → End the comment with the marker: on its own line.\n" + " → STOP. Do NOT proceed until a human replies. The system will re-dispatch you once a human responds.\n" + " - Never assume or guess. Always prefer asking over guessing.\n\n" + "CRITICAL INSTRUCTIONS:\n" + f"1. The repo is already cloned locally at '{self.repo_path}'. DO NOT create a new repository.\n" + f" The repository directory is your current working directory (cwd). You can verify this via the `get_working_directory` tool.\n" + "2. Always start from master: `git checkout master && git pull origin master`\n" + "3. Branch from master: `git checkout -b /issue--`\n" + " Branch names MUST include a descriptive name (words/hyphens), not just the issue number.\n" + " Types: feat, fix, chore, docs, style, refactor, test, build, ci, perf\n" + "4. Use `edit_file`/`write_file` for code changes, then `git add` and `git commit` via `run_command`.\n" + "5. Push: `git push origin ` via `run_command`.\n" + "6. Create PR: Use the `create_pull_request` tool (do NOT use Gitea's `tea` CLI or GitHub's `gh` CLI in run_command, as they can hang/freeze interactively).\n" + " PR description MUST include the Gitea automation template with 'closes #'.\n" + "7. IMPORTANT: After successfully creating the pull request, you MUST comment on the issue (using the `add_comment_to_issue` tool) with the PR number, PR link, and summary.\n" + "8. Check AGENTS.md in repo root for project conventions and verification steps.\n" + "9. Grade severity: Critical/High = fix, Medium = review/fix, Low = skip.\n" + "An issue is DONE when the connected PR is merged (you cannot merge yourself).\n" + "DO NOT edit .git files unless explicitly resolving a git issue.\n" + "DO NOT work on non-meeks organization repos." + ) + + async def process(self, attempt_limit: int) -> str: + # Check if there is an existing PR for the issue + existing_pr = _find_pr_for_issue_helper(self.client, self.repo, self.item.task_number) + is_wip = False + has_request_changes = False + + if existing_pr: + if existing_pr.title: + title_upper = existing_pr.title.strip().upper() + is_wip = title_upper.startswith("WIP:") or "[WIP]" in title_upper + + try: + reviews = self.client.get_pr_reviews(self.owner, self.repo_name, existing_pr.number) + has_request_changes = any(r.get("state") == "REQUEST_CHANGES" for r in reviews) + except Exception as e: + logger.warning(f"Error checking reviews for PR #{existing_pr.number}: {e}") + + if not is_wip and not has_request_changes: + logger.info(f"Issue #{self.item.task_number} already has open PR #{existing_pr.number}. Skipping.") + return f"SKIP: A pull request (PR #{existing_pr.number}) addressing issue #{self.item.task_number} already exists." + + # Check comments on issue and PR + issue_comments = [] + try: + issue_comments = self.client.get_issue_comments(self.owner, self.repo_name, self.item.task_number) + except Exception: + pass + + pr_comments = [] + if existing_pr: + try: + pr_comments = self.client.get_pull_request_comments(self.owner, self.repo_name, existing_pr.number) + except Exception: + pass + + if _is_awaiting_reply_helper(issue_comments) or _is_awaiting_reply_helper(pr_comments): + logger.info(f"Issue #{self.item.task_number}: awaiting human reply. Skipping.") + return f"SKIP: Awaiting human reply on issue #{self.item.task_number} or PR." + + issue_info = self.item.task_info + assert isinstance(issue_info, IssueModel) + title = issue_info.title + issue_body = issue_info.body or "No description provided." + + issue_comments_str = "\n".join([ + f"- @{c.user.login} ({c.created_at}): {c.body}" for c in issue_comments + ]) if issue_comments else "No comments yet." + + pr_info_str = "No existing PR." + if existing_pr: + pr_comments_str = "\n".join([ + f"- @{c.user.login} ({c.created_at}): {c.body}" for c in pr_comments + ]) if pr_comments else "No PR comments yet." + try: + reviews = self.client.get_pr_reviews(self.owner, self.repo_name, existing_pr.number) + reviews_str = "\n".join([ + f"- @{r.get('user', {}).get('login')} ({r.get('submitted_at')}): [{r.get('state')}] {r.get('body')}" + for r in reviews + ]) if reviews else "No reviews yet." + except Exception: + reviews_str = "No reviews available." + pr_info_str = ( + f"PR Number: #{existing_pr.number}\n" + f"PR Title: {existing_pr.title}\n" + f"PR Branch: {existing_pr.head.get('ref', 'unknown')}\n" + f"PR State: {existing_pr.state}\n" + f"PR Comments:\n{pr_comments_str}\n" + f"PR Reviews:\n{reviews_str}" + ) + + state_analysis_mission = ( + f"Analyzing issue #{self.item.task_number} in '{self.repo}'.\n\n" + f"Issue Title: {title}\n" + f"Description:\n{issue_body}\n\n" + f"Issue Comments:\n{issue_comments_str}\n\n" + f"Existing PR Details:\n{pr_info_str}\n" + ) + + for attempt in range(1, attempt_limit + 1): + try: + coord_tools = CoordinatorTools() + coordinator_agent = CoordinatorAgent(self.model_name) + + logger.info(f"Analyzing conversation state for issue #{self.item.task_number}...") + await coordinator_agent.decide_action(state_analysis_mission, self.planning_tools, coord_tools) + + action = coord_tools.action + logger.info(f"Coordinator Decided Action: {action} (tool_called={coord_tools.tool_called})") + + if action == "PROPOSE_PLAN": + comment_body = coord_tools.arguments.get("comment_body", "") + if not comment_body: + plan = coord_tools.arguments.get("plan", "") + comment_body = ( + f"### Proposed Implementation Plan\n\n" + f"{plan}\n\n" + f"Is this plan ok for implementation or do you have any comments/changes?\n" + f"\n" + f"" + ) + self.client.add_comment(self.owner, self.repo_name, self.item.task_number, comment_body) + return f"POSTED_COMMENT: PROPOSE_PLAN comment posted to issue #{self.item.task_number}." + + elif action == "ANSWER_QUESTION": + comment_body = coord_tools.arguments.get("comment_body", "") + if not comment_body: + answer = coord_tools.arguments.get("answer", "") + comment_body = ( + f"{answer}\n\n" + f"Is this answer satisfactory?\n" + f"\n" + f"" + ) + self.client.add_comment(self.owner, self.repo_name, self.item.task_number, comment_body) + return f"POSTED_COMMENT: ANSWER_QUESTION comment posted to issue #{self.item.task_number}." + + elif action == "CLOSE_ISSUE": + comment = coord_tools.arguments.get("comment", "Closing the issue as resolved.") + self.client.add_comment(self.owner, self.repo_name, self.item.task_number, comment) + self.client.close_issue(self.owner, self.repo_name, self.item.task_number) + return f"CLOSED_ISSUE: Issue #{self.item.task_number} closed." + + elif action == "NO_ACTION": + return f"NO_ACTION: No action taken on issue #{self.item.task_number}." + + elif action == "EXECUTE_PLAN": + approved_plan = coord_tools.arguments.get("approved_plan", "") + pr_to_use = existing_pr + branch_name = "" + + if pr_to_use: + branch_name = pr_to_use.head.get("ref", "") + logger.info(f"Resuming work on existing PR #{pr_to_use.number} on branch '{branch_name}'") + else: + logger.info(f"Creating new WIP PR for issue #{self.item.task_number}") + clean_title = re.sub(r'[^a-zA-Z0-9\s-]', '', title).strip().lower() + title_words = clean_title.split()[:5] + desc_suffix = "-".join(title_words) + if not desc_suffix: + desc_suffix = "fix-issue" + branch_name = f"fix/issue-{self.item.task_number}-{desc_suffix}" + + try: + subprocess.run(["git", "checkout", "master"], cwd=str(self.repo_path), check=True) + subprocess.run(["git", "pull", "origin", "master"], cwd=str(self.repo_path), check=True) + subprocess.run(["git", "branch", "-D", branch_name], cwd=str(self.repo_path), stderr=subprocess.DEVNULL) + subprocess.run(["git", "checkout", "-b", branch_name], cwd=str(self.repo_path), check=True) + subprocess.run(["git", "commit", "--allow-empty", "-m", f"WIP: start implementation for issue #{self.item.task_number}"], cwd=str(self.repo_path), check=True) + subprocess.run(["git", "push", "origin", branch_name], cwd=str(self.repo_path), check=True) + + pr_title = f"WIP: {title}" + pr_description = f"Work in progress for issue #{self.item.task_number}." + pr_to_use = self.client.create_pull_request( + self.owner, self.repo_name, head=branch_name, base="master", title=pr_title, description=pr_description + ) + + pr_link = pr_to_use.html_url or f"{self.client.base_url}/{self.repo}/pulls/{pr_to_use.number}" + start_comment = f"Started work on PR #{pr_to_use.number} ({pr_link})." + self.client.add_comment(self.owner, self.repo_name, self.item.task_number, start_comment) + + logger.info(f"Successfully created WIP PR #{pr_to_use.number} on branch '{branch_name}'") + except Exception as e: + logger.error(f"Failed to create WIP PR for issue #{self.item.task_number}: {e}") + return f"FAILED to create WIP PR: {e}" + + base_mission = self._build_issue_mission(issue_info, branch_name) + coding_mission = ( + f"PHASE 2: EXECUTION/CODING PHASE\n\n" + f"You are implementing changes for issue #{self.item.task_number} in repository '{self.repo}'.\n" + f"You are working on the existing Pull Request #{pr_to_use.number} on branch '{branch_name}'.\n\n" + f"--- APPROVED PLAN ---\n{approved_plan}\n--- APPROVED PLAN END ---\n\n" + f"Original Mission details:\n{base_mission}\n\n" + f"DIRECTIONS:\n" + f"1. Checkout the branch '{branch_name}' (it should already be checked out, or run `git checkout {branch_name}`).\n" + f"2. Implement the changes according to the APPROVED PLAN.\n" + f"3. Run verification/tests (check AGENTS.md for conventions).\n" + f"4. Commit and push your changes to origin on the branch '{branch_name}'.\n" + f"5. ONCE COMPLETED SUCCESSFULLY:\n" + f" - Call `update_pull_request(owner='{self.owner}', repo='{self.repo_name}', pull_number={pr_to_use.number}, title='{title}', body='')`.\n" + f" Note: The PR title must not contain 'WIP:'. The PR body must follow the mandatory PR template in CODING AGENT SYSTEM PROMPT.\n" + f" - Call `add_comment_to_issue(owner='{self.owner}', repo='{self.repo_name}', issue_number={self.item.task_number}, body='Work has been completed in PR #{pr_to_use.number}.')`.\n" + f"6. IF YOU ENCOUNTER A BLOCKER OR FAIL:\n" + f" - You are allowed (and encouraged) to comment on the WIP PR #{pr_to_use.number} (using `add_comment`) with any details, logs, or context to help the next agent resume the work.\n\n" + f"AVAILABLE RESEARCH TOOLS:\n" + f" - web_search(query, time_range, categories) — search the web via SearXNG/DuckDuckGo\n" + f" - fetch_url(url) — read any documentation page in full\n" + ) + logger.info(f"Starting Execution/Coding Phase for issue #{self.item.task_number} on branch '{branch_name}'") + coding_agent = CodingAgent(self.model_name) + response = await coding_agent.run_with_tools(coding_mission, self.coding_tools_list) + logger.info(f"Agent response for issue #{self.item.task_number}: {response}") + return response + + except CoordinatorNoToolCalledError as e: + logger.error(f"Coordinator error on issue #{self.item.task_number} (attempt {attempt}/{attempt_limit}): {e}") + if attempt == attempt_limit: + return f"FAILED: Coordinator did not call any tools after {attempt_limit} attempts." + except Exception as e: + logger.error(f"Error processing issue #{self.item.task_number} (attempt {attempt}/{attempt_limit}): {e}") + if attempt == attempt_limit: + return f"FAILED after {attempt_limit} attempts: {str(e)}" + return f"FAILED: Issue #{self.item.task_number} not processed." + + +class AgentDispatcher: + """Dispatches work to a specialized task processor, one repo at a time.""" + + def __init__( + self, + client: GiteaClient, + tools: GiteaTools, + model_name: str = AGENT_MODEL_ID, + max_retries: int = 2, + ) -> None: + self._client = client + self._tools = tools + self._model_name = model_name + self._max_retries = max_retries + + async def dispatch( + self, + repo: str, + work_items: list[WorkItem], + ) -> list[str]: + """Dispatch all work for a single repo by handling specialized processor classes.""" + workspace = WorkspaceManager() + repo_path = workspace.get_repo_path(repo) + + original_cwd = os.getcwd() + changed_dir = False + if os.path.isdir(str(repo_path)): + os.chdir(str(repo_path)) + changed_dir = True + + results: list[str] = [] + try: + # Get authenticated username for reviewer filter + ai_username = "meeks-ai" + try: + user = self._client.get_authenticated_user() + if user: + ai_username = user.login + except Exception: + pass + + for item in work_items: + processor: TaskProcessor + if item.task_type == "pr": + processor = PRTaskProcessor( + client=self._client, + tools=self._tools, + model_name=self._model_name, + repo=repo, + item=item, + ai_username=ai_username, + ) + elif item.task_type == "issue": + processor = IssueTaskProcessor( + client=self._client, + tools=self._tools, + model_name=self._model_name, + repo=repo, + item=item, + ai_username=ai_username, + ) + else: + logger.warning(f"Unknown task type: {item.task_type}") + results.append(f"SKIP: Unknown task type {item.task_type}") + continue + + logger.info(f"Processing {item.task_type} #{item.task_number} via {processor.__class__.__name__}") + result = await processor.process(attempt_limit=self._max_retries) + results.append(result) + + finally: + if changed_dir: + os.chdir(original_cwd) + + return results + + # Backward compatibility helper methods for unit tests + def _find_pr_for_issue(self, repo_full_name: str, issue_number: int) -> PullRequestModel | None: + return _find_pr_for_issue_helper(self._client, repo_full_name, issue_number) + + def _is_awaiting_reply(self, comments: list[CommentModel]) -> bool: + return _is_awaiting_reply_helper(comments) + + def _build_pr_mission(self, item: WorkItem) -> str: + pr_info = item.task_info + assert isinstance(pr_info, PullRequestModel) + processor = PRTaskProcessor( + client=self._client, + tools=self._tools, + model_name=self._model_name, + repo=item.repo_full_name, + item=item, + ai_username="meeks-ai", + ) + return processor._build_pr_mission(pr_info, is_own_pr=False) + + def _build_issue_mission(self, item: WorkItem) -> str: + issue_info = item.task_info + assert isinstance(issue_info, IssueModel) + processor = IssueTaskProcessor( + client=self._client, + tools=self._tools, + model_name=self._model_name, + repo=item.repo_full_name, + item=item, + ai_username="meeks-ai", + ) + return processor._build_issue_mission(issue_info, "dummy-branch") diff --git a/core/factory.py b/core/factory.py index e782ce4..30f2d71 100644 --- a/core/factory.py +++ b/core/factory.py @@ -9,6 +9,8 @@ from core.interfaces import ( ) from core.coding_agent import CodingAgent from core.agent import CavemanAgent +from core.coordinator_agent import CoordinatorAgent +from core.planning_agent import PlanningAgent from gitea.workspace import WorkspaceManager logger: logging.Logger = logging.getLogger("core-factory") @@ -65,6 +67,16 @@ class AgentFactory: logger.info(f"Factory creating CavemanAgent with model: {model_name}") return CavemanAgent(model_name) + @staticmethod + def create_coordinator_agent(model_name: str) -> CoordinatorAgent: + logger.info(f"Factory creating CoordinatorAgent with model: {model_name}") + return CoordinatorAgent(model_name) + + @staticmethod + def create_planning_agent(model_name: str) -> PlanningAgent: + logger.info(f"Factory creating PlanningAgent with model: {model_name}") + return PlanningAgent(model_name) + class WorkspaceFactory: """Factory for creating workspace manager instances.""" diff --git a/core/planning_agent.py b/core/planning_agent.py new file mode 100644 index 0000000..858a510 --- /dev/null +++ b/core/planning_agent.py @@ -0,0 +1,13 @@ +import logging +from core.agent import BaseAgent +from .coding_prompt import CODING_AGENT_SYSTEM_PROMPT + +logger: logging.Logger = logging.getLogger("agent-planning") + + +class PlanningAgent(BaseAgent): + """AI agent that analyzes a PR/issue and builds an implementation plan.""" + + def __init__(self, model_name: str) -> None: + super().__init__(model_name) + self.system_prompt = CODING_AGENT_SYSTEM_PROMPT diff --git a/tests/test_best_practices.py b/tests/test_best_practices.py index 071d536..91a4142 100644 --- a/tests/test_best_practices.py +++ b/tests/test_best_practices.py @@ -94,7 +94,8 @@ def test_run_verification_failure(tmp_path: Path) -> None: @patch("core.dispatcher.CodingAgent") -async def test_dispatch_planning_and_coding_phases(mock_agent_class: MagicMock) -> None: +@patch("core.dispatcher.PlanningAgent") +async def test_dispatch_planning_and_coding_phases(mock_planning_class: MagicMock, mock_coding_class: MagicMock) -> None: mock_client: MagicMock = MagicMock(spec=GiteaClient) mock_tools: MagicMock = MagicMock(spec=GiteaTools) @@ -117,14 +118,14 @@ async def test_dispatch_planning_and_coding_phases(mock_agent_class: MagicMock) mock_client.get_pull_request_files.return_value = [] mock_client.get_pr_reviews.return_value = [] - # Mock CodingAgent instances + # Mock agent instances mock_planning_agent = MagicMock() mock_planning_agent.run_with_tools = AsyncMock(return_value="Plan: Modify file A") + mock_planning_class.return_value = mock_planning_agent mock_coding_agent = MagicMock() mock_coding_agent.run_with_tools = AsyncMock(return_value="PR #1 Created") - - mock_agent_class.side_effect = [mock_planning_agent, mock_coding_agent] + mock_coding_class.return_value = mock_coding_agent dispatcher = AgentDispatcher(client=mock_client, tools=mock_tools) diff --git a/tests/test_dispatcher.py b/tests/test_dispatcher.py index 206a00f..65f065d 100644 --- a/tests/test_dispatcher.py +++ b/tests/test_dispatcher.py @@ -38,8 +38,8 @@ async def test_dispatch_skips_issue_with_existing_pr() -> None: mock_client.list_repo_pull_requests.assert_called_once_with("meeks", "repo1") -@patch("core.dispatcher.CodingAgent") -async def test_dispatch_processes_issue_without_pr(mock_agent_class: MagicMock) -> None: +@patch("core.dispatcher.CoordinatorAgent") +async def test_dispatch_processes_issue_without_pr(mock_coord_class: MagicMock) -> None: mock_client: MagicMock = MagicMock(spec=GiteaClient) mock_tools: MagicMock = MagicMock(spec=GiteaTools) @@ -52,17 +52,14 @@ async def test_dispatch_processes_issue_without_pr(mock_agent_class: MagicMock) mock_client.list_repo_pull_requests.return_value = [pr] mock_client.get_issue_comments.return_value = [] - # Mock CodingAgent run_with_tools - mock_agent_instance = MagicMock() - mock_agent_instance.run_with_tools = AsyncMock(return_value="""```json -{ - "action": "PROPOSE_PLAN", - "reasoning": "Plan needs to be proposed first.", - "comment_body": "### Proposed Plan\\n- change X", - "approved_plan": "" -} -```""") - mock_agent_class.return_value = mock_agent_instance + # Mock CoordinatorAgent invoking propose_plan tool + async def mock_decide(mission: str, planning_tools: list, coord_tools) -> str: + coord_tools.propose_plan(plan="- change X", issue_number=42) + return "Agent proposed plan." + + mock_coord_instance = MagicMock() + mock_coord_instance.decide_action = AsyncMock(side_effect=mock_decide) + mock_coord_class.return_value = mock_coord_instance dispatcher = AgentDispatcher(client=mock_client, tools=mock_tools) @@ -260,8 +257,8 @@ def test_is_awaiting_reply_no_marker_not_detected() -> None: assert dispatcher._is_awaiting_reply(comments) is False -@patch("core.dispatcher.CodingAgent") -async def test_dispatch_proposes_plan(mock_agent_class: MagicMock) -> None: +@patch("core.dispatcher.CoordinatorAgent") +async def test_dispatch_proposes_plan(mock_coord_class: MagicMock) -> None: mock_client: MagicMock = MagicMock(spec=GiteaClient) mock_tools: MagicMock = MagicMock(spec=GiteaTools) @@ -269,15 +266,12 @@ async def test_dispatch_proposes_plan(mock_agent_class: MagicMock) -> None: mock_client.get_issue_comments.return_value = [] mock_client.get_authenticated_user.return_value = UserModel(login="meeks-ai") - mock_agent_instance = MagicMock() - mock_agent_instance.run_with_tools = AsyncMock(return_value="""```json -{ - "action": "PROPOSE_PLAN", - "reasoning": "We need to add a new endpoint.", - "comment_body": "### Proposed Plan\\n- Add endpoint\\n\\n" -} -```""") - mock_agent_class.return_value = mock_agent_instance + async def mock_decide(mission: str, planning_tools: list, coord_tools) -> str: + coord_tools.propose_plan(plan="- Add endpoint\n\n", issue_number=42) + return "Agent proposed plan." + mock_coord_instance = MagicMock() + mock_coord_instance.decide_action = AsyncMock(side_effect=mock_decide) + mock_coord_class.return_value = mock_coord_instance dispatcher = AgentDispatcher(client=mock_client, tools=mock_tools) work_item = WorkItem( @@ -291,11 +285,11 @@ async def test_dispatch_proposes_plan(mock_agent_class: MagicMock) -> None: results = await dispatcher.dispatch("meeks/repo1", [work_item]) assert len(results) == 1 assert "POSTED_COMMENT: PROPOSE_PLAN" in results[0] - mock_client.add_comment.assert_called_once_with("meeks", "repo1", 42, "### Proposed Plan\n- Add endpoint\n\n") + mock_client.add_comment.assert_called_once_with("meeks", "repo1", 42, "### Proposed Implementation Plan\n\n- Add endpoint\n\n\n\nIs this plan ok for implementation or do you have any comments/changes?\n\n") -@patch("core.dispatcher.CodingAgent") -async def test_dispatch_answers_question(mock_agent_class: MagicMock) -> None: +@patch("core.dispatcher.CoordinatorAgent") +async def test_dispatch_answers_question(mock_coord_class: MagicMock) -> None: mock_client: MagicMock = MagicMock(spec=GiteaClient) mock_tools: MagicMock = MagicMock(spec=GiteaTools) @@ -303,15 +297,12 @@ async def test_dispatch_answers_question(mock_agent_class: MagicMock) -> None: mock_client.get_issue_comments.return_value = [] mock_client.get_authenticated_user.return_value = UserModel(login="meeks-ai") - mock_agent_instance = MagicMock() - mock_agent_instance.run_with_tools = AsyncMock(return_value="""```json -{ - "action": "ANSWER_QUESTION", - "reasoning": "This is a question about how X works.", - "comment_body": "X works by doing Y.\\n\\n" -} -```""") - mock_agent_class.return_value = mock_agent_instance + async def mock_decide(mission: str, planning_tools: list, coord_tools) -> str: + coord_tools.answer_question(answer="X works by doing Y.\n\n", issue_number=42) + return "Agent answered question." + mock_coord_instance = MagicMock() + mock_coord_instance.decide_action = AsyncMock(side_effect=mock_decide) + mock_coord_class.return_value = mock_coord_instance dispatcher = AgentDispatcher(client=mock_client, tools=mock_tools) work_item = WorkItem( @@ -325,11 +316,11 @@ async def test_dispatch_answers_question(mock_agent_class: MagicMock) -> None: results = await dispatcher.dispatch("meeks/repo1", [work_item]) assert len(results) == 1 assert "POSTED_COMMENT: ANSWER_QUESTION" in results[0] - mock_client.add_comment.assert_called_once_with("meeks", "repo1", 42, "X works by doing Y.\n\n") + mock_client.add_comment.assert_called_once_with("meeks", "repo1", 42, "X works by doing Y.\n\n\n\nIs this answer satisfactory?\n\n") -@patch("core.dispatcher.CodingAgent") -async def test_dispatch_closes_issue_on_satisfaction(mock_agent_class: MagicMock) -> None: +@patch("core.dispatcher.CoordinatorAgent") +async def test_dispatch_closes_issue_on_satisfaction(mock_coord_class: MagicMock) -> None: mock_client: MagicMock = MagicMock(spec=GiteaClient) mock_tools: MagicMock = MagicMock(spec=GiteaTools) @@ -342,15 +333,12 @@ async def test_dispatch_closes_issue_on_satisfaction(mock_agent_class: MagicMock _make_comment("michael", "Yes, thanks! That makes sense.") ] - mock_agent_instance = MagicMock() - mock_agent_instance.run_with_tools = AsyncMock(return_value="""```json -{ - "action": "CLOSE_ISSUE", - "reasoning": "User is satisfied.", - "comment_body": "Closing the issue now. Let me know if you need anything else!" -} -```""") - mock_agent_class.return_value = mock_agent_instance + async def mock_decide(mission: str, planning_tools: list, coord_tools) -> str: + coord_tools.close_issue(comment="Closing the issue now. Let me know if you need anything else!", issue_number=42) + return "Agent closed issue." + mock_coord_instance = MagicMock() + mock_coord_instance.decide_action = AsyncMock(side_effect=mock_decide) + mock_coord_class.return_value = mock_coord_instance dispatcher = AgentDispatcher(client=mock_client, tools=mock_tools) work_item = WorkItem( @@ -370,7 +358,8 @@ async def test_dispatch_closes_issue_on_satisfaction(mock_agent_class: MagicMock @patch("subprocess.run") @patch("core.dispatcher.CodingAgent") -async def test_dispatch_executes_approved_plan_and_creates_wip_pr(mock_agent_class: MagicMock, mock_run: MagicMock) -> None: +@patch("core.dispatcher.CoordinatorAgent") +async def test_dispatch_executes_approved_plan_and_creates_wip_pr(mock_coord_class: MagicMock, mock_coding_class: MagicMock, mock_run: MagicMock) -> None: mock_client: MagicMock = MagicMock(spec=GiteaClient) mock_tools: MagicMock = MagicMock(spec=GiteaTools) @@ -386,20 +375,13 @@ async def test_dispatch_executes_approved_plan_and_creates_wip_pr(mock_agent_cla mock_client.create_pull_request.return_value = mock_pr # Mock planning agent deciding EXECUTE_PLAN - mock_agent_instance1 = MagicMock() - mock_agent_instance1.run_with_tools = AsyncMock(return_value="""```json -{ - "action": "EXECUTE_PLAN", - "reasoning": "Plan was approved.", - "approved_plan": "Step 1. Code X" -} -```""") + async def mock_decide(mission: str, planning_tools: list, coord_tools) -> str: + coord_tools.start_implementation(approved_plan="Step 1. Code X", issue_number=42) + return "Agent decided execute plan." + mock_coord_class.return_value.decide_action = AsyncMock(side_effect=mock_decide) # Mock coding agent executing plan - mock_agent_instance2 = MagicMock() - mock_agent_instance2.run_with_tools = AsyncMock(return_value="PR Completed Successfully.") - - mock_agent_class.side_effect = [mock_agent_instance1, mock_agent_instance2] + mock_coding_class.return_value.run_with_tools = AsyncMock(return_value="PR Completed Successfully.") dispatcher = AgentDispatcher(client=mock_client, tools=mock_tools) work_item = WorkItem( @@ -427,7 +409,8 @@ async def test_dispatch_executes_approved_plan_and_creates_wip_pr(mock_agent_cla @patch("subprocess.run") @patch("core.dispatcher.CodingAgent") -async def test_dispatch_resumes_wip_pr(mock_agent_class: MagicMock, mock_run: MagicMock) -> None: +@patch("core.dispatcher.CoordinatorAgent") +async def test_dispatch_resumes_wip_pr(mock_coord_class: MagicMock, mock_coding_class: MagicMock, mock_run: MagicMock) -> None: mock_client: MagicMock = MagicMock(spec=GiteaClient) mock_tools: MagicMock = MagicMock(spec=GiteaTools) @@ -450,20 +433,13 @@ async def test_dispatch_resumes_wip_pr(mock_agent_class: MagicMock, mock_run: Ma mock_client.get_pull_request_comments.return_value = [] # Mock planning agent deciding EXECUTE_PLAN - mock_agent_instance1 = MagicMock() - mock_agent_instance1.run_with_tools = AsyncMock(return_value="""```json -{ - "action": "EXECUTE_PLAN", - "reasoning": "WIP PR exists, resume coding.", - "approved_plan": "Step 1. Resume coding" -} -```""") + async def mock_decide(mission: str, planning_tools: list, coord_tools) -> str: + coord_tools.start_implementation(approved_plan="Step 1. Resume coding", issue_number=42) + return "Agent decided execute plan." + mock_coord_class.return_value.decide_action = AsyncMock(side_effect=mock_decide) # Mock coding agent executing plan - mock_agent_instance2 = MagicMock() - mock_agent_instance2.run_with_tools = AsyncMock(return_value="PR Updated Successfully.") - - mock_agent_class.side_effect = [mock_agent_instance1, mock_agent_instance2] + mock_coding_class.return_value.run_with_tools = AsyncMock(return_value="PR Updated Successfully.") dispatcher = AgentDispatcher(client=mock_client, tools=mock_tools) work_item = WorkItem( @@ -528,8 +504,8 @@ def test_coordinator_tools_registration() -> None: assert tools.arguments == {"approved_plan": "my approved plan", "issue_number": 42} -@patch("core.dispatcher.CodingAgent") -async def test_dispatch_uses_coordinator_tool_calling(mock_agent_class: MagicMock) -> None: +@patch("core.dispatcher.CoordinatorAgent") +async def test_dispatch_uses_coordinator_tool_calling(mock_coord_class: MagicMock) -> None: mock_client: MagicMock = MagicMock(spec=GiteaClient) mock_tools: MagicMock = MagicMock(spec=GiteaTools) @@ -538,15 +514,13 @@ async def test_dispatch_uses_coordinator_tool_calling(mock_agent_class: MagicMoc mock_client.get_authenticated_user.return_value = UserModel(login="meeks-ai") # Mock agent invoking propose_plan tool - async def mock_run_tools(mission: str, tools: list[any]) -> str: - for t in tools: - if getattr(t, "__name__", "") == "propose_plan": - t(plan="Step 1. Code X", issue_number=42) + async def mock_decide(mission: str, planning_tools: list, coord_tools) -> str: + coord_tools.propose_plan(plan="Step 1. Code X", issue_number=42) return "Agent finished turn after tool calling." - mock_agent_instance = MagicMock() - mock_agent_instance.run_with_tools = AsyncMock(side_effect=mock_run_tools) - mock_agent_class.return_value = mock_agent_instance + mock_coord_instance = MagicMock() + mock_coord_instance.decide_action = AsyncMock(side_effect=mock_decide) + mock_coord_class.return_value = mock_coord_instance dispatcher = AgentDispatcher(client=mock_client, tools=mock_tools) work_item = WorkItem(