From 5702a9d45256529a63f2862113d10287e83802e5 Mon Sep 17 00:00:00 2001 From: Johannes du Plessis Date: Tue, 26 May 2026 16:24:34 -0700 Subject: [PATCH] feat: reconcile reviewer comment lifecycle (#1332) * feat: reconcile reviewer comment lifecycle Track GitHub review threads for reviewer findings so re-reviews can resolve or reply to existing comments, and collect thumbs feedback on new review comments in LangSmith. * fix: clarify reviewer comment lifecycle * chore: apply reviewer formatting --------- Co-authored-by: open-swe[bot] --- agent/reviewer.py | 27 ++- agent/reviewer_findings.py | 14 ++ agent/reviewer_publish.py | 68 +++++++- agent/reviewer_reconcile.py | 150 ++++++++++++++++ agent/tools/__init__.py | 4 + agent/tools/publish_review.py | 127 +++++++++++--- agent/tools/reply_to_finding_thread.py | 80 +++++++++ agent/tools/resolve_finding_thread.py | 98 +++++++++++ agent/utils/github_feedback.py | 233 +++++++++++++++++++++++++ agent/webapp.py | 15 +- tests/test_github_feedback.py | 184 +++++++++++++++++++ tests/test_reviewer_findings.py | 4 + tests/test_reviewer_publish.py | 55 +++++- tests/test_reviewer_reconcile.py | 75 ++++++++ 14 files changed, 1094 insertions(+), 40 deletions(-) create mode 100644 agent/reviewer_reconcile.py create mode 100644 agent/tools/reply_to_finding_thread.py create mode 100644 agent/tools/resolve_finding_thread.py create mode 100644 agent/utils/github_feedback.py create mode 100644 tests/test_github_feedback.py create mode 100644 tests/test_reviewer_reconcile.py diff --git a/agent/reviewer.py b/agent/reviewer.py index 0d4996d2..acc99807 100644 --- a/agent/reviewer.py +++ b/agent/reviewer.py @@ -43,6 +43,7 @@ from .reviewer_findings import ( list_findings as list_findings_async, ) from .reviewer_publish import fetch_pr_review_threads +from .reviewer_reconcile import reconcile_findings_with_review_threads from .server import ( DEFAULT_LLM_MAX_TOKENS, DEFAULT_RECURSION_LIMIT, @@ -56,6 +57,8 @@ from .tools import ( http_request, list_findings, publish_review, + reply_to_finding_thread, + resolve_finding_thread, update_finding, web_search, ) @@ -87,13 +90,22 @@ Clone the repo so you can grep for full file context: GH_TOKEN=dummy gh repo clone {repo_owner}/{repo_name} && cd {repo_name} && git checkout ``` -Tools: `add_finding`, `update_finding`, `list_findings`, `publish_review`. +Tools: `add_finding`, `update_finding`, `list_findings`, `publish_review`, +`resolve_finding_thread`, `reply_to_finding_thread`. Call `publish_review` once at the end. Re-review: for each open finding, `update_finding(id, status="resolved")` if fixed, `update_finding` with new fields + `note` if changed, otherwise do nothing. Add net-new findings with `add_finding`. +If a human reply shows one of your published findings is invalid, call +`resolve_finding_thread(finding_id, status="dismissed")` after verifying the +claim. If the finding is fixed by code, use `update_finding(..., +status="resolved")`; `publish_review` will close the GitHub thread. Reply with +`reply_to_finding_thread` only when the user directly asks a question or a short +clarification is needed after pushback. Bias strongly toward resolving/dismissing +without replying. + # The bar: file a finding only if it passes these criteria 1. You can anchor it to a specific changed line and quote that line. @@ -341,7 +353,11 @@ def _build_re_review_context( f"For each open finding above, decide whether the new commits resolved " f'it (`update_finding(id, status="resolved")`), left it unchanged ' f"(no action), or changed it materially (`update_finding` with new " - f"fields + a `note`). Then add any net-new findings introduced by the " + f"fields + a `note`). If a human reply on a finding explains why your " + f"comment was invalid, verify that analysis, then call " + f'`resolve_finding_thread(id, status="dismissed")` to close it. ' + f"Reply only when directly asked or when a concise clarification is " + f"necessary. Then add any net-new findings introduced by the " f"new diff — but skip anything already covered by an existing PR " f"review thread above (your own prior threads, another reviewer's, or " f"one a human has already replied to). Call `publish_review` once at " @@ -468,6 +484,10 @@ def _format_existing_findings(findings: list[dict]) -> str: f"- [{f.get('id')}] ({f.get('severity')}, {f.get('category')}) " f"{location} — {f.get('description', '').strip()}" ) + human_reply = f.get("last_human_reply_body") + if isinstance(human_reply, str) and human_reply: + author = f.get("last_human_reply_author") or "human" + lines.append(f" Human reply from {author}: {human_reply}") return "\n".join(lines) if lines else "_(no open findings)_" @@ -533,6 +553,7 @@ async def get_reviewer_agent(config: RunnableConfig) -> Pregel: pr_number=pr_number, token=github_token, ) + await reconcile_findings_with_review_threads(thread_id, threads) existing_threads_block = _format_pr_review_threads(threads) if existing_threads_block: logger.info( @@ -650,6 +671,8 @@ async def get_reviewer_agent(config: RunnableConfig) -> Pregel: update_finding, list_findings, publish_review, + resolve_finding_thread, + reply_to_finding_thread, web_search, fetch_url, http_request, diff --git a/agent/reviewer_findings.py b/agent/reviewer_findings.py index fddd7e03..127d64cf 100644 --- a/agent/reviewer_findings.py +++ b/agent/reviewer_findings.py @@ -77,6 +77,13 @@ class Finding(TypedDict, total=False): first_seen_sha: str last_confirmed_sha: str github_review_comment_id: int | None + github_review_thread_id: str | None + github_review_run_id: str | None + github_thread_resolved: bool + last_human_reply_at: str | None + last_human_reply_author: str | None + last_human_reply_body: str | None + last_reconciliation_note: str | None diff_hunk: str | None @@ -135,6 +142,13 @@ def new_finding( "first_seen_sha": sha, "last_confirmed_sha": sha, "github_review_comment_id": None, + "github_review_thread_id": None, + "github_review_run_id": None, + "github_thread_resolved": False, + "last_human_reply_at": None, + "last_human_reply_author": None, + "last_human_reply_body": None, + "last_reconciliation_note": None, "diff_hunk": diff_hunk, } diff --git a/agent/reviewer_publish.py b/agent/reviewer_publish.py index 7e6b6172..53879782 100644 --- a/agent/reviewer_publish.py +++ b/agent/reviewer_publish.py @@ -20,6 +20,7 @@ the GraphQL ``resolveReviewThread`` mutation (REST doesn't expose this). from __future__ import annotations +import json import logging from typing import Any @@ -51,10 +52,19 @@ def render_inline_comment_body(finding: Finding) -> str: Multi-line suggestions just become multi-line ```suggestion``` blocks. """ description = finding.get("description", "") or "" + marker_payload = { + "id": finding.get("id", ""), + "file_path": finding.get("file", ""), + "start_line": finding.get("start_line"), + "end_line": finding.get("end_line"), + "side": finding.get("side", "RIGHT"), + } + marker = f"" suggestion = finding.get("suggestion") - if not suggestion: - return description - return f"{description}\n\n```suggestion\n{suggestion}\n```" + body = f"{marker}\n\n{description}" + if suggestion: + body = f"{body}\n\n```suggestion\n{suggestion}\n```" + return body def render_inline_comment_payload(finding: Finding) -> dict[str, Any] | None: @@ -230,6 +240,7 @@ async def fetch_pr_review_threads( reviewThreads(first: 50, after: $cursor) { pageInfo { hasNextPage endCursor } nodes { + id isResolved isOutdated path @@ -237,7 +248,9 @@ async def fetch_pr_review_threads( originalLine comments(first: $perThread) { nodes { + databaseId author { login } + authorAssociation body createdAt } @@ -297,7 +310,13 @@ async def fetch_pr_review_threads( login = author_block.get("login") if isinstance(author_block, dict) else None comments.append( { + "id": c.get("databaseId") + if isinstance(c.get("databaseId"), int) + else None, "author": login if isinstance(login, str) else "unknown", + "author_association": c.get("authorAssociation", "") + if isinstance(c.get("authorAssociation"), str) + else "", "body": c.get("body", "") if isinstance(c.get("body"), str) else "", "created_at": c.get("createdAt", "") if isinstance(c.get("createdAt"), str) @@ -306,6 +325,7 @@ async def fetch_pr_review_threads( ) out.append( { + "id": thread.get("id") if isinstance(thread.get("id"), str) else "", "path": thread.get("path", "") if isinstance(thread.get("path"), str) else "", @@ -432,6 +452,48 @@ async def resolve_review_thread(*, thread_node_id: str, token: str) -> bool: return bool(thread.get("isResolved")) +async def reply_to_review_comment( + *, + owner: str, + repo: str, + pr_number: int, + review_comment_id: int, + body: str, + token: str, +) -> dict[str, Any] | None: + """Reply to an existing pull request review comment thread.""" + url = ( + f"{_GITHUB_API_BASE}/repos/{owner}/{repo}/pulls/" + f"{pr_number}/comments/{review_comment_id}/replies" + ) + async with httpx.AsyncClient() as client: + try: + response = await client.post( + url, + headers=_github_headers(token), + json={"body": body}, + timeout=30, + ) + if response.status_code == 401: + raise GitHubAuthError( + f"GitHub returned 401 replying to review comment {review_comment_id}" + ) + response.raise_for_status() + except GitHubAuthError: + raise + except httpx.HTTPError: + logger.exception( + "Failed to reply to review comment %s on %s/%s#%s", + review_comment_id, + owner, + repo, + pr_number, + ) + return None + data = response.json() + return data if isinstance(data, dict) else None + + def _github_headers(token: str) -> dict[str, str]: return { "Authorization": f"Bearer {token}", diff --git a/agent/reviewer_reconcile.py b/agent/reviewer_reconcile.py new file mode 100644 index 00000000..51609779 --- /dev/null +++ b/agent/reviewer_reconcile.py @@ -0,0 +1,150 @@ +from __future__ import annotations + +from typing import Any + +from .reviewer_findings import Finding, list_findings, replace_findings + +ReviewThread = dict[str, Any] + + +def _human_replies_after_bot_comment( + review_thread: ReviewThread, + *, + bot_comment_id: int, +) -> list[ReviewThread]: + comments = review_thread.get("comments") + if not isinstance(comments, list): + return [] + + seen_bot_comment = False + replies: list[ReviewThread] = [] + for comment in comments: + if not isinstance(comment, dict): + continue + comment_id = comment.get("id") + if comment_id == bot_comment_id: + seen_bot_comment = True + continue + if not seen_bot_comment: + continue + author = comment.get("author") + if author == "open-swe[bot]": + continue + replies.append(comment) + return replies + + +def _index_review_threads( + review_threads: list[ReviewThread], +) -> tuple[dict[str, ReviewThread], dict[int, ReviewThread]]: + by_thread_id = { + thread_id: review_thread + for review_thread in review_threads + if isinstance(thread_id := review_thread.get("id"), str) and thread_id + } + by_comment_id: dict[int, ReviewThread] = {} + for review_thread in review_threads: + comments = review_thread.get("comments") + if not isinstance(comments, list): + continue + for comment in comments: + if not isinstance(comment, dict): + continue + comment_id = comment.get("id") + if isinstance(comment_id, int): + by_comment_id[comment_id] = review_thread + return by_thread_id, by_comment_id + + +def _find_review_thread_for_finding( + finding: Finding, + *, + by_thread_id: dict[str, ReviewThread], + by_comment_id: dict[int, ReviewThread], +) -> ReviewThread | None: + github_thread_id = finding.get("github_review_thread_id") + if isinstance(github_thread_id, str) and github_thread_id: + review_thread = by_thread_id.get(github_thread_id) + if review_thread is not None: + return review_thread + + github_comment_id = finding.get("github_review_comment_id") + if isinstance(github_comment_id, int): + return by_comment_id.get(github_comment_id) + return None + + +def _sync_thread_status(finding: Finding, review_thread: ReviewThread) -> bool: + updated = False + github_thread_id = finding.get("github_review_thread_id") + if not isinstance(github_thread_id, str) or not github_thread_id: + new_thread_id = review_thread.get("id") + if isinstance(new_thread_id, str) and new_thread_id: + finding["github_review_thread_id"] = new_thread_id + updated = True + + if review_thread.get("is_resolved") or review_thread.get("is_outdated"): + if finding.get("status") == "open": + finding["status"] = "resolved" + finding["last_reconciliation_note"] = "GitHub thread is resolved or outdated." + updated = True + if review_thread.get("is_resolved") and not finding.get("github_thread_resolved"): + finding["github_thread_resolved"] = True + updated = True + return updated + + +def _sync_latest_human_reply(finding: Finding, review_thread: ReviewThread) -> bool: + github_comment_id = finding.get("github_review_comment_id") + if not isinstance(github_comment_id, int): + return False + + replies = _human_replies_after_bot_comment(review_thread, bot_comment_id=github_comment_id) + if not replies: + return False + + latest = replies[-1] + body = latest.get("body") if isinstance(latest.get("body"), str) else "" + if len(body) > 1000: + body = body[:1000] + "\n...[truncated]" + created_at = latest.get("created_at") if isinstance(latest.get("created_at"), str) else "" + if finding.get("last_human_reply_at") == created_at: + return False + + author = latest.get("author") if isinstance(latest.get("author"), str) else "" + finding["last_human_reply_at"] = created_at + finding["last_human_reply_author"] = author + finding["last_human_reply_body"] = body + finding["last_reconciliation_note"] = ( + "Human replied to this review thread; reassess before taking action." + ) + return True + + +async def reconcile_findings_with_review_threads( + reviewer_thread_id: str, + review_threads: list[ReviewThread], +) -> list[Finding]: + """Sync tracked Open SWE findings with the current GitHub review-thread state.""" + findings = await list_findings(reviewer_thread_id) + if not findings: + return findings + + by_thread_id, by_comment_id = _index_review_threads(review_threads) + + updated = False + for finding in findings: + review_thread = _find_review_thread_for_finding( + finding, + by_thread_id=by_thread_id, + by_comment_id=by_comment_id, + ) + if review_thread is None: + continue + + updated = _sync_thread_status(finding, review_thread) or updated + updated = _sync_latest_human_reply(finding, review_thread) or updated + + if updated: + await replace_findings(reviewer_thread_id, findings) + return findings diff --git a/agent/tools/__init__.py b/agent/tools/__init__.py index 144dace3..ec4304e4 100644 --- a/agent/tools/__init__.py +++ b/agent/tools/__init__.py @@ -10,7 +10,9 @@ from .linear_list_teams import linear_list_teams from .linear_update_issue import linear_update_issue from .list_findings import list_findings from .publish_review import publish_review +from .reply_to_finding_thread import reply_to_finding_thread from .request_pr_review import request_pr_review +from .resolve_finding_thread import resolve_finding_thread from .slack_read_thread_messages import slack_read_thread_messages from .slack_thread_reply import slack_thread_reply from .update_finding import update_finding @@ -30,6 +32,8 @@ __all__ = [ "list_findings", "publish_review", "request_pr_review", + "reply_to_finding_thread", + "resolve_finding_thread", "slack_read_thread_messages", "slack_thread_reply", "update_finding", diff --git a/agent/tools/publish_review.py b/agent/tools/publish_review.py index e31dd9f7..4118b241 100644 --- a/agent/tools/publish_review.py +++ b/agent/tools/publish_review.py @@ -20,6 +20,7 @@ from ..reviewer_findings import ( list_findings as list_findings_async, ) from ..reviewer_publish import ( + fetch_pr_review_threads, fetch_review_comments, fetch_review_thread_id_for_comment, post_pull_request_review, @@ -43,25 +44,13 @@ def publish_review( Call this once at the end of a review run, after you have finished adding findings (and, on a re-review, after marking resolved findings via - ``update_finding``). It will: + ``update_finding``). The tool posts one GitHub PR Review for eligible + inline findings, records the GitHub comment/thread IDs for future + re-reviews, resolves GitHub threads for findings now marked resolved, and + advances the reviewer thread's ``last_reviewed_sha``. - 1. Read findings from the reviewer thread. - 2. Filter to status=open and severity ≥ ``severity_threshold``, capped - at ``cap`` to avoid review spam. - 3. POST a single GitHub PR Review with the eligible findings as inline - comments. ``finding.suggestion`` becomes a ```suggestion``` block - (the "Commit suggestion" UX). The review body is a fixed, - host-formatted summary line — you do not write it. On a re-review - run with no new findings to surface, the GitHub Review post is - skipped entirely (resolved threads and ``last_reviewed_sha`` are - still updated). The "no issues found" summary only posts on the - first review of a PR. - 4. Store the returned per-comment IDs back on each finding so a future - re-review can resolve those threads on GitHub when the issues are fixed. - 5. For findings whose status moved ``open`` → ``resolved`` since the last - publish, resolve their existing GitHub review threads via the GraphQL - ``resolveReviewThread`` mutation. - 6. Update ``last_reviewed_sha`` on the thread to the current head SHA. + On a re-review with no new findings to surface, it skips posting a new + GitHub Review but still resolves fixed threads and updates reviewer state. Args: severity_threshold: Lowest severity to surface to GitHub (default @@ -118,6 +107,7 @@ def publish_review( severity_threshold=_cast_severity(severity_threshold), cap=cap, is_re_review=is_re_review, + langgraph_run_id=_current_run_id(config), ) ) except GitHubAuthError as exc: @@ -188,6 +178,7 @@ async def _publish_review_async( severity_threshold: Severity, cap: int, is_re_review: bool, + langgraph_run_id: str | None = None, ) -> dict[str, Any]: thread_id = get_thread_id_from_runtime() findings = await list_findings_async(thread_id) @@ -272,11 +263,24 @@ async def _publish_review_async( review_id=review_id, token=token, ) + if langgraph_run_id is None: + metadata = await get_thread_metadata(thread_id) + current_run_id = metadata.get("current_reviewer_run_id") + if isinstance(current_run_id, str) and current_run_id: + langgraph_run_id = current_run_id await _store_comment_ids_on_findings( thread_id=thread_id, findings=findings, eligible_with_payload=eligible_with_payload, comment_records=comment_records, + langgraph_run_id=langgraph_run_id, + ) + await _store_thread_ids_on_findings( + thread_id=thread_id, + owner=owner, + repo=repo, + pr_number=pr_number, + token=token, ) resolved_thread_count = await _resolve_threads_for_resolved_findings( @@ -348,6 +352,7 @@ async def _store_comment_ids_on_findings( findings: list[dict[str, Any]], eligible_with_payload: list[tuple[dict[str, Any], dict[str, Any]]], comment_records: list[dict[str, Any]], + langgraph_run_id: str | None, ) -> None: """Match returned GitHub comment ids back to the findings that produced them. @@ -386,12 +391,71 @@ async def _store_comment_ids_on_findings( if finding is None: continue finding["github_review_comment_id"] = comment_id + if langgraph_run_id: + finding["github_review_run_id"] = langgraph_run_id updated = True if updated: await replace_findings(thread_id, list(findings_by_id.values())) +async def _store_thread_ids_on_findings( + *, + thread_id: str, + owner: str, + repo: str, + pr_number: int, + token: str, +) -> None: + findings = await list_findings_async(thread_id) + comment_ids_by_finding_id: dict[str, int] = {} + for finding in findings: + finding_id = finding.get("id") + comment_id = finding.get("github_review_comment_id") + if ( + isinstance(finding_id, str) + and isinstance(comment_id, int) + and not isinstance(finding.get("github_review_thread_id"), str) + ): + comment_ids_by_finding_id[finding_id] = comment_id + if not comment_ids_by_finding_id: + return + + threads = await fetch_pr_review_threads( + owner=owner, + repo=repo, + pr_number=pr_number, + token=token, + ) + thread_id_by_comment_id: dict[int, str] = {} + for thread in threads: + github_thread_id = thread.get("id") + if not isinstance(github_thread_id, str) or not github_thread_id: + continue + for comment in thread.get("comments") or []: + if not isinstance(comment, dict): + continue + comment_id = comment.get("id") + if isinstance(comment_id, int): + thread_id_by_comment_id[comment_id] = github_thread_id + + updated = False + for finding in findings: + finding_id = finding.get("id") + if not isinstance(finding_id, str): + continue + comment_id = comment_ids_by_finding_id.get(finding_id) + if not isinstance(comment_id, int): + continue + github_thread_id = thread_id_by_comment_id.get(comment_id) + if github_thread_id: + finding["github_review_thread_id"] = github_thread_id + updated = True + + if updated: + await replace_findings(thread_id, findings) + + async def _resolve_threads_for_resolved_findings( *, owner: str, @@ -418,13 +482,15 @@ async def _resolve_threads_for_resolved_findings( continue if finding.get("github_thread_resolved"): continue - thread_node_id = await fetch_review_thread_id_for_comment( - owner=owner, - repo=repo, - pr_number=pr_number, - review_comment_id=comment_id, - token=token, - ) + thread_node_id = finding.get("github_review_thread_id") + if not isinstance(thread_node_id, str) or not thread_node_id: + thread_node_id = await fetch_review_thread_id_for_comment( + owner=owner, + repo=repo, + pr_number=pr_number, + review_comment_id=comment_id, + token=token, + ) if not thread_node_id: continue ok = await resolve_review_thread(thread_node_id=thread_node_id, token=token) @@ -438,3 +504,14 @@ async def _resolve_threads_for_resolved_findings( await replace_findings(thread_id, findings) return resolved_count + + +def _current_run_id(config: dict[str, Any]) -> str | None: + candidates = [config.get("run_id")] + configurable = config.get("configurable") + if isinstance(configurable, dict): + candidates.append(configurable.get("run_id")) + for candidate in candidates: + if isinstance(candidate, str) and candidate: + return candidate + return None diff --git a/agent/tools/reply_to_finding_thread.py b/agent/tools/reply_to_finding_thread.py new file mode 100644 index 00000000..b48f5ae9 --- /dev/null +++ b/agent/tools/reply_to_finding_thread.py @@ -0,0 +1,80 @@ +from __future__ import annotations + +import asyncio +from typing import Any + +from langgraph.config import get_config + +from ..reviewer_findings import get_finding, get_thread_id_from_runtime, update_finding_fields +from ..reviewer_publish import reply_to_review_comment +from ..utils.github_token import get_github_token + + +def reply_to_finding_thread(finding_id: str, body: str) -> dict[str, Any]: + """Reply to the GitHub review thread for a tracked finding.""" + if not body.strip(): + return {"success": False, "error": "Reply body is required"} + + config = get_config() + configurable = config.get("configurable", {}) if isinstance(config, dict) else {} + repo_config = configurable.get("repo") if isinstance(configurable, dict) else None + pr_number = configurable.get("pr_number") if isinstance(configurable, dict) else None + if ( + not isinstance(repo_config, dict) + or not repo_config.get("owner") + or not repo_config.get("name") + or not isinstance(pr_number, int) + ): + return {"success": False, "error": "Missing repo or PR info in run config"} + + token = get_github_token() + if not token: + return {"success": False, "error": "No GitHub token available"} + + return asyncio.run( + _reply_to_finding_thread_async( + finding_id=finding_id, + body=body, + owner=str(repo_config["owner"]), + repo=str(repo_config["name"]), + pr_number=pr_number, + token=token, + ) + ) + + +async def _reply_to_finding_thread_async( + *, + finding_id: str, + body: str, + owner: str, + repo: str, + pr_number: int, + token: str, +) -> dict[str, Any]: + thread_id = get_thread_id_from_runtime() + finding = await get_finding(thread_id, finding_id) + if finding is None: + return {"success": False, "error": f"No finding found with id {finding_id}"} + + comment_id = finding.get("github_review_comment_id") + if not isinstance(comment_id, int): + return {"success": False, "error": "Finding has no GitHub review comment mapping"} + + response = await reply_to_review_comment( + owner=owner, + repo=repo, + pr_number=pr_number, + review_comment_id=comment_id, + body=body.strip(), + token=token, + ) + if response is None: + return {"success": False, "error": "GitHub did not accept the reply"} + + reply_id = response.get("id") + updates: dict[str, Any] = {"last_reconciliation_note": "Replied to GitHub review thread."} + if isinstance(reply_id, int): + updates["last_review_reply_comment_id"] = reply_id + updated = await update_finding_fields(thread_id, finding_id, updates) + return {"success": True, "finding": updated, "reply_id": reply_id} diff --git a/agent/tools/resolve_finding_thread.py b/agent/tools/resolve_finding_thread.py new file mode 100644 index 00000000..8fbb5e2b --- /dev/null +++ b/agent/tools/resolve_finding_thread.py @@ -0,0 +1,98 @@ +from __future__ import annotations + +import asyncio +from typing import Any + +from langgraph.config import get_config + +from ..reviewer_findings import get_finding, get_thread_id_from_runtime, update_finding_fields +from ..reviewer_publish import fetch_review_thread_id_for_comment, resolve_review_thread +from ..utils.github_token import get_github_token + + +def resolve_finding_thread( + finding_id: str, + status: str = "dismissed", + note: str | None = None, +) -> dict[str, Any]: + """Resolve the GitHub review thread for a tracked Open SWE finding. + + Use ``status="resolved"`` when the code now fixes the issue. Use + ``status="dismissed"`` when analysis shows the original review comment was + not valid. + """ + if status not in {"resolved", "dismissed"}: + return {"success": False, "error": f"Invalid status: {status}"} + + config = get_config() + configurable = config.get("configurable", {}) if isinstance(config, dict) else {} + repo_config = configurable.get("repo") if isinstance(configurable, dict) else None + pr_number = configurable.get("pr_number") if isinstance(configurable, dict) else None + if ( + not isinstance(repo_config, dict) + or not repo_config.get("owner") + or not repo_config.get("name") + or not isinstance(pr_number, int) + ): + return {"success": False, "error": "Missing repo or PR info in run config"} + + token = get_github_token() + if not token: + return {"success": False, "error": "No GitHub token available"} + + return asyncio.run( + _resolve_finding_thread_async( + finding_id=finding_id, + status=status, + note=note, + owner=str(repo_config["owner"]), + repo=str(repo_config["name"]), + pr_number=pr_number, + token=token, + ) + ) + + +async def _resolve_finding_thread_async( + *, + finding_id: str, + status: str, + note: str | None, + owner: str, + repo: str, + pr_number: int, + token: str, +) -> dict[str, Any]: + thread_id = get_thread_id_from_runtime() + finding = await get_finding(thread_id, finding_id) + if finding is None: + return {"success": False, "error": f"No finding found with id {finding_id}"} + + github_thread_id = finding.get("github_review_thread_id") + if not isinstance(github_thread_id, str) or not github_thread_id: + comment_id = finding.get("github_review_comment_id") + if not isinstance(comment_id, int): + return {"success": False, "error": "Finding has no GitHub review thread mapping"} + github_thread_id = await fetch_review_thread_id_for_comment( + owner=owner, + repo=repo, + pr_number=pr_number, + review_comment_id=comment_id, + token=token, + ) + if not github_thread_id: + return {"success": False, "error": "Could not resolve GitHub review thread id"} + + ok = await resolve_review_thread(thread_node_id=github_thread_id, token=token) + if not ok: + return {"success": False, "error": "GitHub did not resolve the review thread"} + + updates: dict[str, Any] = { + "status": status, + "github_review_thread_id": github_thread_id, + "github_thread_resolved": True, + } + if note: + updates["last_reconciliation_note"] = note + updated = await update_finding_fields(thread_id, finding_id, updates) + return {"success": True, "finding": updated} diff --git a/agent/utils/github_feedback.py b/agent/utils/github_feedback.py new file mode 100644 index 00000000..d6e4f9b8 --- /dev/null +++ b/agent/utils/github_feedback.py @@ -0,0 +1,233 @@ +from __future__ import annotations + +import asyncio +import logging +import os +import re +import uuid +from typing import Any + +from langgraph_sdk import get_client +from langgraph_sdk.client import LangGraphClient + +from ..reviewer_findings import list_findings +from .langsmith import create_langsmith_feedback, delete_langsmith_feedback + +logger = logging.getLogger(__name__) + +LANGGRAPH_URL = os.environ.get("LANGGRAPH_URL") or os.environ.get( + "LANGGRAPH_URL_PROD", "http://localhost:2024" +) + +GITHUB_FEEDBACK_REACTIONS: dict[str, float] = { + "+1": 1.0, + "-1": 0.0, +} + +_REACTION_STATE_NAMESPACE = "github_reaction_state" +_REACTION_EVENT_NAMESPACE = "github_reaction_events" +_PULL_URL_RE = re.compile(r"/pulls/(\d+)\Z") + + +def _reviewer_thread_id(owner: str, repo: str, pr_number: int) -> str: + return str(uuid.uuid5(uuid.NAMESPACE_URL, f"{owner}/{repo}/pr/{pr_number}/reviewer")) + + +def _read_active_reactions(item: dict[str, Any] | None) -> set[str]: + if not item: + return set() + value = item.get("value") + if not isinstance(value, dict): + return set() + reactions = value.get("reactions") + if not isinstance(reactions, list): + return set() + return {reaction for reaction in reactions if isinstance(reaction, str)} + + +def _reaction_state_key(run_id: str, user_login: str, comment_id: int) -> str: + return f"{run_id}:{user_login}:{comment_id}" + + +def _feedback_key(owner: str, repo: str, user_login: str, comment_id: int) -> str: + return f"github_reaction:{owner}/{repo}:{user_login}:{comment_id}" + + +def _score_reactions(reactions: set[str]) -> float | None: + scores = { + GITHUB_FEEDBACK_REACTIONS[reaction] + for reaction in reactions + if reaction in GITHUB_FEEDBACK_REACTIONS + } + if len(scores) != 1: + return None + return next(iter(scores)) + + +def _extract_pr_number(payload: dict[str, Any]) -> int | None: + pull_request = payload.get("pull_request") + if isinstance(pull_request, dict) and isinstance(pull_request.get("number"), int): + return pull_request["number"] + + comment = payload.get("comment") + if isinstance(comment, dict): + url = comment.get("pull_request_url") + if isinstance(url, str): + match = _PULL_URL_RE.search(url) + if match: + return int(match.group(1)) + return None + + +async def _event_was_processed( + langgraph_client: LangGraphClient, repo_key: str, event_id: str +) -> bool: + if not event_id: + return False + item = await langgraph_client.store.get_item((_REACTION_EVENT_NAMESPACE, repo_key), event_id) + return bool(item) + + +async def _mark_event_processed( + langgraph_client: LangGraphClient, repo_key: str, event_id: str +) -> None: + if not event_id: + return + await langgraph_client.store.put_item( + (_REACTION_EVENT_NAMESPACE, repo_key), event_id, {"event_id": event_id} + ) + + +async def _update_reaction_state( + langgraph_client: LangGraphClient, + *, + repo_key: str, + run_id: str, + user_login: str, + comment_id: int, + reaction: str, + added: bool, +) -> set[str]: + namespace = (_REACTION_STATE_NAMESPACE, repo_key) + key = _reaction_state_key(run_id, user_login, comment_id) + item = await langgraph_client.store.get_item(namespace, key) + active_reactions = _read_active_reactions(item) + if added: + active_reactions.add(reaction) + else: + active_reactions.discard(reaction) + if not active_reactions: + await langgraph_client.store.delete_item(namespace, key) + return active_reactions + await langgraph_client.store.put_item( + namespace, + key, + { + "run_id": run_id, + "user_login": user_login, + "comment_id": comment_id, + "reactions": sorted(active_reactions), + }, + ) + return active_reactions + + +async def process_github_reaction( + payload: dict[str, Any], + *, + delivery_id: str = "", + added: bool, +) -> None: + reaction = payload.get("reaction") + content = reaction.get("content") if isinstance(reaction, dict) else None + if not isinstance(content, str) or content not in GITHUB_FEEDBACK_REACTIONS: + return + + comment = payload.get("comment") + comment_id = comment.get("id") if isinstance(comment, dict) else None + if not isinstance(comment_id, int): + return + + repo = payload.get("repository") + owner = repo.get("owner", {}).get("login") if isinstance(repo, dict) else None + repo_name = repo.get("name") if isinstance(repo, dict) else None + pr_number = _extract_pr_number(payload) + sender = payload.get("sender") + user_login = sender.get("login") if isinstance(sender, dict) else None + if not ( + isinstance(owner, str) + and owner + and isinstance(repo_name, str) + and repo_name + and isinstance(pr_number, int) + and isinstance(user_login, str) + and user_login + ): + return + + langgraph_client = get_client(url=LANGGRAPH_URL) + repo_key = f"{owner}/{repo_name}" + if await _event_was_processed(langgraph_client, repo_key, delivery_id): + return + + thread_id = _reviewer_thread_id(owner, repo_name, pr_number) + findings = await list_findings(thread_id) + finding = next( + ( + candidate + for candidate in findings + if candidate.get("github_review_comment_id") == comment_id + ), + None, + ) + if finding is None: + logger.debug("No tracked finding for GitHub review comment id %s", comment_id) + return + + run_id = finding.get("github_review_run_id") + if not isinstance(run_id, str) or not run_id: + logger.debug("Finding %s has no LangSmith run id for feedback", finding.get("id")) + return + + active_reactions = await _update_reaction_state( + langgraph_client, + repo_key=repo_key, + run_id=run_id, + user_login=user_login, + comment_id=comment_id, + reaction=content, + added=added, + ) + + key = _feedback_key(owner, repo_name, user_login, comment_id) + source_info = { + "source": "github_review_reaction", + "owner": owner, + "repo": repo_name, + "pr_number": pr_number, + "comment_id": comment_id, + "finding_id": finding.get("id"), + "user_login": user_login, + } + score = _score_reactions(active_reactions) + if score is None: + success = await asyncio.to_thread(delete_langsmith_feedback, run_id, key) + else: + success = await asyncio.to_thread( + create_langsmith_feedback, + run_id, + key, + score=score, + comment=f"GitHub review reaction feedback from {user_login}", + source_info={**source_info, "reactions": sorted(active_reactions)}, + ) + if success: + await _mark_event_processed(langgraph_client, repo_key, delivery_id) + + +async def process_github_reaction_added(payload: dict[str, Any], delivery_id: str = "") -> None: + await process_github_reaction(payload, delivery_id=delivery_id, added=True) + + +async def process_github_reaction_removed(payload: dict[str, Any], delivery_id: str = "") -> None: + await process_github_reaction(payload, delivery_id=delivery_id, added=False) diff --git a/agent/webapp.py b/agent/webapp.py index 31928145..b1b2b484 100644 --- a/agent/webapp.py +++ b/agent/webapp.py @@ -1592,16 +1592,23 @@ async def trigger_pr_review_from_ref( return {"success": queued, "queued": queued, "thread_id": thread_id, "pr_url": pr_url} logger.info("Creating reviewer run for thread %s from %s PR review request", thread_id, source) - await langgraph_client.runs.create( + run = await langgraph_client.runs.create( thread_id, "reviewer", input={"messages": [{"role": "user", "content": prompt}]}, config={"configurable": configurable, "metadata": _AGENT_VERSION_METADATA}, if_not_exists="create", ) + await _store_current_reviewer_run_id(thread_id, run) return {"success": True, "queued": False, "thread_id": thread_id, "pr_url": pr_url} +async def _store_current_reviewer_run_id(thread_id: str, run: Any) -> None: + run_id = run.get("run_id") if isinstance(run, dict) else None + if isinstance(run_id, str) and run_id: + await set_reviewer_thread_metadata(thread_id, extra={"current_reviewer_run_id": run_id}) + + def _build_reviewer_configurable( *, source: str, @@ -1732,13 +1739,14 @@ async def _dispatch_first_review_from_pr_payload(payload: dict[str, Any], *, sou return logger.info("Creating reviewer run for thread %s (source=%s)", thread_id, source) - await langgraph_client.runs.create( + run = await langgraph_client.runs.create( thread_id, "reviewer", input={"messages": [{"role": "user", "content": prompt}]}, config={"configurable": configurable, "metadata": _AGENT_VERSION_METADATA}, if_not_exists="create", ) + await _store_current_reviewer_run_id(thread_id, run) logger.info("Reviewer run created for thread %s (source=%s)", thread_id, source) @@ -2064,13 +2072,14 @@ async def process_github_push_event(payload: dict[str, Any]) -> None: return logger.info("Creating push re-review run for thread %s", thread_id) - await langgraph_client.runs.create( + run = await langgraph_client.runs.create( thread_id, "reviewer", input={"messages": [{"role": "user", "content": re_review_prompt}]}, config={"configurable": configurable, "metadata": _AGENT_VERSION_METADATA}, if_not_exists="create", ) + await _store_current_reviewer_run_id(thread_id, run) async def _refresh_thread_github_token_after_401(thread_id: str, email: str) -> str | None: diff --git a/tests/test_github_feedback.py b/tests/test_github_feedback.py new file mode 100644 index 00000000..0408a290 --- /dev/null +++ b/tests/test_github_feedback.py @@ -0,0 +1,184 @@ +from __future__ import annotations + +import json +from typing import Any + +import pytest + +from agent import webapp +from agent.utils import github_feedback +from agent.utils.github_feedback import ( + process_github_reaction_added, + process_github_reaction_removed, +) + + +class _FakeStore: + def __init__(self) -> None: + self.items: dict[tuple[tuple[str, ...], str], dict[str, Any]] = {} + + async def get_item(self, namespace: tuple[str, ...], key: str) -> dict[str, Any] | None: + return self.items.get((namespace, key)) + + async def put_item(self, namespace: tuple[str, ...], key: str, value: dict[str, Any]) -> None: + self.items[(namespace, key)] = {"value": value} + + async def delete_item(self, namespace: tuple[str, ...], key: str) -> None: + self.items.pop((namespace, key), None) + + +class _FakeClient: + def __init__(self) -> None: + self.store = _FakeStore() + + +class _FakeBackgroundTasks: + def __init__(self) -> None: + self.tasks: list[tuple[Any, tuple[Any, ...]]] = [] + + def add_task(self, func: Any, *args: Any) -> None: + self.tasks.append((func, args)) + + +class _FakeRequest: + def __init__(self, payload: dict[str, Any]) -> None: + self.headers = { + "X-GitHub-Event": "reaction", + "X-GitHub-Delivery": "delivery-1", + "X-Hub-Signature-256": "sig", + } + self._body = json.dumps(payload).encode() + + async def body(self) -> bytes: + return self._body + + +def _reaction_payload(content: str = "+1", action: str = "created") -> dict[str, Any]: + return { + "action": action, + "reaction": {"content": content}, + "repository": {"owner": {"login": "langchain-ai"}, "name": "open-swe"}, + "pull_request": {"number": 7}, + "comment": {"id": 123, "pull_request_url": "https://api.github.com/repos/o/r/pulls/7"}, + "sender": {"login": "reviewer"}, + } + + +@pytest.mark.asyncio +async def test_github_reaction_added_creates_langsmith_feedback( + monkeypatch: pytest.MonkeyPatch, +) -> None: + client = _FakeClient() + created: dict[str, Any] = {} + + def fake_create_feedback( + run_id: str, + key: str, + *, + score: float, + comment: str | None = None, + source_info: dict[str, Any] | None = None, + ) -> bool: + created.update( + { + "run_id": run_id, + "key": key, + "score": score, + "comment": comment, + "source_info": source_info, + } + ) + return True + + monkeypatch.setattr(github_feedback, "get_client", lambda url: client) + + async def fake_list_findings(thread_id: str) -> list[dict[str, Any]]: + return [ + { + "id": "f1", + "github_review_comment_id": 123, + "github_review_run_id": "run-1", + } + ] + + monkeypatch.setattr( + github_feedback, + "list_findings", + fake_list_findings, + ) + monkeypatch.setattr(github_feedback, "create_langsmith_feedback", fake_create_feedback) + + await process_github_reaction_added(_reaction_payload(), delivery_id="delivery-1") + + assert created["run_id"] == "run-1" + assert created["key"] == "github_reaction:langchain-ai/open-swe:reviewer:123" + assert created["score"] == 1.0 + assert created["source_info"]["finding_id"] == "f1" + assert (("github_reaction_events", "langchain-ai/open-swe"), "delivery-1") in client.store.items + + +@pytest.mark.asyncio +async def test_github_reaction_removed_deletes_langsmith_feedback( + monkeypatch: pytest.MonkeyPatch, +) -> None: + client = _FakeClient() + client.store.items[ + (("github_reaction_state", "langchain-ai/open-swe"), "run-1:reviewer:123") + ] = { + "value": { + "run_id": "run-1", + "user_login": "reviewer", + "comment_id": 123, + "reactions": ["+1"], + } + } + deleted: dict[str, str] = {} + + def fake_delete_feedback(run_id: str, key: str) -> bool: + deleted["run_id"] = run_id + deleted["key"] = key + return True + + monkeypatch.setattr(github_feedback, "get_client", lambda url: client) + + async def fake_list_findings(thread_id: str) -> list[dict[str, Any]]: + return [ + { + "id": "f1", + "github_review_comment_id": 123, + "github_review_run_id": "run-1", + } + ] + + monkeypatch.setattr( + github_feedback, + "list_findings", + fake_list_findings, + ) + monkeypatch.setattr(github_feedback, "delete_langsmith_feedback", fake_delete_feedback) + + await process_github_reaction_removed( + _reaction_payload(action="deleted"), delivery_id="delivery-2" + ) + + assert deleted == { + "run_id": "run-1", + "key": "github_reaction:langchain-ai/open-swe:reviewer:123", + } + assert ( + ("github_reaction_state", "langchain-ai/open-swe"), + "run-1:reviewer:123", + ) not in client.store.items + + +@pytest.mark.asyncio +async def test_github_webhook_ignores_reaction_event(monkeypatch: pytest.MonkeyPatch) -> None: + payload = _reaction_payload() + background_tasks = _FakeBackgroundTasks() + + monkeypatch.setattr(webapp, "verify_github_signature", lambda *args, **kwargs: True) + + response = await webapp.github_webhook(_FakeRequest(payload), background_tasks) + + assert response == {"status": "ignored", "reason": "Unsupported event type: reaction"} + assert background_tasks.tasks == [] diff --git a/tests/test_reviewer_findings.py b/tests/test_reviewer_findings.py index 20f40a6b..8522cf5b 100644 --- a/tests/test_reviewer_findings.py +++ b/tests/test_reviewer_findings.py @@ -49,6 +49,10 @@ def test_new_finding_defaults() -> None: assert finding["first_seen_sha"] == "abc123" assert finding["last_confirmed_sha"] == "abc123" assert finding["github_review_comment_id"] is None + assert finding["github_review_thread_id"] is None + assert finding["github_review_run_id"] is None + assert finding["github_thread_resolved"] is False + assert finding["last_human_reply_at"] is None assert finding["suggestion"] is None diff --git a/tests/test_reviewer_publish.py b/tests/test_reviewer_publish.py index 3d494327..b2c1a5e9 100644 --- a/tests/test_reviewer_publish.py +++ b/tests/test_reviewer_publish.py @@ -14,6 +14,7 @@ from agent.reviewer_publish import ( render_inline_comment_body, render_inline_comment_payload, render_review_body, + reply_to_review_comment, resolve_review_thread, ) @@ -35,7 +36,10 @@ def _f(**overrides: Any) -> Finding: def test_render_inline_comment_body_without_suggestion() -> None: body = render_inline_comment_body(_f(description="just text")) - assert body == "just text" + assert "