mirror of
https://github.com/Sea-Haven-Industries/pr-reviewer.git
synced 2026-09-30 03:23:14 +00:00
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.
This commit is contained in:
parent
46cd856fb3
commit
15a52bff40
11 changed files with 996 additions and 51 deletions
12
.env.example
12
.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
|
||||
|
|
|
|||
4
.gitignore
vendored
4
.gitignore
vendored
|
|
@ -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__/
|
||||
|
|
|
|||
31
README.md
31
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
|
||||
```
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
118
app/main.py
118
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))
|
||||
|
|
|
|||
157
app/store.py
Normal file
157
app/store.py
Normal file
|
|
@ -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
|
||||
232
app/worker.py
Normal file
232
app/worker.py
Normal file
|
|
@ -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()
|
||||
|
|
@ -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 @@
|
|||
<h1>PR Review Desk</h1>
|
||||
<span class="meta" id="cfgMeta">loading config...</span>
|
||||
<span class="spacer"></span>
|
||||
<button id="refreshBtn">Refresh queue</button>
|
||||
<span class="meta" id="syncMeta"></span>
|
||||
<button id="refreshBtn">Refresh now</button>
|
||||
</header>
|
||||
|
||||
<div class="layout">
|
||||
<div class="queue" id="queue"><div class="empty">Click Refresh to load PRs.</div></div>
|
||||
<div class="detail" id="detail"><div class="empty">Select a PR to review.</div></div>
|
||||
<div class="queue" id="queue"><div class="empty spin">Loading queue...</div></div>
|
||||
<div class="detail" id="detail"><div class="empty">Select a PR. Reviews are prepared automatically in the background.</div></div>
|
||||
</div>
|
||||
|
||||
<div id="toastHost"></div>
|
||||
|
||||
<script>
|
||||
const state = { prs: [], me: "", active: null, reviews: {} };
|
||||
const state = { items: [], byKey: {}, me: "", active: null, mentions: [], shownKey: null, pollMs: 8000 };
|
||||
|
||||
function key(pr){ return `${pr.owner}/${pr.repo}#${pr.number}`; }
|
||||
|
||||
|
|
@ -139,63 +144,108 @@ async function loadConfig(){
|
|||
state.mentions = c.mention_authors;
|
||||
}
|
||||
|
||||
async function refresh(){
|
||||
const q = document.getElementById("queue");
|
||||
q.innerHTML = '<div class="empty spin">Loading PRs...</div>';
|
||||
try {
|
||||
const data = await j("/api/prs");
|
||||
state.prs = data.prs; state.me = data.me;
|
||||
renderQueue();
|
||||
} catch(e){ q.innerHTML = `<div class="empty">Failed: ${e.message}</div>`; }
|
||||
// The reviewer identity (for the post confirmation) comes from /api/prs, which
|
||||
// also returns the authenticated login. Best-effort, non-blocking.
|
||||
async function loadMe(){
|
||||
try { const data = await j("/api/prs"); state.me = data.me || ""; } catch(e){ /* ignore */ }
|
||||
}
|
||||
|
||||
// Force an immediate background poll cycle, then refresh the view shortly after.
|
||||
async function refreshNow(){
|
||||
try { await fetch("/api/refresh", {method:"POST"}); toast("Refresh scheduled."); }
|
||||
catch(e){ toast("Could not schedule refresh: "+e.message, true); }
|
||||
setTimeout(poll, 1500);
|
||||
}
|
||||
|
||||
// Poll the cached reviews and repaint the queue. Runs on a timer.
|
||||
async function poll(){
|
||||
let data;
|
||||
try { data = await j("/api/reviews"); }
|
||||
catch(e){
|
||||
document.getElementById("syncMeta").textContent = "sync failed";
|
||||
return;
|
||||
}
|
||||
state.items = data.reviews || [];
|
||||
state.byKey = {};
|
||||
state.items.forEach(it => { state.byKey[key(it)] = it; });
|
||||
if(data.poll_interval){ document.getElementById("cfgMeta").dataset.poll = data.poll_interval; }
|
||||
const now = new Date();
|
||||
document.getElementById("syncMeta").textContent =
|
||||
"synced " + now.toLocaleTimeString();
|
||||
renderQueue();
|
||||
// Upgrade a placeholder to the full review once it becomes ready.
|
||||
if(state.active){
|
||||
const it = state.byKey[state.active];
|
||||
if(it && it.status === "ready" && it.review && state.shownKey !== state.active){
|
||||
renderReview(it, it.review);
|
||||
} else if(it && it.status !== "ready" && state.shownKey === state.active){
|
||||
// it went back to reviewing/error after being shown
|
||||
state.shownKey = null; selectPR(state.active);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const STATUS_LABEL = { reviewing: "reviewing", ready: "ready", error: "error", closed: "closed" };
|
||||
|
||||
function renderQueue(){
|
||||
const q = document.getElementById("queue");
|
||||
if(!state.prs.length){ q.innerHTML = '<div class="empty">No PRs match the filter.</div>'; return; }
|
||||
if(!state.items.length){ q.innerHTML = '<div class="empty">No PRs in the queue yet. The worker reviews new PRs automatically.</div>'; return; }
|
||||
q.innerHTML = "";
|
||||
state.prs.forEach(pr => {
|
||||
const k = key(pr);
|
||||
state.items.forEach(it => {
|
||||
const k = key(it);
|
||||
const div = document.createElement("div");
|
||||
div.className = "queue-item" + (state.active===k ? " active":"");
|
||||
const mention = (state.mentions||[]).includes((pr.author||"").toLowerCase());
|
||||
div.className = "queue-item" + (state.active===k ? " active":"") + (it.status==="closed" ? " closed":"");
|
||||
const mention = (state.mentions||[]).includes((it.author||"").toLowerCase());
|
||||
const statusBadge = `<span class="badge ${it.status}">${STATUS_LABEL[it.status]||it.status}</span>`;
|
||||
div.innerHTML = `
|
||||
<div class="repo">${pr.owner}/${pr.repo} #${pr.number}</div>
|
||||
<div class="title">${escapeHtml(pr.title)}</div>
|
||||
<div class="repo">${it.owner}/${it.repo} #${it.number}</div>
|
||||
<div class="title">${escapeHtml(it.title||"(no title)")}</div>
|
||||
<div class="sub">
|
||||
<span>@${escapeHtml(pr.author)}</span>
|
||||
<span>@${escapeHtml(it.author||"")}</span>
|
||||
${mention ? '<span class="badge mention">will @mention</span>':''}
|
||||
${state.reviews[k] ? '<span class="badge reviewed">reviewed</span>':''}
|
||||
${statusBadge}
|
||||
</div>`;
|
||||
div.onclick = ()=>selectPR(k);
|
||||
q.appendChild(div);
|
||||
});
|
||||
}
|
||||
|
||||
function findPR(k){ return state.prs.find(p=>key(p)===k); }
|
||||
|
||||
async function selectPR(k){
|
||||
function selectPR(k){
|
||||
state.active = k; renderQueue();
|
||||
const pr = findPR(k);
|
||||
const it = state.byKey[k];
|
||||
if(!it) return;
|
||||
const d = document.getElementById("detail");
|
||||
if(state.reviews[k]){ renderReview(pr, state.reviews[k]); return; }
|
||||
if(it.status === "ready" && it.review){ renderReview(it, it.review); return; }
|
||||
state.shownKey = null;
|
||||
let inner;
|
||||
if(it.status === "reviewing"){
|
||||
inner = '<span class="spin">Reviewing with Fireworks in the background...</span>';
|
||||
} else if(it.status === "error"){
|
||||
inner = `<div class="cat block"><h3>REVIEW FAILED</h3><div>${escapeHtml(it.error||"unknown error")} (attempt ${it.attempts||0})</div></div>
|
||||
<div class="toolbar"><button class="primary" id="runBtn">Retry now</button></div>`;
|
||||
} else if(it.status === "closed"){
|
||||
inner = '<div class="spin">This PR left the queue (merged or closed).</div>';
|
||||
} else {
|
||||
inner = '<div class="toolbar"><button class="primary" id="runBtn">Run review now</button></div>';
|
||||
}
|
||||
d.innerHTML = `
|
||||
<h2>${escapeHtml(pr.title)}</h2>
|
||||
<a class="prlink" href="${pr.url}" target="_blank">${pr.owner}/${pr.repo} #${pr.number}</a>
|
||||
<div class="toolbar"><button class="primary" id="runBtn">Run review</button></div>`;
|
||||
document.getElementById("runBtn").onclick = ()=>runReview(pr);
|
||||
<h2>${escapeHtml(it.title||"(no title)")}</h2>
|
||||
<a class="prlink" href="${it.url||"#"}" target="_blank">${it.owner}/${it.repo} #${it.number}</a>
|
||||
<div style="margin-top:16px">${inner}</div>`;
|
||||
const runBtn = document.getElementById("runBtn");
|
||||
if(runBtn) runBtn.onclick = ()=>runReview(it);
|
||||
}
|
||||
|
||||
async function runReview(pr){
|
||||
const d = document.getElementById("detail");
|
||||
d.querySelector(".toolbar").innerHTML = '<span class="spin">Reviewing with Fireworks...</span>';
|
||||
d.querySelector(".toolbar") && (d.querySelector(".toolbar").innerHTML = '<span class="spin">Reviewing with Fireworks...</span>');
|
||||
try {
|
||||
const res = await j("/api/review", {
|
||||
method:"POST", headers:{"Content-Type":"application/json"},
|
||||
body: JSON.stringify(pr)
|
||||
body: JSON.stringify({owner:pr.owner, repo:pr.repo, number:pr.number, title:pr.title||"", author:pr.author||""})
|
||||
});
|
||||
state.reviews[key(pr)] = res.review;
|
||||
renderQueue();
|
||||
renderReview(pr, res.review);
|
||||
poll();
|
||||
} catch(e){
|
||||
toast("Review failed: "+e.message, true);
|
||||
selectPR(key(pr));
|
||||
|
|
@ -204,6 +254,7 @@ async function runReview(pr){
|
|||
|
||||
function renderReview(pr, rv){
|
||||
const d = document.getElementById("detail");
|
||||
state.shownKey = key(pr);
|
||||
const cats = [["block","BLOCK"],["fix","FIX"],["nit","NIT"],["question","QUESTION"]];
|
||||
let catHtml = "";
|
||||
cats.forEach(([k,label])=>{
|
||||
|
|
@ -251,9 +302,8 @@ async function revise(pr){
|
|||
try {
|
||||
const res = await j("/api/revise", {
|
||||
method:"POST", headers:{"Content-Type":"application/json"},
|
||||
body: JSON.stringify({...pr, notes})
|
||||
body: JSON.stringify({owner:pr.owner, repo:pr.repo, number:pr.number, title:pr.title||"", author:pr.author||"", notes})
|
||||
});
|
||||
state.reviews[key(pr)] = res.review;
|
||||
renderReview(pr, res.review);
|
||||
toast("Review revised.");
|
||||
} catch(e){ toast("Revision failed: "+e.message, true); btn.disabled=false; btn.textContent="Request revision"; }
|
||||
|
|
@ -279,8 +329,11 @@ async function postReview(pr){
|
|||
function escapeHtml(s){ return (s||"").replace(/[&<>"']/g, c=>({"&":"&","<":"<",">":">",'"':""","'":"'"}[c])); }
|
||||
function mdInline(s){ return escapeHtml(s).replace(/`([^`]+)`/g, '<code>$1</code>'); }
|
||||
|
||||
document.getElementById("refreshBtn").onclick = refresh;
|
||||
document.getElementById("refreshBtn").onclick = refreshNow;
|
||||
loadConfig();
|
||||
loadMe();
|
||||
poll();
|
||||
setInterval(poll, state.pollMs);
|
||||
</script>
|
||||
</body>
|
||||
</html>
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
93
tests/test_store.py
Normal file
93
tests/test_store.py
Normal file
|
|
@ -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
|
||||
180
tests/test_worker.py
Normal file
180
tests/test_worker.py
Normal file
|
|
@ -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
|
||||
Loading…
Add table
Reference in a new issue