pr-reviewer/app/worker.py

260 lines
8.8 KiB
Python
Raw Normal View History

"""Background auto-review worker.
A daemon thread polls the PR queue on an interval and pre-reviews any new or
changed non-draft PR into the SQLite cache, so a review is ready by the time you
open it in the dashboard. Nothing is ever posted automatically.
Change detection is two-level to keep token spend down:
1. cheap gate - skip if the PR's ``updated_at`` is unchanged since the last
ready review (no diff fetched at all);
2. precise gate - if it did change, fetch the diff and skip the LLM call when
the diff's SHA-256 matches (e.g. the PR was only commented on).
``run_cycle`` takes its GitHub client, reviewer, and store as parameters so it is
directly testable with fakes; the ``Worker`` class only owns the thread/timing.
"""
from __future__ import annotations
import hashlib
import json
import logging
import random
import threading
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from typing import Any, Callable
# Keep a review for a closed/merged PR this long before hard-deleting it, so an
# open dashboard doesn't lose a review the moment the PR merges.
CLOSED_GRACE_SECONDS = 600
_log = logging.getLogger("pr_reviewer.worker")
def _sha256(text: str) -> str:
return hashlib.sha256(text.encode("utf-8", "ignore")).hexdigest()
def _is_rate_limited(exc: Exception) -> bool:
resp = getattr(exc, "response", None)
return getattr(resp, "status_code", None) == 429
def _review_with_backoff(
reviewer: Any,
pr: dict[str, Any],
diff: str,
*,
retries: int = 3,
base: float = 2.0,
sleep: Callable[[float], None] = time.sleep,
) -> dict[str, Any]:
"""Call the reviewer, backing off on HTTP 429 with exponential delay +
jitter. Non-429 errors (and the final 429) propagate to the caller."""
for attempt in range(retries + 1):
try:
return reviewer.review(pr, diff)
except Exception as e: # noqa: BLE001 - re-raised below unless retriable
if _is_rate_limited(e) and attempt < retries:
sleep(base * (2**attempt) + random.uniform(0, 1))
continue
raise
def _review_one(
gh: Any,
reviewer: Any,
store: Any,
pr: dict[str, Any],
sleep: Callable[[float], None],
) -> str:
"""Review a single PR and persist the outcome. Returns an outcome label and
never raises (failures are stored as ``error`` rows)."""
owner, repo, number = pr["owner"], pr["repo"], pr["number"]
meta = {
"pr_title": pr.get("title"),
"pr_author": pr.get("author"),
"pr_url": pr.get("url"),
"pr_updated_at": pr.get("updated_at"),
"pr_created_at": pr.get("created_at"),
"pr_node_id": pr.get("node_id"),
}
try:
diff = gh.pr_diff(owner, repo, number)
diff_hash = _sha256(diff)
existing = store.get(owner, repo, number)
if (
existing
and existing.get("diff_hash") == diff_hash
and existing.get("review_json")
):
# Diff unchanged (updated_at moved for a non-code reason) and we
# already have a review: keep it, just mark ready and refresh metadata
# so the cheap gate catches it next time. (Status is "reviewing" here
# because run_cycle marks queued PRs before handing them off.)
store.upsert(owner, repo, number, status="ready", **meta)
return "unchanged"
review = _review_with_backoff(reviewer, pr, diff, sleep=sleep)
store.upsert(
owner,
repo,
number,
status="ready",
diff_hash=diff_hash,
review_json=json.dumps(review),
error=None,
attempts=0,
**meta,
)
return "reviewed"
except Exception as e: # noqa: BLE001 - stored as an error row, surfaced in UI
existing = store.get(owner, repo, number)
attempts = (existing.get("attempts", 0) if existing else 0) + 1
store.upsert(
owner,
repo,
number,
status="error",
error=str(e),
attempts=attempts,
**meta,
)
return "error"
def run_cycle(
gh: Any,
reviewer: Any,
store: Any,
*,
concurrency: int = 2,
max_attempts: int = 3,
grace_seconds: float = CLOSED_GRACE_SECONDS,
sleep: Callable[[float], None] = time.sleep,
handbook: Any = None,
) -> dict[str, int]:
"""One poll cycle: refresh handbook guidance if stale, refresh the queue,
close/purge departed PRs, and review every new or changed non-draft PR
(bounded concurrency). Returns per-outcome counts."""
if handbook is not None:
handbook.refresh_if_stale() # daily in-process trigger; no-op when fresh
prs = gh.search_prs()
active = [p for p in prs if not p.get("draft")]
live_keys = {(p["owner"], p["repo"], p["number"]) for p in active}
store.mark_missing_closed(live_keys)
store.purge_closed(grace_seconds)
todo: list[dict[str, Any]] = []
for p in active:
existing = store.get(p["owner"], p["repo"], p["number"])
if existing:
same_update = existing.get("pr_updated_at") == p.get("updated_at")
if existing.get("status") == "ready" and same_update:
# Cheap gate: no re-review. But backfill lightweight PR metadata
# (e.g. node_id/created_at added in a later release) when it
# actually differs, so features like auto-merge have current data
# without re-running the model. No-op write when nothing differs.
if existing.get("pr_node_id") != p.get("node_id") or existing.get(
"pr_created_at"
) != p.get("created_at"):
store.upsert(
p["owner"],
p["repo"],
p["number"],
status="ready",
pr_node_id=p.get("node_id"),
pr_created_at=p.get("created_at"),
)
continue
if (
existing.get("status") == "error"
and same_update
and existing.get("attempts", 0) >= max_attempts
):
continue # keeps failing: wait until the PR itself changes
todo.append(p)
# Mark queued PRs as reviewing up front so the UI reflects in-flight work.
for p in todo:
store.upsert(
p["owner"],
p["repo"],
p["number"],
status="reviewing",
pr_title=p.get("title"),
pr_author=p.get("author"),
pr_url=p.get("url"),
pr_updated_at=p.get("updated_at"),
pr_created_at=p.get("created_at"),
pr_node_id=p.get("node_id"),
)
results = {
"reviewed": 0,
"unchanged": 0,
"error": 0,
"skipped": len(active) - len(todo),
}
if todo:
with ThreadPoolExecutor(max_workers=max(1, concurrency)) as ex:
futures = [
ex.submit(_review_one, gh, reviewer, store, p, sleep) for p in todo
]
for fut in as_completed(futures):
results[fut.result()] += 1
return results
class Worker:
"""Owns the poll thread and its timing; delegates the work to ``run_cycle``."""
def __init__(
self, cfg: Any, gh: Any, reviewer: Any, store: Any, handbook: Any = None
) -> None:
self.cfg = cfg
self.gh = gh
self.reviewer = reviewer
self.store = store
self.handbook = handbook
self._stop = threading.Event()
self._wake = threading.Event()
self._thread: threading.Thread | None = None
def start(self) -> None:
n = self.store.reset_stale_reviewing()
if n:
_log.info("reset %d interrupted review(s) from a prior run", n)
self._thread = threading.Thread(
target=self._loop, name="review-worker", daemon=True
)
self._thread.start()
def trigger(self) -> None:
"""Ask the loop to run a cycle now (non-blocking)."""
self._wake.set()
def stop(self) -> None:
self._stop.set()
self._wake.set()
if self._thread:
self._thread.join(timeout=5)
def _loop(self) -> None:
while not self._stop.is_set():
try:
res = run_cycle(
self.gh,
self.reviewer,
self.store,
concurrency=self.cfg.WORKER_CONCURRENCY,
max_attempts=self.cfg.MAX_REVIEW_ATTEMPTS,
handbook=self.handbook,
)
_log.info("review cycle: %s", res)
except Exception:
_log.exception("review cycle failed")
self._wake.wait(timeout=self.cfg.POLL_INTERVAL)
self._wake.clear()