diff --git a/app.js b/app.js index 42bf37c7..eba49b44 100644 --- a/app.js +++ b/app.js @@ -4,7 +4,7 @@ import bodyParser from 'body-parser'; import expressSession from 'express-session'; import authManager from './lib/authentication.js'; import socketAuthenticator from './lib/socket-auth.js'; -import refresh from './lib/refresh.js'; +import refresh, { startRecentPullsSweep } from './lib/refresh.js'; import pullManager from './lib/pull-manager.js'; import git from './lib/git-manager.js'; import dbManager from './lib/db-manager.js'; @@ -82,9 +82,17 @@ app.get('/api/v1/pulls', apiAuth, apiController.getPulls); // the lookup. git.getBotLogin(); -debug('Loading all recent pulls from the DB'); +// Read before webhooks and the startup refresh write any pull: the first sweep +// reaches back to it, so events dropped while the server was down get re-read. +let latestPullUpdate = null; + dbManager - .getRecentPulls(pullManager.getOldestAllowedPullTimestamp()) + .getLatestPullUpdate() + .then(function (latest) { + latestPullUpdate = latest; + debug('Loading all recent pulls from the DB'); + return dbManager.getRecentPulls(pullManager.getOldestAllowedPullTimestamp()); + }) .then(function (pulls) { debug('Loaded %s pulls', pulls.length); pullQueue.pause(); @@ -96,6 +104,7 @@ dbManager .then(function () { debug('Refreshing all open pulls from the API'); refresh.openPulls(); + startRecentPullsSweep(refresh, config.repos, latestPullUpdate); }) .done(); diff --git a/lib/db-manager.js b/lib/db-manager.js index fa10d41d..e8d437a5 100644 --- a/lib/db-manager.js +++ b/lib/db-manager.js @@ -401,6 +401,17 @@ const dbManager = { return db.query('SELECT repo, number FROM pulls WHERE state = ?', ['open']); }, + /** + * Returns a promise which resolves to the newest `date_updated` (epoch + * seconds) of any pull in the DB, or null when there are no pulls. + */ + getLatestPullUpdate: function () { + dbDebug('Calling getLatestPullUpdate'); + return db + .query('SELECT MAX(date_updated) AS latest FROM pulls') + .then(rows => rows[0].latest); + }, + /** * Returns a promise which resolves to a pull's number for the given head * commit sha. diff --git a/lib/git-manager.js b/lib/git-manager.js index 722d9b11..5a793bbe 100644 --- a/lib/git-manager.js +++ b/lib/git-manager.js @@ -170,6 +170,34 @@ export default { ); }, + /** + * Resolves to `{ repo, number, updatedAt }` for every pull, open or closed, + * updated at or after `since` in any repo `owner` holds, configured or not. + * `updatedAt` is GitHub's current updated_at, which can be newer than the + * time the search matched on. + */ + searchUpdatedPulls: function (owner, since) { + const updatedSince = since.toISOString().slice(0, 19) + 'Z'; + return logErrors( + github + .paginate(githubRest.search.issuesAndPullRequests, { + q: `user:${owner} is:pr updated:>=${updatedSince}`, + per_page: 100, + }) + .then(items => + items.map(item => ({ + // Search results carry only the API URL, .../repos/owner/repo + repo: item.repository_url.split('/').slice(-2).join('/'), + number: item.number, + updatedAt: item.updated_at, + })) + ), + 'Searching pulls in %s updated since %s', + owner, + updatedSince + ); + }, + /** * Get *all* pull requests for a repo. * diff --git a/lib/refresh.js b/lib/refresh.js index e09fec5e..198e7ce0 100644 --- a/lib/refresh.js +++ b/lib/refresh.js @@ -4,6 +4,7 @@ import utils from './utils.js'; import NotifyQueue from 'notify-queue'; import debug from './debug.js'; import Promise from 'bluebird'; +import _ from 'underscore'; import { createPacer, noopPacer } from './pacer.js'; const refreshDebug = debug('pulldasher:refresh'); @@ -79,11 +80,11 @@ export function createRefresh({ pacer = noopPacer } = {}) { /////// Pulls ///////// - pull: function refreshPull(repo, number) { + pull: function refreshPull(repo, number, onFailure = null) { refreshDebug('refresh pull %s', number); return gitManager .getPull(repo, number) - .then(pushOnQueue(pullQueue)) + .then(pushOnQueue(pullQueue, onFailure)) .catch(function (err) { // getPull's rejection (a dead/expired bot token, a 5xx, ...) lands // here before pushOnQueue ever runs -- log it so a single-item @@ -141,6 +142,92 @@ export function createRefresh({ pacer = noopPacer } = {}) { // bins build their own paced instance via createPacedRefresh. export default createRefresh(); +const SWEEP_INTERVAL_MS = 20 * 60 * 1000; +// Search can index an update late: on 2026-09-30 a comment on iFixit/ops#1119 +// took about 15 minutes to become findable by its updated_at. So each sweep +// reaches this far back into windows earlier sweeps covered, and skips hits it +// already re-read at the same updated_at, which keeps the overlap from costing +// re-reads. +const SWEEP_OVERLAP_MS = 60 * 60 * 1000; +// How far back the first sweep after a restart may reach. A day of changed +// pulls is a few hundred re-reads, well inside the hourly quota. +const SWEEP_MAX_CATCH_UP_MS = 24 * 60 * 60 * 1000; + +/** + * GitHub never resends a dropped webhook. Searching by owner also reaches repos + * missing from config.repos, which the startup refresh skips. + * + * `lastSeen` is the newest pull `date_updated` (epoch seconds) the DB held at + * boot, which is about when the previous run stopped hearing from GitHub. + */ +export function startRecentPullsSweep(refreshApi, repos, lastSeen = null) { + const sweep = createRecentPullsSweep(refreshApi, repoOwners(repos), Date.now, lastSeen); + const scheduleNext = () => setTimeout(() => sweep().then(scheduleNext), SWEEP_INTERVAL_MS); + scheduleNext(); +} + +export function repoOwners(repos) { + return _.uniq(repos.map(repo => repo.name.split('/')[0].toLowerCase())); +} + +/** + * The sweep never rejects. A failed search keeps its window for the next sweep. + * A pull that fails to fetch, parse or save is logged and retried by the next + * sweep that finds it. + */ +export function createRecentPullsSweep(refreshApi, owners, now = Date.now, lastSeen = null) { + // The first sweep also covers a restart shorter than one interval. After a + // longer outage it reaches back to `lastSeen`, but no further than a day. + let lastStart = now() - SWEEP_INTERVAL_MS; + if (lastSeen !== null) { + lastStart = Math.max(now() - SWEEP_MAX_CATCH_UP_MS, Math.min(lastStart, lastSeen * 1000)); + } + // `repo#number` -> the updated_at a hit had when the sweep last re-read it. + const reread = new Map(); + return async function sweep() { + const startedAt = now(); + const since = new Date(lastStart - SWEEP_OVERLAP_MS); + let pulls = []; + try { + for (const owner of owners) { + pulls = pulls.concat(await gitManager.searchUpdatedPulls(owner, since)); + } + } catch (err) { + console.error( + 'Failed to search for pulls updated since %s: %s', + since.toISOString(), + (err && err.message) || err + ); + return; + } + const due = pulls.filter(pull => { + const key = `${pull.repo}#${pull.number}`; + return !reread.has(key) || reread.get(key) !== pull.updatedAt; + }); + refreshDebug( + 'refreshing %s of %s pulls updated since %s', + due.length, + pulls.length, + since.toISOString() + ); + for (const { repo, number, updatedAt } of due) { + // refreshApi.pull resolves even when the parse or save fails. + let saved = true; + await refreshApi.pull(repo, number, () => (saved = false)).then( + () => saved && reread.set(`${repo}#${number}`, updatedAt), + () => {} + ); + } + // A hit updated before this window only comes back with a new updated_at. + for (const [key, updatedAt] of reread) { + if (Date.parse(updatedAt) < since.getTime()) { + reread.delete(key); + } + } + lastStart = startedAt; + }; +} + /** * Build the refresh API a CLI backfill bin runs against. One per-process pacer * is installed as a rate-limit observer on the GitHub client (so every @@ -183,14 +270,15 @@ export function makeQueueConsumer(pacer, processItem, deps) { * push its first argument to the specified Queue * and return a promise that is fulfilled when the item is fully processed. * - * Single-item refreshes (webhooks, socket) have no end-of-run report, so no - * failure collector is attached — and they aren't quota-paced, so live traffic - * is never delayed. + * Webhook and socket refreshes have no end-of-run report, so they attach no + * failure collector; the recent-pulls sweep passes one to learn whether the + * re-read saved. Single-item refreshes aren't quota-paced, so live traffic is + * never delayed. */ -function pushOnQueue(queue) { +function pushOnQueue(queue, onFailure = null) { return function (githubResponse) { return new Promise(function (resolve) { - queue.push({ response: githubResponse, onFailure: null }); + queue.push({ response: githubResponse, onFailure: onFailure }); queue.push(resolve); }); }; diff --git a/test/recent-pulls-sweep.test.js b/test/recent-pulls-sweep.test.js new file mode 100644 index 00000000..ce7b32b5 --- /dev/null +++ b/test/recent-pulls-sweep.test.js @@ -0,0 +1,226 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import gitManager from "../lib/git-manager.js"; +import dbManager from "../lib/db-manager.js"; +import { createRecentPullsSweep, createRefresh, repoOwners } from "../lib/refresh.js"; + +const MINUTE = 60 * 1000; + +function fakeClock(start) { + let time = start; + const now = () => time; + now.advance = (ms) => (time += ms); + return now; +} + +test("repoOwners lists each owner once, ignoring case", () => { + const repos = [ + { name: "iFixit/ifixit" }, + { name: "ifixit/expo" }, + { name: "other/repo" }, + ]; + assert.deepEqual(repoOwners(repos), ["ifixit", "other"]); +}); + +// Each window starts an hour before the previous sweep started, and the first +// one reaches back a full interval so it covers a short restart. +test("the sweep refreshes every pull the search finds and advances its window", async (t) => { + const searches = []; + t.mock.method(gitManager, "searchUpdatedPulls", (owner, since) => { + searches.push(`${owner} ${since.toISOString()}`); + // A new updated_at on every search, so each sweep re-reads both. + const updatedAt = new Date(now()).toISOString(); + return Promise.resolve( + owner === "a" + ? [{ repo: "a/one", number: 1, updatedAt }] + : [{ repo: "b/two", number: 2, updatedAt }] + ); + }); + const refreshed = []; + const refreshApi = { + pull: (repo, number) => { + refreshed.push(`${repo}#${number}`); + return Promise.resolve(); + }, + }; + const now = fakeClock(Date.parse("2026-09-29T12:00:00Z")); + const sweep = createRecentPullsSweep(refreshApi, ["a", "b"], now); + + await sweep(); + now.advance(20 * MINUTE); + await sweep(); + + assert.deepEqual(searches, [ + "a 2026-09-29T10:40:00.000Z", + "b 2026-09-29T10:40:00.000Z", + "a 2026-09-29T11:00:00.000Z", + "b 2026-09-29T11:00:00.000Z", + ]); + assert.deepEqual(refreshed, ["a/one#1", "b/two#2", "a/one#1", "b/two#2"]); +}); + +test("a failed search keeps its window for the next sweep", async (t) => { + const searches = []; + let fail = true; + t.mock.method(gitManager, "searchUpdatedPulls", (owner, since) => { + searches.push(since.toISOString()); + return fail ? Promise.reject(new Error("search 503")) : Promise.resolve([]); + }); + t.mock.method(console, "error", () => {}); + const now = fakeClock(Date.parse("2026-09-29T12:00:00Z")); + const sweep = createRecentPullsSweep( + { pull: () => Promise.resolve() }, + ["a"], + now + ); + + await sweep(); + fail = false; + now.advance(20 * MINUTE); + await sweep(); + + assert.deepEqual(searches, [ + "2026-09-29T10:40:00.000Z", + "2026-09-29T10:40:00.000Z", + ]); +}); + +// After an outage, the first window starts an hour before the newest update +// the DB held at boot, and its start is never more than a day back. +test("the first sweep after a restart reaches back to the newest update the DB held", async (t) => { + const searches = []; + t.mock.method(gitManager, "searchUpdatedPulls", (owner, since) => { + searches.push(since.toISOString()); + return Promise.resolve([]); + }); + const refreshApi = { pull: () => Promise.resolve() }; + const now = fakeClock(Date.parse("2026-09-29T15:00:00Z")); + const at = (iso) => Date.parse(iso) / 1000; + + // The 2026-09-29 cominor outage: nothing written after 10:28:56Z. + await createRecentPullsSweep(refreshApi, ["a"], now, at("2026-09-29T10:28:56Z"))(); + // Down for three days: capped at a day. + await createRecentPullsSweep(refreshApi, ["a"], now, at("2026-09-26T15:00:00Z"))(); + // Updated a minute before the restart: the usual 20-minute window. + await createRecentPullsSweep(refreshApi, ["a"], now, at("2026-09-29T14:59:00Z"))(); + + assert.deepEqual(searches, [ + "2026-09-29T09:28:56.000Z", + "2026-09-28T14:00:00.000Z", + "2026-09-29T13:40:00.000Z", + ]); +}); + +// The overlap finds the same hits again; only a new updated_at re-reads one. +test("a hit already re-read at the same updated_at is skipped", async (t) => { + let hits = [ + { repo: "a/one", number: 1, updatedAt: "2026-09-29T11:50:00Z" }, + { repo: "a/one", number: 2, updatedAt: "2026-09-29T11:50:00Z" }, + ]; + t.mock.method(gitManager, "searchUpdatedPulls", () => Promise.resolve(hits)); + const refreshed = []; + const refreshApi = { + pull: (repo, number) => { + refreshed.push(number); + return Promise.resolve(); + }, + }; + const now = fakeClock(Date.parse("2026-09-29T12:00:00Z")); + const sweep = createRecentPullsSweep(refreshApi, ["a"], now); + + await sweep(); + hits = [hits[0], { ...hits[1], updatedAt: "2026-09-29T12:10:00Z" }]; + now.advance(20 * MINUTE); + await sweep(); + + assert.deepEqual(refreshed, [1, 2, 2]); +}); + +test("a pull that failed to refresh is retried by the next sweep", async (t) => { + t.mock.method(gitManager, "searchUpdatedPulls", () => + Promise.resolve([{ repo: "a/one", number: 1, updatedAt: "2026-09-29T11:50:00Z" }]) + ); + let fail = true; + const attempts = []; + const refreshApi = { + pull: (repo, number) => { + attempts.push(number); + return fail ? Promise.reject(new Error("transient 500")) : Promise.resolve(); + }, + }; + const now = fakeClock(Date.parse("2026-09-29T12:00:00Z")); + const sweep = createRecentPullsSweep(refreshApi, ["a"], now); + + await sweep(); + fail = false; + now.advance(20 * MINUTE); + await sweep(); + now.advance(20 * MINUTE); + await sweep(); + + assert.deepEqual(attempts, [1, 1]); +}); + +// refreshApi.pull resolves after a failed parse or save, so the sweep learns +// about it only through the onFailure it passes. +test("a pull that failed to save is retried by the next sweep", async (t) => { + t.mock.method(gitManager, "searchUpdatedPulls", () => + Promise.resolve([{ repo: "a/one", number: 1, updatedAt: "2026-09-29T11:50:00Z" }]) + ); + let fail = true; + const attempts = []; + const refreshApi = { + pull: (repo, number, onFailure) => { + attempts.push(number); + if (fail) onFailure(repo, number); + return Promise.resolve(); + }, + }; + const now = fakeClock(Date.parse("2026-09-29T12:00:00Z")); + const sweep = createRecentPullsSweep(refreshApi, ["a"], now); + + await sweep(); + fail = false; + now.advance(20 * MINUTE); + await sweep(); + now.advance(20 * MINUTE); + await sweep(); + + assert.deepEqual(attempts, [1, 1]); +}); + +test("refresh.pull reports a failed save to its onFailure and still resolves", async (t) => { + t.mock.method(gitManager, "getPull", (repo, number) => + Promise.resolve({ number, base: { repo: { full_name: repo } } }) + ); + t.mock.method(gitManager, "parse", (response) => Promise.resolve(response)); + t.mock.method(dbManager, "updateAllPullData", () => Promise.reject(new Error("db down"))); + t.mock.method(console, "error", () => {}); + const failures = []; + + await createRefresh().pull("a/one", 1, (repo, number) => failures.push(`${repo}#${number}`)); + + assert.deepEqual(failures, ["a/one#1"]); +}); + +test("one pull failing to refresh doesn't stop the sweep", async (t) => { + t.mock.method(gitManager, "searchUpdatedPulls", () => + Promise.resolve([ + { repo: "a/one", number: 1 }, + { repo: "a/one", number: 2 }, + ]) + ); + const refreshed = []; + const refreshApi = { + pull: (repo, number) => { + refreshed.push(number); + return number === 1 + ? Promise.reject(new Error("transient 500")) + : Promise.resolve(); + }, + }; + + await createRecentPullsSweep(refreshApi, ["a"])(); + + assert.deepEqual(refreshed, [1, 2]); +});