open-swe/agent/ci_autofix.py
Adam Moussa b3fc62da80
refactor: split webapp.py into api/ + per-source webhook routes
Plan step C4 (docs/upstream-sync/domain-reorg/reorg-build-plan.md, approved
decisions 1-2): split the 2,590-line agent/webapp.py monolith into
agent/webhooks/common.py (shared verify/dispatch helpers), agent/api/app.py
(composition), agent/api/health.py (/health + /webhooks/run-complete), and
per-source {github,linear,slack,jira,confluence}_routes.py. Atlassian
Connect lifecycle + descriptor routes (/connect/*) fold into
confluence_routes.py; webapp.py becomes the upstream-shaped compatibility
shim (from .api.app import app). langgraph.json http.app stays
agent.webapp:app via the shim.

Fork content, upstream layout: linear/slack route files verified
content-identical to upstream 8356eb34 and taken verbatim; github_routes is
upstream + the fork's CI auto-fix trigger wiring; jira/confluence routes are
fork-only, transformed to the same common.X / service.X module-attribute
style. All signature verification (GitHub HMAC, Slack, Linear
timestamp-freshness, verify_jira_secret + opt-in HMAC/timestamp/IP
allowlist, Connect JWT/qsh), token-attribution gating, TID-COLLIDE-01 repo
binding, _is_repo_auto_review_enabled gates, and public-repo org gate move
unchanged.

Handlers rewired from webapp.X to common.X; test monkeypatch sites across
26 files + conftest.py + e2e/harness.py retargeted to
webhook_common/handler/route modules per upstream's pattern. Residual
agent.webapp importers: only the shim, langgraph.json http.app, Makefile
uvicorn target, and docs (doc-path updates land in C7).

Gates: ruff check + format, pytest --co, full unit (1637 passed), full
Playwright E2E vs real langgraph dev (9/9), residual-importer sweep.
2026-07-17 14:30:05 -04:00

616 lines
23 KiB
Python

"""Auto-fix CI failures and review feedback on agent-authored pull requests.
This is the shared core for "PR babysitting": when a CI check fails (or a
reviewer leaves actionable feedback) on a PR that Open SWE opened, locate the
originating agent thread and dispatch a confidence-gated fix run on it.
Both the GitHub webhook path (:mod:`agent.webhooks.github_routes`) and the polling fallback
(:mod:`agent.ci_monitor`) call into here, so all the skip-rules, dedupe, and
loop-capping live in one place. Skip-rules mirror Cursor/Claude Code:
* Only PRs Open SWE authored (an agent thread with this ``pr_url`` exists).
* Skip failures inherited from the base branch.
* Skip when the latest commit was authored by a human (don't fight pushes).
* Dedupe per head SHA; cap total attempts.
* Honor the per-user ``auto_fix_ci`` profile flag and the per-PR opt-out.
"""
from __future__ import annotations
import logging
from typing import Any
from langgraph_sdk import get_client
from .dashboard.agent_overrides import load_profile, resolve_login_from_email_async
from .dashboard.autofix_state import is_pr_autofix_disabled
from .dashboard.enabled_repos import is_review_repo_enabled
from .dispatch import dispatch_agent_run
from .review.findings import REVIEWER_THREAD_KIND
from .utils.dashboard_links import dashboard_thread_url
from .utils.github_app import get_github_app_installation_token
from .utils.github_checks import post_autofix_status_check
from .utils.github_ci import (
fetch_open_pr_for_branch,
fetch_pr,
has_repo_write_permission,
head_commit_author_login,
list_failing_check_runs,
list_failing_statuses,
names_failing_on_base,
)
from .utils.github_org_membership import INTERNAL_BOT_LOGINS
from .utils.thread_ops import (
get_thread_active_status,
langgraph_client,
)
logger = logging.getLogger(__name__)
# Hard cap on auto-fix follow-ups per PR so a failure the agent can't resolve
# doesn't loop forever (Cursor caps at 10).
MAX_AUTOFIX_ATTEMPTS = 10
# Keep the dedupe list bounded on thread metadata.
_MAX_HANDLED_KEYS = 30
# Store location for batched auto-fix events, consumed by the message-queue middleware.
_PENDING_AUTOFIX_NS = "autofix"
_PENDING_AUTOFIX_KEY = "pending_event"
def _dedupe_key(head_sha: str) -> str:
return head_sha
async def _user_autofix_enabled(github_login: str, user_email: str = "") -> bool:
"""Check the per-user ``auto_fix_ci`` profile flag (defaults to True)."""
login = github_login.strip() if isinstance(github_login, str) else ""
if not login and user_email:
login = await resolve_login_from_email_async(user_email)
if not login:
return True
profile = await load_profile(login)
if not isinstance(profile, dict):
return True
value = profile.get("auto_fix_ci")
return value if isinstance(value, bool) else True
async def find_agent_thread_for_pr(pr_url: str) -> tuple[str, dict[str, Any]] | None:
"""Return ``(thread_id, metadata)`` of the agent thread that opened ``pr_url``.
Reviewer threads are skipped — only the coding-agent thread can push fixes.
"""
if not pr_url:
return None
client = get_client()
try:
threads = await client.threads.search(metadata={"pr_url": pr_url}, limit=10)
except Exception: # noqa: BLE001
logger.debug("Could not search threads for PR %s", pr_url, exc_info=True)
return None
for thread in threads or []:
metadata = thread.get("metadata") if isinstance(thread, dict) else None
if not isinstance(metadata, dict):
continue
if metadata.get("kind") == REVIEWER_THREAD_KIND:
continue
if metadata.get("agent_kind") != "agent":
continue
thread_id = thread.get("thread_id") or thread.get("id")
if isinstance(thread_id, str) and thread_id:
return thread_id, metadata
return None
def _build_ci_fix_prompt(
*,
owner: str,
repo: str,
pr_number: int,
pr_url: str,
branch: str,
head_sha: str,
failing_checks: list[dict[str, Any]],
) -> str:
lines = []
for check in failing_checks:
name = check.get("name", "check")
conclusion = check.get("conclusion", "failure")
details = check.get("details_url") or ""
suffix = f" — {details}" if details else ""
lines.append(f"- {name} ({conclusion}){suffix}")
failing_block = "\n".join(lines)
return (
"An automated CI check failed on a pull request you opened. Please "
"investigate and fix it.\n\n"
f"## Repository: {owner}/{repo}\n\n"
f"## Pull Request: {pr_url} (#{pr_number})\n\n"
f"## Branch: {branch}\n\n"
f"## Head commit: {head_sha}\n\n"
f"## Failing checks:\n{failing_block}\n\n"
"Instructions:\n"
"1. Make sure you are on the PR branch, then read the failing logs "
"(e.g. `GH_TOKEN=dummy gh pr checks` and `GH_TOKEN=dummy gh run view "
"<run-id> --log-failed`).\n"
"2. Confidence gating — fix autonomously ONLY when the cause is clear "
"and deterministic (lint/format, type errors, missing imports, failed "
"assertions, snapshot updates, build errors). Commit and push to the "
"existing branch; do NOT open a new PR.\n"
"3. If the failure is ambiguous, flaky, infrastructure-related, appears "
"pre-existing, or needs an architectural/design decision, do NOT guess. "
"Post a short PR comment explaining what you found and what input you "
"need, then stop.\n"
"4. Never force-push. Never weaken or delete test assertions just to go "
"green unless the behavior change is intentional and correct.\n"
"5. Before finishing, re-check the PR's latest CI status and review "
"comments. Address any newly failed checks or unhandled actionable "
"comments that arrived while you were working.\n"
"6. After you push, CI re-runs automatically — you don't need to merge."
)
def _build_review_feedback_prompt(
*,
owner: str,
repo: str,
pr_number: int,
pr_url: str,
reviewer: str,
body: str,
) -> str:
return (
"A reviewer left feedback on a pull request you opened. Please respond.\n\n"
f"## Repository: {owner}/{repo}\n\n"
f"## Pull Request: {pr_url} (#{pr_number})\n\n"
f"## Reviewer: {reviewer}\n\n"
f"## Feedback:\n{body}\n\n"
"Instructions:\n"
"1. If the requested change is unambiguous (rename, typo, missing null "
"check, small refactor, add a test), make it, commit, and push to the "
"existing branch.\n"
"2. If the comment is ambiguous, opinion-based, or needs a design "
"decision, reply on the PR asking for clarification instead of guessing.\n"
"3. Before finishing, re-check the PR's latest review comments and CI "
"status. Address any newly arrived actionable comments or failed checks "
"that are clear and deterministic.\n"
"4. Never force-push. Reply to the reviewer on GitHub to explain what "
"you changed."
)
async def _thread_autofix_state(metadata: dict[str, Any]) -> tuple[int, list[str], str, str]:
attempts = metadata.get("autofix_attempts")
attempts = attempts if isinstance(attempts, int) and attempts >= 0 else 0
handled = metadata.get("autofix_handled")
handled = [h for h in handled if isinstance(h, str)] if isinstance(handled, list) else []
github_login = metadata.get("github_login")
github_login = github_login if isinstance(github_login, str) else ""
user_email = metadata.get("triggering_user_email")
user_email = user_email if isinstance(user_email, str) else ""
return attempts, handled, github_login, user_email
async def _record_attempt(
thread_id: str, *, attempts: int, handled: list[str], dedupe_key: str, head_sha: str
) -> None:
new_handled = [*handled, dedupe_key][-_MAX_HANDLED_KEYS:]
try:
await get_client().threads.update(
thread_id=thread_id,
metadata={
"autofix_attempts": attempts + 1,
"autofix_handled": new_handled,
"autofix_last_head_sha": head_sha,
},
)
except Exception: # noqa: BLE001
logger.debug("Failed to record auto-fix attempt for thread %s", thread_id, exc_info=True)
# Run sources the agent's GitHub-token resolver knows how to authenticate.
_AUTH_RESOLVABLE_SOURCES = frozenset(["github", "slack", "dashboard", "linear", "schedule"])
def _run_configurable(
metadata: dict[str, Any], *, repo_config: dict[str, str], pr_number: int
) -> dict[str, Any]:
"""Build the run config for a fix run by reusing the PR thread's identity.
The agent's GitHub-token resolver only authenticates known sources, so a
bespoke ``github_ci`` source would fail in non-bot-token deployments. Reuse
the originating thread's ``source`` + login/email so auth resolves exactly
as it did for the run that opened the PR.
"""
source = metadata.get("source")
if source not in _AUTH_RESOLVABLE_SOURCES:
source = "github"
configurable: dict[str, Any] = {
"source": source,
"repo": repo_config,
"pr_number": pr_number,
}
login = metadata.get("github_login")
if isinstance(login, str) and login:
configurable["github_login"] = login
email = metadata.get("triggering_user_email")
if isinstance(email, str) and email:
configurable["user_email"] = email
return configurable
async def _mark_pending_autofix_event(thread_id: str, reason: str, detail: str = "") -> None:
"""Record a batched auto-fix event in the store the message-queue middleware reads.
Uses the same store namespace mechanism as ``queue_message_for_thread`` so the
in-flight run picks it up in-process at its next ``before_model`` step — no
per-step thread fetch. ``detail`` (e.g. a reviewer's comment) is accumulated so
specifics aren't lost when several events batch against one busy run.
"""
client = langgraph_client()
namespace = (_PENDING_AUTOFIX_NS, thread_id)
try:
details: list[str] = []
try:
existing = await client.store.get_item(namespace, _PENDING_AUTOFIX_KEY)
if existing and existing.get("value"):
prior = existing["value"].get("details")
if isinstance(prior, list):
details = [d for d in prior if isinstance(d, str)]
except Exception: # noqa: BLE001
logger.debug("No existing pending auto-fix event for thread %s", thread_id)
if detail and detail not in details:
details.append(detail)
await client.store.put_item(
namespace, _PENDING_AUTOFIX_KEY, {"reason": reason, "details": details}
)
except Exception: # noqa: BLE001
logger.debug(
"Failed to record pending auto-fix event for thread %s", thread_id, exc_info=True
)
async def _dispatch_or_batch(
thread_id: str, prompt: str, *, configurable: dict[str, Any], reason: str, detail: str = ""
) -> str:
# Deliberate skip-rule: batch auto-fix events while the agent thread is
# actively running so we don't interrupt an in-progress fix. ``interrupt``
# is fine for human follow-ups but undesirable for autofix, so we keep the
# busy-check here even though the webhook hot-path no longer needs one.
if await get_thread_active_status(thread_id) is True:
logger.info("Agent thread %s busy; batching auto-fix event %s", thread_id, reason)
await _mark_pending_autofix_event(thread_id, reason, detail)
return "batched"
# The busy-check above has a TOCTOU window (the dedupe SHA is only recorded
# after dispatch), so a burst of near-simultaneous CI events for one head SHA
# can all pass the gate. Dispatch with ``reject`` — matching ``dev``'s prior
# platform default — so the platform drops the duplicate concurrent creates
# instead of letting them interrupt each other.
await dispatch_agent_run(
thread_id,
prompt,
configurable,
source=str(configurable.get("source") or "github_autofix"),
multitask_strategy="reject",
)
logger.info(
"Created auto-fix run for thread %s (source=%s)", thread_id, configurable.get("source")
)
return "dispatched"
async def handle_ci_failure(
*,
repo_config: dict[str, str],
branch: str,
head_sha: str,
token: str | None = None,
source: str = "github_ci",
failing_checks: list[dict[str, Any]] | None = None,
pr: dict[str, Any] | None = None,
) -> str:
"""Auto-fix failing CI on an agent-authored PR. Returns a status string."""
owner = repo_config.get("owner", "")
repo = repo_config.get("name", "")
if not owner or not repo:
return "missing_repo"
if not await is_review_repo_enabled(owner, repo):
return "repo_not_enabled"
if token is None:
token = await get_github_app_installation_token()
if not token:
logger.warning("No GitHub App token for CI auto-fix on %s/%s", owner, repo)
return "no_token"
if pr is None:
if not branch:
return "no_branch"
pr = await fetch_open_pr_for_branch(owner=owner, repo=repo, branch=branch, token=token)
if not pr:
return "no_open_pr"
pr_number = pr.get("number")
if not isinstance(pr_number, int):
return "no_pr_number"
pr_url = pr.get("html_url") or pr.get("url") or ""
base_sha = (pr.get("base") or {}).get("sha", "")
branch = branch or (pr.get("head") or {}).get("ref", "")
head_sha = head_sha or (pr.get("head") or {}).get("sha", "")
if not head_sha:
return "no_head_sha"
if await is_pr_autofix_disabled(owner, repo, pr_number):
return "pr_disabled"
found = await find_agent_thread_for_pr(pr_url)
if found is None:
return "no_agent_thread"
thread_id, metadata = found
attempts, handled, github_login, user_email = await _thread_autofix_state(metadata)
if not await _user_autofix_enabled(github_login, user_email):
return "autofix_disabled_user"
if attempts >= MAX_AUTOFIX_ATTEMPTS:
await post_autofix_status_check(
owner=owner,
repo=repo,
head_sha=head_sha,
token=token,
title="Auto-fix limit reached",
summary=(
f"Open SWE has attempted {attempts} auto-fixes on this PR and "
"stopped to avoid a loop. Push a commit or comment to continue."
),
details_url=dashboard_thread_url(thread_id),
)
return "max_attempts"
if failing_checks is None:
runs = await list_failing_check_runs(owner=owner, repo=repo, ref=head_sha, token=token)
statuses = await list_failing_statuses(owner=owner, repo=repo, ref=head_sha, token=token)
if runs is None and statuses is None:
return "ci_read_failed"
failing_checks = (runs or []) + (statuses or [])
if not failing_checks:
return "no_failing_checks"
base_failing = await names_failing_on_base(
owner=owner, repo=repo, base_sha=base_sha, token=token
)
actionable = [c for c in failing_checks if c.get("name") not in base_failing]
if not actionable:
return "all_failing_on_base"
dedupe_key = _dedupe_key(head_sha)
if dedupe_key in handled:
return "already_handled"
author_login = await head_commit_author_login(owner=owner, repo=repo, sha=head_sha, token=token)
if (
author_login is not None
and author_login not in INTERNAL_BOT_LOGINS
and (not github_login or author_login.lower() != github_login.lower())
):
logger.info(
"Skipping CI auto-fix on %s/%s#%s: head commit %s authored by human %s",
owner,
repo,
pr_number,
head_sha,
author_login,
)
return "human_commit"
prompt = _build_ci_fix_prompt(
owner=owner,
repo=repo,
pr_number=pr_number,
pr_url=pr_url,
branch=branch,
head_sha=head_sha,
failing_checks=actionable,
)
result = await _dispatch_or_batch(
thread_id,
prompt,
configurable=_run_configurable(
metadata, repo_config={"owner": owner, "name": repo}, pr_number=pr_number
),
reason="ci_failure",
)
if result == "dispatched":
# Only burn an attempt / mark the SHA handled on a real dispatch. A batched
# event is just a nudge to the in-flight run; if that run ends before
# consuming it, leaving the SHA un-handled lets a later webhook or the sweep
# re-dispatch instead of silently dropping the failure.
await _record_attempt(
thread_id, attempts=attempts, handled=handled, dedupe_key=dedupe_key, head_sha=head_sha
)
await post_autofix_status_check(
owner=owner,
repo=repo,
head_sha=head_sha,
token=token,
title=f"Auto-fixing {len(actionable)} failing check(s)",
summary=(
"Open SWE is investigating the failing checks and will push a fix if "
"the cause is clear. Track progress in the linked run."
),
details_url=dashboard_thread_url(thread_id),
)
return result
async def handle_review_feedback(
*,
repo_config: dict[str, str],
pr_number: int,
pr_url: str,
reviewer: str,
body: str,
token: str | None = None,
source: str = "github_review",
) -> str:
"""Auto-respond to a human review comment on an agent-authored PR."""
owner = repo_config.get("owner", "")
repo = repo_config.get("name", "")
if not owner or not repo or not pr_url:
return "missing_repo"
if not await is_review_repo_enabled(owner, repo):
return "repo_not_enabled"
if await is_pr_autofix_disabled(owner, repo, pr_number):
return "pr_disabled"
found = await find_agent_thread_for_pr(pr_url)
if found is None:
return "no_agent_thread"
thread_id, metadata = found
_, _, github_login, user_email = await _thread_autofix_state(metadata)
if not await _user_autofix_enabled(github_login, user_email):
return "autofix_disabled_user"
if token is None:
token = await get_github_app_installation_token()
if not token:
return "no_token"
if not await has_repo_write_permission(owner=owner, repo=repo, username=reviewer, token=token):
logger.info(
"Skipping auto-fix review feedback on %s/%s#%s: %s lacks write access",
owner,
repo,
pr_number,
reviewer or "<unknown>",
)
return "reviewer_no_write_permission"
prompt = _build_review_feedback_prompt(
owner=owner,
repo=repo,
pr_number=pr_number,
pr_url=pr_url,
reviewer=reviewer,
body=body,
)
return await _dispatch_or_batch(
thread_id,
prompt,
configurable=_run_configurable(
metadata, repo_config={"owner": owner, "name": repo}, pr_number=pr_number
),
reason="review_feedback",
detail=f"Reviewer {reviewer or 'unknown'} commented: {body.strip()}"
if body.strip()
else "",
)
async def sweep_open_prs() -> dict[str, int]:
"""Poll open agent-authored PRs and auto-fix failing CI / flag conflicts.
The polling fallback for deployments without reliable CI webhooks, and the
only path that can react to base-branch merge conflicts (GitHub emits no
webhook for those).
"""
counts = {"scanned": 0, "dispatched": 0, "batched": 0, "conflicts": 0}
token = await get_github_app_installation_token()
if not token:
logger.warning("CI monitor sweep: no GitHub App token")
return counts
client = get_client()
try:
threads = await client.threads.search(
metadata={"agent_kind": "agent", "pr_state": "open"}, limit=100
)
except Exception: # noqa: BLE001
logger.warning("CI monitor sweep: thread search failed", exc_info=True)
return counts
for thread in threads or []:
metadata = thread.get("metadata") if isinstance(thread, dict) else None
if not isinstance(metadata, dict):
continue
repo = metadata.get("repo")
pr_number = metadata.get("pr_number")
branch = metadata.get("branch_name")
if not isinstance(repo, dict) or not isinstance(pr_number, int):
continue
owner = repo.get("owner", "")
name = repo.get("name", "")
if not owner or not name:
continue
counts["scanned"] += 1
pr = await fetch_pr(owner=owner, repo=name, pr_number=pr_number, token=token)
if not pr:
continue
head_sha = (pr.get("head") or {}).get("sha", "")
branch = (pr.get("head") or {}).get("ref", "") or (
branch if isinstance(branch, str) else ""
)
if pr.get("mergeable_state") == "dirty":
counts["conflicts"] += 1
await _flag_merge_conflict(
owner=owner,
repo=name,
pr_number=pr_number,
pr_url=pr.get("html_url") or "",
head_sha=head_sha,
token=token,
)
continue
result = await handle_ci_failure(
repo_config={"owner": owner, "name": name},
branch=branch,
head_sha=head_sha,
token=token,
source="ci_monitor",
pr=pr,
)
if result == "dispatched":
counts["dispatched"] += 1
elif result == "batched":
counts["batched"] += 1
logger.info("CI monitor sweep complete: %s", counts)
return counts
async def _flag_merge_conflict(
*, owner: str, repo: str, pr_number: int, pr_url: str, head_sha: str, token: str
) -> None:
"""Ask the agent to rebase a PR that has merge conflicts with its base."""
if await is_pr_autofix_disabled(owner, repo, pr_number):
return
found = await find_agent_thread_for_pr(pr_url)
if found is None:
return
thread_id, metadata = found
_, _, github_login, user_email = await _thread_autofix_state(metadata)
if not await _user_autofix_enabled(github_login, user_email):
return
if metadata.get("autofix_conflict_head") == head_sha:
return
prompt = (
f"The pull request you opened (#{pr_number}, {pr_url}) now has merge "
"conflicts with its base branch. Rebase or merge the base branch into "
"the PR branch, resolve the conflicts carefully, and push. If a "
"conflict resolution is ambiguous, comment on the PR and ask before "
"guessing. Never force-push over commits already on the remote."
)
await _dispatch_or_batch(
thread_id,
prompt,
configurable=_run_configurable(
metadata, repo_config={"owner": owner, "name": repo}, pr_number=pr_number
),
reason="merge_conflict",
)
try:
await get_client().threads.update(
thread_id=thread_id, metadata={"autofix_conflict_head": head_sha}
)
except Exception: # noqa: BLE001
logger.debug("Failed to record conflict head for thread %s", thread_id, exc_info=True)