Skip to content

Commit 277e6ba

Browse files
ralyodioclaude
andcommitted
Believe a 429 once for the whole host, instead of once per feed
Reported as the crawler being slow again: 5,000 feeds/hour down to 2,000. It was neither the database nor a knob, and the diagnostic order in the notes settled that in three minutes — an empty write transaction measured 0.39s (the old saturation was 29-300s), and the poller's own `started` event showed all four tuning variables already correct. The live log said it instead. 130 of 136 feed events in five minutes were `www.reddit.com`, and 106 of those were http-429. **A bulk import put 50,099 reddit feeds in the directory** — 30,000 on 08-28 and 19,999 on 08-27, two submissions, which is exactly when throughput halved. 49,174 of them were due at once: 41% of the entire due queue on a single hostname. Over one hour, reddit took 1,074 crawl attempts and was refused 1,015 times, while the whole rest of the directory got 218. The crawler had not slowed down. It was spending 84% of itself being told to go away, and had read 650 of those 50,099 feeds successfully in 36 hours of trying. Two things combine to make one host able to do that, and neither is wrong on its own. A host's feeds are crawled strictly in series — that is the politeness guarantee. And `spreadHosts` deliberately lifts its per-host cap rather than under-fill a batch, because the alternative throttles a legitimate bulk import to eight feeds a tick. So when the due set genuinely is one host, one worker gets most of a 900-feed batch on that host and walks it into the same rate limit nine hundred times. **The fix is to stop asking.** A 429 is a fact about the host, not about the feed that happened to be at the front of its queue, so the worker now believes it once and leaves. The rest of that host's queue is put back rather than walked. Three things that decide whether this is right: * **The held-back feeds are not blamed.** `markThrottled` already exists so a rate limit is not recorded as feed health — ten consecutive failures marks a feed dead, so counting 429s against feeds would retire a whole platform for our own crawl rate. `markHostThrottled` extends that to feeds we did not even send a request for, where there is even less to judge. * **They come back spread across an hour, not at one instant.** Handing a thousand held-back feeds a single timestamp just moves the pile-up one interval into the future, where it walks into the same wall together. * **It is bucketed by minute, not one statement per feed.** This database gives the cluster one writer, and a batch of nine hundred statements holds it long enough to push the jobs behind it past their deadline — the failure `TURSO_QUEUE_GROUP_STATEMENTS` was lowered to 10 to avoid. Bucketing makes the statement count a property of the window (~60) rather than of the queue. A 404 still walks the queue to the end: that is evidence about one feed and says nothing about the host. `crawlDue` gained a `crawl` option so a whole batch can be driven without the network — the throttle path only exists *across* a host's queue, so it was not reachable from `crawlFeed`'s tests at all. The tick summary gains `throttled`, which is what distinguishes a slow tick that was one publisher rate-limiting us from the crawler struggling; `ms` alone cannot tell those apart. Also done, in production rather than in code: the 50,099 reddit feeds were rescheduled across fourteen days, taking that host from 49,174 feeds due at once to ~145 an hour. Nothing was deleted and the update only ever pushed schedules out. Whether a directory that describes itself as the small web rather than the platforms should carry 50k subreddits at all is a separate question, and not one to answer by deleting someone's upload. Suite green: 1,206 tests, 0 failures across 11 packages. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gV7AyoWsiihaakh7RYZbP
1 parent a972f07 commit 277e6ba

4 files changed

Lines changed: 353 additions & 6 deletions

File tree

‎apps/poller/src/index.js‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -289,7 +289,7 @@ async function tick() {
289289
// stdout: they are the content of a log somebody is watching, and twenty-five
290290
// a minute forever is not what Railway's viewer is for. The tick summary
291291
// below is the durable record of the same work.
292-
const { crawled, failed, items, unchanged, hosts } = await crawlDue(
292+
const { crawled, failed, items, unchanged, throttled, hosts } = await crawlDue(
293293
db,
294294
batchSize,
295295
concurrency,
@@ -309,6 +309,12 @@ async function tick() {
309309
// entirely up to that server, so this cannot be predicted from here and
310310
// has to be watched.
311311
unchanged,
312+
// Hosts that answered 429 and had the rest of their queue put back
313+
// rather than walked into the same wall. This is the number that says a
314+
// slow tick was one publisher rate-limiting us and not the crawler
315+
// struggling — the two are indistinguishable from `ms` alone, and
316+
// telling them apart is what this whole change is for.
317+
throttled,
312318
// How many distinct hosts the batch touched. The number that says
313319
// whether the worker pool had anything to do: `crawled` and `ms`
314320
// together look identical for a batch spread over 80 hosts and one

‎packages/db/src/queries.js‎

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3022,6 +3022,70 @@ export async function markThrottled(db, id, minutes) {
30223022
});
30233023
}
30243024

3025+
/**
3026+
* Come back later for every feed still queued on a host that just throttled us.
3027+
*
3028+
* `markThrottled` answers the one feed that was refused. This answers the rest
3029+
* of that publisher's queue, and it exists because the one-feed version is the
3030+
* wrong shape of reply to a 429: a throttle is a fact about the *host*, not
3031+
* about the feed that happened to be at the front of it. Measured on
3032+
* 2026-08-29, a bulk import put 50,099 feeds on `www.reddit.com` — 41% of the
3033+
* whole due queue — and because a host's feeds are crawled strictly in series,
3034+
* one worker walked that queue into the same 429 over and over: 84% of every
3035+
* crawl attempt in an hour went to reddit and 1,015 of 1,074 were refused,
3036+
* while the rest of the directory got 218 crawls. Stopping at the first refusal
3037+
* turns that whole wasted walk into a single request.
3038+
*
3039+
* **The spread is the point, not a detail.** Rescheduling a thousand held-back
3040+
* feeds to one instant just moves the pile-up, and they would come due together
3041+
* and walk into the same wall in one tick's time. Each feed is given its own
3042+
* offset across `spreadMinutes`, so the host's queue returns as a trickle that
3043+
* its rate limit can actually absorb.
3044+
*
3045+
* Like `markThrottled`, only the schedule moves: `status`, `error_count`,
3046+
* `last_error`, `last_success_at` and `last_fetched_at` are all left alone.
3047+
* These feeds were never even asked, so there is nothing to judge them on —
3048+
* recording a failure here is how a rate limit would retire a whole platform.
3049+
*
3050+
* @param {Client} db
3051+
* @param {string[]} ids feeds left unread on the throttled host
3052+
* @param {number} minutes how long before the first of them is tried again
3053+
* @param {number} [spreadMinutes] window to scatter them across
3054+
* @returns {Promise<number>} how many were rescheduled
3055+
*/
3056+
export async function markHostThrottled(db, ids, minutes, spreadMinutes = 60) {
3057+
const queued = (Array.isArray(ids) ? ids : []).filter((id) => typeof id === 'string' && id);
3058+
if (queued.length === 0) return 0;
3059+
3060+
const wait = Math.max(1, Math.round(Number(minutes) || 30));
3061+
const spread = Math.max(0, Math.round(Number(spreadMinutes) || 0));
3062+
const now = nowIso();
3063+
3064+
// Bucketed by minute rather than one statement per feed, and that is a
3065+
// constraint rather than a tidiness preference. A throttled host's queue here
3066+
// runs to hundreds of feeds, and this database gives the whole cluster one
3067+
// writer: a batch of nine hundred statements holds it long enough to push the
3068+
// jobs behind it past their deadline, which is the failure
3069+
// `TURSO_QUEUE_GROUP_STATEMENTS` was lowered to 10 to avoid. Bucketing makes
3070+
// the statement count a property of the *window* -- at most `spread` + 1,
3071+
// sixty-odd -- no matter how many feeds are held back.
3072+
const buckets = new Map();
3073+
for (let index = 0; index < queued.length; index += 1) {
3074+
const offset = spread === 0 ? 0 : Math.round((index / queued.length) * spread);
3075+
const bucket = buckets.get(offset);
3076+
if (bucket) bucket.push(queued[index]);
3077+
else buckets.set(offset, [queued[index]]);
3078+
}
3079+
3080+
const statements = [...buckets.entries()].map(([offset, members]) => ({
3081+
sql: `update feeds set next_fetch_at = ?, updated_at = ? where id in (${members.map(() => '?').join(',')})`,
3082+
args: [nowIso((wait + offset) * 60_000), now, ...members],
3083+
}));
3084+
3085+
await db.batch(statements, 'write');
3086+
return queued.length;
3087+
}
3088+
30253089
/* ------------------------------------------------------------- feed cards */
30263090

30273091
/**

‎packages/ingest/src/crawl.js‎

Lines changed: 95 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -228,7 +228,15 @@ export async function crawlFeed(db, feed, opts = {}) {
228228
// So: come back when asked, leave every health column exactly as it was.
229229
if (resolved.throttled) {
230230
await q.markThrottled(db, id, throttleMinutes(resolved.retryAfter));
231-
return { ok: false, newItems: 0, throttled: true, error: resolved.error };
231+
// `retryAfter` travels with the result so the caller can hold the rest of
232+
// this host's queue back by the same interval the server itself named.
233+
return {
234+
ok: false,
235+
newItems: 0,
236+
throttled: true,
237+
retryAfter: resolved.retryAfter ?? null,
238+
error: resolved.error,
239+
};
232240
}
233241

234242
if (!resolved.ok) {
@@ -575,7 +583,9 @@ function hostOf(feed) {
575583
* @param {number} [batchSize]
576584
* @param {number} [concurrency] hosts crawled at once
577585
* @param {((event: { at: string, event: 'feed', status: 'ok'|'error', subject: string, slug: string|null, amount: number|null, detail: string|null, ms: number }) => void)|null} [onEvent]
578-
* @param {{ perHost?: number }} [opts]
586+
* @param {{ perHost?: number, crawl?: { resolve?: Function, scrape?: Function } }} [opts] `crawl` is
587+
* handed to each `crawlFeed`, which is what makes a whole batch testable without
588+
* the network — the throttle path in particular only exists across a host's queue.
579589
* @returns {Promise<{ crawled: number, failed: number, items: number, hosts: number }>} items being posts stored, not posts seen
580590
*/
581591
export async function crawlDue(db, batchSize = 25, concurrency = 8, onEvent = null, opts = {}) {
@@ -604,6 +614,9 @@ export async function crawlDue(db, batchSize = 25, concurrency = 8, onEvent = nu
604614
// out there -- it is entirely up to other people's servers, so it cannot be
605615
// predicted and has to be measured.
606616
let unchanged = 0;
617+
// Hosts abandoned mid-queue because they answered 429. Distinct from `failed`:
618+
// nothing is wrong with these feeds and they were mostly never even asked.
619+
let throttled = 0;
607620
let next = 0;
608621

609622
const worker = async () => {
@@ -612,14 +625,17 @@ export async function crawlDue(db, batchSize = 25, concurrency = 8, onEvent = nu
612625
next += 1;
613626
if (index >= queues.length) return;
614627

615-
for (const feed of queues[index]) {
628+
const queue = queues[index];
629+
630+
for (let position = 0; position < queue.length; position += 1) {
631+
const feed = queue[position];
616632
const started = Date.now();
617633

618634
// One feed that throws — a write that times out, a URL that breaks the
619635
// parser — must not reject the whole batch and take the other workers'
620636
// completed crawls down with it.
621637
try {
622-
const res = await crawlFeed(db, feed);
638+
const res = await crawlFeed(db, feed, opts.crawl);
623639
if (res.ok) {
624640
crawled += 1;
625641
items += Number(res.newItems ?? 0);
@@ -631,6 +647,20 @@ export async function crawlDue(db, batchSize = 25, concurrency = 8, onEvent = nu
631647
amount: res.ok ? Number(res.newItems ?? 0) : null,
632648
detail: res.ok ? null : (res.error ?? 'unknown'),
633649
});
650+
651+
// A 429 is the host speaking for all of its feeds, so believe it once
652+
// and leave. The rest of this queue is the same server, and asking it
653+
// the same question another eight hundred times is precisely what it
654+
// just told us to stop doing — see `markHostThrottled` for the
655+
// measurement that put this here.
656+
if (res.throttled) {
657+
throttled += 1;
658+
const rest = queue.slice(position + 1);
659+
if (rest.length > 0) {
660+
await holdBackHost(db, rest, throttleMinutes(res.retryAfter), onEvent, feed);
661+
}
662+
break;
663+
}
634664
} catch (err) {
635665
failed += 1;
636666
report(onEvent, feed, started, {
@@ -649,7 +679,67 @@ export async function crawlDue(db, batchSize = 25, concurrency = 8, onEvent = nu
649679
// `hosts` is what says whether the spread is working: a batch of 300 feeds
650680
// across 4 hosts cannot go faster than its biggest queue however many workers
651681
// are pointed at it, and the number is otherwise invisible from outside.
652-
return { crawled, failed, items, unchanged, hosts: queues.length };
682+
return { crawled, failed, items, unchanged, throttled, hosts: queues.length };
683+
}
684+
685+
/**
686+
* Put back the feeds a throttled host never got asked about.
687+
*
688+
* Separated from the worker so the worker reads as the policy — believe the
689+
* 429, leave — rather than as the bookkeeping. Two things it must not do, and
690+
* both are why it exists as its own function:
691+
*
692+
* * it must not fail the batch. These feeds are fine; the reschedule is an
693+
* optimisation and the crawl has already done its useful work. If the write
694+
* times out they keep their existing `next_fetch_at` and are simply offered
695+
* again next tick, which is the behaviour this replaces.
696+
* * it must not report the held-back feeds as errors. Nothing was wrong with
697+
* them and nothing was even sent, so a per-feed line would put hundreds of
698+
* healthy feeds on the failure panel. One line for the host says it.
699+
*
700+
* @param {import('@libsql/client').Client} db
701+
* @param {Array<{ id: string }>} rest feeds left unread on this host
702+
* @param {number} minutes how long before the first is tried again
703+
* @param {((event: object) => void)|null} onEvent
704+
* @param {{ feed_url: string, title?: unknown }} feed the one that was refused
705+
* @returns {Promise<void>}
706+
*/
707+
async function holdBackHost(db, rest, minutes, onEvent, feed) {
708+
const started = Date.now();
709+
710+
try {
711+
// Spread across an hour rather than returned all at once: a thousand feeds
712+
// handed back to the same instant is the same pile-up one tick later.
713+
await q.markHostThrottled(
714+
db,
715+
rest.map((f) => f.id),
716+
minutes,
717+
60,
718+
);
719+
} catch {
720+
// The schedule is unchanged, so they come back next tick exactly as they
721+
// would have without this. Losing the batch over it would be worse.
722+
return;
723+
}
724+
725+
if (typeof onEvent !== 'function') return;
726+
727+
try {
728+
onEvent({
729+
at: new Date().toISOString(),
730+
event: 'host-throttled',
731+
status: 'info',
732+
subject: hostOf(feed),
733+
slug: null,
734+
amount: rest.length,
735+
// Not `message`: `toEntry` reads any row carrying one as an error, and
736+
// this is the crawler behaving correctly. See the daemon error panel note.
737+
detail: `held back ${rest.length} feeds for ${minutes}m`,
738+
ms: Date.now() - started,
739+
});
740+
} catch {
741+
// A broken listener loses its line and nothing else.
742+
}
653743
}
654744

655745
/**

0 commit comments

Comments
 (0)