"""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()