"""Top-level coordinator: polls Gitea, queues work, dispatches to agent.""" import asyncio import logging from pathlib import Path from typing import Any from core.queue import WorkQueue, WorkItem from core.dispatcher import AgentDispatcher from gitea.client import GiteaClient from gitea.models import IssueModel, PullRequestModel from gitea.tools.gitea_tools import GiteaTools from gitea.config import AGENT_MODEL_ID, AGENT_MAX_RETRIES from gitea.workspace import WorkspaceManager logger: logging.Logger = logging.getLogger("agent-orchestrator") class AgentOrchestrator: """Top-level coordinator: polls Gitea, queues work, dispatches to agent.""" def __init__( self, client: GiteaClient, tools: GiteaTools, model_name: str = AGENT_MODEL_ID, max_retries: int = AGENT_MAX_RETRIES, ) -> None: self._client = client self._tools = tools self._model_name = model_name self._work_queue = WorkQueue() self._dispatcher = AgentDispatcher(client, tools, model_name, max_retries) async def poll_and_dispatch(self) -> None: """Poll Gitea for tasks, enqueue them, and dispatch to agent.""" issues: list[IssueModel] = self._client.list_assigned_issues() prs: list[PullRequestModel] = self._client.list_assigned_pull_requests() if issues: logger.info(f"Found {len(issues)} assigned issues") self._enqueue_tasks("issue", issues) else: logger.info("No assigned issues found.") if prs: logger.info(f"Found {len(prs)} assigned PRs") self._enqueue_tasks("pr", prs) else: logger.info("No assigned PRs found.") if not self._work_queue.is_empty: await self._process_work() def _enqueue_tasks(self, task_type: str, tasks: list[IssueModel] | list[PullRequestModel]) -> None: for task in tasks: repo_full_name: str | None = task.repository.full_name if task.repository else None task_number: int = task.number if not repo_full_name or not task_number: continue item = WorkItem( repo_full_name=repo_full_name, task_type=task_type, task_number=task_number, task_info=task, priority=0, ) self._work_queue.enqueue(item) logger.info(f"Enqueued {task_type} #{task_number} from {repo_full_name}") async def _process_work(self) -> None: """Process all queued work, repo by repo.""" while not self._work_queue.is_empty: repo: str | None = self._work_queue.get_next_repo() if not repo: break work_items: list[WorkItem] = self._work_queue.get_repo_work(repo) self._work_queue.remove_repo_work(repo) # Ensure the workspace repository is cloned and sanitized workspace = WorkspaceManager() repo_path = workspace.get_repo_path(repo) if not repo_path.exists(): workspace.clone_repo(repo) logger.info(f"Cloned {repo} to {repo_path}") else: workspace.sanitize_repo(repo, repo_path) logger.info(f"Sanitized existing repo at {repo_path}") logger.info(f"Dispatching {len(work_items)} tasks for {repo}") results: list[str] = await self._dispatcher.dispatch(repo, work_items) for i, result in enumerate(results): item = work_items[i] logger.info(f"Completed {item.task_type} #{item.task_number}: {result[:200]}")