From e426186e4690812b5c91f88a89f1df8ad00b1b00 Mon Sep 17 00:00:00 2001 From: can1357 Date: Wed, 15 Jul 2026 00:46:08 +0200 Subject: [PATCH] feat(robomp): added local issue indexing and commit-search tool support - Added `ROBOMP_ISSUE_INDEX_SYNC_SECONDS` configuration and lifecycle-managed issue indexing through startup/shutdown hooks. - Added issue/PR index tables with FTS5 triggers plus upsert and keyword/filter search helpers for indexed records. - Added GitHub backend/proxy support for `IssueIndexEntry` and issue-index page retrieval, including webhook ingestion and periodic watermark-driven sync. - Updated `gh_search_issues` to prefer local index queries when synchronized and added `search_commits` host tool with query modes and validation. --- python/robomp/.env.example | 10 + python/robomp/src/config.py | 6 + python/robomp/src/db.py | 179 +++++++++++++++ python/robomp/src/github_backend.py | 9 + python/robomp/src/github_client.py | 106 +++++++++ python/robomp/src/host_tools.py | 187 +++++++++++++-- python/robomp/src/issue_index.py | 251 +++++++++++++++++++++ python/robomp/src/prompts/host_tools.toml | 13 +- python/robomp/src/prompts/system_append.md | 4 +- python/robomp/src/proxy/server.py | 16 ++ python/robomp/src/proxy_client.py | 36 +++ python/robomp/src/server.py | 17 +- python/robomp/tests/conftest.py | 4 + python/robomp/tests/test_host_tools.py | 95 +++++++- python/robomp/tests/test_issue_index.py | 189 ++++++++++++++++ python/robomp/tests/test_server.py | 36 +++ 16 files changed, 1129 insertions(+), 29 deletions(-) create mode 100644 python/robomp/src/issue_index.py create mode 100644 python/robomp/tests/test_issue_index.py diff --git a/python/robomp/.env.example b/python/robomp/.env.example index 692206eec..a338cdf82 100644 --- a/python/robomp/.env.example +++ b/python/robomp/.env.example @@ -164,6 +164,16 @@ ROBOMP_QUESTION_AUTOCLOSE_HOURS=4 # multi-hour close window. ROBOMP_QUESTION_AUTOCLOSE_SCAN_SECONDS=60 +# ============================================================================= +# --- Local issue search index --- +# ============================================================================= +# `gh_search_issues` answers from a local SQLite FTS mirror of every issue/PR +# in the allowlisted repos. Webhooks keep it fresh in real time; this interval +# controls the periodic GitHub reconcile (first pass backfills each repo). +# Set <= 0 to disable the reconciler — the tool then uses the GitHub search +# API until a repo has been backfilled. +ROBOMP_ISSUE_INDEX_SYNC_SECONDS=900 + # Path or command name for the omp binary inside the container. The shipped # image installs a shim that invokes Bun against the mounted pi checkout. ROBOMP_OMP_COMMAND=omp diff --git a/python/robomp/src/config.py b/python/robomp/src/config.py index 765507905..046e7c656 100644 --- a/python/robomp/src/config.py +++ b/python/robomp/src/config.py @@ -147,6 +147,12 @@ class Settings(BaseSettings): question_autoclose_enabled: bool = Field(True, alias="ROBOMP_QUESTION_AUTOCLOSE_ENABLED") question_autoclose_hours: float = Field(4.0, alias="ROBOMP_QUESTION_AUTOCLOSE_HOURS") question_autoclose_scan_seconds: float = Field(60.0, alias="ROBOMP_QUESTION_AUTOCLOSE_SCAN_SECONDS") + # Local issue/PR search index. Webhooks keep it fresh in real time; this + # interval controls the periodic GitHub reconcile (first pass = full + # backfill of every allowlisted repo). <= 0 disables the reconciler — + # `gh_search_issues` then falls back to the remote search API until the + # repo has a sync watermark. + issue_index_sync_seconds: float = Field(900.0, alias="ROBOMP_ISSUE_INDEX_SYNC_SECONDS") # pi-natives build-output cache. Hardlinks pre-built # `packages/natives/native/*.node` (and its companions) into new diff --git a/python/robomp/src/db.py b/python/robomp/src/db.py index ffbbf2d2c..7fa13940d 100644 --- a/python/robomp/src/db.py +++ b/python/robomp/src/db.py @@ -12,6 +12,8 @@ from datetime import UTC, datetime, timedelta from pathlib import Path from typing import Any, Literal +from robomp.github_client import IssueIndexEntry + EventState = Literal["queued", "running", "done", "failed", "skipped"] INACTIVE_EVENT_STATES: tuple[EventState, ...] = ("done", "failed", "skipped") @@ -113,6 +115,53 @@ CREATE TABLE IF NOT EXISTS pending_closures ( ); CREATE INDEX IF NOT EXISTS pending_closures_state_close_at ON pending_closures(state, close_at); + +-- Local mirror of every issue/PR in allowlisted repos, kept fresh by webhook +-- upserts plus the periodic `IssueIndexSync` reconciler. `gh_search_issues` +-- serves from here so triage lookups cost no GitHub API calls. +CREATE TABLE IF NOT EXISTS issue_index ( + repo TEXT NOT NULL, + number INTEGER NOT NULL, + is_pr INTEGER NOT NULL DEFAULT 0, + title TEXT NOT NULL DEFAULT '', + body TEXT NOT NULL DEFAULT '', + state TEXT NOT NULL DEFAULT 'open', + state_reason TEXT NOT NULL DEFAULT '', + merged_at TEXT NOT NULL DEFAULT '', + author TEXT NOT NULL DEFAULT '', + labels_json TEXT NOT NULL DEFAULT '[]', + comments INTEGER NOT NULL DEFAULT 0, + created_at TEXT NOT NULL DEFAULT '', + updated_at TEXT NOT NULL DEFAULT '', + html_url TEXT NOT NULL DEFAULT '', + PRIMARY KEY (repo, number) +); +CREATE INDEX IF NOT EXISTS issue_index_repo_updated + ON issue_index(repo, updated_at); + +CREATE VIRTUAL TABLE IF NOT EXISTS issue_index_fts USING fts5( + title, body, content='issue_index', content_rowid='rowid' +); + +CREATE TRIGGER IF NOT EXISTS issue_index_ai AFTER INSERT ON issue_index BEGIN + INSERT INTO issue_index_fts(rowid, title, body) VALUES (new.rowid, new.title, new.body); +END; +CREATE TRIGGER IF NOT EXISTS issue_index_ad AFTER DELETE ON issue_index BEGIN + INSERT INTO issue_index_fts(issue_index_fts, rowid, title, body) + VALUES ('delete', old.rowid, old.title, old.body); +END; +CREATE TRIGGER IF NOT EXISTS issue_index_au AFTER UPDATE ON issue_index BEGIN + INSERT INTO issue_index_fts(issue_index_fts, rowid, title, body) + VALUES ('delete', old.rowid, old.title, old.body); + INSERT INTO issue_index_fts(rowid, title, body) VALUES (new.rowid, new.title, new.body); +END; + +-- Per-repo reconcile watermark: the max `updated_at` the sync has fully +-- ingested. Absent row = repo never backfilled. +CREATE TABLE IF NOT EXISTS issue_index_sync ( + repo TEXT PRIMARY KEY, + last_synced TEXT NOT NULL +); """ @@ -1163,6 +1212,136 @@ class Database: ).fetchone() return _pending_closure_from_row(row) if row is not None else None + # ---- issue search index ---- + def upsert_issue_index(self, entry: IssueIndexEntry) -> None: + """Insert or refresh one issue/PR in the local search index.""" + with self._lock: + self._conn.execute( + """ + INSERT INTO issue_index + (repo, number, is_pr, title, body, state, state_reason, merged_at, + author, labels_json, comments, created_at, updated_at, html_url) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(repo, number) DO UPDATE SET + is_pr = excluded.is_pr, + title = excluded.title, + body = excluded.body, + state = excluded.state, + state_reason = excluded.state_reason, + merged_at = excluded.merged_at, + author = excluded.author, + labels_json = excluded.labels_json, + comments = excluded.comments, + created_at = excluded.created_at, + updated_at = excluded.updated_at, + html_url = excluded.html_url + """, + ( + entry.repo, + entry.number, + 1 if entry.is_pull_request else 0, + entry.title, + entry.body, + entry.state, + entry.state_reason, + entry.merged_at, + entry.author, + json.dumps(list(entry.labels), separators=(",", ":")), + entry.comments, + entry.created_at, + entry.updated_at, + entry.html_url, + ), + ) + + def search_issue_index( + self, + repo: str, + *, + keywords: Iterable[str] = (), + is_pr: bool | None = None, + state: str | None = None, + merged: bool | None = None, + label: str | None = None, + author: str | None = None, + limit: int = 10, + ) -> list[IssueIndexEntry]: + """Query the local index. Keywords go through FTS5 (bm25-ranked, AND + semantics); the remaining filters are exact. With no keywords, results + order by `updated_at` descending. + """ + conds = ["i.repo = ?"] + params: list[Any] = [repo] + if is_pr is not None: + conds.append("i.is_pr = ?") + params.append(1 if is_pr else 0) + if state is not None: + conds.append("i.state = ?") + params.append(state) + if merged is not None: + conds.append("i.merged_at != ''" if merged else "i.merged_at = ''") + if label is not None: + conds.append("EXISTS (SELECT 1 FROM json_each(i.labels_json) WHERE json_each.value = ?)") + params.append(label) + if author is not None: + conds.append("i.author = ?") + params.append(author) + terms = [t for t in keywords if t.strip()] + limit = max(1, min(int(limit), 50)) + with self._lock: + if terms: + # Quote every term so reporter text can never inject FTS5 syntax. + match = " ".join('"' + t.replace('"', '""') + '"' for t in terms) + sql = ( + "SELECT i.* FROM issue_index_fts f JOIN issue_index i ON i.rowid = f.rowid " + f"WHERE issue_index_fts MATCH ? AND {' AND '.join(conds)} " + "ORDER BY bm25(issue_index_fts) LIMIT ?" + ) + rows = self._conn.execute(sql, (match, *params, limit)).fetchall() + else: + sql = f"SELECT i.* FROM issue_index i WHERE {' AND '.join(conds)} ORDER BY i.updated_at DESC LIMIT ?" + rows = self._conn.execute(sql, (*params, limit)).fetchall() + return [_index_entry_from_row(row) for row in rows] + + def issue_index_watermark(self, repo: str) -> str | None: + """Max `updated_at` fully ingested for `repo`; None = never backfilled.""" + with self._lock: + row = self._conn.execute("SELECT last_synced FROM issue_index_sync WHERE repo = ?", (repo,)).fetchone() + return str(row["last_synced"]) if row is not None else None + + def set_issue_index_watermark(self, repo: str, last_synced: str) -> None: + with self._lock: + self._conn.execute( + """ + INSERT INTO issue_index_sync (repo, last_synced) VALUES (?, ?) + ON CONFLICT(repo) DO UPDATE SET last_synced = excluded.last_synced + """, + (repo, last_synced), + ) + + +def _index_entry_from_row(row: sqlite3.Row) -> IssueIndexEntry: + try: + labels = tuple(str(x) for x in json.loads(row["labels_json"])) + except (ValueError, TypeError): + labels = () + return IssueIndexEntry( + repo=str(row["repo"]), + number=int(row["number"]), + is_pull_request=bool(row["is_pr"]), + title=str(row["title"]), + body=str(row["body"]), + state=str(row["state"]), + state_reason=str(row["state_reason"]), + merged_at=str(row["merged_at"]), + author=str(row["author"]), + labels=labels, + comments=int(row["comments"]), + created_at=str(row["created_at"]), + updated_at=str(row["updated_at"]), + html_url=str(row["html_url"]), + ) + _DB_SINGLETON: Database | None = None _DB_LOCK = threading.Lock() diff --git a/python/robomp/src/github_backend.py b/python/robomp/src/github_backend.py index b9b09ba6f..38e78c19b 100644 --- a/python/robomp/src/github_backend.py +++ b/python/robomp/src/github_backend.py @@ -13,6 +13,7 @@ from typing import Any, Protocol from robomp.github_client import ( CommentInfo, + IssueIndexEntry, IssueInfo, IssueSummary, PullRequestFileInfo, @@ -47,6 +48,14 @@ class GitHubBackend(Protocol): ) -> list[IssueSummary]: ... async def search_issues(self, repo: str, query: str, *, limit: int = 10) -> list[IssueSummary]: ... + async def list_issue_index_entries( + self, + repo: str, + *, + since: str | None = None, + page: int = 1, + per_page: int = 100, + ) -> list[IssueIndexEntry]: ... async def list_comments(self, repo: str, number: int) -> list[CommentInfo]: ... diff --git a/python/robomp/src/github_client.py b/python/robomp/src/github_client.py index 6ad2bf222..f6b03d088 100644 --- a/python/robomp/src/github_client.py +++ b/python/robomp/src/github_client.py @@ -122,6 +122,30 @@ class IssueSummary: is_pull_request: bool = False +@dataclass(slots=True, frozen=True) +class IssueIndexEntry: + """Full projection of an issue/PR for the local search index (includes body). + + Produced by `GitHubClient.list_issue_index_entries` / webhook payloads and + stored verbatim in the orchestrator's `issue_index` table. + """ + + repo: str + number: int + is_pull_request: bool + title: str + body: str + state: str # open | closed + state_reason: str # completed | not_planned | reopened | "" + merged_at: str # ISO timestamp for merged PRs; "" otherwise + author: str + labels: tuple[str, ...] + comments: int + created_at: str + updated_at: str + html_url: str + + @dataclass(slots=True, frozen=True) class ReactionInfo: """A reaction on an issue/comment. @@ -359,6 +383,31 @@ class GitHubClient: items = (data or {}).get("items") or [] return [_summary_from_item(repo, item) for item in items] + async def list_issue_index_entries( + self, + repo: str, + *, + since: str | None = None, + page: int = 1, + per_page: int = 100, + ) -> list[IssueIndexEntry]: + """One page of issues AND PRs (with bodies) for the local search index. + + `since` is GitHub's ISO `updated_at` lower bound; omit for a full + backfill. Callers page from 1 until a short page comes back. + """ + params: dict[str, Any] = { + "state": "all", + "per_page": max(1, min(int(per_page), 100)), + "page": max(1, int(page)), + "sort": "updated", + "direction": "asc", + } + if since: + params["since"] = since + data = await self.request("GET", f"/repos/{repo}/issues", params=params) + return [index_entry_from_issue_object(repo, item) for item in (data or [])] + async def list_comments(self, repo: str, number: int) -> list[CommentInfo]: data = await self.request("GET", f"/repos/{repo}/issues/{number}/comments", params={"per_page": 100}) return [_comment_from_payload(item) for item in (data or [])] @@ -602,6 +651,60 @@ def _summary_from_item(repo: str, item: Mapping[str, Any]) -> IssueSummary: ) +def index_entry_from_issue_object(repo: str, item: Mapping[str, Any]) -> IssueIndexEntry: + """Build an `IssueIndexEntry` from a REST *issue-shaped* object. + + Accepts both plain issues and the issue representation of a PR (webhook + `issues`/`issue_comment` payloads, `/repos/{repo}/issues` items): PRs carry + a `pull_request` sub-object holding `merged_at`. + """ + user = item.get("user") or {} + labels_raw = item.get("labels") or [] + pr_obj = item.get("pull_request") + is_pr = pr_obj is not None + merged_at = str(pr_obj.get("merged_at") or "") if isinstance(pr_obj, Mapping) else "" + return IssueIndexEntry( + repo=repo, + number=int(item["number"]), + is_pull_request=is_pr, + title=str(item.get("title") or ""), + body=str(item.get("body") or ""), + state=str(item.get("state") or "open"), + state_reason=str(item.get("state_reason") or ""), + merged_at=merged_at, + author=str(user.get("login") or ""), + labels=tuple(str(lbl["name"]) if isinstance(lbl, dict) else str(lbl) for lbl in labels_raw), + comments=int(item.get("comments") or 0), + created_at=str(item.get("created_at") or ""), + updated_at=str(item.get("updated_at") or ""), + html_url=str(item.get("html_url") or ""), + ) + + +def index_entry_from_pr_object(repo: str, item: Mapping[str, Any]) -> IssueIndexEntry: + """Build an `IssueIndexEntry` from a REST *pull-request-shaped* object + (webhook `pull_request*` payloads), where `merged_at` sits at the top level. + """ + user = item.get("user") or {} + labels_raw = item.get("labels") or [] + return IssueIndexEntry( + repo=repo, + number=int(item["number"]), + is_pull_request=True, + title=str(item.get("title") or ""), + body=str(item.get("body") or ""), + state=str(item.get("state") or "open"), + state_reason="", + merged_at=str(item.get("merged_at") or ""), + author=str(user.get("login") or ""), + labels=tuple(str(lbl["name"]) if isinstance(lbl, dict) else str(lbl) for lbl in labels_raw), + comments=int(item.get("comments") or 0), + created_at=str(item.get("created_at") or ""), + updated_at=str(item.get("updated_at") or ""), + html_url=str(item.get("html_url") or ""), + ) + + def _pr_file_from_payload(data: Mapping[str, Any]) -> PullRequestFileInfo: return PullRequestFileInfo( path=str(data.get("filename") or data.get("path") or ""), @@ -663,6 +766,7 @@ __all__ = [ "CommentInfo", "GitHubClient", "GitHubError", + "IssueIndexEntry", "IssueInfo", "IssueSummary", "PullRequestFileInfo", @@ -671,5 +775,7 @@ __all__ = [ "ReactionInfo", "RepoInfo", "ReviewCommentInfo", + "index_entry_from_issue_object", + "index_entry_from_pr_object", "parse_issue_payload", ] diff --git a/python/robomp/src/host_tools.py b/python/robomp/src/host_tools.py index ac4bf3633..92cdd42be 100644 --- a/python/robomp/src/host_tools.py +++ b/python/robomp/src/host_tools.py @@ -27,6 +27,7 @@ from robomp.db import Database, IssueState, issue_key from robomp.git_ops import GitCommandError, HeadDriftError from robomp.github_backend import GitHubBackend from robomp.github_client import GitHubError, IssueInfo, PullRequestFileInfo, RepoInfo +from robomp.issue_index import parse_search_query from robomp.sandbox import ( GitTransport, Workspace, @@ -1260,11 +1261,25 @@ def _build_fetch_thread(bindings: ToolBindings) -> HostTool[Any, Any]: _REPO_QUALIFIER_RE = re.compile(r"(?i)\brepo:") +def _render_search_matches( + query: str, repo: str, rows: list[tuple[bool, int, str, str, str, tuple[str, ...], str]] +) -> str: + """Render (is_pr, number, state_display, title, author, labels, updated) rows.""" + lines = [f"# {len(rows)} match(es) for {query!r} in {repo}"] + for is_pr, number, state, title, author, labels, updated in rows: + kind = "PR" if is_pr else "issue" + label_sfx = f" [{', '.join(labels)}]" if labels else "" + lines.append(f"- #{number} ({kind}, {state}) {title} — @{author}, updated {updated[:10]}{label_sfx}") + return "\n".join(lines) + + def _build_search_issues(bindings: ToolBindings) -> HostTool[Any, Any]: - """Read-only issue/PR search scoped to the current repo. + """Issue/PR search scoped to the current repo, served from the local index. Exists so triage can find duplicates and already-merged fixes instead of - classifying blind; the inbound issue itself is filtered out of results. + classifying blind. Queries hit the webhook-fed SQLite FTS index (zero API + cost); the GitHub search API is only used before the repo's first + reconcile completes. The inbound issue is filtered out of results. """ def execute(args: dict[str, Any], _ctx: HostToolContext[Any]) -> str: @@ -1280,28 +1295,61 @@ def _build_search_issues(bindings: ToolBindings) -> HostTool[Any, Any]: _raise_command(msg) limit_raw = args.get("limit") limit = max(1, min(int(limit_raw), 20)) if isinstance(limit_raw, int) else 10 - try: - results = _run_coro( - bindings.loop, - bindings.github.search_issues(bindings.repo.full_name, query, limit=limit), + repo = bindings.repo.full_name + + rows: list[tuple[bool, int, str, str, str, tuple[str, ...], str]] + if bindings.db.issue_index_watermark(repo) is not None: + parsed = parse_search_query(query) + entries = bindings.db.search_issue_index( + repo, + keywords=parsed.keywords, + is_pr=parsed.is_pr, + state=parsed.state, + merged=parsed.merged, + label=parsed.label, + author=parsed.author, + limit=limit + 1, # headroom for the self-filter below ) - except GitHubError as exc: - _audit(bindings, "gh_search_issues", args, error=str(exc)) - _raise_command(f"GitHub search failed: {exc.status} {exc.message}") - results = [s for s in results if s.is_pull_request or s.number != bindings.issue.number] - if not results: - _audit(bindings, "gh_search_issues", args, result={"matches": 0}) - return f"No issues or PRs in {bindings.repo.full_name} match {query!r}." - lines = [f"# {len(results)} match(es) for {query!r} in {bindings.repo.full_name}"] - for s in results: - kind = "PR" if s.is_pull_request else "issue" - state = f"{s.state} ({s.state_reason})" if s.state_reason else s.state - labels = f" [{', '.join(s.labels)}]" if s.labels else "" - lines.append( - f"- #{s.number} ({kind}, {state}) {s.title} — @{s.author}, updated {s.updated_at[:10]}{labels}" - ) - _audit(bindings, "gh_search_issues", args, result={"matches": len(results)}) - return "\n".join(lines) + entries = [e for e in entries if e.is_pull_request or e.number != bindings.issue.number][:limit] + rows = [] + for e in entries: + if e.is_pull_request and e.merged_at: + state = "merged" + elif e.state_reason: + state = f"{e.state} ({e.state_reason})" + else: + state = e.state + rows.append((e.is_pull_request, e.number, state, e.title, e.author, e.labels, e.updated_at)) + source = "local" + else: + # Index not backfilled yet — fall through to the GitHub search API. + try: + found = _run_coro( + bindings.loop, + bindings.github.search_issues(repo, query, limit=limit), + ) + except GitHubError as exc: + _audit(bindings, "gh_search_issues", args, error=str(exc)) + _raise_command(f"GitHub search failed: {exc.status} {exc.message}") + found = [s for s in found if s.is_pull_request or s.number != bindings.issue.number] + rows = [ + ( + s.is_pull_request, + s.number, + f"{s.state} ({s.state_reason})" if s.state_reason else s.state, + s.title, + s.author, + s.labels, + s.updated_at, + ) + for s in found + ] + source = "remote" + if not rows: + _audit(bindings, "gh_search_issues", args, result={"matches": 0, "source": source}) + return f"No issues or PRs in {repo} match {query!r}." + _audit(bindings, "gh_search_issues", args, result={"matches": len(rows), "source": source}) + return _render_search_matches(query, repo, rows) return host_tool( name="gh_search_issues", @@ -1325,6 +1373,98 @@ def _build_search_issues(bindings: ToolBindings) -> HostTool[Any, Any]: ) +# ---------- search_commits ---------- +_COMMIT_SEARCH_TIMEOUT_SECONDS = 120.0 + + +def _build_search_commits(bindings: ToolBindings) -> HostTool[Any, Any]: + """Local `git log` search over the default branch's history. + + Two modes: `message` greps commit subjects/bodies (case-insensitive + regex), `patch` runs the pickaxe (`-S`) to find commits whose diff adds or + removes the literal string — the sharp tool for "was this already fixed". + The search interface (query in, ranked commits out) is deliberately opaque + about its backend so a semantic index can replace git plumbing later. + """ + + def execute(args: dict[str, Any], _ctx: HostToolContext[Any]) -> str: + query = args.get("query") + if not isinstance(query, str) or not query.strip(): + msg = "search_commits requires a non-empty 'query'." + _audit(bindings, "search_commits", args, error=msg) + _raise_command(msg) + query = query.strip() + mode = args.get("mode") or "message" + if mode not in ("message", "patch"): + msg = "search_commits 'mode' must be 'message' or 'patch'." + _audit(bindings, "search_commits", args, error=msg) + _raise_command(msg) + limit_raw = args.get("limit") + limit = max(1, min(int(limit_raw), 30)) if isinstance(limit_raw, int) else 10 + paths = [p for p in (args.get("paths") or ()) if isinstance(p, str) and p.strip()] + + rev = f"origin/{bindings.repo.default_branch}" + probe = _run_repo_command(bindings, ["git", "rev-parse", "--verify", "--quiet", rev], timeout=30.0) + if probe.returncode != 0: + rev = "HEAD" + cmd = ["git", "log", rev, "-n", str(limit), "--date=short", "--pretty=format:%h %ad %an — %s"] + if mode == "message": + cmd += [f"--grep={query}", "--regexp-ignore-case"] + else: + cmd += ["-S", query] + if paths: + cmd += ["--", *paths] + try: + proc = _run_repo_command(bindings, cmd, timeout=_COMMIT_SEARCH_TIMEOUT_SECONDS) + except subprocess.TimeoutExpired: + msg = f"search_commits timed out after {_COMMIT_SEARCH_TIMEOUT_SECONDS:.0f}s; narrow with 'paths' or a shorter history window." + _audit(bindings, "search_commits", args, error=msg) + _raise_command(msg) + if proc.returncode != 0: + msg = f"git log failed: {(proc.stderr or proc.stdout).strip()[:500]}" + _audit(bindings, "search_commits", args, error=msg) + _raise_command(msg) + out = proc.stdout.strip() + if not out: + _audit(bindings, "search_commits", args, result={"matches": 0}) + return f"No commits on {rev} match {query!r} (mode={mode})." + matches = out.splitlines() + _audit(bindings, "search_commits", args, result={"matches": len(matches)}) + header = f"# {len(matches)} commit(s) on {rev} matching {query!r} (mode={mode})" + return "\n".join([header, *matches]) + + return host_tool( + name="search_commits", + description=persona.host_tool_description("search_commits"), + parameters={ + "type": "object", + "properties": { + "query": { + "type": "string", + "description": persona.host_tool_parameter_description("search_commits", "query"), + }, + "mode": { + "type": "string", + "enum": ["message", "patch"], + "description": persona.host_tool_parameter_description("search_commits", "mode"), + }, + "paths": { + "type": "array", + "items": {"type": "string"}, + "description": persona.host_tool_parameter_description("search_commits", "paths"), + }, + "limit": { + "type": "integer", + "description": persona.host_tool_parameter_description("search_commits", "limit"), + }, + }, + "required": ["query"], + "additionalProperties": False, + }, + execute=execute, + ) + + _PRIMARY_TYPES = ("bug", "enhancement", "question", "proposal", "documentation", "wontfix", "invalid", "duplicate") _AUTO_PR_CLASSIFICATIONS = frozenset({"bug", "documentation"}) _PRIORITIES = ("prio:p0", "prio:p1", "prio:p2", "prio:p3") @@ -1905,6 +2045,7 @@ def build(bindings: ToolBindings) -> tuple[HostTool[Any, Any], ...]: _build_abort_task(bindings), _build_fetch_thread(bindings), _build_search_issues(bindings), + _build_search_commits(bindings), ) diff --git a/python/robomp/src/issue_index.py b/python/robomp/src/issue_index.py new file mode 100644 index 000000000..98f966018 --- /dev/null +++ b/python/robomp/src/issue_index.py @@ -0,0 +1,251 @@ +"""Local issue/PR search index: webhook ingest, periodic reconcile, query parsing. + +The `issue_index` table mirrors every issue and PR of the allowlisted repos so +`gh_search_issues` answers from SQLite FTS5 instead of the GitHub search API. +Freshness comes from two directions: + + - `ingest_webhook_payload` upserts on every `issues` / `issue_comment` / + `pull_request*` delivery, keeping the hot path current in real time. + - `IssueIndexSync` reconciles each repo every `issue_index_sync_seconds` + (and backfills on first run) via `/repos/{repo}/issues?since=…`, catching + anything webhooks missed while the orchestrator was down. + +`parse_search_query` translates the GitHub-search-flavored tool query into the +structured filters `Database.search_issue_index` takes, so the agent keeps one +query language whether the lookup is served locally or remotely. +""" + +from __future__ import annotations + +import asyncio +import logging +from collections.abc import Mapping +from dataclasses import dataclass +from datetime import UTC, datetime, timedelta + +from robomp.config import Settings +from robomp.db import Database +from robomp.github_backend import GitHubBackend +from robomp.github_client import ( + GitHubError, + index_entry_from_issue_object, + index_entry_from_pr_object, +) + +log = logging.getLogger(__name__) + +# Overlap subtracted from the watermark on every reconcile so a sync that +# raced a concurrent update can never permanently skip it. +_SYNC_OVERLAP = timedelta(minutes=2) +_PAGE_SIZE = 100 +_MAX_PAGES_PER_TICK = 30 + + +@dataclass(slots=True, frozen=True) +class ParsedSearchQuery: + """Structured form of a GitHub-issue-search style query string.""" + + keywords: tuple[str, ...] + is_pr: bool | None = None + state: str | None = None + merged: bool | None = None + label: str | None = None + author: str | None = None + + +def parse_search_query(query: str) -> ParsedSearchQuery: + """Split a GitHub-search style string into keywords + structured filters. + + Supported qualifiers: `is:pr` / `is:issue` / `is:open` / `is:closed` / + `is:merged`, `label:`, `author:`. Unrecognized `key:value` + qualifiers are dropped rather than fed to FTS5 (a bare `in:title` token + would otherwise be a syntax error). Everything else is a keyword. + """ + keywords: list[str] = [] + is_pr: bool | None = None + state: str | None = None + merged: bool | None = None + label: str | None = None + author: str | None = None + for token in query.split(): + key, sep, value = token.partition(":") + if not sep or not value or " " in key: + keywords.append(token) + continue + key = key.lower() + if key == "is": + v = value.lower() + if v == "pr": + is_pr = True + elif v == "issue": + is_pr = False + elif v in ("open", "closed"): + state = v + elif v == "merged": + is_pr = True + merged = True + elif key == "label": + label = value.strip('"') + elif key == "author": + author = value.lstrip("@") + # Any other qualifier (in:, sort:, created:, …) is intentionally dropped. + return ParsedSearchQuery( + keywords=tuple(keywords), + is_pr=is_pr, + state=state, + merged=merged, + label=label, + author=author, + ) + + +def ingest_webhook_payload(db: Database, repo: str, event_type: str, payload: Mapping[str, object]) -> bool: + """Upsert the issue/PR carried by a webhook delivery into the index. + + Returns True when the payload contained an indexable object. Runs before + routing so even deliveries the router skips (bot comments, unhandled + actions) still refresh the index. + """ + if event_type in ("issues", "issue_comment"): + obj = payload.get("issue") + if isinstance(obj, Mapping) and obj.get("number") is not None: + db.upsert_issue_index(index_entry_from_issue_object(repo, obj)) + return True + return False + if event_type.startswith("pull_request"): + obj = payload.get("pull_request") + if isinstance(obj, Mapping) and obj.get("number") is not None: + db.upsert_issue_index(index_entry_from_pr_object(repo, obj)) + return True + return False + return False + + +def _utcnow_iso() -> str: + return datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%SZ") + + +def _overlapped(watermark: str) -> str: + """Rewind an ISO watermark by the sync overlap; fall back to the raw value.""" + try: + parsed = datetime.strptime(watermark, "%Y-%m-%dT%H:%M:%SZ").replace(tzinfo=UTC) + except ValueError: + return watermark + return (parsed - _SYNC_OVERLAP).strftime("%Y-%m-%dT%H:%M:%SZ") + + +class IssueIndexSync: + """Background reconciler for the local issue index. + + First tick per repo backfills everything (`since=None`); later ticks pull + only issues updated after the stored watermark (minus a small overlap). + Pagination is bounded per tick — a huge backlog finishes across ticks + rather than hogging one. + """ + + def __init__(self, *, settings: Settings, db: Database, github: GitHubBackend) -> None: + self._settings = settings + self._db = db + self._github = github + self._task: asyncio.Task[None] | None = None + self._stop_event: asyncio.Event | None = None + + @property + def enabled(self) -> bool: + return self._settings.issue_index_sync_seconds > 0 + + async def start(self) -> None: + """Spawn the background loop. No-op when disabled.""" + if not self.enabled: + log.info("issue index sync disabled") + return + if self._task is not None: + return + self._stop_event = asyncio.Event() + self._task = asyncio.create_task(self._run(), name="issue-index-sync") + log.info( + "issue index sync started", + extra={"interval_seconds": self._settings.issue_index_sync_seconds}, + ) + + async def stop(self) -> None: + """Signal the loop to exit and await its termination.""" + if self._task is None: + return + assert self._stop_event is not None + self._stop_event.set() + try: + await asyncio.wait_for(self._task, timeout=5.0) + except TimeoutError: + self._task.cancel() + try: + await self._task + except (asyncio.CancelledError, Exception): + pass + finally: + self._task = None + self._stop_event = None + + async def _run(self) -> None: + assert self._stop_event is not None + interval = float(self._settings.issue_index_sync_seconds) + while not self._stop_event.is_set(): + try: + await self.tick() + except Exception: + log.exception("issue index sync tick failed") + try: + await asyncio.wait_for(self._stop_event.wait(), timeout=interval) + except TimeoutError: + continue + + async def tick(self) -> None: + """Reconcile every allowlisted repo once.""" + for repo in self._settings.repo_allowlist: + try: + await self.sync_repo(repo) + except GitHubError as exc: + log.warning( + "issue index sync failed; will retry next tick", + extra={"repo": repo, "status": exc.status, "gh_message": exc.message}, + ) + + async def sync_repo(self, repo: str) -> int: + """Pull updated issues/PRs for one repo into the index. Returns count ingested.""" + started_at = _utcnow_iso() + watermark = self._db.issue_index_watermark(repo) + since = _overlapped(watermark) if watermark else None + ingested = 0 + exhausted = False + last_seen = "" + for page in range(1, _MAX_PAGES_PER_TICK + 1): + batch = await self._github.list_issue_index_entries(repo, since=since, page=page, per_page=_PAGE_SIZE) + for entry in batch: + self._db.upsert_issue_index(entry) + if entry.updated_at > last_seen: + last_seen = entry.updated_at + ingested += len(batch) + if len(batch) < _PAGE_SIZE: + exhausted = True + break + if exhausted: + # Everything up to tick start is ingested; later updates are the + # next tick's problem (or a webhook's). + self._db.set_issue_index_watermark(repo, started_at) + elif last_seen: + # Page budget hit mid-backfill: advance the watermark only to the + # newest updated_at actually ingested, so the next tick resumes there. + self._db.set_issue_index_watermark(repo, last_seen) + log.info( + "issue index synced", + extra={"repo": repo, "ingested": ingested, "backfill": watermark is None, "complete": exhausted}, + ) + return ingested + + +__all__ = [ + "IssueIndexSync", + "ParsedSearchQuery", + "ingest_webhook_payload", + "parse_search_query", +] diff --git a/python/robomp/src/prompts/host_tools.toml b/python/robomp/src/prompts/host_tools.toml index 80b0aa211..82a84a32d 100644 --- a/python/robomp/src/prompts/host_tools.toml +++ b/python/robomp/src/prompts/host_tools.toml @@ -73,12 +73,21 @@ reason = "Internal diagnosis for the operator. Concrete, specific, blameless. NE description = "Refetch the originating issue and its comments. Use sparingly." [gh_search_issues] -description = "Search issues AND pull requests in the current repo (GitHub issue-search syntax; repo scope applied automatically). Use during triage to find duplicates and to check whether a merged PR already fixed the reported problem before classifying." +description = "Search issues AND pull requests in the current repo. Served from a local index of every issue/PR (webhook-fed, periodically reconciled), so calls are free — search liberally during triage to find duplicates and already-merged fixes before classifying." [gh_search_issues.parameters] -query = "GitHub issue-search syntax: bare keywords plus qualifiers like `is:pr`, `is:closed`, `is:merged`, `label:bug`, `in:title`, `author:`. NEVER include a `repo:` qualifier — scope is applied automatically." +query = "GitHub issue-search syntax: bare keywords plus qualifiers like `is:pr`, `is:closed`, `is:merged`, `label:bug`, `author:`. NEVER include a `repo:` qualifier — scope is applied automatically." limit = "Max results, 1-20. Default 10." +[search_commits] +description = "Search the default branch's commit history. `mode=message` greps commit messages (case-insensitive regex); `mode=patch` finds commits whose diff adds/removes the literal query string — use it to check whether the broken code path was already touched by a fix." + +[search_commits.parameters] +query = "Regex for `mode=message`; literal string for `mode=patch`." +mode = "`message` (default) searches commit messages; `patch` pickaxe-searches diff content." +paths = "Optional path filters (files or directories) to narrow the walk." +limit = "Max commits, 1-30. Default 10." + [set_issue_labels] description = "Append labels to the originating issue/PR. NEVER removes existing labels." diff --git a/python/robomp/src/prompts/system_append.md b/python/robomp/src/prompts/system_append.md index 5212c9ed9..9bbced981 100644 --- a/python/robomp/src/prompts/system_append.md +++ b/python/robomp/src/prompts/system_append.md @@ -25,10 +25,10 @@ Pick exactly ONE primary label per issue: ## Duplicate & already-fixed check -Before `classify_issue`, run `gh_search_issues` with the report's key terms (retry with synonyms and an `is:pr` variant — one search proves nothing): +Before `classify_issue`, run `gh_search_issues` with the report's key terms (retry with synonyms and an `is:pr` variant — searches are served from a local index and cost nothing; one search proves nothing): - **Prior issue on the same problem** → `duplicate`, cite it. A prior closure as not-planned/`wontfix` on the same complaint is binding precedent — adopt that verdict; NEVER relitigate it. -- **Already fixed.** Your worktree is the CURRENT default branch; reporters often run older releases. When the reported version lags the latest release (topmost released section of the relevant `packages/*/CHANGELOG.md`), check the changelog and merged PRs (`is:pr is:merged `) for an existing fix, and try the repro against the worktree — failing on the reporter's version but passing here means it is already fixed. Classify `duplicate`: cite the fixing PR, name the release carrying it (or say it ships in the next release when still under `[Unreleased]`), and tell the reporter to update. NEVER re-fix what main already fixed. +- **Already fixed.** Your worktree is the CURRENT default branch; reporters often run older releases. When the reported version lags the latest release (topmost released section of the relevant `packages/*/CHANGELOG.md`), check the changelog, merged PRs (`is:pr is:merged `), and recent commits (`search_commits` — `mode=message` for symptom keywords, `mode=patch` for the exact broken code) for an existing fix, and try the repro against the worktree — failing on the reporter's version but passing here means it is already fixed. Classify `duplicate`: cite the fixing PR/commit, name the release carrying it (or say it ships in the next release when still under `[Unreleased]`), and tell the reporter to update. NEVER re-fix what main already fixed. ## Merit gate — `bug` vs `wontfix` vs `enhancement` diff --git a/python/robomp/src/proxy/server.py b/python/robomp/src/proxy/server.py index 923753b40..2d3bd96a6 100644 --- a/python/robomp/src/proxy/server.py +++ b/python/robomp/src/proxy/server.py @@ -523,6 +523,22 @@ def create_proxy_app(settings: Settings) -> FastAPI: return _gh_error_response(exc) return JSONResponse({"items": [_serialize(s) for s in items]}) + @app.get("/gh/v1/issue_index_entries") + async def list_issue_index_entries( + request: Request, + repo: str, + since: str | None = None, + page: int = 1, + per_page: int = 100, + ) -> JSONResponse: + await _authenticate(request) + github: GitHubClient = request.app.state.github + try: + items = await github.list_issue_index_entries(repo, since=since, page=page, per_page=per_page) + except GitHubError as exc: + return _gh_error_response(exc) + return JSONResponse({"items": [_serialize(s) for s in items]}) + @app.get("/gh/v1/comments") async def list_comments(request: Request, repo: str, number: int) -> JSONResponse: await _authenticate(request) diff --git a/python/robomp/src/proxy_client.py b/python/robomp/src/proxy_client.py index a2069c950..ddf0c3229 100644 --- a/python/robomp/src/proxy_client.py +++ b/python/robomp/src/proxy_client.py @@ -25,6 +25,7 @@ from robomp.git_ops import GitCommandError, HeadDriftError, PushResult from robomp.github_client import ( CommentInfo, GitHubError, + IssueIndexEntry, IssueInfo, IssueSummary, PullRequestFileInfo, @@ -212,6 +213,20 @@ class GitHubProxyClient: ) return [_issue_summary_from(item) for item in (data.get("items") if isinstance(data, dict) else None) or []] + async def list_issue_index_entries( + self, + repo: str, + *, + since: str | None = None, + page: int = 1, + per_page: int = 100, + ) -> list[IssueIndexEntry]: + params: dict[str, Any] = {"repo": repo, "page": page, "per_page": per_page} + if since: + params["since"] = since + data = await self._request("GET", "/gh/v1/issue_index_entries", params=params) + return [_index_entry_from(item) for item in (data.get("items") if isinstance(data, dict) else None) or []] + async def list_comments(self, repo: str, number: int) -> list[CommentInfo]: data = await self._request("GET", "/gh/v1/comments", params={"repo": repo, "number": number}) return [_comment_from(item) for item in (data.get("items") if isinstance(data, dict) else None) or []] @@ -511,6 +526,27 @@ def _issue_summary_from(data: Any) -> IssueSummary: ) +def _index_entry_from(data: Any) -> IssueIndexEntry: + if not isinstance(data, dict): + raise GitHubError(500, "proxy returned malformed issue index payload") + return IssueIndexEntry( + repo=str(data["repo"]), + number=int(data["number"]), + is_pull_request=bool(data.get("is_pull_request")), + title=str(data.get("title") or ""), + body=str(data.get("body") or ""), + state=str(data.get("state") or ""), + state_reason=str(data.get("state_reason") or ""), + merged_at=str(data.get("merged_at") or ""), + author=str(data.get("author") or ""), + labels=tuple(str(x) for x in (data.get("labels") or [])), + comments=int(data.get("comments") or 0), + created_at=str(data.get("created_at") or ""), + updated_at=str(data.get("updated_at") or ""), + html_url=str(data.get("html_url") or ""), + ) + + def _comment_from(data: Any) -> CommentInfo: if not isinstance(data, dict): raise GitHubError(500, "proxy returned malformed comment payload") diff --git a/python/robomp/src/server.py b/python/robomp/src/server.py index 6a9371f96..36aa3d5a5 100644 --- a/python/robomp/src/server.py +++ b/python/robomp/src/server.py @@ -14,7 +14,7 @@ from fastapi import Body, FastAPI, Header, HTTPException, Request, status from fastapi.responses import HTMLResponse, JSONResponse from fastapi.staticfiles import StaticFiles -from robomp import github_events +from robomp import github_events, issue_index from robomp.autoclose import AutocloseScheduler from robomp.config import Settings, get_settings from robomp.dashboard import render_index, static_dir, tail_jsonl @@ -29,6 +29,7 @@ from robomp.db import ( ) from robomp.github_backend import GitHubBackend from robomp.github_client import GitHubError, IssueSummary +from robomp.issue_index import IssueIndexSync from robomp.manual_triage import ( InvalidIssueRef, ManualTriageConflict, @@ -252,6 +253,7 @@ def _build_state(settings: Settings) -> dict[str, Any]: ) pool = WorkerPool(settings=settings, db=db, github=github, sandbox=sandbox, git_transport=git_transport) autoclose = AutocloseScheduler(settings=settings, db=db, github=github) + index_sync = IssueIndexSync(settings=settings, db=db, github=github) return { "settings": settings, "db": db, @@ -262,6 +264,7 @@ def _build_state(settings: Settings) -> dict[str, Any]: "pool": pool, "issue_browse_cache": _IssueBrowseCache(), "autoclose": autoclose, + "issue_index_sync": index_sync, } @@ -278,9 +281,12 @@ def create_app(settings: Settings | None = None) -> FastAPI: await pool.start() autoclose: AutocloseScheduler = app.state.bag["autoclose"] await autoclose.start() + index_sync: IssueIndexSync = app.state.bag["issue_index_sync"] + await index_sync.start() try: yield finally: + await index_sync.stop() await autoclose.stop() await pool.stop( drain_timeout=cfg.shutdown_drain_timeout_seconds, @@ -328,6 +334,15 @@ def create_app(settings: Settings | None = None) -> FastAPI: payload=payload, allowlist=cfg.repo_allowlist, ) + # Keep the local search index fresh from every delivery that carries an + # issue/PR object — including ones the router will skip. + if x_github_event in ("issues", "issue_comment") or x_github_event.startswith("pull_request"): + repo_full = str((payload.get("repository") or {}).get("full_name") or "") + if repo_full and repo_full in cfg.repo_allowlist: + try: + issue_index.ingest_webhook_payload(db, repo_full, x_github_event, payload) + except Exception: + log.exception("issue index webhook ingest failed", extra={"repo": repo_full}) def _resolve(repo_full: str, pr_number: int) -> str | None: row = db.find_issue_by_pr(repo_full, pr_number) diff --git a/python/robomp/tests/conftest.py b/python/robomp/tests/conftest.py index 89959e5fb..6fc4395ab 100644 --- a/python/robomp/tests/conftest.py +++ b/python/robomp/tests/conftest.py @@ -97,6 +97,10 @@ def _baseline_env(tmp_path: Path) -> dict[str, str]: # cache flip `ROBOMP_NATIVES_CACHE_ENABLED=true` explicitly. "ROBOMP_NATIVES_CACHE_ROOT": str(tmp_path / "natives-cache"), "ROBOMP_NATIVES_CACHE_ENABLED": "false", + # Same reasoning for the issue-index reconciler: its first tick would + # spin connect-retries against the .invalid proxy URL inside server + # tests. Tests that want it construct IssueIndexSync directly. + "ROBOMP_ISSUE_INDEX_SYNC_SECONDS": "0", } diff --git a/python/robomp/tests/test_host_tools.py b/python/robomp/tests/test_host_tools.py index 102074c90..3e2f2f541 100644 --- a/python/robomp/tests/test_host_tools.py +++ b/python/robomp/tests/test_host_tools.py @@ -4,6 +4,7 @@ from __future__ import annotations import asyncio import json +import subprocess import threading from pathlib import Path from typing import Any @@ -14,7 +15,7 @@ from omp_rpc import HostToolContext, RpcCommandError from robomp import host_tools from robomp.db import Database -from robomp.github_client import GitHubClient, IssueInfo, RepoInfo +from robomp.github_client import GitHubClient, IssueIndexEntry, IssueInfo, RepoInfo from robomp.host_tools import AbortController, ToolBindings, build from robomp.sandbox import LocalGitTransport, Workspace @@ -937,6 +938,98 @@ def test_gh_search_issues_rejects_repo_qualifier_and_empty_query(db: Database, t _stop_loop(loop, t) +def test_gh_search_issues_serves_from_local_index_once_synced(db: Database, tmp_path: Path) -> None: + """With a sync watermark present the tool answers from SQLite: qualifiers + become filters, merged PRs render as `merged`, and NO GitHub call happens.""" + + def handler(_request: httpx.Request) -> httpx.Response: + raise AssertionError("local-index search must not call GitHub") + + bindings, loop, t = _bindings(db, tmp_path, httpx.MockTransport(handler)) + db.set_issue_index_watermark("octo/widget", "2026-07-01T00:00:00Z") + db.upsert_issue_index( + IssueIndexEntry( + repo="octo/widget", + number=31, + is_pull_request=True, + title="fix: resize crash", + body="handles narrow terminals", + state="closed", + state_reason="", + merged_at="2026-06-02T00:00:00Z", + author="bot", + labels=(), + comments=1, + created_at="2026-06-02T00:00:00Z", + updated_at="2026-06-02T00:00:00Z", + html_url="https://example/pull/31", + ) + ) + db.upsert_issue_index( + IssueIndexEntry( + repo="octo/widget", + number=30, + is_pull_request=False, + title="resize crash report", + body="", + state="closed", + state_reason="not_planned", + merged_at="", + author="bob", + labels=("wontfix",), + comments=3, + created_at="2026-05-01T00:00:00Z", + updated_at="2026-06-01T00:00:00Z", + html_url="https://example/30", + ) + ) + try: + tool = next(x for x in build(bindings) if x.name == "gh_search_issues") + result = tool.execute({"query": "resize crash"}, _ctx()) + pr_only = tool.execute({"query": "resize crash is:merged"}, _ctx()) + finally: + _stop_loop(loop, t) + assert "#30 (issue, closed (not_planned))" in result + assert "#31 (PR, merged)" in result + assert "#31" in pr_only and "#30" not in pr_only + + +def _git_repo_with_commits(bindings) -> None: + """Turn the stub workspace repo_dir into a git repo with two commits.""" + repo = str(bindings.workspace.repo_dir) + ident = ["-c", "user.name=t", "-c", "user.email=t@example.invalid"] + subprocess.run(["git", "init", "-q", "-b", "main", repo], check=True) + Path(repo, "a.txt").write_text("plain start\n", encoding="utf-8") + subprocess.run(["git", "-C", repo, "add", "."], check=True) + subprocess.run(["git", "-C", repo, *ident, "commit", "-q", "-m", "feat: initial import"], check=True) + Path(repo, "a.txt").write_text("plain start\nsplitPathAndSel guard\n", encoding="utf-8") + subprocess.run(["git", "-C", repo, "add", "."], check=True) + subprocess.run( + ["git", "-C", repo, *ident, "commit", "-q", "-m", "fix(tools): colon selector literal paths"], + check=True, + ) + + +def test_search_commits_message_and_patch_modes(db: Database, tmp_path: Path) -> None: + """message mode greps commit messages; patch mode pickaxes diff content. + Without an origin ref the search falls back to HEAD instead of failing.""" + bindings, loop, t = _bindings(db, tmp_path, httpx.MockTransport(lambda r: httpx.Response(500))) + _git_repo_with_commits(bindings) + try: + tool = next(x for x in build(bindings) if x.name == "search_commits") + by_message = tool.execute({"query": "colon selector"}, _ctx()) + by_patch = tool.execute({"query": "splitPathAndSel", "mode": "patch"}, _ctx()) + none = tool.execute({"query": "nonexistent-topic"}, _ctx()) + with pytest.raises(RpcCommandError): + tool.execute({"query": "x", "mode": "bogus"}, _ctx()) + finally: + _stop_loop(loop, t) + assert "fix(tools): colon selector literal paths" in by_message + assert "feat: initial import" not in by_message + assert "fix(tools): colon selector literal paths" in by_patch + assert none.startswith("No commits") + + def test_classify_issue_rejects_bug_without_priority(db: Database, tmp_path: Path) -> None: bindings, loop, t = _bindings(db, tmp_path, httpx.MockTransport(lambda r: httpx.Response(500))) try: diff --git a/python/robomp/tests/test_issue_index.py b/python/robomp/tests/test_issue_index.py new file mode 100644 index 000000000..d3da7021f --- /dev/null +++ b/python/robomp/tests/test_issue_index.py @@ -0,0 +1,189 @@ +"""Local issue index: query parsing, webhook ingest, FTS search, reconcile sync.""" + +from __future__ import annotations + +from pathlib import Path + +from robomp.db import Database +from robomp.github_client import IssueIndexEntry +from robomp.issue_index import IssueIndexSync, ingest_webhook_payload, parse_search_query + + +def _entry(number: int, **overrides) -> IssueIndexEntry: + base = { + "repo": "octo/widget", + "number": number, + "is_pull_request": False, + "title": f"issue {number}", + "body": "", + "state": "open", + "state_reason": "", + "merged_at": "", + "author": "alice", + "labels": (), + "comments": 0, + "created_at": "2026-01-01T00:00:00Z", + "updated_at": "2026-01-01T00:00:00Z", + "html_url": f"https://example/{number}", + } + base.update(overrides) + return IssueIndexEntry(**base) + + +# ---- parse_search_query ---- + + +def test_parse_search_query_extracts_supported_qualifiers() -> None: + parsed = parse_search_query("colon selector is:pr is:merged label:bug author:@alice in:title") + assert parsed.keywords == ("colon", "selector") # `in:title` dropped, not fed to FTS + assert parsed.is_pr is True + assert parsed.merged is True + assert parsed.label == "bug" + assert parsed.author == "alice" + + +def test_parse_search_query_state_and_issue_kind() -> None: + parsed = parse_search_query("is:issue is:closed crash") + assert parsed.is_pr is False + assert parsed.state == "closed" + assert parsed.keywords == ("crash",) + + +# ---- db index: upsert + search ---- + + +def test_search_issue_index_matches_body_text_and_ranks(db: Database) -> None: + db.upsert_issue_index(_entry(1, title="TUI crash on resize", body="stack trace mentions overlay")) + db.upsert_issue_index(_entry(2, title="unrelated docs typo", body="readme wording")) + found = db.search_issue_index("octo/widget", keywords=("resize", "crash")) + assert [e.number for e in found] == [1] + # body-only terms also hit + found = db.search_issue_index("octo/widget", keywords=("overlay",)) + assert [e.number for e in found] == [1] + + +def test_search_issue_index_filters(db: Database) -> None: + db.upsert_issue_index( + _entry(1, title="fix crash", is_pull_request=True, merged_at="2026-02-01T00:00:00Z", state="closed") + ) + db.upsert_issue_index( + _entry(2, title="crash report", state="closed", state_reason="not_planned", labels=("wontfix",)) + ) + db.upsert_issue_index(_entry(3, title="crash report open", state="open")) + + merged_prs = db.search_issue_index("octo/widget", keywords=("crash",), is_pr=True, merged=True) + assert [e.number for e in merged_prs] == [1] + wontfixed = db.search_issue_index("octo/widget", keywords=("crash",), label="wontfix") + assert [e.number for e in wontfixed] == [2] + open_only = db.search_issue_index("octo/widget", keywords=("crash",), state="open") + assert [e.number for e in open_only] == [3] + + +def test_upsert_refreshes_fts_so_stale_text_stops_matching(db: Database) -> None: + """The UPDATE trigger must swap FTS content, not accumulate it.""" + db.upsert_issue_index(_entry(1, title="original scrollback wipe")) + db.upsert_issue_index(_entry(1, title="renamed: alternate screen request", state="closed")) + assert db.search_issue_index("octo/widget", keywords=("scrollback",)) == [] + found = db.search_issue_index("octo/widget", keywords=("alternate",)) + assert len(found) == 1 and found[0].state == "closed" + + +def test_search_issue_index_quotes_fts_metacharacters(db: Database) -> None: + """Reporter text like `"AND (` must never raise an FTS5 syntax error.""" + db.upsert_issue_index(_entry(1, title='crash with "quoted" AND (parens)')) + found = db.search_issue_index("octo/widget", keywords=('"quoted"', "AND", "(parens)")) + assert [e.number for e in found] == [1] + + +def test_issue_index_watermark_roundtrip(db: Database) -> None: + assert db.issue_index_watermark("octo/widget") is None + db.set_issue_index_watermark("octo/widget", "2026-07-01T00:00:00Z") + assert db.issue_index_watermark("octo/widget") == "2026-07-01T00:00:00Z" + db.set_issue_index_watermark("octo/widget", "2026-07-02T00:00:00Z") + assert db.issue_index_watermark("octo/widget") == "2026-07-02T00:00:00Z" + + +# ---- webhook ingest ---- + + +def test_ingest_webhook_issue_and_pr_payloads(db: Database) -> None: + ingested = ingest_webhook_payload( + db, + "octo/widget", + "issues", + {"issue": {"number": 5, "title": "boom", "body": "b", "state": "open", "user": {"login": "alice"}}}, + ) + assert ingested + # PR-flavored issue payload (issue_comment on a PR) carries pull_request.merged_at. + ingest_webhook_payload( + db, + "octo/widget", + "issue_comment", + { + "issue": { + "number": 6, + "title": "fixes boom", + "state": "closed", + "user": {"login": "bob"}, + "pull_request": {"merged_at": "2026-03-01T00:00:00Z"}, + } + }, + ) + # Native pull_request payload: merged_at at top level. + ingest_webhook_payload( + db, + "octo/widget", + "pull_request", + {"pull_request": {"number": 7, "title": "another fix", "state": "closed", "merged_at": "2026-04-01T00:00:00Z"}}, + ) + assert not ingest_webhook_payload(db, "octo/widget", "push", {"ref": "refs/heads/main"}) + + boom = db.search_issue_index("octo/widget", keywords=("boom",)) + assert {e.number for e in boom} == {5, 6} + pr6 = next(e for e in boom if e.number == 6) + assert pr6.is_pull_request and pr6.merged_at == "2026-03-01T00:00:00Z" + pr7 = db.search_issue_index("octo/widget", keywords=("another",))[0] + assert pr7.is_pull_request and pr7.merged_at == "2026-04-01T00:00:00Z" + + +# ---- reconcile sync ---- + + +class _FakeBackend: + """Pages of index entries keyed by page number; records `since` per call.""" + + def __init__(self, pages: dict[int, list[IssueIndexEntry]]) -> None: + self.pages = pages + self.calls: list[tuple[str | None, int]] = [] + + async def list_issue_index_entries( + self, repo: str, *, since: str | None = None, page: int = 1, per_page: int = 100 + ) -> list[IssueIndexEntry]: + self.calls.append((since, page)) + return self.pages.get(page, []) + + +class _SyncSettings: + issue_index_sync_seconds = 900.0 + repo_allowlist = frozenset({"octo/widget"}) + + +async def test_sync_repo_backfills_pages_and_sets_watermark(db: Database, tmp_path: Path) -> None: + full_page = [_entry(n, updated_at=f"2026-06-{n:02d}T00:00:00Z") for n in range(1, 101)] + short_page = [_entry(101, updated_at="2026-07-01T00:00:00Z")] + backend = _FakeBackend({1: full_page, 2: short_page}) + sync = IssueIndexSync(settings=_SyncSettings(), db=db, github=backend) # type: ignore[arg-type] + + ingested = await sync.sync_repo("octo/widget") + assert ingested == 101 + # First run is a backfill: no `since` on any call, pages walked in order. + assert backend.calls == [(None, 1), (None, 2)] + watermark = db.issue_index_watermark("octo/widget") + assert watermark is not None + assert db.search_issue_index("octo/widget", keywords=("issue",), limit=5) + + # Second run is incremental: `since` derives from the stored watermark. + backend.calls.clear() + backend.pages = {1: []} + await sync.sync_repo("octo/widget") + assert backend.calls and backend.calls[0][0] is not None diff --git a/python/robomp/tests/test_server.py b/python/robomp/tests/test_server.py index f4c492b5b..7141fddc3 100644 --- a/python/robomp/tests/test_server.py +++ b/python/robomp/tests/test_server.py @@ -958,6 +958,42 @@ def test_webhook_incoming_pr_comment_without_directive_skips_without_counting_bu assert states == ["queued", "queued", "skipped"] +def test_webhook_delivery_populates_issue_index(settings: Settings) -> None: + """Every issue-carrying delivery upserts the local search index — including + ones the router skips (here: a conversation comment on an incoming PR).""" + app = create_app(settings) + with TestClient(app) as client: + payload = { + "action": "opened", + "issue": { + "number": 501, + "title": "grep misses colon filenames", + "body": "read tool peels the selector suffix", + "state": "open", + "user": {"login": "alice"}, + "author_association": "NONE", + }, + "repository": {"full_name": "octo/widget"}, + } + body = json.dumps(payload).encode() + resp = client.post( + "/webhook/github", + content=body, + headers=_signed_headers("test-webhook-secret", body, event="issues", delivery="idx-1"), + ) + assert resp.status_code == 202 + + skipped = _post_pr_issue_comment(client, delivery="idx-2", user="stranger", pr_number=502) + assert skipped.json()["state"] == "skipped" + + db = get_database(settings.sqlite_path) + by_body = db.search_issue_index("octo/widget", keywords=("selector", "suffix")) + pr_row = db.search_issue_index("octo/widget", is_pr=True) + close_database() + assert [e.number for e in by_body] == [501] + assert [e.number for e in pr_row] == [502] + + def test_webhook_contributor_gets_higher_cap(rate_limited_settings: Settings) -> None: app = create_app(rate_limited_settings) with TestClient(app) as client: