Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 12 additions & 3 deletions app.js
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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();
Expand All @@ -96,6 +104,7 @@ dbManager
.then(function () {
debug('Refreshing all open pulls from the API');
refresh.openPulls();
startRecentPullsSweep(refresh, config.repos, latestPullUpdate);
})
.done();

Expand Down
11 changes: 11 additions & 0 deletions lib/db-manager.js
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
28 changes: 28 additions & 0 deletions lib/git-manager.js
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
Expand Down
102 changes: 95 additions & 7 deletions lib/refresh.js
Original file line number Diff line number Diff line change
Expand Up @@ -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');
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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);
});
};
Expand Down
Loading
Loading