754 lines
38 KiB
Python
754 lines
38 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.gitea_tools import GiteaTools
|
|
from gitea.client import GiteaClient
|
|
from core.coordinator_tools import CoordinatorTools
|
|
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.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]) -> 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 "<!-- 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,
|
|
tools: GiteaTools,
|
|
model_name: str,
|
|
repo: str,
|
|
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
|
|
|
|
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,
|
|
]
|
|
|
|
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,
|
|
]
|
|
|
|
@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:
|
|
pass
|
|
|
|
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:
|
|
pass
|
|
|
|
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:
|
|
pass
|
|
|
|
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:
|
|
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: <!-- 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:
|
|
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
|
|
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:
|
|
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,
|
|
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)
|
|
|
|
results: list[str] = []
|
|
# 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)
|
|
|
|
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
|
|
if not isinstance(pr_info, PullRequestModel):
|
|
raise TypeError("Expected task_info to be a 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
|
|
if not isinstance(issue_info, IssueModel):
|
|
raise TypeError("Expected task_info to be an 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")
|