From a132925f9b933de033f97e58f457d268553cbf18 Mon Sep 17 00:00:00 2001 From: Claude Lin & Lay Date: Thu, 23 Apr 2026 20:35:39 +0900 Subject: [PATCH] feat(ingest): index issue comments, PR reviews, PR inline comments for memory DB MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 判断履歴の主要容器であるコメント欄を ingest 対象に追加する。issue_comment / pull_request_review / pull_request_review_comment を新 type として扱い、 bot (`sender.login` が `[bot]` で終わる) と trim 後 10 文字未満の body は ingest 時点で除外する。埋め込み入力は `author\n\n{body}` とし、speaker 文脈を dense ベクトルに残す。search_issues の type enum に 3 値を追加し、scan mode も新規 recent endpoints を叩いて aggregation する。poller には per-repo 上位 20 parent issue/PR について comments / reviews / review_comments を backfill する経路を 追加し、API fan-out と embedding 回数を別々に cap する。 Closes #106 --- README.ja.md | 30 ++- README.md | 32 ++- src/mcp.ts | 212 +++++++++++++++++++- src/pipeline.ts | 474 ++++++++++++++++++++++++++++++++++++++++++++ src/poller.ts | 223 +++++++++++++++++++++ src/store.ts | 516 +++++++++++++++++++++++++++++++++++++++++++++++- src/types.ts | 68 ++++++- src/webhook.ts | 284 ++++++++++++++++++++++++++ 8 files changed, 1823 insertions(+), 16 deletions(-) diff --git a/README.ja.md b/README.ja.md index 4431f5c..29b8e86 100644 --- a/README.ja.md +++ b/README.ja.md @@ -85,7 +85,7 @@ GitHub webhooks + GitHub API ### `search_issues` -GitHub の issue / pull request / release / documentation / commit diff を対象にした統合検索ツールです。 +GitHub の issue / pull request / release / documentation / commit diff / comment 系 (issue と PR の top-level comment、PR review 本文、PR インラインレビューコメント) を対象にした統合検索ツールです。 `query` と `sort` の組み合わせで、以下の 3 モードを切り替えます。 @@ -95,6 +95,8 @@ GitHub の issue / pull request / release / documentation / commit diff を対 structured filter (`repo` / `state` / `labels` / `milestone` / `assignee` / `type`) はすべてのモードで有効です。 +bot (`sender.login` が `[bot]` で終わる) と trim 後 10 文字未満の body は ingest 時点で除外されます。`LGTM` / `+1` / CI ノイズなどは retrieval 面に残りません。 + #### パラメータ | 名前 | 型 | 説明 | @@ -105,7 +107,7 @@ structured filter (`repo` / `state` / `labels` / `milestone` / `assignee` / `typ | `labels` | string[] | label 名で AND 絞り込み。 | | `milestone` | string | milestone title で絞り込み。 | | `assignee` | string | assignee login で絞り込み。 | -| `type` | `"issue"` / `"pull_request"` / `"release"` / `"doc"` / `"diff"` / `"all"` | type で絞り込み (既定 `all`)。 | +| `type` | 下記参照 | type で絞り込み (既定 `all`)。 | | `top_k` | number | 最大件数 (既定 10、上限 50)。 | | `fusion` | `"rrf"` / `"dense_only"` / `"sparse_only"` | fusion 戦略 (既定 `rrf`)。scan モードでは無視。 | | `rerank` | boolean | cross-encoder rerank (既定 `true`)。scan モードでは無視。 | @@ -114,6 +116,20 @@ structured filter (`repo` / `state` / `labels` / `milestone` / `assignee` / `typ | `until` | ISO 8601 文字列 | `updated_at < until` の結果だけを残します。 | | `include_content` | boolean | 上位 doc 結果に本文を inline する (既定 `false`)。 | +#### `type` 値 + +| 値 | 対象 | +|----|------| +| `"issue"` | GitHub issue (title + body)。 | +| `"pull_request"` | pull request 本文 (title + body)。 | +| `"release"` | release notes (name + body)。 | +| `"doc"` | Markdown documentation。 | +| `"diff"` | commit の per-file diff (commit message + file path + patch)。 | +| `"issue_comment"` | issue と PR の top-level コメント。 | +| `"pr_review"` | PR レビュー本文 (`APPROVED` / `CHANGES_REQUESTED` / `COMMENTED`)。 | +| `"pr_review_comment"` | PR の per-line インラインレビューコメント。 | +| `"all"` | 上記すべての union (既定)。 | + #### 使用例 特定トピックの意味検索: @@ -147,6 +163,16 @@ structured filter (`repo` / `state` / `labels` / `milestone` / `assignee` / `typ } ``` +特定トピックに関する過去の PR review 判断を検索: + +```json +{ + "query": "rerank threshold tuning", + "type": "pr_review", + "top_k": 5 +} +``` + ## Repository Structure ```text diff --git a/README.md b/README.md index 31e00f7..2eb4495 100644 --- a/README.md +++ b/README.md @@ -85,16 +85,18 @@ This MCP server exposes a single consolidated tool. All retrieval modes — sema ### `search_issues` -Unified search across GitHub issues, pull requests, releases, repository documentation, and commit diffs. +Unified search across GitHub issues, pull requests, releases, repository documentation, commit diffs, and comment / review surfaces (top-level comments on issues and PRs, PR review bodies, and PR inline review comments). Three modes are selected by the combination of `query` and `sort`: 1. **Hybrid semantic search (default)** — dense BGE-M3 over Vectorize + sparse BM25 over D1 FTS5, fused via Reciprocal Rank Fusion (RRF, k=60), then re-scored with the `@cf/baai/bge-reranker-base` cross-encoder. Pass a natural-language `query`. -2. **Time-ordered activity scan** — omit or leave `query` empty and set `sort` to `"updated_desc"` or `"created_desc"`. Optionally narrow with `since` / `until` to list recent issue / PR / release / doc / diff activity. This subsumes the previous `list_recent_activity` tool. +2. **Time-ordered activity scan** — omit or leave `query` empty and set `sort` to `"updated_desc"` or `"created_desc"`. Optionally narrow with `since` / `until` to list recent activity across every type. This subsumes the previous `list_recent_activity` tool. 3. **Doc content fetch** — set `include_content: true`. For result rows whose `type` is `"doc"`, the raw file content is fetched from the GitHub contents API and inlined as a `content` field. Capped at the first few doc rows to bound API fan-out. This subsumes the previous `get_doc_content` tool. Structured filters (`repo`, `state`, `labels`, `milestone`, `assignee`, `type`) apply in every mode. +Bot-authored comments (`sender.login` ending in `[bot]`) and comments shorter than 10 characters (trimmed) are filtered out at ingest time so noise such as `LGTM`, `+1`, or CI chatter does not dilute the retrieval surface. + #### Parameters | Name | Type | Description | @@ -105,7 +107,7 @@ Structured filters (`repo`, `state`, `labels`, `milestone`, `assignee`, `type`) | `labels` | string[] | Filter by label names (AND). | | `milestone` | string | Filter by milestone title. | | `assignee` | string | Filter by assignee login. | -| `type` | `"issue"` \| `"pull_request"` \| `"release"` \| `"doc"` \| `"diff"` \| `"all"` | Filter by type (default `all`). | +| `type` | see below | Filter by type (default `all`). | | `top_k` | number | Max results (default 10, max 50). | | `fusion` | `"rrf"` \| `"dense_only"` \| `"sparse_only"` | Fusion strategy (default `rrf`). Ignored in scan mode. | | `rerank` | boolean | Cross-encoder rerank (default `true`). Ignored in scan mode. | @@ -114,6 +116,20 @@ Structured filters (`repo`, `state`, `labels`, `milestone`, `assignee`, `type`) | `until` | ISO 8601 string | Keep only results with `updated_at < until`. | | `include_content` | boolean | Inline raw content on top doc results (default `false`). | +#### `type` values + +| Value | Surface | +|-------|---------| +| `"issue"` | GitHub issues (title + body). | +| `"pull_request"` | Pull request descriptions (title + body). | +| `"release"` | Release notes (name + body). | +| `"doc"` | Markdown documentation files. | +| `"diff"` | Per-file commit diffs (commit message + file path + patch). | +| `"issue_comment"` | Top-level comments on issues and PRs. | +| `"pr_review"` | PR review bodies (`APPROVED` / `CHANGES_REQUESTED` / `COMMENTED`). | +| `"pr_review_comment"` | PR inline review comments (per-line diff comments). | +| `"all"` | Union of every type above (default). | + #### Examples Semantic search for a specific topic: @@ -147,6 +163,16 @@ Semantic search with inline doc content on the top doc hits: } ``` +Search past PR review judgments about a specific topic: + +```json +{ + "query": "rerank threshold tuning", + "type": "pr_review", + "top_k": 5 +} +``` + ## Repository Structure ```text diff --git a/src/mcp.ts b/src/mcp.ts index 4795429..745a8fc 100644 --- a/src/mcp.ts +++ b/src/mcp.ts @@ -15,7 +15,17 @@ import { McpAgent } from "agents/mcp"; import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import { z } from "zod"; -import type { Env, IssueRecord, ReleaseRecord, DocRecord, DiffRecord, VectorMetadata } from "./types.js"; +import type { + Env, + IssueRecord, + ReleaseRecord, + DocRecord, + DiffRecord, + IssueCommentRecord, + PRReviewRecord, + PRReviewCommentRecord, + VectorMetadata, +} from "./types.js"; import type { GitHubUserProps } from "./oauth.js"; import { queryFts, @@ -81,18 +91,21 @@ export class RagMcpAgent extends McpAgent { // ── search_issues ────────────────────────────────────────── this.server.tool( "search_issues", - "Unified search across GitHub issues, PRs, releases, repository documentation, and commit diffs. " + + "Unified search across GitHub issues, PRs, releases, repository documentation, commit diffs, " + + "issue/PR top-level comments, PR reviews, and PR inline review comments. " + "Three modes via the query / sort axes:\n" + " 1. Hybrid semantic search (default): dense BGE-M3 over Vectorize + sparse BM25 over D1 FTS5, " + "fused via Reciprocal Rank Fusion (RRF, k=60), then re-scored with a cross-encoder " + "(@cf/baai/bge-reranker-base; set rerank: false to skip).\n" + " 2. Time-ordered activity scan: pass an empty (or omitted) query with sort=\"updated_desc\" or \"created_desc\"; " + - "optionally narrow via since / until to list recent issue / PR / release / doc / diff activity.\n" + + "optionally narrow via since / until to list recent activity across every type.\n" + " 3. Doc content fetch: pass include_content: true to inline the raw file content of top doc results " + "(fetched from the GitHub contents API; capped at the first few doc rows).\n" + "Optional metadata filters (repo, state, labels, milestone, assignee, type) apply across all modes. " + "Use type: \"diff\" to retrieve judgment history preserved in commit diffs — including changes to deleted files " + - "and non-.md files that are not present in the live document index.", + "and non-.md files that are not present in the live document index. " + + "Use type: \"issue_comment\" / \"pr_review\" / \"pr_review_comment\" to retrieve comment-level judgment history " + + "(Master's feedback, AI responses, self-review now/later/accepted classifications).", { query: z .string() @@ -124,10 +137,26 @@ export class RagMcpAgent extends McpAgent { .optional() .describe("Filter by assignee login"), type: z - .enum(["issue", "pull_request", "release", "doc", "diff", "all"]) + .enum([ + "issue", + "pull_request", + "release", + "doc", + "diff", + "issue_comment", + "pr_review", + "pr_review_comment", + "all", + ]) .optional() .default("all") - .describe("Filter by type (default: all). Use \"diff\" to search per-file commit diffs."), + .describe( + "Filter by type (default: all). " + + "\"diff\" = per-file commit diffs. " + + "\"issue_comment\" = top-level comments on issues and PRs. " + + "\"pr_review\" = PR review bodies (approve / request_changes / comment). " + + "\"pr_review_comment\" = inline per-line review comments on PR diffs.", + ), top_k: z .number() .min(1) @@ -236,7 +265,15 @@ export class RagMcpAgent extends McpAgent { }; type ScanRow = { - type: "issue" | "pull_request" | "release" | "doc" | "diff"; + type: + | "issue" + | "pull_request" + | "release" + | "doc" + | "diff" + | "issue_comment" + | "pr_review" + | "pr_review_comment"; repo: string; number: number; title: string; @@ -255,6 +292,14 @@ export class RagMcpAgent extends McpAgent { file_status?: string; commit_date?: string; commit_author?: string; + /** Comment / review author login */ + author?: string; + /** GitHub comment / review id (comment rows only) */ + comment_id?: number; + /** GitHub review id (pr_review rows only) */ + review_id?: number; + /** Inline review-comment line number (pr_review_comment rows only) */ + line?: number; }; const rows: ScanRow[] = []; @@ -393,6 +438,108 @@ export class RagMcpAgent extends McpAgent { } } + // Issue / PR top-level comments + if (wantType("issue_comment")) { + try { + const res = await store.fetch( + new Request( + `http://store/recent-comments?${buildParams().toString()}`, + ), + ); + if (res.ok) { + const records = (await res.json()) as IssueCommentRecord[]; + for (const c of records) { + rows.push({ + type: "issue_comment", + repo: c.repo, + number: c.number, + title: `${c.author} on #${c.number}`, + state: "active", + labels: [], + milestone: "", + assignees: [], + url: `https://github.com/${c.repo}/issues/${c.number}#issuecomment-${c.commentId}`, + updated_at: c.updatedAt, + created_at: c.createdAt, + author: c.author, + comment_id: c.commentId, + }); + } + } + } catch { + // Non-critical. + } + } + + // PR reviews (approve / request_changes / comment body) + if (wantType("pr_review")) { + try { + const res = await store.fetch( + new Request( + `http://store/recent-reviews?${buildParams().toString()}`, + ), + ); + if (res.ok) { + const records = (await res.json()) as PRReviewRecord[]; + for (const r of records) { + rows.push({ + type: "pr_review", + repo: r.repo, + number: r.number, + title: `${r.author} ${r.state} on #${r.number}`, + state: r.state, + labels: [], + milestone: "", + assignees: [], + url: `https://github.com/${r.repo}/pull/${r.number}#pullrequestreview-${r.reviewId}`, + updated_at: r.updatedAt, + created_at: r.submittedAt, + author: r.author, + review_id: r.reviewId, + }); + } + } + } catch { + // Non-critical. + } + } + + // PR inline review comments + if (wantType("pr_review_comment")) { + try { + const res = await store.fetch( + new Request( + `http://store/recent-review-comments?${buildParams().toString()}`, + ), + ); + if (res.ok) { + const records = (await res.json()) as PRReviewCommentRecord[]; + for (const rc of records) { + rows.push({ + type: "pr_review_comment", + repo: rc.repo, + number: rc.number, + title: `${rc.author} @ ${rc.filePath}:${rc.line}`, + state: "active", + labels: [], + milestone: "", + assignees: [], + url: `https://github.com/${rc.repo}/pull/${rc.number}#discussion_r${rc.commentId}`, + updated_at: rc.updatedAt, + created_at: rc.createdAt, + author: rc.author, + comment_id: rc.commentId, + file_path: rc.filePath, + line: rc.line, + commit_sha: rc.commitId, + }); + } + } + } catch { + // Non-critical. + } + } + // State / milestone / assignee / labels post-filters (best-effort // over the metadata we have; assignees / labels are already arrays). let filteredRows = rows; @@ -792,6 +939,10 @@ export class RagMcpAgent extends McpAgent { file_status?: string; commit_date?: string; commit_author?: string; + author?: string; + comment_id?: number; + review_id?: number; + line?: number; content?: string; }; @@ -815,6 +966,10 @@ export class RagMcpAgent extends McpAgent { const fileStatus = meta?.file_status ?? ftsRow?.fileStatus ?? ""; const commitDate = meta?.commit_date ?? ftsRow?.commitDate ?? ""; const commitAuthor = meta?.commit_author ?? ftsRow?.commitAuthor ?? ""; + const author = meta?.author ?? ""; + const commentId = meta?.comment_id ?? 0; + const reviewId = meta?.review_id ?? 0; + const line = meta?.line ?? 0; let url: string; if (itemType === "release" && tagName) { @@ -823,6 +978,12 @@ export class RagMcpAgent extends McpAgent { url = `https://github.com/${itemRepo}/blob/main/${docPath}`; } else if (itemType === "diff" && commitSha) { url = `https://github.com/${itemRepo}/commit/${commitSha}`; + } else if (itemType === "issue_comment" && commentId) { + url = `https://github.com/${itemRepo}/issues/${number}#issuecomment-${commentId}`; + } else if (itemType === "pr_review" && reviewId) { + url = `https://github.com/${itemRepo}/pull/${number}#pullrequestreview-${reviewId}`; + } else if (itemType === "pr_review_comment" && commentId) { + url = `https://github.com/${itemRepo}/pull/${number}#discussion_r${commentId}`; } else { url = `https://github.com/${itemRepo}/issues/${number}`; } @@ -857,6 +1018,27 @@ export class RagMcpAgent extends McpAgent { commit_author: commitAuthor, } : {}), + ...(itemType === "issue_comment" + ? { + author, + comment_id: commentId, + } + : {}), + ...(itemType === "pr_review" + ? { + author, + review_id: reviewId, + } + : {}), + ...(itemType === "pr_review_comment" + ? { + author, + comment_id: commentId, + file_path: filePath, + line, + commit_sha: commitSha, + } + : {}), }; }); @@ -887,6 +1069,22 @@ export class RagMcpAgent extends McpAgent { const sha = item.commit_sha ?? ""; const shortSha = sha ? sha.slice(0, 7) : ""; item.title = [shortSha, fp].filter(Boolean).join(" "); + } else if (item.type === "issue_comment") { + // Title = "{author} on #{number}" — enough context to skim results. + const author = item.author ?? ""; + item.title = author ? `${author} on #${item.number}` : `comment on #${item.number}`; + } else if (item.type === "pr_review") { + // Title = "{author} {state} on #{number}" — review state gives + // the classification at a glance (APPROVED / CHANGES_REQUESTED / COMMENTED). + const author = item.author ?? ""; + const state = item.state || ""; + item.title = author ? `${author} ${state} on #${item.number}` : `review on #${item.number}`; + } else if (item.type === "pr_review_comment") { + // Title = "{author} @ {file_path}:{line}" — inline comment location. + const author = item.author ?? ""; + const fp = item.file_path ?? ""; + const line = item.line ?? 0; + item.title = author ? `${author} @ ${fp}:${line}` : `inline on #${item.number}`; } else if (item.repo && item.number) { try { const res = await store.fetch( diff --git a/src/pipeline.ts b/src/pipeline.ts index f8791bd..ac269a4 100644 --- a/src/pipeline.ts +++ b/src/pipeline.ts @@ -13,9 +13,41 @@ import type { DocRecord, DiffRecord, DiffFileStatus, + IssueCommentRecord, + PRReviewRecord, + PRReviewCommentRecord, } from "./types.js"; import { upsertFtsRow, deleteFtsRow } from "./fts.js"; +// ── Ingest filters (shared by webhook + poller paths) ──────── + +/** + * Minimum trimmed body length for comment / review ingest. + * Filters out "LGTM", "+1", emoji-only reactions, etc. + */ +export const MIN_COMMENT_BODY_CHARS = 10; + +/** + * Returns true when the login looks like a GitHub App / bot account. + * + * Bot accounts end in the `[bot]` suffix on the sender.login field. We + * filter them out because bot-authored comments (CI notes, dependabot + * summaries, auto-merge status) add noise without judgment history. + */ +export function isBotSender(login: string | null | undefined): boolean { + if (!login) return false; + return /\[bot\]$/.test(login); +} + +/** + * Returns true when the body is too short (or empty) to carry judgment + * history. Trim first so whitespace-only payloads count as empty. + */ +export function isBodyTooShort(body: string | null | undefined): boolean { + if (!body) return true; + return body.trim().length < MIN_COMMENT_BODY_CHARS; +} + // ── Constants ──────────────────────────────────────────────── /** Maximum characters for embedding input (BGE-M3 context limit ~8192 tokens, conservative char limit) */ @@ -211,6 +243,39 @@ export function diffVectorId( return stableVectorId("c", repo, commitSha, filePath); } +/** + * Build Vectorize vector ID for an issue / PR top-level comment. + * Deterministic SHA-256-based ID under the "ic" prefix. + */ +export function issueCommentVectorId( + repo: string, + commentId: number, +): Promise { + return stableVectorId("ic", repo, String(commentId)); +} + +/** + * Build Vectorize vector ID for a PR review (approve / request_changes / comment body). + * Deterministic SHA-256-based ID under the "pv" prefix. + */ +export function prReviewVectorId( + repo: string, + reviewId: number, +): Promise { + return stableVectorId("pv", repo, String(reviewId)); +} + +/** + * Build Vectorize vector ID for a PR inline review comment (per-line diff comment). + * Deterministic SHA-256-based ID under the "pc" prefix. + */ +export function prReviewCommentVectorId( + repo: string, + commentId: number, +): Promise { + return stableVectorId("pc", repo, String(commentId)); +} + // ── Per-item upsert functions ──────────────────────────────── /** GitHub API issue/PR response shape (subset of fields we need) */ @@ -1004,3 +1069,412 @@ export async function processAndUpsertCommitDiff( return { embedded, skipped, failed, batches }; } + +// ── Comment / review ingest surface ────────────────────────── + +/** GitHub API issue/PR comment shape (subset we need) */ +export interface GitHubCommentData { + id: number; + body: string | null; + user: { login: string } | null; + created_at: string; + updated_at: string; +} + +/** GitHub API PR review shape (subset we need) */ +export interface GitHubPRReviewData { + id: number; + body: string | null; + user: { login: string } | null; + state: string; + submitted_at: string | null; +} + +/** GitHub API PR inline review comment shape (subset we need) */ +export interface GitHubPRReviewCommentData { + id: number; + body: string | null; + user: { login: string } | null; + path: string | null; + line: number | null; + original_line?: number | null; + commit_id: string | null; + created_at: string; + updated_at: string; +} + +/** + * Build the embedding input for a comment / review body. + * Format: "{author}\n\n{body}", truncated to MAX_EMBEDDING_INPUT_CHARS. + * The author prefix supplies speaker context so the dense embedding can + * distinguish the same body authored by different reviewers. + */ +export function prepareCommentEmbeddingInput( + author: string, + body: string, +): string { + const text = `${author}\n\n${body}`; + if (text.length <= MAX_EMBEDDING_INPUT_CHARS) return text; + return text.slice(0, MAX_EMBEDDING_INPUT_CHARS); +} + +/** Result of a comment / review ingest operation */ +export interface CommentUpsertResult { + embedded: boolean; + skippedUnchanged: boolean; + /** True when the item was filtered out (bot author or body too short). */ + filtered: boolean; + failed: boolean; +} + +/** + * Process and upsert a single issue/PR top-level comment. + * + * Flow mirrors processAndUpsertIssue: bot / short-body filter, hash-based + * change detection, embedding, Vectorize upsert, FTS5 upsert, IssueStore record. + */ +export async function ingestIssueComment( + env: Env, + storeStub: DurableObjectStub, + repo: string, + parentNumber: number, + comment: GitHubCommentData, +): Promise { + const author = comment.user?.login ?? ""; + const body = comment.body ?? ""; + + if (isBotSender(author) || isBodyTooShort(body)) { + return { embedded: false, skippedUnchanged: false, filtered: true, failed: false }; + } + + const bodyHash = await computeBodyHash(author, body); + + // Change detection: compare stored hash + const existingResp = await storeStub.fetch( + new Request( + `http://store/comment?repo=${encodeURIComponent(repo)}&comment_id=${comment.id}`, + ), + ); + if (existingResp.ok) { + const existing = (await existingResp.json()) as IssueCommentRecord; + if (existing.bodyHash === bodyHash) { + return { embedded: false, skippedUnchanged: true, filtered: false, failed: false }; + } + } + + const embeddingInput = prepareCommentEmbeddingInput(author, body); + + let embeddingSucceeded = false; + try { + const embedding = await generateEmbedding(env.AI, embeddingInput); + + const metadata: Record = { + repo, + number: parentNumber, + type: "issue_comment", + state: "active", + labels: "", + milestone: "", + assignees: "", + updated_at: comment.updated_at, + author, + comment_id: comment.id, + }; + + const vid = await issueCommentVectorId(repo, comment.id); + await env.VECTORIZE.upsert([{ id: vid, values: embedding, metadata }]); + + try { + await upsertFtsRow(env.DB_FTS, { + vectorId: vid, + repo, + type: "issue_comment", + state: "active", + labels: "", + milestone: "", + assignees: "", + updatedAt: comment.updated_at, + number: parentNumber, + content: embeddingInput, + }); + } catch (ftsErr) { + console.error( + `Failed to upsert FTS5 row for comment ${repo}#${comment.id}:`, + ftsErr instanceof Error ? ftsErr.message : String(ftsErr), + ); + } + + embeddingSucceeded = true; + } catch (err) { + console.error( + `Failed to embed comment ${repo}#${comment.id}:`, + err instanceof Error ? err.message : String(err), + ); + } + + const record: IssueCommentRecord = { + repo, + commentId: comment.id, + number: parentNumber, + author, + bodyHash: embeddingSucceeded ? bodyHash : "", + createdAt: comment.created_at, + updatedAt: comment.updated_at, + }; + + await storeStub.fetch( + new Request("http://store/upsert-comment", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(record), + }), + ); + + return { + embedded: embeddingSucceeded, + skippedUnchanged: false, + filtered: false, + failed: !embeddingSucceeded, + }; +} + +/** + * Process and upsert a single PR review (approve / request_changes / comment body). + * + * Reviews without a body (approve-only, no prose) pass the min-length + * filter and are skipped. Reviews with meaningful prose go through the + * normal embed + upsert flow. + */ +export async function ingestPRReview( + env: Env, + storeStub: DurableObjectStub, + repo: string, + parentNumber: number, + review: GitHubPRReviewData, +): Promise { + const author = review.user?.login ?? ""; + const body = review.body ?? ""; + + if (isBotSender(author) || isBodyTooShort(body)) { + return { embedded: false, skippedUnchanged: false, filtered: true, failed: false }; + } + + const bodyHash = await computeBodyHash(author + "\n\n" + review.state, body); + + const existingResp = await storeStub.fetch( + new Request( + `http://store/review?repo=${encodeURIComponent(repo)}&review_id=${review.id}`, + ), + ); + if (existingResp.ok) { + const existing = (await existingResp.json()) as PRReviewRecord; + if (existing.bodyHash === bodyHash) { + return { embedded: false, skippedUnchanged: true, filtered: false, failed: false }; + } + } + + const submittedAt = review.submitted_at ?? new Date().toISOString(); + const embeddingInput = prepareCommentEmbeddingInput(author, body); + + let embeddingSucceeded = false; + try { + const embedding = await generateEmbedding(env.AI, embeddingInput); + + const metadata: Record = { + repo, + number: parentNumber, + type: "pr_review", + // Store the GitHub review state verbatim (APPROVED / CHANGES_REQUESTED / COMMENTED ...) + state: review.state, + labels: "", + milestone: "", + assignees: "", + updated_at: submittedAt, + author, + review_id: review.id, + }; + + const vid = await prReviewVectorId(repo, review.id); + await env.VECTORIZE.upsert([{ id: vid, values: embedding, metadata }]); + + try { + await upsertFtsRow(env.DB_FTS, { + vectorId: vid, + repo, + type: "pr_review", + state: review.state, + labels: "", + milestone: "", + assignees: "", + updatedAt: submittedAt, + number: parentNumber, + content: embeddingInput, + }); + } catch (ftsErr) { + console.error( + `Failed to upsert FTS5 row for PR review ${repo}#${review.id}:`, + ftsErr instanceof Error ? ftsErr.message : String(ftsErr), + ); + } + + embeddingSucceeded = true; + } catch (err) { + console.error( + `Failed to embed PR review ${repo}#${review.id}:`, + err instanceof Error ? err.message : String(err), + ); + } + + const record: PRReviewRecord = { + repo, + reviewId: review.id, + number: parentNumber, + author, + state: review.state, + bodyHash: embeddingSucceeded ? bodyHash : "", + submittedAt, + updatedAt: submittedAt, + }; + + await storeStub.fetch( + new Request("http://store/upsert-review", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(record), + }), + ); + + return { + embedded: embeddingSucceeded, + skippedUnchanged: false, + filtered: false, + failed: !embeddingSucceeded, + }; +} + +/** + * Process and upsert a single PR inline review comment (per-line diff comment). + * + * Inline comments carry extra diff context: file path, line, commit SHA. + * We surface these in the Vectorize metadata and FTS5 row so query-time + * filters can narrow to a specific file or commit. + */ +export async function ingestPRReviewComment( + env: Env, + storeStub: DurableObjectStub, + repo: string, + parentNumber: number, + comment: GitHubPRReviewCommentData, +): Promise { + const author = comment.user?.login ?? ""; + const body = comment.body ?? ""; + + if (isBotSender(author) || isBodyTooShort(body)) { + return { embedded: false, skippedUnchanged: false, filtered: true, failed: false }; + } + + // Line numbers can be null (outdated) or appear on original_line only; fall back. + const line = comment.line ?? comment.original_line ?? 0; + const filePath = comment.path ?? ""; + const commitId = comment.commit_id ?? ""; + + const bodyHash = await computeBodyHash( + `${author}\n${filePath}:${line}`, + body, + ); + + const existingResp = await storeStub.fetch( + new Request( + `http://store/review-comment?repo=${encodeURIComponent(repo)}&comment_id=${comment.id}`, + ), + ); + if (existingResp.ok) { + const existing = (await existingResp.json()) as PRReviewCommentRecord; + if (existing.bodyHash === bodyHash) { + return { embedded: false, skippedUnchanged: true, filtered: false, failed: false }; + } + } + + const embeddingInput = prepareCommentEmbeddingInput(author, body); + + let embeddingSucceeded = false; + try { + const embedding = await generateEmbedding(env.AI, embeddingInput); + + const metadata: Record = { + repo, + number: parentNumber, + type: "pr_review_comment", + state: "active", + labels: "", + milestone: "", + assignees: "", + updated_at: comment.updated_at, + author, + comment_id: comment.id, + file_path: filePath, + line, + commit_sha: commitId, + }; + + const vid = await prReviewCommentVectorId(repo, comment.id); + await env.VECTORIZE.upsert([{ id: vid, values: embedding, metadata }]); + + try { + await upsertFtsRow(env.DB_FTS, { + vectorId: vid, + repo, + type: "pr_review_comment", + state: "active", + labels: "", + milestone: "", + assignees: "", + updatedAt: comment.updated_at, + number: parentNumber, + filePath, + commitSha: commitId, + content: embeddingInput, + }); + } catch (ftsErr) { + console.error( + `Failed to upsert FTS5 row for PR review comment ${repo}#${comment.id}:`, + ftsErr instanceof Error ? ftsErr.message : String(ftsErr), + ); + } + + embeddingSucceeded = true; + } catch (err) { + console.error( + `Failed to embed PR review comment ${repo}#${comment.id}:`, + err instanceof Error ? err.message : String(err), + ); + } + + const record: PRReviewCommentRecord = { + repo, + commentId: comment.id, + number: parentNumber, + author, + filePath, + line, + commitId, + bodyHash: embeddingSucceeded ? bodyHash : "", + createdAt: comment.created_at, + updatedAt: comment.updated_at, + }; + + await storeStub.fetch( + new Request("http://store/upsert-review-comment", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(record), + }), + ); + + return { + embedded: embeddingSucceeded, + skippedUnchanged: false, + filtered: false, + failed: !embeddingSucceeded, + }; +} diff --git a/src/poller.ts b/src/poller.ts index 8e5a56a..4d8efde 100644 --- a/src/poller.ts +++ b/src/poller.ts @@ -13,8 +13,14 @@ import { processAndUpsertIssue, processAndUpsertRelease, processAndUpsertDoc, + ingestIssueComment, + ingestPRReview, + ingestPRReviewComment, type GitHubIssueData, type GitHubReleaseData, + type GitHubCommentData, + type GitHubPRReviewData, + type GitHubPRReviewCommentData, } from "./pipeline.js"; import { deleteFtsRow } from "./fts.js"; @@ -33,6 +39,15 @@ const MAX_EMBEDDINGS_PER_RUN = 50; * so the next cron continues from where it left off. */ const MAX_PAGES_PER_RUN = 2; +/** Maximum number of recent parent issues/PRs to fan out comment backfill over per repo. + * Keeps the API fan-out bounded on large, active repos so a single cron run + * does not exhaust the rate budget. Older parents are left to the next cron. */ +const MAX_COMMENT_BACKFILL_PARENTS = 20; + +/** Maximum number of comment-level items embedded per repo per run. + * Workers AI embed calls are the dominant cost for the comment surface. */ +const MAX_COMMENTS_EMBEDDED_PER_REPO = 30; + /** Sentinel value indicating GitHub returned 304 Not Modified */ const NOT_MODIFIED = Symbol("NOT_MODIFIED"); @@ -744,6 +759,205 @@ async function pollDocs( ); } +// ── Comment / review backfill ──────────────────────────────── + +/** Identify whether an issue record represents a pull request (has the PR surface) */ +function isPullRequestRecord(record: IssueRecord): boolean { + return record.type === "pull_request"; +} + +/** Fetch top-level comments for a single issue/PR. Returns [] on transient failures. */ +async function fetchIssueComments( + repo: string, + number: number, + token: string, +): Promise { + const url = `https://api.github.com/repos/${repo}/issues/${number}/comments?per_page=100`; + const resp = await fetch(url, { + headers: { + Authorization: `Bearer ${token}`, + Accept: "application/vnd.github+json", + "X-GitHub-Api-Version": "2022-11-28", + "User-Agent": "github-rag-mcp/0.1.0", + }, + cache: "no-store", + } as RequestInit); + + if (!resp.ok) { + throw new Error(`GitHub Issues Comments API error ${resp.status} for ${repo}#${number}`); + } + + return (await resp.json()) as GitHubCommentData[]; +} + +/** Fetch PR reviews for a single PR. */ +async function fetchPRReviews( + repo: string, + number: number, + token: string, +): Promise { + const url = `https://api.github.com/repos/${repo}/pulls/${number}/reviews?per_page=100`; + const resp = await fetch(url, { + headers: { + Authorization: `Bearer ${token}`, + Accept: "application/vnd.github+json", + "X-GitHub-Api-Version": "2022-11-28", + "User-Agent": "github-rag-mcp/0.1.0", + }, + cache: "no-store", + } as RequestInit); + + if (!resp.ok) { + throw new Error(`GitHub PR Reviews API error ${resp.status} for ${repo}#${number}`); + } + + return (await resp.json()) as GitHubPRReviewData[]; +} + +/** Fetch PR inline review comments for a single PR. */ +async function fetchPRReviewComments( + repo: string, + number: number, + token: string, +): Promise { + const url = `https://api.github.com/repos/${repo}/pulls/${number}/comments?per_page=100`; + const resp = await fetch(url, { + headers: { + Authorization: `Bearer ${token}`, + Accept: "application/vnd.github+json", + "X-GitHub-Api-Version": "2022-11-28", + "User-Agent": "github-rag-mcp/0.1.0", + }, + cache: "no-store", + } as RequestInit); + + if (!resp.ok) { + throw new Error(`GitHub PR Review Comments API error ${resp.status} for ${repo}#${number}`); + } + + return (await resp.json()) as GitHubPRReviewCommentData[]; +} + +/** + * Backfill comments, reviews, and review comments for a repo. + * + * Strategy: iterate over the most recently updated issues/PRs in the store + * (capped at MAX_COMMENT_BACKFILL_PARENTS), fetch their comment lists, + * and ingest each comment via the shared pipeline (bot / min-length filter + * + hash-based skip handle deduplication and noise). + * + * Embedding count is capped at MAX_COMMENTS_EMBEDDED_PER_REPO to stay within + * Workers AI rate budgets. Remaining items are picked up on the next cron. + */ +async function pollComments( + repo: string, + env: Env, + storeStub: DurableObjectStub, +): Promise { + // Pull the most recent issues/PRs from the store; these are the most likely + // to have fresh comments. Limit keeps fan-out bounded on busy repos. + const recentResp = await storeStub.fetch( + new Request( + `http://store/issues?repo=${encodeURIComponent(repo)}&limit=${MAX_COMMENT_BACKFILL_PARENTS}`, + ), + ); + if (!recentResp.ok) { + console.warn(`pollComments: unable to list recent issues for ${repo}`); + return; + } + + const parents = (await recentResp.json()) as IssueRecord[]; + if (parents.length === 0) { + console.log(`${repo} comments: no parents to backfill`); + return; + } + + let commentsEmbedded = 0; + let commentsSkipped = 0; + let commentsFiltered = 0; + let reviewsEmbedded = 0; + let reviewsSkipped = 0; + let reviewsFiltered = 0; + let reviewCommentsEmbedded = 0; + let reviewCommentsSkipped = 0; + let reviewCommentsFiltered = 0; + let fetchFailures = 0; + + const embedBudget = (): boolean => + commentsEmbedded + reviewsEmbedded + reviewCommentsEmbedded < MAX_COMMENTS_EMBEDDED_PER_REPO; + + for (const parent of parents) { + if (!embedBudget()) break; + + // Top-level comments (issues and PRs both route through /issues/{N}/comments) + try { + const comments = await fetchIssueComments(repo, parent.number, env.GITHUB_TOKEN); + for (const c of comments) { + if (!embedBudget()) break; + const result = await ingestIssueComment(env, storeStub, repo, parent.number, c); + if (result.embedded) commentsEmbedded++; + else if (result.skippedUnchanged) commentsSkipped++; + else if (result.filtered) commentsFiltered++; + } + } catch (err) { + fetchFailures++; + console.error( + `pollComments: failed to fetch comments for ${repo}#${parent.number}:`, + err instanceof Error ? err.message : String(err), + ); + } + + // PR-only: review bodies + inline review comments + if (!isPullRequestRecord(parent)) continue; + + if (!embedBudget()) break; + + try { + const reviews = await fetchPRReviews(repo, parent.number, env.GITHUB_TOKEN); + for (const r of reviews) { + if (!embedBudget()) break; + const result = await ingestPRReview(env, storeStub, repo, parent.number, r); + if (result.embedded) reviewsEmbedded++; + else if (result.skippedUnchanged) reviewsSkipped++; + else if (result.filtered) reviewsFiltered++; + } + } catch (err) { + fetchFailures++; + console.error( + `pollComments: failed to fetch reviews for ${repo}#${parent.number}:`, + err instanceof Error ? err.message : String(err), + ); + } + + if (!embedBudget()) break; + + try { + const inline = await fetchPRReviewComments(repo, parent.number, env.GITHUB_TOKEN); + for (const rc of inline) { + if (!embedBudget()) break; + const result = await ingestPRReviewComment(env, storeStub, repo, parent.number, rc); + if (result.embedded) reviewCommentsEmbedded++; + else if (result.skippedUnchanged) reviewCommentsSkipped++; + else if (result.filtered) reviewCommentsFiltered++; + } + } catch (err) { + fetchFailures++; + console.error( + `pollComments: failed to fetch review comments for ${repo}#${parent.number}:`, + err instanceof Error ? err.message : String(err), + ); + } + } + + console.log( + `${repo} comments: scanned ${parents.length} parents, ` + + `top-level [embedded=${commentsEmbedded}, skipped=${commentsSkipped}, filtered=${commentsFiltered}], ` + + `reviews [embedded=${reviewsEmbedded}, skipped=${reviewsSkipped}, filtered=${reviewsFiltered}], ` + + `inline [embedded=${reviewCommentsEmbedded}, skipped=${reviewCommentsSkipped}, filtered=${reviewCommentsFiltered}], ` + + `fetch_failures=${fetchFailures}`, + ); +} + /** * Main scheduled handler — called by Cron Trigger hourly as fallback. * Polls all configured repositories for issue/PR updates. @@ -810,5 +1024,14 @@ export async function handleScheduled( err instanceof Error ? err.message : String(err), ); } + + try { + await pollComments(repo, env, storeStub); + } catch (err) { + console.error( + `Failed to poll comments for ${repo}:`, + err instanceof Error ? err.message : String(err), + ); + } } } diff --git a/src/store.ts b/src/store.ts index 69ee90e..06fed7e 100644 --- a/src/store.ts +++ b/src/store.ts @@ -12,6 +12,9 @@ import type { DocRecord, DiffRecord, DiffFileStatus, + IssueCommentRecord, + PRReviewRecord, + PRReviewCommentRecord, PollWatermark, } from "./types.js"; @@ -75,6 +78,86 @@ type WatermarkRow = { etag: string; }; +/** Row shape returned by SQLite for the issue_comments table */ +type IssueCommentRow = { + [key: string]: SqlStorageValue; + repo: string; + comment_id: number; + number: number; + author: string; + body_hash: string; + created_at: string; + updated_at: string; +}; + +/** Row shape returned by SQLite for the pr_reviews table */ +type PRReviewRow = { + [key: string]: SqlStorageValue; + repo: string; + review_id: number; + number: number; + author: string; + state: string; + body_hash: string; + submitted_at: string; + updated_at: string; +}; + +/** Row shape returned by SQLite for the pr_review_comments table */ +type PRReviewCommentRow = { + [key: string]: SqlStorageValue; + repo: string; + comment_id: number; + number: number; + author: string; + file_path: string; + line: number; + commit_id: string; + body_hash: string; + created_at: string; + updated_at: string; +}; + +function rowToIssueCommentRecord(row: IssueCommentRow): IssueCommentRecord { + return { + repo: row.repo, + commentId: row.comment_id, + number: row.number, + author: row.author, + bodyHash: row.body_hash, + createdAt: row.created_at, + updatedAt: row.updated_at, + }; +} + +function rowToPRReviewRecord(row: PRReviewRow): PRReviewRecord { + return { + repo: row.repo, + reviewId: row.review_id, + number: row.number, + author: row.author, + state: row.state, + bodyHash: row.body_hash, + submittedAt: row.submitted_at, + updatedAt: row.updated_at, + }; +} + +function rowToPRReviewCommentRecord(row: PRReviewCommentRow): PRReviewCommentRecord { + return { + repo: row.repo, + commentId: row.comment_id, + number: row.number, + author: row.author, + filePath: row.file_path, + line: row.line, + commitId: row.commit_id, + bodyHash: row.body_hash, + createdAt: row.created_at, + updatedAt: row.updated_at, + }; +} + function rowToReleaseRecord(row: ReleaseRow): ReleaseRecord { return { repo: row.repo, @@ -237,6 +320,82 @@ export class IssueStore implements DurableObject { // Column already exists — ignore } + // Top-level comments on issues and PRs share the same number space, so we + // key on (repo, comment_id) directly. `number` stores the parent issue/PR + // number so we can filter by parent or reindex quickly. + this.sql.exec(` + CREATE TABLE IF NOT EXISTS issue_comments ( + repo TEXT NOT NULL, + comment_id INTEGER NOT NULL, + number INTEGER NOT NULL, + author TEXT NOT NULL DEFAULT '', + body_hash TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (repo, comment_id) + ); + `); + + this.sql.exec(` + CREATE INDEX IF NOT EXISTS idx_issue_comments_parent + ON issue_comments (repo, number, updated_at DESC); + `); + + this.sql.exec(` + CREATE INDEX IF NOT EXISTS idx_issue_comments_recent + ON issue_comments (updated_at DESC); + `); + + this.sql.exec(` + CREATE TABLE IF NOT EXISTS pr_reviews ( + repo TEXT NOT NULL, + review_id INTEGER NOT NULL, + number INTEGER NOT NULL, + author TEXT NOT NULL DEFAULT '', + state TEXT NOT NULL DEFAULT '', + body_hash TEXT NOT NULL DEFAULT '', + submitted_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (repo, review_id) + ); + `); + + this.sql.exec(` + CREATE INDEX IF NOT EXISTS idx_pr_reviews_parent + ON pr_reviews (repo, number, submitted_at DESC); + `); + + this.sql.exec(` + CREATE INDEX IF NOT EXISTS idx_pr_reviews_recent + ON pr_reviews (updated_at DESC); + `); + + this.sql.exec(` + CREATE TABLE IF NOT EXISTS pr_review_comments ( + repo TEXT NOT NULL, + comment_id INTEGER NOT NULL, + number INTEGER NOT NULL, + author TEXT NOT NULL DEFAULT '', + file_path TEXT NOT NULL DEFAULT '', + line INTEGER NOT NULL DEFAULT 0, + commit_id TEXT NOT NULL DEFAULT '', + body_hash TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (repo, comment_id) + ); + `); + + this.sql.exec(` + CREATE INDEX IF NOT EXISTS idx_pr_review_comments_parent + ON pr_review_comments (repo, number, updated_at DESC); + `); + + this.sql.exec(` + CREATE INDEX IF NOT EXISTS idx_pr_review_comments_recent + ON pr_review_comments (updated_at DESC); + `); + // Index for recent-activity queries (updated_at descending scan) this.sql.exec(` CREATE INDEX IF NOT EXISTS idx_issues_updated @@ -549,6 +708,198 @@ export class IssueStore implements DurableObject { return [...cursor].map(rowToDiffRecord); } + // ---- Issue comment CRUD ---- + + upsertIssueComment(record: IssueCommentRecord): void { + this.sql.exec( + `INSERT INTO issue_comments (repo, comment_id, number, author, body_hash, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (repo, comment_id) DO UPDATE SET + number = excluded.number, + author = excluded.author, + body_hash = excluded.body_hash, + updated_at = excluded.updated_at`, + record.repo, + record.commentId, + record.number, + record.author, + record.bodyHash, + record.createdAt, + record.updatedAt, + ); + } + + getIssueComment(repo: string, commentId: number): IssueCommentRecord | null { + const cursor = this.sql.exec( + `SELECT * FROM issue_comments WHERE repo = ? AND comment_id = ?`, + repo, + commentId, + ); + const rows = [...cursor]; + if (rows.length === 0) return null; + return rowToIssueCommentRecord(rows[0]); + } + + deleteIssueComment(repo: string, commentId: number): void { + this.sql.exec( + `DELETE FROM issue_comments WHERE repo = ? AND comment_id = ?`, + repo, + commentId, + ); + } + + getRecentIssueComments( + opts?: { since?: string; limit?: number; repo?: string }, + ): IssueCommentRecord[] { + const limit = opts?.limit ?? 20; + const since = opts?.since ?? new Date(Date.now() - 24 * 60 * 60 * 1000).toISOString(); + + let query: string; + let params: (string | number)[]; + + if (opts?.repo) { + query = `SELECT * FROM issue_comments WHERE repo = ? AND updated_at >= ? ORDER BY updated_at DESC LIMIT ?`; + params = [opts.repo, since, limit]; + } else { + query = `SELECT * FROM issue_comments WHERE updated_at >= ? ORDER BY updated_at DESC LIMIT ?`; + params = [since, limit]; + } + + const cursor = this.sql.exec(query, ...params); + return [...cursor].map(rowToIssueCommentRecord); + } + + // ---- PR review CRUD ---- + + upsertPRReview(record: PRReviewRecord): void { + this.sql.exec( + `INSERT INTO pr_reviews (repo, review_id, number, author, state, body_hash, submitted_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (repo, review_id) DO UPDATE SET + number = excluded.number, + author = excluded.author, + state = excluded.state, + body_hash = excluded.body_hash, + submitted_at = excluded.submitted_at, + updated_at = excluded.updated_at`, + record.repo, + record.reviewId, + record.number, + record.author, + record.state, + record.bodyHash, + record.submittedAt, + record.updatedAt, + ); + } + + getPRReview(repo: string, reviewId: number): PRReviewRecord | null { + const cursor = this.sql.exec( + `SELECT * FROM pr_reviews WHERE repo = ? AND review_id = ?`, + repo, + reviewId, + ); + const rows = [...cursor]; + if (rows.length === 0) return null; + return rowToPRReviewRecord(rows[0]); + } + + deletePRReview(repo: string, reviewId: number): void { + this.sql.exec( + `DELETE FROM pr_reviews WHERE repo = ? AND review_id = ?`, + repo, + reviewId, + ); + } + + getRecentPRReviews( + opts?: { since?: string; limit?: number; repo?: string }, + ): PRReviewRecord[] { + const limit = opts?.limit ?? 20; + const since = opts?.since ?? new Date(Date.now() - 24 * 60 * 60 * 1000).toISOString(); + + let query: string; + let params: (string | number)[]; + + if (opts?.repo) { + query = `SELECT * FROM pr_reviews WHERE repo = ? AND updated_at >= ? ORDER BY updated_at DESC LIMIT ?`; + params = [opts.repo, since, limit]; + } else { + query = `SELECT * FROM pr_reviews WHERE updated_at >= ? ORDER BY updated_at DESC LIMIT ?`; + params = [since, limit]; + } + + const cursor = this.sql.exec(query, ...params); + return [...cursor].map(rowToPRReviewRecord); + } + + // ---- PR review comment CRUD ---- + + upsertPRReviewComment(record: PRReviewCommentRecord): void { + this.sql.exec( + `INSERT INTO pr_review_comments (repo, comment_id, number, author, file_path, line, commit_id, body_hash, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (repo, comment_id) DO UPDATE SET + number = excluded.number, + author = excluded.author, + file_path = excluded.file_path, + line = excluded.line, + commit_id = excluded.commit_id, + body_hash = excluded.body_hash, + updated_at = excluded.updated_at`, + record.repo, + record.commentId, + record.number, + record.author, + record.filePath, + record.line, + record.commitId, + record.bodyHash, + record.createdAt, + record.updatedAt, + ); + } + + getPRReviewComment(repo: string, commentId: number): PRReviewCommentRecord | null { + const cursor = this.sql.exec( + `SELECT * FROM pr_review_comments WHERE repo = ? AND comment_id = ?`, + repo, + commentId, + ); + const rows = [...cursor]; + if (rows.length === 0) return null; + return rowToPRReviewCommentRecord(rows[0]); + } + + deletePRReviewComment(repo: string, commentId: number): void { + this.sql.exec( + `DELETE FROM pr_review_comments WHERE repo = ? AND comment_id = ?`, + repo, + commentId, + ); + } + + getRecentPRReviewComments( + opts?: { since?: string; limit?: number; repo?: string }, + ): PRReviewCommentRecord[] { + const limit = opts?.limit ?? 20; + const since = opts?.since ?? new Date(Date.now() - 24 * 60 * 60 * 1000).toISOString(); + + let query: string; + let params: (string | number)[]; + + if (opts?.repo) { + query = `SELECT * FROM pr_review_comments WHERE repo = ? AND updated_at >= ? ORDER BY updated_at DESC LIMIT ?`; + params = [opts.repo, since, limit]; + } else { + query = `SELECT * FROM pr_review_comments WHERE updated_at >= ? ORDER BY updated_at DESC LIMIT ?`; + params = [since, limit]; + } + + const cursor = this.sql.exec(query, ...params); + return [...cursor].map(rowToPRReviewCommentRecord); + } + // ---- Full re-embed reset ---- /** @@ -571,6 +922,9 @@ export class IssueStore implements DurableObject { releaseHashesReset: number; docsDeleted: number; diffsDeleted: number; + issueCommentHashesReset: number; + prReviewHashesReset: number; + prReviewCommentHashesReset: number; watermarksDeleted: number; } { const issuesCursor = this.sql.exec( @@ -593,12 +947,32 @@ export class IssueStore implements DurableObject { repo, ); - // Delete all three watermark namespaces: issues (repo), releases (releases:{repo}), docs (docs:{repo}) + const issueCommentsCursor = this.sql.exec( + `UPDATE issue_comments SET body_hash = '' WHERE repo = ? AND body_hash != ''`, + repo, + ); + + const prReviewsCursor = this.sql.exec( + `UPDATE pr_reviews SET body_hash = '' WHERE repo = ? AND body_hash != ''`, + repo, + ); + + const prReviewCommentsCursor = this.sql.exec( + `UPDATE pr_review_comments SET body_hash = '' WHERE repo = ? AND body_hash != ''`, + repo, + ); + + // Delete all watermark namespaces: issues (repo), releases (releases:{repo}), + // docs (docs:{repo}), and the three comment/review surfaces + // (comments:{repo}, reviews:{repo}, review_comments:{repo}). const watermarksCursor = this.sql.exec( - `DELETE FROM watermarks WHERE repo IN (?, ?, ?)`, + `DELETE FROM watermarks WHERE repo IN (?, ?, ?, ?, ?, ?)`, repo, `releases:${repo}`, `docs:${repo}`, + `comments:${repo}`, + `reviews:${repo}`, + `review_comments:${repo}`, ); return { @@ -606,6 +980,9 @@ export class IssueStore implements DurableObject { releaseHashesReset: releasesCursor.rowsWritten, docsDeleted: docsCursor.rowsWritten, diffsDeleted: diffsCursor.rowsWritten, + issueCommentHashesReset: issueCommentsCursor.rowsWritten, + prReviewHashesReset: prReviewsCursor.rowsWritten, + prReviewCommentHashesReset: prReviewCommentsCursor.rowsWritten, watermarksDeleted: watermarksCursor.rowsWritten, }; } @@ -874,6 +1251,141 @@ export class IssueStore implements DurableObject { return Response.json(items); } + // ── Issue comment endpoints ─────────────────────────────── + + // POST /upsert-comment — upsert a single issue/PR top-level comment + if (request.method === "POST" && path === "/upsert-comment") { + const record = (await request.json()) as IssueCommentRecord; + this.upsertIssueComment(record); + return new Response("ok", { status: 200 }); + } + + // GET /comment?repo=...&comment_id=... — get a single comment + if (request.method === "GET" && path === "/comment") { + const repo = url.searchParams.get("repo"); + const commentId = url.searchParams.get("comment_id"); + if (!repo || !commentId) { + return new Response("missing repo or comment_id", { status: 400 }); + } + const item = this.getIssueComment(repo, parseInt(commentId, 10)); + if (!item) return new Response("not found", { status: 404 }); + return Response.json(item); + } + + // DELETE /comment?repo=...&comment_id=... — delete a comment + if (request.method === "DELETE" && path === "/comment") { + const repo = url.searchParams.get("repo"); + const commentId = url.searchParams.get("comment_id"); + if (!repo || !commentId) { + return new Response("missing repo or comment_id", { status: 400 }); + } + this.deleteIssueComment(repo, parseInt(commentId, 10)); + return new Response("ok", { status: 200 }); + } + + // GET /recent-comments?since=...&limit=...&repo=... — recent comments + if (request.method === "GET" && path === "/recent-comments") { + const since = url.searchParams.get("since") ?? undefined; + const limit = url.searchParams.get("limit"); + const repo = url.searchParams.get("repo") ?? undefined; + const items = this.getRecentIssueComments({ + since, + limit: limit ? parseInt(limit, 10) : undefined, + repo, + }); + return Response.json(items); + } + + // ── PR review endpoints ─────────────────────────────────── + + // POST /upsert-review — upsert a single PR review + if (request.method === "POST" && path === "/upsert-review") { + const record = (await request.json()) as PRReviewRecord; + this.upsertPRReview(record); + return new Response("ok", { status: 200 }); + } + + // GET /review?repo=...&review_id=... — get a single review + if (request.method === "GET" && path === "/review") { + const repo = url.searchParams.get("repo"); + const reviewId = url.searchParams.get("review_id"); + if (!repo || !reviewId) { + return new Response("missing repo or review_id", { status: 400 }); + } + const item = this.getPRReview(repo, parseInt(reviewId, 10)); + if (!item) return new Response("not found", { status: 404 }); + return Response.json(item); + } + + // DELETE /review?repo=...&review_id=... — delete a review + if (request.method === "DELETE" && path === "/review") { + const repo = url.searchParams.get("repo"); + const reviewId = url.searchParams.get("review_id"); + if (!repo || !reviewId) { + return new Response("missing repo or review_id", { status: 400 }); + } + this.deletePRReview(repo, parseInt(reviewId, 10)); + return new Response("ok", { status: 200 }); + } + + // GET /recent-reviews?since=...&limit=...&repo=... — recent PR reviews + if (request.method === "GET" && path === "/recent-reviews") { + const since = url.searchParams.get("since") ?? undefined; + const limit = url.searchParams.get("limit"); + const repo = url.searchParams.get("repo") ?? undefined; + const items = this.getRecentPRReviews({ + since, + limit: limit ? parseInt(limit, 10) : undefined, + repo, + }); + return Response.json(items); + } + + // ── PR review comment endpoints ─────────────────────────── + + // POST /upsert-review-comment — upsert a single PR inline review comment + if (request.method === "POST" && path === "/upsert-review-comment") { + const record = (await request.json()) as PRReviewCommentRecord; + this.upsertPRReviewComment(record); + return new Response("ok", { status: 200 }); + } + + // GET /review-comment?repo=...&comment_id=... — get a single review comment + if (request.method === "GET" && path === "/review-comment") { + const repo = url.searchParams.get("repo"); + const commentId = url.searchParams.get("comment_id"); + if (!repo || !commentId) { + return new Response("missing repo or comment_id", { status: 400 }); + } + const item = this.getPRReviewComment(repo, parseInt(commentId, 10)); + if (!item) return new Response("not found", { status: 404 }); + return Response.json(item); + } + + // DELETE /review-comment?repo=...&comment_id=... — delete a review comment + if (request.method === "DELETE" && path === "/review-comment") { + const repo = url.searchParams.get("repo"); + const commentId = url.searchParams.get("comment_id"); + if (!repo || !commentId) { + return new Response("missing repo or comment_id", { status: 400 }); + } + this.deletePRReviewComment(repo, parseInt(commentId, 10)); + return new Response("ok", { status: 200 }); + } + + // GET /recent-review-comments?since=...&limit=...&repo=... — recent review comments + if (request.method === "GET" && path === "/recent-review-comments") { + const since = url.searchParams.get("since") ?? undefined; + const limit = url.searchParams.get("limit"); + const repo = url.searchParams.get("repo") ?? undefined; + const items = this.getRecentPRReviewComments({ + since, + limit: limit ? parseInt(limit, 10) : undefined, + repo, + }); + return Response.json(items); + } + return new Response("not found", { status: 404 }); } catch (err) { const message = err instanceof Error ? err.message : String(err); diff --git a/src/types.ts b/src/types.ts index cff81bc..a601776 100644 --- a/src/types.ts +++ b/src/types.ts @@ -37,6 +37,54 @@ export interface DocRecord { updatedAt: string; } +/** + * Stored top-level comment record (issues + PRs) in Durable Object SQLite. + * `number` is the parent issue or PR number (issues and PRs share the same number space). + */ +export interface IssueCommentRecord { + repo: string; + commentId: number; + number: number; + author: string; + bodyHash: string; + createdAt: string; + updatedAt: string; +} + +/** + * Stored PR review record (approve / request_changes / comment bodies). + * `state` carries the GitHub review state enum verbatim. + */ +export interface PRReviewRecord { + repo: string; + reviewId: number; + number: number; + author: string; + /** GitHub review state: APPROVED | CHANGES_REQUESTED | COMMENTED | DISMISSED | PENDING */ + state: string; + bodyHash: string; + submittedAt: string; + updatedAt: string; +} + +/** + * Stored PR inline review comment record (per-line comments on a diff). + * `filePath` / `line` pinpoint the diff location; `commitId` ties the + * comment to the reviewed commit SHA. + */ +export interface PRReviewCommentRecord { + repo: string; + commentId: number; + number: number; + author: string; + filePath: string; + line: number; + commitId: string; + bodyHash: string; + createdAt: string; + updatedAt: string; +} + /** File change status reported by GitHub for a file inside a commit */ export type DiffFileStatus = | "added" @@ -72,8 +120,16 @@ export interface PollWatermark { export interface VectorMetadata { repo: string; number: number; - type: "issue" | "pull_request" | "release" | "doc" | "diff"; - state: "open" | "closed" | "published" | "active"; + type: + | "issue" + | "pull_request" + | "release" + | "doc" + | "diff" + | "issue_comment" + | "pr_review" + | "pr_review_comment"; + state: string; labels: string; milestone: string; assignees: string; @@ -96,6 +152,14 @@ export interface VectorMetadata { blob_sha_before?: string; /** Git blob SHA after the commit, empty when file was removed (diffs only) */ blob_sha_after?: string; + /** Comment / review author login (comments + reviews + review comments only) */ + author?: string; + /** GitHub comment id (issue_comment + pr_review_comment only) */ + comment_id?: number; + /** GitHub review id (pr_review only) */ + review_id?: number; + /** Review-comment inline line number (pr_review_comment only) */ + line?: number; /** * Expanded label fields for Vectorize pre-filtering (first 4 labels, sorted). * Empty string when slot is unused. diff --git a/src/webhook.ts b/src/webhook.ts index 4e1da6b..e76fba2 100644 --- a/src/webhook.ts +++ b/src/webhook.ts @@ -13,12 +13,21 @@ import { processAndUpsertRelease, processAndUpsertDoc, processAndUpsertCommitDiff, + ingestIssueComment, + ingestPRReview, + ingestPRReviewComment, vectorId, releaseVectorId, docVectorId, + issueCommentVectorId, + prReviewVectorId, + prReviewCommentVectorId, type GitHubIssueData, type GitHubReleaseData, type GitHubCommitDetail, + type GitHubCommentData, + type GitHubPRReviewData, + type GitHubPRReviewCommentData, } from "./pipeline.js"; import { deleteFtsRow } from "./fts.js"; @@ -96,6 +105,12 @@ export async function handleWebhook(request: Request, env: Env): Promise, + env: Env, + storeStub: DurableObjectStub, +): Promise { + const action = payload.action as string; + const repo = (payload.repository as { full_name: string }).full_name; + const raw = payload.comment as Record; + const parent = payload.issue as Record; + + const commentId = raw.id as number; + const parentNumber = parent.number as number; + + // Deletion: drop from Vectorize + FTS5 + store + if (action === "deleted") { + const vid = await issueCommentVectorId(repo, commentId); + try { + await env.VECTORIZE.deleteByIds([vid]); + } catch (err) { + console.error( + `Failed to delete comment vector ${vid}:`, + err instanceof Error ? err.message : String(err), + ); + } + try { + await deleteFtsRow(env.DB_FTS, vid); + } catch (err) { + console.error( + `Failed to delete FTS5 row ${vid}:`, + err instanceof Error ? err.message : String(err), + ); + } + try { + await storeStub.fetch( + new Request( + `http://store/comment?repo=${encodeURIComponent(repo)}&comment_id=${commentId}`, + { method: "DELETE" }, + ), + ); + } catch (err) { + console.error( + `Failed to delete comment record ${repo}#${commentId}:`, + err instanceof Error ? err.message : String(err), + ); + } + return jsonResponse(202, { + received: true, + event: "issue_comment", + action, + repo, + number: parentNumber, + commentId, + result: "deleted", + }); + } + + const comment: GitHubCommentData = { + id: commentId, + body: (raw.body as string | null) ?? null, + user: (raw.user as { login: string } | null) ?? null, + created_at: raw.created_at as string, + updated_at: raw.updated_at as string, + }; + + const result = await ingestIssueComment(env, storeStub, repo, parentNumber, comment); + + return jsonResponse(202, { + received: true, + event: "issue_comment", + action, + repo, + number: parentNumber, + commentId, + result: { + embedded: result.embedded, + skippedUnchanged: result.skippedUnchanged, + filtered: result.filtered, + failed: result.failed, + }, + }); +} + +/** + * Handle `pull_request_review.*` webhook events. + * + * For dismissed: remove the review from Vectorize + FTS5 + store (the + * dismissed body is no longer useful judgment history). + * For submitted / edited: ingest via the shared review pipeline. + */ +async function handlePRReviewEvent( + payload: Record, + env: Env, + storeStub: DurableObjectStub, +): Promise { + const action = payload.action as string; + const repo = (payload.repository as { full_name: string }).full_name; + const raw = payload.review as Record; + const parent = payload.pull_request as Record; + + const reviewId = raw.id as number; + const parentNumber = parent.number as number; + + if (action === "dismissed") { + const vid = await prReviewVectorId(repo, reviewId); + try { + await env.VECTORIZE.deleteByIds([vid]); + } catch (err) { + console.error( + `Failed to delete review vector ${vid}:`, + err instanceof Error ? err.message : String(err), + ); + } + try { + await deleteFtsRow(env.DB_FTS, vid); + } catch (err) { + console.error( + `Failed to delete FTS5 row ${vid}:`, + err instanceof Error ? err.message : String(err), + ); + } + try { + await storeStub.fetch( + new Request( + `http://store/review?repo=${encodeURIComponent(repo)}&review_id=${reviewId}`, + { method: "DELETE" }, + ), + ); + } catch (err) { + console.error( + `Failed to delete review record ${repo}#${reviewId}:`, + err instanceof Error ? err.message : String(err), + ); + } + return jsonResponse(202, { + received: true, + event: "pull_request_review", + action, + repo, + number: parentNumber, + reviewId, + result: "deleted", + }); + } + + const review: GitHubPRReviewData = { + id: reviewId, + body: (raw.body as string | null) ?? null, + user: (raw.user as { login: string } | null) ?? null, + state: (raw.state as string) ?? "", + submitted_at: (raw.submitted_at as string | null) ?? null, + }; + + const result = await ingestPRReview(env, storeStub, repo, parentNumber, review); + + return jsonResponse(202, { + received: true, + event: "pull_request_review", + action, + repo, + number: parentNumber, + reviewId, + result: { + embedded: result.embedded, + skippedUnchanged: result.skippedUnchanged, + filtered: result.filtered, + failed: result.failed, + }, + }); +} + +/** + * Handle `pull_request_review_comment.*` webhook events (inline diff comments). + * + * For deleted: remove from Vectorize + FTS5 + store. + * For created / edited: ingest via the shared review-comment pipeline. + */ +async function handlePRReviewCommentEvent( + payload: Record, + env: Env, + storeStub: DurableObjectStub, +): Promise { + const action = payload.action as string; + const repo = (payload.repository as { full_name: string }).full_name; + const raw = payload.comment as Record; + const parent = payload.pull_request as Record; + + const commentId = raw.id as number; + const parentNumber = parent.number as number; + + if (action === "deleted") { + const vid = await prReviewCommentVectorId(repo, commentId); + try { + await env.VECTORIZE.deleteByIds([vid]); + } catch (err) { + console.error( + `Failed to delete review comment vector ${vid}:`, + err instanceof Error ? err.message : String(err), + ); + } + try { + await deleteFtsRow(env.DB_FTS, vid); + } catch (err) { + console.error( + `Failed to delete FTS5 row ${vid}:`, + err instanceof Error ? err.message : String(err), + ); + } + try { + await storeStub.fetch( + new Request( + `http://store/review-comment?repo=${encodeURIComponent(repo)}&comment_id=${commentId}`, + { method: "DELETE" }, + ), + ); + } catch (err) { + console.error( + `Failed to delete review comment record ${repo}#${commentId}:`, + err instanceof Error ? err.message : String(err), + ); + } + return jsonResponse(202, { + received: true, + event: "pull_request_review_comment", + action, + repo, + number: parentNumber, + commentId, + result: "deleted", + }); + } + + const comment: GitHubPRReviewCommentData = { + id: commentId, + body: (raw.body as string | null) ?? null, + user: (raw.user as { login: string } | null) ?? null, + path: (raw.path as string | null) ?? null, + line: (raw.line as number | null) ?? null, + original_line: (raw.original_line as number | null) ?? null, + commit_id: (raw.commit_id as string | null) ?? null, + created_at: raw.created_at as string, + updated_at: raw.updated_at as string, + }; + + const result = await ingestPRReviewComment(env, storeStub, repo, parentNumber, comment); + + return jsonResponse(202, { + received: true, + event: "pull_request_review_comment", + action, + repo, + number: parentNumber, + commentId, + result: { + embedded: result.embedded, + skippedUnchanged: result.skippedUnchanged, + filtered: result.filtered, + failed: result.failed, + }, + }); +}