Skip to content

Commit 889f04c

Browse files
ralyodioclaude
andauthored
Drain an upload even while the crawler is catching up (#135)
A list of 109,474 URLs was submitted at 09:21 today. Every entry staged into import_entries within 55 seconds, and then nothing happened at all: queued_count stayed 0, and the page the submitter was given said "0% of 0 of 109,474 feeds — working". It would have said that for ever. The cause was position. `drainImport` sat eleven lines below `if (catchupOnly) return`, and CRAWL_CATCHUP=1 is set on the poller. Nothing errored and nothing logged, because from the tick's point of view nothing went wrong — it simply never reached the queue somebody was waiting on. The only way to find it was to read this file. Catch-up mode is right about everything else it skips. Discovery finds feeds on its own schedule, cards and clusters decorate what is already indexed, and none of it has a person watching a page. An import is the opposite: somebody handed over a list and was given a URL to follow. A slice a tick is a small enough price that it never had to be deferred. `notifyFinishedSubmissions` moves with it, because the two are a pair — this file already said so, that the daemon draining the queue is the one that tells the submitter it drained. Draining in catch-up mode while leaving that below would finish an upload and never say so, which is a stranger failure than not draining. It is now guarded, too: down there a throw cost only the housekeeping that followed, up here it would cost the crawl. `notifyFinishedDiscoveries` deliberately stays below. Nobody is waiting on it. The test is a source-order assertion, which is a compromise worth naming: index.js connects to a database and starts timers at import, so `tick()` cannot be loaded by a test without a refactor far larger than the bug. Checking the order of three lines is crude and it does hold the one invariant whose violation is invisible from outside. Verified it fails three ways with the drain moved back below the return. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 36595fd commit 889f04c

2 files changed

Lines changed: 130 additions & 26 deletions

File tree

‎apps/poller/src/index.js‎

Lines changed: 52 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -310,6 +310,58 @@ async function tick() {
310310
}
311311
}
312312

313+
// An upload hands over its list and leaves; this is where the list becomes
314+
// feeds. One slice a tick rather than a whole submission, so a very large
315+
// catalogue cannot hold the crawl above hostage while it queues — the two
316+
// share the process and the import is the one that can wait.
317+
//
318+
// Above the catch-up return, and deliberately the only thing that is.
319+
// Everything below is work nobody asked for by name: discovery finds feeds
320+
// on its own schedule, cards and clusters decorate what is already indexed,
321+
// and none of it has somebody watching a page. An import is the opposite —
322+
// a person handed us a list, was given a URL to follow, and that page says
323+
// "working". In catch-up mode it said "working" for ever: a 109,474-entry
324+
// upload on 2026-08-19 staged all its entries in 55 seconds and then sat at
325+
// 0 queued, because this block was eleven lines below a `return`. Nothing
326+
// errored, nothing logged, and the only way to find out was to read the
327+
// poller source.
328+
//
329+
// Catch-up mode exists to stop unrelated writes competing with a deep
330+
// first-crawl backlog, and that reasoning holds for everything else here.
331+
// It does not hold for the one queue a user is actively waiting on, and a
332+
// slice a tick is a small enough price that it never had to.
333+
try {
334+
const drained = await drainImport(db);
335+
if (drained.ran) {
336+
log('import-drain', {
337+
submission: drained.submissionId,
338+
queued: drained.queued,
339+
skipped: drained.skipped,
340+
remaining: drained.remaining,
341+
finished: drained.finished,
342+
});
343+
}
344+
} catch (err) {
345+
// One bad slice must not take the crawl down with it; the entries are
346+
// still staged, so the next tick tries again.
347+
log('import-drain-error', { message: String(err?.message ?? err) });
348+
}
349+
350+
// Above the return for the same reason, and because the two are a pair: the
351+
// daemon that drains the queue is the one that tells the submitter it
352+
// drained. Draining in catch-up mode while leaving this below would finish
353+
// somebody's upload and never say so, which is a stranger failure than not
354+
// draining at all. A no-op when no mail provider is configured.
355+
//
356+
// Guarded, unlike where it used to sit. Down there a throw only cost the
357+
// housekeeping that followed it; up here it would cost the crawl.
358+
try {
359+
const notified = await notifyFinishedSubmissions(db);
360+
if (notified.sent || notified.failed) log('notified', notified);
361+
} catch (err) {
362+
log('notify-error', { message: String(err?.message ?? err) });
363+
}
364+
313365
// All work below is resumable enrichment or housekeeping. In recovery mode
314366
// the deep first-crawl queue is the job, and returning here lets the next
315367
// minute tick begin immediately instead of waiting behind unrelated writes.
@@ -360,32 +412,6 @@ async function tick() {
360412
}
361413
}
362414

363-
// An upload hands over its list and leaves; this is where the list becomes
364-
// feeds. One slice a tick rather than a whole submission, so a very large
365-
// catalogue cannot hold the crawl above hostage while it queues — the two
366-
// share the process and the import is the one that can wait.
367-
try {
368-
const drained = await drainImport(db);
369-
if (drained.ran) {
370-
log('import-drain', {
371-
submission: drained.submissionId,
372-
queued: drained.queued,
373-
skipped: drained.skipped,
374-
remaining: drained.remaining,
375-
finished: drained.finished,
376-
});
377-
}
378-
} catch (err) {
379-
// One bad slice must not take the crawl down with it; the entries are
380-
// still staged, so the next tick tries again.
381-
log('import-drain-error', { message: String(err?.message ?? err) });
382-
}
383-
384-
// Queued submissions finish long after the upload, so the daemon that
385-
// drains the queue is also what tells the submitter it drained. A no-op
386-
// when no mail provider is configured.
387-
const notified = await notifyFinishedSubmissions(db);
388-
if (notified.sent || notified.failed) log('notified', notified);
389415

390416
const toldSearchers = await notifyFinishedDiscoveries(db);
391417
if (toldSearchers.sent || toldSearchers.failed) log('notified-discovery', toldSearchers);
Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,78 @@
1+
import assert from 'node:assert/strict';
2+
import { test } from 'node:test';
3+
import { readFile } from 'node:fs/promises';
4+
import { dirname, join } from 'node:path';
5+
import { fileURLToPath } from 'node:url';
6+
7+
/**
8+
* That catch-up mode cannot strand somebody's upload again.
9+
*
10+
* On 2026-08-19 a 109,474-entry list was submitted. Every entry staged in 55
11+
* seconds and then nothing happened at all: `queued_count` stayed 0, the
12+
* submitter's page said "working", and it would have said so for ever. The
13+
* cause was position — `drainImport` sat eleven lines *below*
14+
* `if (catchupOnly) return`, and `CRAWL_CATCHUP=1` is set on the poller.
15+
*
16+
* Nothing errored and nothing logged, because from the tick's point of view
17+
* nothing went wrong; it simply never reached the queue a person was waiting
18+
* on. The only way to find it was to read this file.
19+
*
20+
* This is a source-order assertion rather than a real test of `tick()`, and
21+
* that is a compromise worth naming: `apps/poller/src/index.js` connects to a
22+
* database and starts timers at import, so it cannot be loaded by a test
23+
* without a refactor far larger than the bug. Checking the order of three
24+
* lines in the source is crude, and it does hold the one invariant whose
25+
* violation is invisible from the outside.
26+
*/
27+
28+
const POLLER = join(dirname(fileURLToPath(import.meta.url)), '..', 'src', 'index.js');
29+
const source = await readFile(POLLER, 'utf8');
30+
31+
/**
32+
* Where a snippet appears, asserted to be present exactly once so a rename
33+
* fails loudly here rather than silently passing on the wrong line.
34+
*
35+
* @param {string} needle
36+
* @returns {number}
37+
*/
38+
function positionOf(needle) {
39+
const first = source.indexOf(needle);
40+
assert.notEqual(first, -1, `expected to find ${JSON.stringify(needle)} in the poller`);
41+
assert.equal(
42+
source.indexOf(needle, first + 1),
43+
-1,
44+
`${JSON.stringify(needle)} appears more than once; this check needs a better anchor`,
45+
);
46+
return first;
47+
}
48+
49+
test('the catch-up early return still exists to be checked against', () => {
50+
// If this ever goes away the rest of the file is asserting nothing, and a
51+
// silently vacuous test is worse than no test.
52+
positionOf('if (catchupOnly) return;');
53+
});
54+
55+
test('an upload is drained even while the crawler is catching up', () => {
56+
assert.ok(
57+
positionOf('await drainImport(db)') < positionOf('if (catchupOnly) return;'),
58+
'drainImport must run before the catch-up return, or a submission sits at 0% for ever',
59+
);
60+
});
61+
62+
test('and the submitter is told, since draining without telling is stranger still', () => {
63+
assert.ok(
64+
positionOf('await notifyFinishedSubmissions(db)') < positionOf('if (catchupOnly) return;'),
65+
'notifyFinishedSubmissions must run before the catch-up return: the daemon that drains ' +
66+
'the queue is the one that says it drained',
67+
);
68+
});
69+
70+
test('discovery stays below the return, because nobody is waiting on it', () => {
71+
// The other half of the rule. Everything catch-up mode skips is work the
72+
// crawler gave itself; an import is work a person handed over and was given a
73+
// URL to watch. If this ever moves up, the mode has stopped meaning anything.
74+
assert.ok(
75+
positionOf('await notifyFinishedDiscoveries(db)') > positionOf('if (catchupOnly) return;'),
76+
'discovery notifications are housekeeping and belong after the catch-up return',
77+
);
78+
});

0 commit comments

Comments
 (0)