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.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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]: ...
|
||||
|
||||
|
||||
@@ -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",
|
||||
]
|
||||
|
||||
+164
-23
@@ -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),
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -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:<name>`, `author:<login>`. 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",
|
||||
]
|
||||
@@ -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:<login>`. 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:<login>`. 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."
|
||||
|
||||
|
||||
@@ -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 <keywords>`) 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 <keywords>`), 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`
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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",
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user