diff --git a/roboco/enforcement/task_lifecycle.py b/roboco/enforcement/task_lifecycle.py index dceefff7..b1187935 100644 --- a/roboco/enforcement/task_lifecycle.py +++ b/roboco/enforcement/task_lifecycle.py @@ -321,6 +321,10 @@ class GitContext: # `_is_coordination_task` in the orchestrator and the `_ensure_branch_for_task` # short-circuit in TaskService. is_coordination: bool = False + # An inbound external-PR review task reviews someone else's PR read-only; it + # does no git work of its own and never gets a branch, so it is exempt from + # the claimed->in_progress branch gate (same rationale as is_coordination). + is_external_review: bool = False def validate_git_requirements( @@ -387,6 +391,7 @@ def validate_git_requirements( transition == ("claimed", "in_progress") and not git_ctx.branch_name and not git_ctx.is_coordination + and not git_ctx.is_external_review ): raise GitRequirementError( transition=transition, diff --git a/roboco/services/gateway/choreographer/__init__.py b/roboco/services/gateway/choreographer/__init__.py index fcf0bca9..8a41b4c8 100644 --- a/roboco/services/gateway/choreographer/__init__.py +++ b/roboco/services/gateway/choreographer/__init__.py @@ -22,10 +22,13 @@ from roboco.services.gateway.choreographer._impl import ( ) from roboco.services.gateway.choreographer.board import BoardMixin from roboco.services.gateway.choreographer.doc import DocMixin +from roboco.services.gateway.choreographer.pr_review import PRReviewerMixin from roboco.services.gateway.choreographer.qa import QAMixin -class Choreographer(BoardMixin, DocMixin, QAMixin, _LegacyChoreographer): +class Choreographer( + BoardMixin, DocMixin, QAMixin, PRReviewerMixin, _LegacyChoreographer +): """Composed choreographer. MRO walks left-to-right: extracted mixins resolve first, then the diff --git a/roboco/services/gateway/choreographer/_verb_runner.py b/roboco/services/gateway/choreographer/_verb_runner.py index 83bf4134..bb0707e7 100644 --- a/roboco/services/gateway/choreographer/_verb_runner.py +++ b/roboco/services/gateway/choreographer/_verb_runner.py @@ -152,6 +152,11 @@ class VerbRunner: "create_subtask requires DelegateInputs; verb body owns dispatch" ) + async def _do_pr_review_done(self, task: Any, agent: Any, ctx: spec.Context) -> Any: + return await self.task_service.complete_review( + agent.id, task.id, ctx.notes or "" + ) + @classmethod def _atomic_handlers(cls) -> dict[str, _AtomicHandler]: return { @@ -170,6 +175,7 @@ class VerbRunner: "unblock": cls._do_unblock, "resume": cls._do_resume, "create_subtask": cls._do_create_subtask, + "pr_review_done": cls._do_pr_review_done, } # -- Side-effect handlers --------------------------------------------- diff --git a/roboco/services/gateway/choreographer/pr_review.py b/roboco/services/gateway/choreographer/pr_review.py new file mode 100644 index 00000000..037ecc26 --- /dev/null +++ b/roboco/services/gateway/choreographer/pr_review.py @@ -0,0 +1,315 @@ +"""PR-reviewer verbs (inbound external/fork PR review). + +Mixin for ``claim_pr_review`` and ``post_pr_review`` — distinct from QA's +surface (per the locked design). The reviewer reviews PRs the org did NOT +author. The review is **read-only**: the diff is fetched from the GitHub API +(``git.get_pr_diff``), never checked out, and the contributor's code is never +run here. ``post_pr_review`` posts exactly one change-request to the PR and +finishes the review task. + +Inherits ``ChoreographerHelpers`` under ``TYPE_CHECKING`` only so mypy resolves +``self.task`` etc.; at runtime the composed ``Choreographer`` supplies the real +attributes via MRO (same pattern as ``QAMixin``). +""" + +from __future__ import annotations + +from types import SimpleNamespace +from typing import TYPE_CHECKING, Any + +import structlog + +from roboco.foundation.policy import lifecycle as spec_module +from roboco.foundation.policy import tracing as _tr +from roboco.services.gateway.envelope import Envelope + +if TYPE_CHECKING: + from uuid import UUID + + from roboco.services.gateway.choreographer._protocol import ChoreographerHelpers + + _Base = ChoreographerHelpers +else: + _Base = object + +logger = structlog.get_logger() + + +class PRReviewerMixin(_Base): + """PR-reviewer-role verbs (inbound external/fork PRs).""" + + async def claim_pr_review(self, reviewer_agent_id: UUID, task_id: UUID) -> Envelope: + """Reviewer claims an external-PR review task and starts work. + + Spec gate enforces role (pr_reviewer) + the composed ``claim`` action's + source-status (PENDING). The composed ``claim``+``start`` runs through + the verb runner (pending -> claimed -> in_progress; the review task is + branch-gate exempt). The response carries the contributor's unified diff + INLINE — fetched read-only via the GitHub API; the fork code is never + checked out or run — so the reviewer inspects it before posting. + """ + t = await self.task.get(task_id) + if t is None: + return await self._emit_rejection( + Envelope.not_found(message=f"task {task_id} not found"), + agent_id=reviewer_agent_id, + task_id=task_id, + verb="claim_pr_review", + ) + agent = await self.task.agent_for(reviewer_agent_id) + role_str = str(agent.role) if agent is not None else "pr_reviewer" + briefing = await self._briefing_for(reviewer_agent_id, task_id) + role_or_rejection = await self._resolve_role( + t, role_str, briefing, reviewer_agent_id, task_id, "claim_pr_review" + ) + if isinstance(role_or_rejection, Envelope): + return role_or_rejection + role = role_or_rejection + spec_ctx = spec_module.Context( + actor_id=reviewer_agent_id, + actor_slug=getattr(agent, "slug", None) if agent is not None else None, + ) + decision = spec_module.can_invoke_intent(role, "claim_pr_review", t, spec_ctx) + if not decision.allowed: + return await self._emit_rejection( + Envelope.from_decision(decision, briefing=briefing).with_introspection( + task=t, role=role_str + ), + agent_id=reviewer_agent_id, + task_id=task_id, + verb="claim_pr_review", + ) + guard = await self._run_claim_guards(agent_id=reviewer_agent_id, task=t) + if guard: + guard.with_introspection(task=t, role=role_str) + return await self._emit_rejection( + self._with_briefing(guard, briefing), + agent_id=reviewer_agent_id, + task_id=task_id, + verb="claim_pr_review", + ) + runner = self._verb_runner() + try: + t = await runner.run_intent("claim_pr_review", t, agent, spec_ctx) + except Exception as exc: + return await self._runner_failure( + exc, + t, + role_str, + briefing, + reviewer_agent_id, + task_id, + "claim_pr_review", + ) + evidence = await self._build_pr_review_evidence(t) + return Envelope.ok( + status=str(t.status), + task_id=str(task_id), + next=spec_module._INTENT_VERBS["claim_pr_review"].next_hint(t), + evidence=evidence, + context_briefing=briefing, + ).with_introspection(task=t, role=role_str) + + async def post_pr_review( + self, + reviewer_agent_id: UUID, + task_id: UUID, + body: str, + event: str = "REQUEST_CHANGES", + ) -> Envelope: + """Post ONE change-request to the PR and finish the review task. + + Spec gate enforces role + the ``pr_review_done`` source-status + (IN_PROGRESS). The tracing gate requires a journal:learning entry. The + composed ``pr_review_done`` runs through the verb runner + (in_progress -> completed); then — mirroring QA's ``a2a.send`` pattern — + the review is posted to GitHub from the verb body, after the transition. + """ + t = await self.task.get(task_id) + if t is None: + return await self._emit_rejection( + Envelope.not_found(message=f"task {task_id} not found"), + agent_id=reviewer_agent_id, + task_id=task_id, + verb="post_pr_review", + ) + pre = await self._post_pr_review_preflight( + t, reviewer_agent_id, task_id, body + ) + if isinstance(pre, Envelope): + return pre + agent, role_str, briefing, spec_ctx = pre + slug = await self._project_slug_for(t) + pr_number = t.pr_number + runner = self._verb_runner() + try: + t = await runner.run_intent("post_pr_review", t, agent, spec_ctx) + except Exception as exc: + return await self._runner_failure( + exc, t, role_str, briefing, reviewer_agent_id, task_id, "post_pr_review" + ) + # GitHub side-effect AFTER the DB transition (a2a.send pattern). Best- + # effort: a posting failure is logged, not rolled back — the review task + # is complete; a missed post can be re-driven manually. + if slug and pr_number: + try: + await self.git.post_pr_review(slug, pr_number, body, event=event) + except Exception: + logger.exception( + "post_pr_review GitHub post failed", task_id=str(task_id) + ) + return Envelope.ok( + status=str(t.status), + task_id=str(task_id), + next=spec_module._INTENT_VERBS["post_pr_review"].next_hint(t), + context_briefing=briefing, + ).with_introspection(task=t, role=role_str) + + # -- helpers ---------------------------------------------------------- + + async def _post_pr_review_preflight( + self, t: Any, reviewer_agent_id: UUID, task_id: UUID, body: str + ) -> Any: + """Pre-runner guards for post_pr_review. + + Returns a rejection ``Envelope`` or the context tuple + ``(agent, role_str, briefing, spec_ctx)`` on pass. + """ + agent = await self.task.agent_for(reviewer_agent_id) + role_str = str(agent.role) if agent is not None else "pr_reviewer" + briefing = await self._briefing_for(reviewer_agent_id, task_id) + if not body or not body.strip(): + return await self._emit_rejection( + Envelope.invalid_state( + message="post_pr_review requires a non-empty review body", + remediate=( + "pass body=''" + ), + context_briefing=briefing, + ).with_introspection(task=t, role=role_str), + agent_id=reviewer_agent_id, + task_id=task_id, + verb="post_pr_review", + ) + role_or_rejection = await self._resolve_role( + t, role_str, briefing, reviewer_agent_id, task_id, "post_pr_review" + ) + if isinstance(role_or_rejection, Envelope): + return role_or_rejection + role = role_or_rejection + spec_ctx = spec_module.Context( + actor_id=reviewer_agent_id, + actor_slug=getattr(agent, "slug", None) if agent is not None else None, + notes=body, + ) + decision = spec_module.can_invoke_intent(role, "post_pr_review", t, spec_ctx) + if not decision.allowed: + return await self._emit_rejection( + Envelope.from_decision(decision, briefing=briefing).with_introspection( + task=t, role=role_str + ), + agent_id=reviewer_agent_id, + task_id=task_id, + verb="post_pr_review", + ) + gate = await self._pr_review_tracing_gate( + reviewer_agent_id, task_id, t, role_str + ) + if gate is not None: + return gate + return (agent, role_str, briefing, spec_ctx) + + async def _resolve_role( + self, + t: Any, + role_str: str, + briefing: dict[str, Any], + agent_id: UUID, + task_id: UUID, + verb: str, + ) -> Any: + """Parse the role enum, or return a not_authorized rejection Envelope.""" + try: + return spec_module.Role(role_str) + except ValueError: + return await self._emit_rejection( + Envelope.not_authorized( + message=f"unknown role '{role_str}'", + remediate="role is not declared in the lifecycle spec", + context_briefing=briefing, + ).with_introspection(task=t, role=role_str), + agent_id=agent_id, + task_id=task_id, + verb=verb, + ) + + async def _runner_failure( + self, + exc: Exception, + t: Any, + role_str: str, + briefing: dict[str, Any], + agent_id: UUID, + task_id: UUID, + verb: str, + ) -> Envelope: + """Shared rejection for a verb-runner failure.""" + return await self._emit_rejection( + Envelope.invalid_state( + message=f"verb runner failed: {exc}", + remediate="retry; if persistent, unclaim and notify the CEO", + context_briefing=briefing, + ).with_introspection(task=t, role=role_str), + agent_id=agent_id, + task_id=task_id, + verb=verb, + ) + + async def _build_pr_review_evidence(self, t: Any) -> dict[str, Any]: + """Inline evidence for claim_pr_review: the PR's unified diff (read-only).""" + slug = await self._project_slug_for(t) + diff = "" + if slug and t.pr_number: + diff = await self.git.get_pr_diff(slug, t.pr_number) + return { + "pr_number": t.pr_number, + "pr_url": t.pr_url, + "pr_diff": diff, + "is_external_pr": True, + } + + async def _pr_review_tracing_gate( + self, reviewer_agent_id: UUID, task_id: UUID, t: Any, role_str: str + ) -> Envelope | None: + """post_pr_review requires a journal:learning entry (parity with QA).""" + has_learning = await self.journal.has_learning_for_task( + reviewer_agent_id, task_id + ) + ctx = _tr.GateContext(journal_learning_present=has_learning) + result = _tr.check_requirements( + task=SimpleNamespace(), + requirements=list(_tr.requirements_for("post_pr_review")), + ctx=ctx, + ) + if result.passed: + return None + return await self._emit_rejection( + ( + await self._build_tracing_gap( + reviewer_agent_id, task_id, result.missing + ) + ).with_introspection(task=t, role=role_str), + agent_id=reviewer_agent_id, + task_id=task_id, + verb="post_pr_review", + ) + + async def _project_slug_for(self, t: Any) -> str | None: + """Resolve the project slug for a task (for read-only PR API calls).""" + from roboco.services.project import get_project_service + + if t.project_id is None: + return None + project = await get_project_service(self.task.session).get(t.project_id) + return project.slug if project is not None else None diff --git a/roboco/services/git.py b/roboco/services/git.py index 0f91dc82..2fc674af 100644 --- a/roboco/services/git.py +++ b/roboco/services/git.py @@ -1696,12 +1696,59 @@ class GitService(BaseService): raise GitError(f"PR not found: #{pr_number} on {owner}/{repo}", details) if not resp.is_success: raise GitError( - f"GitHub API refused PR review ({resp.status_code}): " - f"{resp.text[:200]}", + f"GitHub API refused PR review ({resp.status_code}): {resp.text[:200]}", details, ) return cast("dict[str, Any]", resp.json()) + async def get_pr_diff(self, project_slug: str, pr_number: int) -> str: + """Fetch a PR's unified diff READ-ONLY via the GitHub API. + + ``GET /pulls/{n}`` with the diff media type returns the unified diff + text without checking out or running any of the contributor's code — + the review is read-only; untrusted fork code never executes here. + Returns ``""`` on a missing token / unparseable remote / GitHub error + so the reviewer's claim still returns context instead of crashing. + """ + project = await get_project_service(self.session).get_by_slug(project_slug) + if project is None or not project.git_url: + return "" + try: + owner, repo = self._parse_git_url(project.git_url) + except GitError: + return "" + git_token = await self._token_for_project(project_slug) + if not git_token: + return "" + api_base = settings.github_api_base_url.rstrip("/") + try: + async with httpx.AsyncClient(timeout=_default_git_timeout()) as client: + resp = await client.get( + f"{api_base}/repos/{owner}/{repo}/pulls/{pr_number}", + headers={ + "Authorization": f"Bearer {git_token}", + "Accept": "application/vnd.github.v3.diff", + "X-GitHub-Api-Version": "2022-11-28", + }, + ) + except httpx.HTTPError as e: + self.log.warning( + "get_pr_diff request failed", + project=project_slug, + pr=pr_number, + error=str(e), + ) + return "" + if not resp.is_success: + self.log.warning( + "get_pr_diff non-2xx", + project=project_slug, + pr=pr_number, + status=resp.status_code, + ) + return "" + return resp.text + async def update_pr_for_task( self, task_id: UUID, diff --git a/roboco/services/task.py b/roboco/services/task.py index f02bbfac..78a9d23d 100644 --- a/roboco/services/task.py +++ b/roboco/services/task.py @@ -399,6 +399,9 @@ class TaskService(BaseService): # can reach in_progress and delegate. product_id is a plain column # (no lazy load). is_coordination=(task.project_id is None and task.product_id is not None), + # An external-PR review task reviews someone else's PR read-only — + # no branch of its own — so it is branch-gate exempt. + is_external_review=(getattr(task, "source", "manual") == "external_pr"), ) validate_git_requirements(current, target, git_ctx) @@ -672,6 +675,38 @@ class TaskService(BaseService): await self.session.flush() return task + async def complete_review( + self, reviewer_agent_id: UUID, task_id: UUID, notes: str | None = None + ) -> TaskTable | None: + """Mark an external-PR review task complete (in_progress -> completed). + + The terminal for the pr_reviewer's ``post_pr_review`` verb: the review + has been posted, so the review task is done. Attributed to the reviewer. + Mirrors ``qa_pass``'s validated-transition shape; the ``pr_review_done`` + transition (in_progress -> completed, role pr_reviewer) is defined in the + lifecycle spec and is git-gate exempt (the review task has no branch). + """ + task = await self.get(task_id) + if task is None: + return None + if task.status != TaskStatus.IN_PROGRESS: + return None + if notes: + task.qa_notes = notes + reviewer_id = to_python_uuid(task.claimed_by) or reviewer_agent_id + task.assigned_to = None + task.claimed_by = None + task.active_claimant_id = cast("Any", None) + self._validate_and_set_status( + task, + TaskStatus.COMPLETED, + "pr_reviewer", + audit_agent_id=reviewer_id, + ) + await self.session.flush() + self.log.info("External PR review complete", task_id=str(task_id)) + return task + async def _inherit_parent_session( self, task_id: UUID,