Skip to content

Commit 4158e9d

Browse files
feat(poller): add historical commit diff backfill via 2-phase polling (#82) (#119)
webhook path が新規 push 限定なため、過去 commit diff が retrieval から欠落していた。poller 側に forward + backward の 2-phase pollDiffs を追加し、`#80` の `processAndUpsertCommitDiff` と `fetchCommitDetail` を共用する形で歴史 backfill を段階的に進める。 - pipeline.ts: `fetchCommitDetail` を export 化して webhook / poller 共用のユーティリティに移動。webhook.ts 側は import 切替のみ。 - poller.ts: `pollDiffs(repo, env, storeStub)` を追加。 forward phase は `since=lastPolledAt`(初回 1h ago)で webhook redundancy、backward phase は `until=oldestUnprocessedDate` (初回 now)で履歴遡行。各 phase 10 commits/run/repo で cap。 `diffs:${repo}` と `diffs_backfill:${repo}` の 2 namespace で watermark を分離。`handleScheduled` から pollComments 後に呼び出す。 - docs/0-requirements.*: Cron Poller 節に commit diff surface と 2-phase 進行 / watermark 方針を追記。 upsert は `(repo, commit_sha, file_path)` で idempotent なため、 webhook / 両 phase 間 overlap は副作用なし。 Closes #82
1 parent 26b8b71 commit 4158e9d

5 files changed

Lines changed: 289 additions & 35 deletions

File tree

‎docs/0-requirements.ja.md‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -111,11 +111,18 @@ Responsibilities:
111111

112112
- webhook 取りこぼしを補償する
113113
- 新しい repository の backfill を行う
114-
- issue、pull request、release、docs の変更を再取得する
114+
- issue、pull request、release、docs、issue/PR comments、commit diff の変更を再取得する
115115
- 一時障害後も store を収束させる
116116

117117
現在の deployment では hourly で実行する。
118118

119+
commit diff poller は 2-phase 構成:
120+
121+
- **forward phase** — `since=lastPolledAt` で直近 commit を再取得する(webhook 取りこぼし時の redundancy)。watermark namespace は `diffs:${repo}`。
122+
- **backward phase** — `until=oldestUnprocessedDate` で履歴を徐々に遡行する(新規 deployment や webhook 起動前の commit を backfill する経路)。watermark namespace は `diffs_backfill:${repo}`。
123+
124+
1 run あたり上限は forward / backward それぞれ 10 commits。`processAndUpsertCommitDiff` の upsert は `(repo, commit_sha, file_path)` で idempotent なので、webhook / 両 phase 間で overlap しても副作用はない。
125+
119126
### 4. Embedding Pipeline
120127

121128
embedding model:

‎docs/0-requirements.md‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -111,11 +111,18 @@ Responsibilities:
111111

112112
- repair missed webhook deliveries
113113
- backfill new repositories
114-
- refresh changed issues, pull requests, releases, and docs
114+
- refresh changed issues, pull requests, releases, docs, issue/PR comments, and commit diffs
115115
- keep the stores converged even after transient failures
116116

117117
The poller runs hourly in the current deployment.
118118

119+
The commit-diff poller runs in two phases:
120+
121+
- **forward phase** — re-fetches recent commits using `since=lastPolledAt`, acting as redundancy when webhook delivery has stalled. Watermark namespace: `diffs:${repo}`.
122+
- **backward phase** — walks backward through history using `until=oldestUnprocessedDate`, backfilling commits that predate the webhook or a fresh deployment. Watermark namespace: `diffs_backfill:${repo}`.
123+
124+
Each phase is capped at 10 commits per repo per run. Upserts through `processAndUpsertCommitDiff` are idempotent on `(repo, commit_sha, file_path)`, so overlap between webhook and either phase is safe.
125+
119126
### 4. Embedding Pipeline
120127

121128
Embedding model:

‎src/pipeline.ts‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -841,6 +841,39 @@ export interface DiffUpsertResult {
841841
batches: number;
842842
}
843843

844+
/**
845+
* Fetch a single commit with per-file patches via the GitHub REST API.
846+
* Returns the commit detail including `files[]` with inline `patch` fields.
847+
* Throws on non-2xx responses. Shared between webhook (new-commit path) and
848+
* poller (historical backfill path).
849+
*/
850+
export async function fetchCommitDetail(
851+
repo: string,
852+
sha: string,
853+
token: string,
854+
): Promise<GitHubCommitDetail> {
855+
const url = `https://api.github.com/repos/${repo}/commits/${sha}`;
856+
857+
const resp = await fetch(url, {
858+
headers: {
859+
Authorization: `Bearer ${token}`,
860+
Accept: "application/vnd.github+json",
861+
"X-GitHub-Api-Version": "2022-11-28",
862+
"User-Agent": "github-rag-mcp/0.1.0",
863+
},
864+
cache: "no-store",
865+
} as RequestInit);
866+
867+
if (!resp.ok) {
868+
const text = await resp.text();
869+
throw new Error(
870+
`GitHub Commits API error ${resp.status} for ${repo}@${sha}: ${text}`,
871+
);
872+
}
873+
874+
return (await resp.json()) as GitHubCommitDetail;
875+
}
876+
844877
/**
845878
* Normalise GitHub's file status string to our DiffFileStatus union.
846879
* Unknown values fall through to "changed" (the generic GitHub bucket).

‎src/poller.ts‎

Lines changed: 239 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,8 @@ import {
1313
processAndUpsertIssue,
1414
processAndUpsertRelease,
1515
processAndUpsertDoc,
16+
processAndUpsertCommitDiff,
17+
fetchCommitDetail,
1618
ingestIssueComment,
1719
ingestPRReview,
1820
ingestPRReviewComment,
@@ -48,6 +50,19 @@ const MAX_COMMENT_BACKFILL_PARENTS = 20;
4850
* Workers AI embed calls are the dominant cost for the comment surface. */
4951
const MAX_COMMENTS_EMBEDDED_PER_REPO = 30;
5052

53+
/** Maximum number of commits fetched in the forward (webhook-redundancy) phase
54+
* of the diff poller per repo per run.
55+
* Forward is normally a no-op because the webhook path already indexes new
56+
* commits; this cap bounds the work when webhook delivery has stalled. */
57+
const MAX_DIFF_COMMITS_FORWARD_PER_RUN = 10;
58+
59+
/** Maximum number of commits fetched in the backward (historical backfill) phase
60+
* of the diff poller per repo per run.
61+
* Backfill walks backward through repo history one hourly run at a time; the
62+
* cap keeps per-run API and embedding cost bounded so the total sweep spreads
63+
* over many runs (e.g. 10 commits/run × 24 runs/day = 240 commits/day). */
64+
const MAX_DIFF_COMMITS_BACKWARD_PER_RUN = 10;
65+
5166
/** Sentinel value indicating GitHub returned 304 Not Modified */
5267
const NOT_MODIFIED = Symbol("NOT_MODIFIED");
5368

@@ -958,6 +973,221 @@ async function pollComments(
958973
);
959974
}
960975

976+
/** GitHub API commit list item — subset used by the diff poller. */
977+
interface GitHubCommitSummary {
978+
sha: string;
979+
commit: {
980+
message?: string;
981+
author?: { date?: string | null } | null;
982+
committer?: { date?: string | null } | null;
983+
};
984+
}
985+
986+
/**
987+
* Fetch a single page of commits from `GET /repos/{repo}/commits`.
988+
* Supports `since` (inclusive lower bound on committer date) and `until`
989+
* (inclusive upper bound) filters; the two are combined by GitHub with AND.
990+
* Results are ordered newest-first by committer date.
991+
* Throws on non-2xx responses so the caller can log and fall back.
992+
*/
993+
async function fetchRepoCommits(
994+
repo: string,
995+
token: string,
996+
opts: { since?: string; until?: string; per_page: number },
997+
): Promise<GitHubCommitSummary[]> {
998+
const url = new URL(`https://api.github.com/repos/${repo}/commits`);
999+
url.searchParams.set("per_page", String(opts.per_page));
1000+
if (opts.since) url.searchParams.set("since", opts.since);
1001+
if (opts.until) url.searchParams.set("until", opts.until);
1002+
1003+
const resp = await fetch(url.toString(), {
1004+
headers: {
1005+
Authorization: `Bearer ${token}`,
1006+
Accept: "application/vnd.github+json",
1007+
"X-GitHub-Api-Version": "2022-11-28",
1008+
"User-Agent": "github-rag-mcp/0.1.0",
1009+
},
1010+
cache: "no-store",
1011+
} as RequestInit);
1012+
1013+
if (!resp.ok) {
1014+
const text = await resp.text();
1015+
throw new Error(
1016+
`GitHub Commits list API error ${resp.status} for ${repo}: ${text}`,
1017+
);
1018+
}
1019+
1020+
return (await resp.json()) as GitHubCommitSummary[];
1021+
}
1022+
1023+
/** Read a watermark record by its namespaced key; returns null when absent. */
1024+
async function readWatermark(
1025+
storeStub: DurableObjectStub,
1026+
key: string,
1027+
): Promise<{ lastPolledAt: string } | null> {
1028+
const resp = await storeStub.fetch(
1029+
new Request(`http://store/watermark?repo=${encodeURIComponent(key)}`),
1030+
);
1031+
if (!resp.ok) return null;
1032+
const wm = (await resp.json()) as { repo: string; lastPolledAt: string };
1033+
return { lastPolledAt: wm.lastPolledAt };
1034+
}
1035+
1036+
/** Upsert a watermark record under the given namespaced key. */
1037+
async function writeWatermark(
1038+
storeStub: DurableObjectStub,
1039+
key: string,
1040+
lastPolledAt: string,
1041+
): Promise<void> {
1042+
await storeStub.fetch(
1043+
new Request("http://store/watermark", {
1044+
method: "POST",
1045+
headers: { "Content-Type": "application/json" },
1046+
body: JSON.stringify({ repo: key, lastPolledAt }),
1047+
}),
1048+
);
1049+
}
1050+
1051+
/** Extract the best-available ISO timestamp from a commit summary. */
1052+
function commitDateOf(summary: GitHubCommitSummary): string | undefined {
1053+
return (
1054+
summary.commit.author?.date ??
1055+
summary.commit.committer?.date ??
1056+
undefined
1057+
);
1058+
}
1059+
1060+
/**
1061+
* Poll historical and recent commit diffs for a repository and upsert them
1062+
* through the shared commit-diff pipeline.
1063+
*
1064+
* Two phases run per cron tick:
1065+
*
1066+
* 1. **Forward** (webhook redundancy): fetch commits with `since=lastPolledAt`
1067+
* so the poller re-covers any commits missed while webhook delivery was
1068+
* stalled. The first run uses "one hour ago" as the initial since so the
1069+
* initial fetch stays bounded; subsequent runs advance the forward
1070+
* watermark to the current poll start time unconditionally.
1071+
*
1072+
* 2. **Backward** (historical backfill): fetch commits with
1073+
* `until=oldestUnprocessedDate` so the poller walks backward through the
1074+
* repo's history one tick at a time. The first run uses "now" as the
1075+
* initial until; subsequent runs advance the backward watermark to the
1076+
* commit_date of the oldest commit processed in this run. When the repo's
1077+
* history is exhausted the API returns 0 commits and the watermark stops
1078+
* advancing — subsequent runs will repeatedly return 0 commits, which is
1079+
* acceptable idle-state behavior.
1080+
*
1081+
* Each phase is capped at a small commit count (see MAX_DIFF_COMMITS_*) to
1082+
* spread cost across many cron ticks. `processAndUpsertCommitDiff` upserts
1083+
* on the (repo, commit_sha, file_path) primary key, so overlap with webhook
1084+
* or with the opposite phase is idempotent.
1085+
*/
1086+
export async function pollDiffs(
1087+
repo: string,
1088+
env: Env,
1089+
storeStub: DurableObjectStub,
1090+
): Promise<void> {
1091+
const pollStartTime = new Date().toISOString();
1092+
1093+
// ── Forward phase ───────────────────────────────────────────
1094+
const fwdKey = `diffs:${repo}`;
1095+
const fwdWm = await readWatermark(storeStub, fwdKey);
1096+
// First run: start one hour ago so the initial forward sweep covers the
1097+
// last cron interval without pulling the whole history into this phase.
1098+
const sinceFwd =
1099+
fwdWm?.lastPolledAt ??
1100+
new Date(Date.now() - 60 * 60 * 1000).toISOString();
1101+
1102+
let fwdProcessed = 0;
1103+
let fwdFailed = 0;
1104+
try {
1105+
const fwdCommits = await fetchRepoCommits(repo, env.GITHUB_TOKEN, {
1106+
since: sinceFwd,
1107+
per_page: MAX_DIFF_COMMITS_FORWARD_PER_RUN,
1108+
});
1109+
for (const summary of fwdCommits) {
1110+
try {
1111+
const detail = await fetchCommitDetail(
1112+
repo,
1113+
summary.sha,
1114+
env.GITHUB_TOKEN,
1115+
);
1116+
await processAndUpsertCommitDiff(env, storeStub, repo, detail);
1117+
fwdProcessed++;
1118+
} catch (err) {
1119+
fwdFailed++;
1120+
console.error(
1121+
`pollDiffs: forward commit ${repo}@${summary.sha} failed:`,
1122+
err instanceof Error ? err.message : String(err),
1123+
);
1124+
}
1125+
}
1126+
} catch (err) {
1127+
console.error(
1128+
`pollDiffs: forward list failed for ${repo}:`,
1129+
err instanceof Error ? err.message : String(err),
1130+
);
1131+
}
1132+
// Advance the forward watermark regardless of partial failures so the next
1133+
// run continues from pollStartTime instead of reprocessing the same window.
1134+
// Upstream upsert is idempotent on (repo, commit_sha, file_path).
1135+
await writeWatermark(storeStub, fwdKey, pollStartTime);
1136+
1137+
// ── Backward phase ──────────────────────────────────────────
1138+
const bwdKey = `diffs_backfill:${repo}`;
1139+
const bwdWm = await readWatermark(storeStub, bwdKey);
1140+
// First run: start walking backward from the current time.
1141+
const untilBwd = bwdWm?.lastPolledAt ?? pollStartTime;
1142+
1143+
let bwdProcessed = 0;
1144+
let bwdFailed = 0;
1145+
let oldestSeenDate: string | undefined;
1146+
try {
1147+
const bwdCommits = await fetchRepoCommits(repo, env.GITHUB_TOKEN, {
1148+
until: untilBwd,
1149+
per_page: MAX_DIFF_COMMITS_BACKWARD_PER_RUN,
1150+
});
1151+
// GitHub returns commits newest-first; the last entry is the oldest in
1152+
// this page and becomes the next-run watermark.
1153+
for (const summary of bwdCommits) {
1154+
try {
1155+
const detail = await fetchCommitDetail(
1156+
repo,
1157+
summary.sha,
1158+
env.GITHUB_TOKEN,
1159+
);
1160+
await processAndUpsertCommitDiff(env, storeStub, repo, detail);
1161+
bwdProcessed++;
1162+
const d = commitDateOf(summary);
1163+
if (d) oldestSeenDate = d;
1164+
} catch (err) {
1165+
bwdFailed++;
1166+
console.error(
1167+
`pollDiffs: backward commit ${repo}@${summary.sha} failed:`,
1168+
err instanceof Error ? err.message : String(err),
1169+
);
1170+
}
1171+
}
1172+
} catch (err) {
1173+
console.error(
1174+
`pollDiffs: backward list failed for ${repo}:`,
1175+
err instanceof Error ? err.message : String(err),
1176+
);
1177+
}
1178+
// Only advance the backward watermark when we actually saw a commit. If the
1179+
// API returned 0 commits the repo's history is exhausted (or the token lost
1180+
// access); leaving the watermark alone avoids silently skipping a window.
1181+
if (oldestSeenDate) {
1182+
await writeWatermark(storeStub, bwdKey, oldestSeenDate);
1183+
}
1184+
1185+
console.log(
1186+
`${repo} diffs: forward [processed=${fwdProcessed}, failed=${fwdFailed}], ` +
1187+
`backward [processed=${bwdProcessed}, failed=${bwdFailed}]`,
1188+
);
1189+
}
1190+
9611191
/**
9621192
* Main scheduled handler — called by Cron Trigger hourly as fallback.
9631193
* Polls all configured repositories for issue/PR updates.
@@ -1033,5 +1263,14 @@ export async function handleScheduled(
10331263
err instanceof Error ? err.message : String(err),
10341264
);
10351265
}
1266+
1267+
try {
1268+
await pollDiffs(repo, env, storeStub);
1269+
} catch (err) {
1270+
console.error(
1271+
`Failed to poll diffs for ${repo}:`,
1272+
err instanceof Error ? err.message : String(err),
1273+
);
1274+
}
10361275
}
10371276
}

0 commit comments

Comments
 (0)