From 15a52bff405fb413d19c9fdd17197aac72e57c6a Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Wed, 1 Jul 2026 14:28:11 -0400 Subject: [PATCH] Add background auto-review worker with SQLite cache Pre-compute PR reviews so they are ready the moment a PR is opened in the queue, instead of waiting on an on-demand Fireworks call each time. A daemon worker polls the queue every POLL_INTERVAL and reviews new or changed non-draft PRs into a local SQLite cache (gitignored). Change detection is two-level: skip when the PR's updated_at is unchanged, and even when it moved, skip the model call when the diff's SHA-256 matches, so comment-only bumps do not burn tokens. Failed reviews retry up to MAX_REVIEW_ATTEMPTS with 429 backoff; departed PRs are closed with a grace window before purge. Singletons are initialized eagerly in the FastAPI lifespan before the worker thread starts to avoid an init race; all cache writes are serialized. New endpoints GET /api/reviews and POST /api/refresh back a dashboard that polls for status (reviewing/ready/error) and opens a ready review instantly. Nothing is posted automatically; the human still decides. Stdlib-only, so no new runtime or test dependencies. --- .env.example | 12 ++- .gitignore | 4 + README.md | 31 +++++- app/config.py | 21 +++- app/main.py | 118 ++++++++++++++++++-- app/store.py | 157 +++++++++++++++++++++++++++ app/worker.py | 232 ++++++++++++++++++++++++++++++++++++++++ static/index.html | 129 +++++++++++++++------- tests/test_endpoints.py | 70 ++++++++++++ tests/test_store.py | 93 ++++++++++++++++ tests/test_worker.py | 180 +++++++++++++++++++++++++++++++ 11 files changed, 996 insertions(+), 51 deletions(-) create mode 100644 app/store.py create mode 100644 app/worker.py create mode 100644 tests/test_store.py create mode 100644 tests/test_worker.py diff --git a/.env.example b/.env.example index e08e991..0c30cc8 100644 --- a/.env.example +++ b/.env.example @@ -3,7 +3,7 @@ GITHUB_TOKEN= GITHUB_ORG=Sea-Haven-Industries PR_SEARCH_FILTER=is:pr state:open archived:false sort:updated-desc org:Sea-Haven-Industries -MAX_PRS=30 +MAX_PRS=100 MAX_DIFF_BYTES=120000 # --- Fireworks --- @@ -21,3 +21,13 @@ MENTION_AUTHORS=openswe HOST=127.0.0.1 PORT=8765 + +# --- Background auto-review worker --- +# Seconds between poll cycles (pre-reviews new/changed non-draft PRs). +POLL_INTERVAL=300 +# How many PRs to review in parallel per cycle (kept low to respect rate limits). +WORKER_CONCURRENCY=2 +# Stop auto-retrying a PR that keeps failing until its diff changes. +MAX_REVIEW_ATTEMPTS=3 +# SQLite cache path. Defaults to pr_cache.db in the repo root (gitignored). +# CACHE_DB=/absolute/path/to/pr_cache.db diff --git a/.gitignore b/.gitignore index 47bae5a..dff70cf 100644 --- a/.gitignore +++ b/.gitignore @@ -1,6 +1,10 @@ # Secrets and local env .env +# Local review cache (SQLite + WAL/SHM sidecars) +pr_cache.db +pr_cache.db-* + # Python virtualenv and caches .venv/ __pycache__/ diff --git a/README.md b/README.md index 3506762..7837162 100644 --- a/README.md +++ b/README.md @@ -2,6 +2,8 @@ A local dashboard that pulls open PRs from your GitHub org, reviews each one with a Fireworks model using the BLOCK / FIX / NIT / QUESTION skill format, and lets you request revisions or post the review to GitHub as yourself. +A background worker pre-reviews non-draft PRs on an interval, so a review is usually ready the moment you open one in the queue. You still decide whether and how to post; nothing is ever posted automatically. + Everything runs on your machine. This is a local, single-user tool. It is not deployed anywhere, so there is no AWS stack, no CI deploy path, and secrets live only in a local `.env` (gitignored). Your GitHub token and Fireworks key stay in the backend and never reach the browser. ## Setup @@ -30,11 +32,18 @@ Open http://127.0.0.1:8765 ## How it works -1. **Refresh queue** runs your filter (`PR_SEARCH_FILTER`, default matches your org filter) and lists the PRs. -2. **Run review** fetches the PR diff and sends it to Fireworks. The result is parsed into the four categories plus a summary line and a recommended event. +1. **Background worker** polls your filter (`PR_SEARCH_FILTER`) every `POLL_INTERVAL` seconds and pre-reviews any new or changed non-draft PR, caching the result. The queue shows each PR's status: `reviewing`, `ready`, `error`, or `closed`. **Refresh now** forces an immediate poll. +2. **Open a PR** — if its review is `ready`, it appears instantly. Otherwise you see its status, and you can **Run review now** on demand. 3. **Request revision** re-runs the review with your notes folded in as a trusted instruction, separate from the untrusted diff. 4. **Post review to GitHub** submits it as a PR review. You confirm the event type (COMMENT / APPROVE / REQUEST_CHANGES) and can edit the body first. +### Auto-review worker + +- Reviews are cached in a local SQLite file (`CACHE_DB`, default `pr_cache.db` in the repo root, gitignored) so they survive restarts and aren't recomputed for unchanged PRs. +- Change detection is two-level: a PR is skipped if its `updated_at` hasn't moved since the last review, and even when it has, the diff's SHA-256 is compared so a comment-only bump doesn't burn tokens. +- Drafts are skipped. Failed reviews are retried on later cycles up to `MAX_REVIEW_ATTEMPTS`, then left until the PR changes. Rate-limit (HTTP 429) responses back off and retry. +- `WORKER_CONCURRENCY` controls how many PRs are reviewed in parallel per cycle (default 2). + ## The @mention rule Any PR whose author login is in `MENTION_AUTHORS` (default `openswe`) gets an `@author` mention prepended to the review summary. Add more logins comma-separated. @@ -42,9 +51,21 @@ Any PR whose author login is in `MENTION_AUTHORS` (default `openswe`) gets an `@ ## Notes - The diff is treated as untrusted input; the model is instructed to ignore any embedded instructions. -- State is in memory and resets on restart. This is a single-user local tool, not a shared service. +- Cached reviews persist in a local SQLite file. This is a single-user local tool, not a shared service. - Large diffs are truncated at `MAX_DIFF_BYTES` to control token cost. +## Configuration + +Set in `.env` (see `.env.example`). Beyond the GitHub/Fireworks keys: + +| Var | Default | Purpose | +|---|---|---| +| `POLL_INTERVAL` | `300` | seconds between background poll cycles | +| `WORKER_CONCURRENCY` | `2` | PRs reviewed in parallel per cycle | +| `MAX_REVIEW_ATTEMPTS` | `3` | error retries before giving up until the PR changes | +| `MAX_PRS` | `100` | cap on PRs pulled per cycle (GitHub search page max) | +| `CACHE_DB` | `pr_cache.db` | SQLite cache path (absolute, repo root by default) | + ## Layout ``` @@ -52,7 +73,9 @@ app/ config.py settings from .env github_client.py search PRs, fetch diffs, post reviews reviewer.py Fireworks call + skill format + markdown rendering - main.py FastAPI endpoints + store.py SQLite cache of pre-computed reviews + worker.py background poll + auto-review (run_cycle) + main.py FastAPI endpoints (incl. /api/reviews, /api/refresh) static/ index.html the dashboard ``` diff --git a/app/config.py b/app/config.py index 83e886e..e810ae3 100644 --- a/app/config.py +++ b/app/config.py @@ -6,11 +6,14 @@ backend; the backend holds the GitHub PAT and Fireworks key. import os from functools import lru_cache +from pathlib import Path from dotenv import load_dotenv load_dotenv() +_REPO_ROOT = Path(__file__).resolve().parent.parent + class Config: # --- GitHub --- @@ -23,8 +26,9 @@ class Config: "PR_SEARCH_FILTER", "is:pr state:open archived:false sort:updated-desc org:Sea-Haven-Industries", ).strip() - # Cap how many PRs we pull per refresh. - MAX_PRS: int = int(os.getenv("MAX_PRS", "30")) + # Cap how many PRs we pull per refresh. GitHub search returns at most 100 + # per page; PRs beyond that are not paginated (fine at current org volume). + MAX_PRS: int = int(os.getenv("MAX_PRS", "100")) # Skip diffs larger than this many bytes (keeps token cost sane). MAX_DIFF_BYTES: int = int(os.getenv("MAX_DIFF_BYTES", "120000")) @@ -54,6 +58,19 @@ class Config: HOST: str = os.getenv("HOST", "127.0.0.1") PORT: int = int(os.getenv("PORT", "8765")) + # --- Background auto-review worker --- + # Seconds between poll cycles. The worker fetches the queue, then pre-reviews + # any new or changed non-draft PR so results are ready before you open them. + POLL_INTERVAL: int = int(os.getenv("POLL_INTERVAL", "300")) + # How many PRs to review in parallel per cycle. Kept low for a single-user + # tool to stay well under Fireworks rate limits. + WORKER_CONCURRENCY: int = int(os.getenv("WORKER_CONCURRENCY", "2")) + # Give up auto-retrying a PR that keeps failing until its diff changes. + MAX_REVIEW_ATTEMPTS: int = int(os.getenv("MAX_REVIEW_ATTEMPTS", "3")) + # SQLite cache of pre-computed reviews. Absolute path resolved from the repo + # root so it lands in the same place regardless of the working directory. + CACHE_DB: str = os.getenv("CACHE_DB", str(_REPO_ROOT / "pr_cache.db")) + @lru_cache def get_config() -> Config: diff --git a/app/main.py b/app/main.py index 39f5c63..4bae2d5 100644 --- a/app/main.py +++ b/app/main.py @@ -2,37 +2,47 @@ Endpoints: GET / -> serves the dashboard - GET /api/prs -> list PRs matching the filter (no review yet) - POST /api/review -> run a Fireworks review for one PR + GET /api/config -> org/filter/model info for the UI + GET /api/prs -> live list of PRs matching the filter + GET /api/reviews -> cached, pre-computed reviews + per-PR status + POST /api/refresh -> ask the background worker to poll now (202) + POST /api/review -> run a Fireworks review for one PR on demand POST /api/revise -> re-run review with the user's revision notes POST /api/post -> submit the review to GitHub (side-effectful) -In-memory store only; state resets on restart. This is a single-user local tool. +A background worker (app/worker.py) pre-reviews non-draft PRs into a SQLite cache +(app/store.py) so reviews are ready when opened. Nothing is ever posted +automatically. Single-user local tool. """ from __future__ import annotations +import hashlib +import json +from contextlib import asynccontextmanager from pathlib import Path from typing import Any from fastapi import FastAPI, HTTPException -from fastapi.responses import FileResponse +from fastapi.responses import FileResponse, JSONResponse from fastapi.staticfiles import StaticFiles from pydantic import BaseModel from .config import get_config from .github_client import GitHubClient, GitHubError from .reviewer import Reviewer +from .store import ReviewStore +from .worker import Worker cfg = get_config() -app = FastAPI(title="PR Review Dashboard") STATIC_DIR = Path(__file__).parent.parent / "static" -app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static") # lazy singletons so a missing token doesn't crash import _gh: GitHubClient | None = None _reviewer: Reviewer | None = None +_store: ReviewStore | None = None +_worker: Worker | None = None def gh() -> GitHubClient: @@ -49,6 +59,35 @@ def reviewer() -> Reviewer: return _reviewer +def store() -> ReviewStore | None: + return _store + + +def worker() -> Worker | None: + return _worker + + +@asynccontextmanager +async def lifespan(app: FastAPI): + # Initialize the singletons eagerly, before the worker thread starts, so the + # worker and request handlers never race to lazily create them. + global _store, _worker + gh() + reviewer() + _store = ReviewStore(cfg.CACHE_DB) + _worker = Worker(cfg, _gh, _reviewer, _store) + _worker.start() + try: + yield + finally: + if _worker is not None: + _worker.stop() + + +app = FastAPI(title="PR Review Dashboard", lifespan=lifespan) +app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static") + + @app.get("/") def index() -> FileResponse: return FileResponse(STATIC_DIR / "index.html") @@ -73,6 +112,50 @@ def api_prs() -> dict[str, Any]: raise HTTPException(400, str(e)) +def _serialize_review(row: dict[str, Any]) -> dict[str, Any]: + review = None + if row.get("review_json"): + try: + review = json.loads(row["review_json"]) + except (json.JSONDecodeError, TypeError): + review = None + return { + "owner": row["owner"], + "repo": row["repo"], + "number": row["number"], + "title": row.get("pr_title"), + "author": row.get("pr_author"), + "url": row.get("pr_url"), + "status": row["status"], + "updated_at": row.get("pr_updated_at"), + "cached_at": row.get("cached_at"), + "attempts": row.get("attempts", 0), + "error": row.get("error"), + "review": review, + } + + +@app.get("/api/reviews") +def api_reviews() -> dict[str, Any]: + s = store() + if s is None: + raise HTTPException(503, "Review worker is not running yet.") + rows = s.list() + return { + "reviews": [_serialize_review(r) for r in rows], + "poll_interval": cfg.POLL_INTERVAL, + } + + +@app.post("/api/refresh") +def api_refresh() -> JSONResponse: + w = worker() + if w is None: + raise HTTPException(503, "Review worker is not running yet.") + w.trigger() # non-blocking: the worker thread runs the cycle + return JSONResponse(status_code=202, content={"ok": True, "scheduled": True}) + + class ReviewReq(BaseModel): owner: str repo: str @@ -82,12 +165,35 @@ class ReviewReq(BaseModel): body: str = "" +def _cache_write_through(req: ReviewReq, diff: str, result: dict[str, Any]) -> None: + """Best-effort: store an on-demand review in the cache so the dashboard and + the worker share one source of truth. Never fails the request.""" + if _store is None: + return + try: + _store.upsert( + req.owner, + req.repo, + req.number, + status="ready", + diff_hash=hashlib.sha256(diff.encode("utf-8", "ignore")).hexdigest(), + review_json=json.dumps(result), + error=None, + attempts=0, + pr_title=req.title, + pr_author=req.author, + ) + except Exception: # noqa: BLE001 - caching is best-effort + pass + + @app.post("/api/review") def api_review(req: ReviewReq) -> dict[str, Any]: try: diff = gh().pr_diff(req.owner, req.repo, req.number) pr = req.model_dump() result = reviewer().review(pr, diff) + _cache_write_through(req, diff, result) return {"review": result} except GitHubError as e: raise HTTPException(400, str(e)) diff --git a/app/store.py b/app/store.py new file mode 100644 index 0000000..63a3b04 --- /dev/null +++ b/app/store.py @@ -0,0 +1,157 @@ +"""SQLite cache of pre-computed PR reviews. + +One row per PR, keyed by ``(owner, repo, number)``. The background worker writes; +the API reads. SQLite allows only a single writer at a time, so every write goes +through one module-level lock and connections use a busy timeout; readers open +their own short-lived connection. The DB file is local and gitignored. + +Status lifecycle: + reviewing -> ready | error (worker reviews a new/changed PR) + ready -> ready (unchanged; refreshed metadata only) + * -> closed (PR left the queue; kept briefly, then purged) +""" + +from __future__ import annotations + +import sqlite3 +import threading +import time +from typing import Any + +# One writer at a time. SQLite serializes writes anyway; this keeps the Python +# side from piling up concurrent write transactions from the worker pool. +_WRITE_LOCK = threading.Lock() + +_SCHEMA = """ +CREATE TABLE IF NOT EXISTS reviews ( + owner TEXT NOT NULL, + repo TEXT NOT NULL, + number INTEGER NOT NULL, + status TEXT NOT NULL, -- reviewing | ready | error | closed + diff_hash TEXT, + review_json TEXT, + error TEXT, + attempts INTEGER NOT NULL DEFAULT 0, + pr_title TEXT, + pr_author TEXT, + pr_url TEXT, + pr_updated_at TEXT, -- PR's updatedAt on GitHub + cached_at REAL NOT NULL, -- when this row was last written + closed_at REAL, -- when the PR left the queue + PRIMARY KEY (owner, repo, number) +); +""" + +# Columns a caller may set via upsert(); owner/repo/number are the key and +# cached_at is stamped automatically. +_UPSERTABLE = ( + "status", + "diff_hash", + "review_json", + "error", + "attempts", + "pr_title", + "pr_author", + "pr_url", + "pr_updated_at", + "closed_at", +) + + +class ReviewStore: + def __init__(self, db_path: str) -> None: + self.db_path = db_path + with _WRITE_LOCK, self._connect() as conn: + conn.executescript(_SCHEMA) + + def _connect(self) -> sqlite3.Connection: + conn = sqlite3.connect(self.db_path, timeout=30) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA journal_mode=WAL") + conn.execute("PRAGMA busy_timeout=30000") + return conn + + # --- reads -------------------------------------------------------------- + def get(self, owner: str, repo: str, number: int) -> dict[str, Any] | None: + with self._connect() as conn: + row = conn.execute( + "SELECT * FROM reviews WHERE owner=? AND repo=? AND number=?", + (owner, repo, number), + ).fetchone() + return dict(row) if row else None + + def list(self, *, include_closed: bool = True) -> list[dict[str, Any]]: + sql = "SELECT * FROM reviews" + if not include_closed: + sql += " WHERE status != 'closed'" + sql += " ORDER BY pr_updated_at DESC" + with self._connect() as conn: + rows = conn.execute(sql).fetchall() + return [dict(r) for r in rows] + + # --- writes (serialized) ------------------------------------------------ + def upsert(self, owner: str, repo: str, number: int, **fields: Any) -> None: + """Insert or update one PR row. Only recognized columns are written; + ``cached_at`` is stamped automatically. Unspecified columns keep their + previous value on update (and their default on insert).""" + bad = set(fields) - set(_UPSERTABLE) + if bad: + raise ValueError(f"unknown review columns: {sorted(bad)}") + + cols = ["owner", "repo", "number", "cached_at", *fields] + vals = [owner, repo, number, time.time(), *fields.values()] + placeholders = ", ".join("?" for _ in cols) + # On conflict, update every provided column plus cached_at. + updates = ", ".join(f"{c}=excluded.{c}" for c in ("cached_at", *fields)) + sql = ( + f"INSERT INTO reviews ({', '.join(cols)}) VALUES ({placeholders}) " + f"ON CONFLICT(owner, repo, number) DO UPDATE SET {updates}" + ) + with _WRITE_LOCK, self._connect() as conn: + conn.execute(sql, vals) + conn.commit() + + def reset_stale_reviewing(self) -> int: + """On startup, any row still marked ``reviewing`` was interrupted by a + crash/restart. Flip it to ``error`` so the next cycle retries it.""" + with _WRITE_LOCK, self._connect() as conn: + cur = conn.execute( + "UPDATE reviews SET status='error', error='interrupted' " + "WHERE status='reviewing'" + ) + conn.commit() + return cur.rowcount + + def mark_missing_closed(self, live_keys: set[tuple[str, str, int]]) -> None: + """Mark rows whose PR is no longer in the queue as ``closed`` (keeping + them briefly so an open dashboard doesn't lose a review mid-read).""" + now = time.time() + with _WRITE_LOCK, self._connect() as conn: + rows = conn.execute( + "SELECT owner, repo, number FROM reviews WHERE status != 'closed'" + ).fetchall() + gone = [ + (r["owner"], r["repo"], r["number"]) + for r in rows + if (r["owner"], r["repo"], r["number"]) not in live_keys + ] + for owner, repo, number in gone: + conn.execute( + "UPDATE reviews SET status='closed', closed_at=? " + "WHERE owner=? AND repo=? AND number=?", + (now, owner, repo, number), + ) + conn.commit() + + def purge_closed(self, older_than_seconds: float) -> int: + """Hard-delete rows that have been ``closed`` longer than the grace + window.""" + cutoff = time.time() - older_than_seconds + with _WRITE_LOCK, self._connect() as conn: + cur = conn.execute( + "DELETE FROM reviews WHERE status='closed' AND closed_at IS NOT NULL " + "AND closed_at < ?", + (cutoff,), + ) + conn.commit() + return cur.rowcount diff --git a/app/worker.py b/app/worker.py new file mode 100644 index 0000000..04ac66c --- /dev/null +++ b/app/worker.py @@ -0,0 +1,232 @@ +"""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"), + } + 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, +) -> dict[str, int]: + """One poll cycle: refresh the queue, close/purge departed PRs, and review + every new or changed non-draft PR (bounded concurrency). Returns per-outcome + counts.""" + 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: + continue # cheap gate: nothing changed since last ready review + 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"), + ) + + 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) -> None: + self.cfg = cfg + self.gh = gh + self.reviewer = reviewer + self.store = store + 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, + ) + _log.info("review cycle: %s", res) + except Exception: + _log.exception("review cycle failed") + self._wake.wait(timeout=self.cfg.POLL_INTERVAL) + self._wake.clear() diff --git a/static/index.html b/static/index.html index 9f88ae9..278982b 100644 --- a/static/index.html +++ b/static/index.html @@ -54,8 +54,12 @@ .queue-item .title { font-weight: 500; margin: 2px 0; } .queue-item .sub { font-size: 12px; color: var(--ink-dim); display: flex; gap: 8px; align-items: center; } .badge { font-size: 10px; font-family: var(--mono); padding: 1px 6px; border-radius: 4px; border: 1px solid var(--line); } - .badge.reviewed { color: var(--nit); border-color: var(--nit); } + .badge.reviewed, .badge.ready { color: var(--nit); border-color: var(--nit); } + .badge.reviewing { color: var(--fix); border-color: var(--fix); } + .badge.error { color: var(--block); border-color: var(--block); } + .badge.closed { color: var(--ink-dim); border-color: var(--line); } .badge.mention { color: var(--accent); border-color: var(--accent); } + .queue-item.closed { opacity: .55; } .detail { overflow-y: auto; padding: 24px 28px; } .detail h2 { font-size: 17px; margin: 0 0 4px; } .detail .prlink { font-family: var(--mono); font-size: 12px; color: var(--accent); text-decoration: none; } @@ -102,18 +106,19 @@

PR Review Desk

loading config... - + +
-
Click Refresh to load PRs.
-
Select a PR to review.
+
Loading queue...
+
Select a PR. Reviews are prepared automatically in the background.
diff --git a/tests/test_endpoints.py b/tests/test_endpoints.py index 702782c..8378ab5 100644 --- a/tests/test_endpoints.py +++ b/tests/test_endpoints.py @@ -7,6 +7,7 @@ already-installed httpx2, so no extra test deps. from __future__ import annotations +import json from typing import Any import pytest @@ -14,6 +15,7 @@ from fastapi.testclient import TestClient from app import main as main_mod from app.github_client import GitHubError +from app.store import ReviewStore @pytest.fixture @@ -203,3 +205,71 @@ def test_api_post_github_error_maps_to_400(client: TestClient, monkeypatch) -> N ) assert r.status_code == 400 assert "Invalid event" in r.json()["detail"] + + +# --------------------------------------------------------------------------- # +# /api/reviews and /api/refresh (background-worker cache) +# --------------------------------------------------------------------------- # + + +def test_api_reviews_returns_cached(client: TestClient, monkeypatch, tmp_path) -> None: + s = ReviewStore(str(tmp_path / "cache.db")) + s.upsert( + "o", + "r", + 1, + status="ready", + review_json=json.dumps({"summary": "hi"}), + pr_title="t", + pr_author="a", + pr_url="u", + pr_updated_at="2026-01-01", + ) + monkeypatch.setattr(main_mod, "_store", s) + r = client.get("/api/reviews") + assert r.status_code == 200 + data = r.json() + assert data["poll_interval"] > 0 + item = data["reviews"][0] + assert item["status"] == "ready" + assert item["review"]["summary"] == "hi" + + +def test_api_reviews_503_without_store(client: TestClient, monkeypatch) -> None: + monkeypatch.setattr(main_mod, "_store", None) + assert client.get("/api/reviews").status_code == 503 + + +def test_api_refresh_triggers_worker(client: TestClient, monkeypatch) -> None: + class FakeWorker: + def __init__(self) -> None: + self.triggered = False + + def trigger(self) -> None: + self.triggered = True + + w = FakeWorker() + monkeypatch.setattr(main_mod, "_worker", w) + r = client.post("/api/refresh") + assert r.status_code == 202 + assert r.json()["scheduled"] is True + assert w.triggered is True + + +def test_api_refresh_503_without_worker(client: TestClient, monkeypatch) -> None: + monkeypatch.setattr(main_mod, "_worker", None) + assert client.post("/api/refresh").status_code == 503 + + +def test_api_review_writes_through_to_cache( + client: TestClient, monkeypatch, tmp_path +) -> None: + s = ReviewStore(str(tmp_path / "cache.db")) + monkeypatch.setattr(main_mod, "_store", s) + _patch(monkeypatch, gh=FakeGH(), reviewer=FakeReviewer()) + r = client.post("/api/review", json=_review_body()) + assert r.status_code == 200 + row = s.get("o", "r", 1) + assert row is not None + assert row["status"] == "ready" + assert json.loads(row["review_json"])["summary"] == "ok" diff --git a/tests/test_store.py b/tests/test_store.py new file mode 100644 index 0000000..4734d01 --- /dev/null +++ b/tests/test_store.py @@ -0,0 +1,93 @@ +"""Tests for the SQLite review cache (stdlib sqlite3, temp DB per test).""" + +from __future__ import annotations + +import json + +import pytest + +from app.store import ReviewStore + + +@pytest.fixture +def store(tmp_path) -> ReviewStore: + return ReviewStore(str(tmp_path / "cache.db")) + + +def test_upsert_insert_then_get(store: ReviewStore) -> None: + store.upsert( + "o", + "r", + 1, + status="ready", + diff_hash="abc", + review_json=json.dumps({"summary": "hi"}), + pr_title="Title", + pr_author="octocat", + pr_url="https://x/1", + pr_updated_at="2026-01-01", + ) + row = store.get("o", "r", 1) + assert row is not None + assert row["status"] == "ready" + assert row["diff_hash"] == "abc" + assert row["pr_title"] == "Title" + assert json.loads(row["review_json"])["summary"] == "hi" + assert row["cached_at"] > 0 + + +def test_get_missing_returns_none(store: ReviewStore) -> None: + assert store.get("o", "r", 999) is None + + +def test_update_preserves_unspecified_columns(store: ReviewStore) -> None: + store.upsert("o", "r", 1, status="ready", diff_hash="h1", review_json="{}") + # Update only status; diff_hash and review_json must survive. + store.upsert("o", "r", 1, status="reviewing") + row = store.get("o", "r", 1) + assert row["status"] == "reviewing" + assert row["diff_hash"] == "h1" + assert row["review_json"] == "{}" + + +def test_unknown_column_raises(store: ReviewStore) -> None: + with pytest.raises(ValueError, match="unknown review columns"): + store.upsert("o", "r", 1, status="ready", bogus="x") + + +def test_list_orders_and_filters_closed(store: ReviewStore) -> None: + store.upsert("o", "r", 1, status="ready", pr_updated_at="2026-01-01") + store.upsert( + "o", "r", 2, status="closed", pr_updated_at="2026-02-02", closed_at=1.0 + ) + all_rows = store.list() + assert [r["number"] for r in all_rows] == [2, 1] # updated_at desc + open_rows = store.list(include_closed=False) + assert [r["number"] for r in open_rows] == [1] + + +def test_reset_stale_reviewing(store: ReviewStore) -> None: + store.upsert("o", "r", 1, status="reviewing") + store.upsert("o", "r", 2, status="ready") + n = store.reset_stale_reviewing() + assert n == 1 + assert store.get("o", "r", 1)["status"] == "error" + assert store.get("o", "r", 1)["error"] == "interrupted" + assert store.get("o", "r", 2)["status"] == "ready" + + +def test_mark_missing_closed_and_purge(store: ReviewStore) -> None: + store.upsert("o", "r", 1, status="ready") + store.upsert("o", "r", 2, status="ready") + # Only PR 1 is still live; PR 2 should be closed. + store.mark_missing_closed({("o", "r", 1)}) + assert store.get("o", "r", 1)["status"] == "ready" + closed = store.get("o", "r", 2) + assert closed["status"] == "closed" + assert closed["closed_at"] is not None + + # A long grace keeps it; a zero grace purges it. + assert store.purge_closed(older_than_seconds=10_000) == 0 + assert store.get("o", "r", 2) is not None + assert store.purge_closed(older_than_seconds=0) == 1 + assert store.get("o", "r", 2) is None diff --git a/tests/test_worker.py b/tests/test_worker.py new file mode 100644 index 0000000..4ae0596 --- /dev/null +++ b/tests/test_worker.py @@ -0,0 +1,180 @@ +"""Tests for the auto-review worker's run_cycle and backoff (no network). + +Uses a real ReviewStore on a temp DB plus fake GitHub/reviewer objects, and +drives run_cycle directly (its deps are injected, so no thread/timer needed). +""" + +from __future__ import annotations + +import json +from types import SimpleNamespace + +import pytest + +from app.store import ReviewStore +from app.worker import _review_with_backoff, run_cycle + + +class FakeGH: + def __init__(self, prs: list[dict], diffs: dict[int, str]) -> None: + self.prs = prs + self.diffs = diffs + + def search_prs(self) -> list[dict]: + return [dict(p) for p in self.prs] # copy: mimic a fresh fetch + + def pr_diff(self, owner: str, repo: str, number: int) -> str: + return self.diffs[number] + + +class CountingReviewer: + def __init__(self) -> None: + self.calls = 0 + + def review(self, pr: dict, diff: str) -> dict: + self.calls += 1 + return {"summary": f"review {pr['number']}", "recommended_event": "COMMENT"} + + +def _pr(number: int, *, updated_at: str, draft: bool = False) -> dict: + return { + "owner": "o", + "repo": "r", + "number": number, + "title": f"PR {number}", + "author": "octocat", + "url": f"https://x/{number}", + "updated_at": updated_at, + "draft": draft, + } + + +@pytest.fixture +def store(tmp_path) -> ReviewStore: + return ReviewStore(str(tmp_path / "cache.db")) + + +def test_reviews_nondraft_skips_draft(store: ReviewStore) -> None: + gh = FakeGH( + [_pr(1, updated_at="d1"), _pr(2, updated_at="d1", draft=True)], + {1: "diff-1", 2: "diff-2"}, + ) + rev = CountingReviewer() + res = run_cycle(gh, rev, store, concurrency=1) + + assert res["reviewed"] == 1 + assert rev.calls == 1 + row = store.get("o", "r", 1) + assert row["status"] == "ready" + assert json.loads(row["review_json"])["summary"] == "review 1" + assert store.get("o", "r", 2) is None # draft never touched + + +def test_cheap_gate_skips_unchanged_updated_at(store: ReviewStore) -> None: + gh = FakeGH([_pr(1, updated_at="v1")], {1: "diff-1"}) + rev = CountingReviewer() + run_cycle(gh, rev, store, concurrency=1) + res2 = run_cycle(gh, rev, store, concurrency=1) + + assert res2["reviewed"] == 0 + assert res2["skipped"] == 1 + assert rev.calls == 1 # not re-reviewed, and no diff even compared + + +def test_precise_gate_skips_when_diff_unchanged(store: ReviewStore) -> None: + gh = FakeGH([_pr(1, updated_at="v1")], {1: "diff-1"}) + rev = CountingReviewer() + run_cycle(gh, rev, store, concurrency=1) + + # updated_at bumped (e.g. a comment) but the diff is identical. + gh.prs[0]["updated_at"] = "v2" + res = run_cycle(gh, rev, store, concurrency=1) + + assert res["unchanged"] == 1 + assert rev.calls == 1 # LLM not called again + assert store.get("o", "r", 1)["pr_updated_at"] == "v2" # metadata refreshed + + +def test_changed_diff_triggers_rereview(store: ReviewStore) -> None: + gh = FakeGH([_pr(1, updated_at="v1")], {1: "diff-1"}) + rev = CountingReviewer() + run_cycle(gh, rev, store, concurrency=1) + + gh.prs[0]["updated_at"] = "v2" + gh.diffs[1] = "diff-1-CHANGED" + res = run_cycle(gh, rev, store, concurrency=1) + + assert res["reviewed"] == 1 + assert rev.calls == 2 + + +def test_error_stored_and_attempts_capped(store: ReviewStore) -> None: + class BoomReviewer: + def review(self, pr, diff): + raise ValueError("fireworks exploded") + + gh = FakeGH([_pr(1, updated_at="v1")], {1: "diff-1"}) + rev = BoomReviewer() + + r1 = run_cycle(gh, rev, store, concurrency=1, max_attempts=2) + assert r1["error"] == 1 + assert store.get("o", "r", 1)["status"] == "error" + assert store.get("o", "r", 1)["attempts"] == 1 + + run_cycle(gh, rev, store, concurrency=1, max_attempts=2) + assert store.get("o", "r", 1)["attempts"] == 2 + + # attempts now >= max and updated_at unchanged -> stop retrying + r3 = run_cycle(gh, rev, store, concurrency=1, max_attempts=2) + assert r3["skipped"] == 1 + assert r3["error"] == 0 + assert store.get("o", "r", 1)["attempts"] == 2 + + +def test_departed_pr_is_closed(store: ReviewStore) -> None: + gh = FakeGH([_pr(1, updated_at="v1")], {1: "diff-1"}) + rev = CountingReviewer() + run_cycle(gh, rev, store, concurrency=1) + + gh.prs = [] # PR merged/closed, no longer in the queue + run_cycle(gh, rev, store, concurrency=1) + assert store.get("o", "r", 1)["status"] == "closed" + + +# --- backoff helper --------------------------------------------------------- # + + +def _http_error(status: int) -> Exception: + e = Exception(f"HTTP {status}") + e.response = SimpleNamespace(status_code=status) + return e + + +def test_backoff_retries_on_429_then_succeeds() -> None: + class Flaky: + def __init__(self) -> None: + self.n = 0 + + def review(self, pr, diff): + self.n += 1 + if self.n < 3: + raise _http_error(429) + return {"summary": "ok"} + + sleeps: list[float] = [] + out = _review_with_backoff( + Flaky(), {}, "d", retries=3, base=0.01, sleep=sleeps.append + ) + assert out["summary"] == "ok" + assert len(sleeps) == 2 # two 429s -> two backoffs + + +def test_backoff_does_not_retry_non_429() -> None: + class Boom: + def review(self, pr, diff): + raise ValueError("nope") + + sleeps: list[float] = [] + with pytest.raises(ValueError): + _review_with_backoff(Boom(), {}, "d", retries=3, sleep=sleeps.append) + assert sleeps == [] # non-429 fails fast