Files
coding-agent-gitea/core/dispatcher.py
T
meeks 3c94c3cfac refactor: remove duplicate add_comment and add_label methods from IssueTools
- Removed duplicate dd_comment method (kept dd_comment_to_issue)
- Removed duplicate dd_label method (kept dd_label_to_issue)
- Updated dispatcher.py to remove duplicate tool registrations
- Updated coding_prompt.py to reference only dd_comment_to_issue
- Removed corresponding duplicate tests from test_issue_tools.py

This addresses section 10.1 (Confusing Naming) in bad_code.md
2026-07-16 14:33:58 +02:00

1020 lines
45 KiB
Python

"""Dispatches work to a specialized task processor, one repo at a time."""
import logging
import re
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 gitea.tools.coding_tools import CodingTools
from gitea.tools.research_tools import ResearchTools
from gitea.tools.issue_tools import IssueTools
from gitea.tools.pr_tools import PRTools
from gitea.tools.file_tools import FileTools
from gitea.tools.git_tools import GitTools
from gitea.client import GiteaClient
from core.coordinator_tools import CoordinatorTools
from gitea.workspace import WorkspaceManager
from gitea.config import AGENT_MODEL_ID, AGENT_USERNAMES
from gitea.models import (
CommentModel,
PullRequestFileModel,
PullRequestModel,
IssueModel,
)
logger: logging.Logger = logging.getLogger("agent-dispatcher")
CLOSE_KEYWORDS_PATTERN: re.Pattern[str] = re.compile(
rf"\b(?:close|closes|closed|fix|fixes|fixed|resolve|resolves|resolved)\s+#(\d+)\b",
re.IGNORECASE,
)
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"(?<!\d){issue_number}(?!\d)", ref):
return pr
body = pr.body or ""
title = pr.title or ""
matches = CLOSE_KEYWORDS_PATTERN.findall(
body
) + CLOSE_KEYWORDS_PATTERN.findall(title)
if any(int(m) == issue_number for m in matches):
return pr
issue_ref_pattern = re.compile(rf"(?<!\w)#{issue_number}\b")
if issue_ref_pattern.search(title) or issue_ref_pattern.search(body):
return pr
except Exception as e:
logger.warning(
f"Error checking PRs for issue #{issue_number} in {repo_full_name}: {e}"
)
return None
def _find_issues_for_pr_helper(pr_body: str) -> 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], ai_username: str) -> 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
agent_usernames = {ai_username, *AGENT_USERNAMES}
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 "<!-- agent:awaiting-reply -->" 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,
issue_tools: IssueTools,
pr_tools: PRTools,
file_tools: FileTools,
git_tools: GitTools,
model_name: str,
repo: str,
item: WorkItem,
ai_username: str,
) -> None:
self.client = client
self.issue_tools = issue_tools
self.pr_tools = pr_tools
self.file_tools = file_tools
self.git_tools = git_tools
self.model_name = model_name
self.repo = repo
self.item = item
self.ai_username = ai_username
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.issue_tools.get_issue,
self.pr_tools.get_pull_request,
self.issue_tools.list_issues,
self.pr_tools.list_pull_requests,
self.file_tools.get_file_content,
self.issue_tools.get_issue_comments,
self.pr_tools.get_pull_request_comments,
self.pr_tools.get_pull_request_diff,
self.pr_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,
]
self.coding_tools_list: list[Callable[..., Any]] = [
self.issue_tools.get_issue,
self.pr_tools.get_pull_request,
self.issue_tools.list_issues,
self.pr_tools.list_pull_requests,
self.file_tools.get_file_content,
self.pr_tools.create_pull_request,
self.pr_tools.update_pull_request,
self.issue_tools.add_label_to_issue,
self.pr_tools.add_label_to_pr,
self.git_tools.create_branch,
self.file_tools.commit_file,
self.issue_tools.create_issue,
self.issue_tools.add_comment_to_issue,
self.issue_tools.close_issue,
self.pr_tools.close_pull_request,
self.issue_tools.get_issue_comments,
self.pr_tools.get_pull_request_comments,
self.file_tools.update_file,
self.pr_tools.get_pull_request_diff,
self.pr_tools.get_pull_request_patch,
self.pr_tools.approve_pull_request,
self.pr_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,
]
@abstractmethod
async def process(self, attempt_limit: int) -> str:
"""Execute the task flow, including planning, coding, or coordination."""
pass
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:
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(
self.owner, self.repo_name, pr_number
)
except Exception as e:
logger.warning(
f"Error fetching files for PR #{pr_number}: {e}", exc_info=True
)
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(
self.owner, self.repo_name, pr_number
)
if not isinstance(comments, list):
comments = []
except Exception as e:
logger.warning(
f"Error fetching comments for PR #{pr_number}: {e}", exc_info=True
)
reviews: list[dict[str, Any]] = []
try:
reviews = self.client.get_pr_reviews(self.owner, self.repo_name, pr_number)
if not isinstance(reviews, list):
reviews = []
except Exception as e:
logger.warning(
f"Error fetching reviews for PR #{pr_number}: {e}", exc_info=True
)
timeline: list[dict[str, Any]] = []
for c in comments:
timeline.append(
{
"timestamp": c.created_at or "",
"user": c.user.login,
"type": "comment",
"body": c.body,
"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")
r_body = r.get("body", "")
r_state = r.get("state", "")
timeline.append(
{
"timestamp": r.get("submitted_at") or r.get("updated_at") or "",
"user": r_user,
"type": "review",
"body": f"[{r_state}] {r_body}",
"by_ai": r_user == self.ai_username,
}
)
timeline.sort(key=lambda x: x["timestamp"])
last_action_by_ai = False
if timeline:
last_action_by_ai = timeline[-1]["by_ai"]
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 = (
"\n".join(
[f"- @{c.user.login} ({c.created_at}): {c.body}" for c in comments]
)
if comments
else "No comments yet."
)
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 = _find_issues_for_pr_helper(pr_body)
if linked_issues:
issues_details = []
for issue_num in linked_issues:
try:
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
]
)
if issue_comments
else " No comments yet."
)
issues_details.append(
f"### Connected Issue #{issue_num}: {issue.title}\n"
f"Author: @{issue.user.login} (created {issue.created_at})\n"
f"Description:\n{issue.body or 'No description'}\n"
f"Discussion:\n{comments_list}"
)
except Exception as e:
logger.warning(f"Could not fetch connected issue #{issue_num}: {e}")
if issues_details:
connected_issues_ctx = (
"\n---\n\n## 📋 CONNECTED ISSUE CONTEXT\n"
+ "\n\n".join(issues_details)
)
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"
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 '{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"
f" 2. Implement the requested fixes or changes on this branch.\n"
f" 3. Verify your fixes and run verification/tests.\n"
f' 4. Commit and push the changes directly: `git add <files> && git commit -m "fix: address feedback" && git push origin {pr_head_branch}`\n'
f" 5. After pushing, comment on the PR (using the `add_comment` tool) with a summary of the fixes implemented."
)
else:
instructions = (
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"
"3. Check for: code quality, security issues, edge cases, test coverage.\n"
"4. If the PR is good: approve it (using approve_pull_request tool) with a meaningful comment.\n"
"5. If the PR has issues: request changes (using request_changes tool) with specific feedback.\n"
"6. Post your review comment on the PR (using add_comment tool).\n"
"IMPORTANT: Never merge the PR yourself - that is handled by humans."
)
return (
f"Your mission is to process PR #{pr_number} in {self.repo}.\n\n"
f"PR: {pr_info.title}\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"
f"Files Changed ({len(pr_files)}):\n{files_summary}\n\n"
f"Timeline Comments:\n{comments_str}\n\n"
f"Reviews:\n{reviews_str}\n\n"
f"{connected_issues_ctx}\n\n"
"BEFORE MAKING ANY CHANGES:\n"
" - Search online for any technology, API, or behavior you are not 100% certain about.\n"
" - Read ALL review comments and change requests carefully.\n"
" - If any review comment is ambiguous or unclear:\n"
" → Post a clarifying comment on the PR (using the `add_comment` tool) with your specific question(s).\n"
" → End the comment with the marker: <!-- agent:awaiting-reply --> on its own line.\n"
" → STOP. Do NOT implement anything until a human replies. The system will re-dispatch you once a human responds.\n"
" - Never assume or guess what a reviewer meant. Always prefer asking over guessing.\n\n"
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 as e:
logger.warning(
f"Error fetching comments for PR #{self.item.task_number}: {e}",
exc_info=True,
)
if _is_awaiting_reply_helper(pr_comments, self.ai_username):
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 as e:
logger.warning(
f"Error fetching comments for issue #{issue_number}: {e}", exc_info=True
)
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: <!-- agent:awaiting-reply --> 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 <type>/issue-<number>-<descriptive-name>`\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 <branch>` 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 #<ISSUE>'.\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 as e:
logger.warning(
f"Error fetching comments for issue #{self.item.task_number}: {e}",
exc_info=True,
)
pr_comments = []
if existing_pr:
try:
pr_comments = self.client.get_pull_request_comments(
self.owner, self.repo_name, existing_pr.number
)
except Exception as e:
logger.warning(
f"Error fetching comments for PR #{existing_pr.number}: {e}",
exc_info=True,
)
if _is_awaiting_reply_helper(
issue_comments, self.ai_username
) or _is_awaiting_reply_helper(pr_comments, self.ai_username):
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
if not isinstance(issue_info, IssueModel):
raise TypeError("Expected task_info to be an 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 as e:
logger.warning(
f"Error fetching reviews for PR #{existing_pr.number}: {e}",
exc_info=True,
)
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"<!-- agent:plan-proposal -->\n"
f"<!-- agent:awaiting-reply -->"
)
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"<!-- agent:question-response -->\n"
f"<!-- agent:awaiting-reply -->"
)
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='<filled PR template>')`.\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,
issue_tools: IssueTools,
pr_tools: PRTools,
file_tools: FileTools,
git_tools: GitTools,
model_name: str = AGENT_MODEL_ID,
max_retries: int = 2,
) -> None:
self._client = client
self._issue_tools = issue_tools
self._pr_tools = pr_tools
self._file_tools = file_tools
self._git_tools = git_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)
results: list[str] = []
# Get authenticated username for reviewer filter
try:
user = self._client.get_authenticated_user()
except Exception as e:
raise RuntimeError("No authenticated user found.") from e
if not user or not user.login:
raise RuntimeError("No authenticated user found.")
ai_username = user.login
for item in work_items:
processor: TaskProcessor
if item.task_type == "pr":
processor = PRTaskProcessor(
client=self._client,
issue_tools=self._issue_tools,
pr_tools=self._pr_tools,
file_tools=self._file_tools,
git_tools=self._git_tools,
model_name=self._model_name,
repo=repo,
item=item,
ai_username=ai_username,
)
elif item.task_type == "issue":
processor = IssueTaskProcessor(
client=self._client,
issue_tools=self._issue_tools,
pr_tools=self._pr_tools,
file_tools=self._file_tools,
git_tools=self._git_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)
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:
user = self._client.get_authenticated_user()
if not user or not user.login:
raise RuntimeError("No authenticated user found.")
return _is_awaiting_reply_helper(comments, user.login)
def _build_pr_mission(self, item: WorkItem) -> str:
pr_info = item.task_info
if not isinstance(pr_info, PullRequestModel):
raise TypeError("Expected task_info to be a PullRequestModel")
user = self._client.get_authenticated_user()
if not user or not user.login:
raise RuntimeError("No authenticated user found.")
processor = PRTaskProcessor(
client=self._client,
issue_tools=self._issue_tools,
pr_tools=self._pr_tools,
file_tools=self._file_tools,
git_tools=self._git_tools,
model_name=self._model_name,
repo=item.repo_full_name,
item=item,
ai_username=user.login,
)
return processor._build_pr_mission(pr_info, is_own_pr=False)
def _build_issue_mission(self, item: WorkItem) -> str:
issue_info = item.task_info
if not isinstance(issue_info, IssueModel):
raise TypeError("Expected task_info to be an IssueModel")
user = self._client.get_authenticated_user()
if not user or not user.login:
raise RuntimeError("No authenticated user found.")
processor = IssueTaskProcessor(
client=self._client,
issue_tools=self._issue_tools,
pr_tools=self._pr_tools,
file_tools=self._file_tools,
git_tools=self._git_tools,
model_name=self._model_name,
repo=item.repo_full_name,
item=item,
ai_username=user.login,
)
return processor._build_issue_mission(issue_info, "dummy-branch")