From 37692dbf6503a855a386303918b294960e736f25 Mon Sep 17 00:00:00 2001 From: Andre Date: Sun, 27 Sep 2026 16:01:09 -0700 Subject: [PATCH 1/2] fix(scraper): keep crawling while one page is slow The crawl processed pages in fixed batches of maxConcurrency and waited for the whole batch before starting the next one. One slow item therefore idled every other slot. On brendangregg.com, Guess/guess.py, guess.pl and guess.pl.txt always answer HTTP 500 (Apache runs them as CGI), and each is retried with 1+2+4 s of backoff, so the crawl stood still for about 7 s at a time. The Jobs page shows the last finished URL, which is why the reporter saw it "stop" at guess.qbas or guess.ps depending on how the batch lined up. The loop now runs a sliding window: a slot is refilled as soon as its item finishes. Discovered links are still admitted in dequeue order, not in completion order, so the queue stays breadth-first, every URL is still first offered at its shortest depth, and the crawl order does not depend on timing. Running items count against the remaining maxPages budget, so the limit is never overshot. processBatch is split into processQueueItem (one item to its outcome) and admitToQueue (deduplication and queue totals). Most of the diff is the re-indentation of that body. Fixes #490 --- .../strategies/BaseScraperStrategy.test.ts | 121 +++ src/scraper/strategies/BaseScraperStrategy.ts | 777 ++++++++++-------- .../strategies/WebScraperStrategy.test.ts | 20 +- 3 files changed, 550 insertions(+), 368 deletions(-) diff --git a/src/scraper/strategies/BaseScraperStrategy.test.ts b/src/scraper/strategies/BaseScraperStrategy.test.ts index e9278a0c..2370d31d 100644 --- a/src/scraper/strategies/BaseScraperStrategy.test.ts +++ b/src/scraper/strategies/BaseScraperStrategy.test.ts @@ -1985,6 +1985,127 @@ describe("BaseScraperStrategy progress counter semantics", () => { }); }); +describe("BaseScraperStrategy concurrent processing", () => { + const opts = (o: Partial = {}): ScraperOptions => ({ + url: "https://example.com", + library: "test", + version: "1.0", + maxPages: 1000, + maxDepth: 3, + maxConcurrency: 3, + ...o, + }); + + const page = (url: string, links: string[] = []) => ({ + url, + links, + status: FetchStatus.SUCCESS, + content: { textContent: "content", chunks: [], links: [], errors: [] }, + }); + + /** A promise the test resolves by hand, to hold one item mid-flight. */ + const gate = () => { + let open: () => void = () => {}; + const closed = new Promise((resolve) => { + open = resolve; + }); + return { closed, open }; + }; + + it("keeps the other slots working while one item is slow", async () => { + // A URL stuck in retry backoff used to hold its whole batch, so the crawl + // looked frozen on whichever page had finished last (#490). + const strategy = new TestScraperStrategy(loadConfig()); + const slow = gate(); + const children = ["slow", "a", "b", "c", "d"].map((p) => `https://example.com/${p}`); + const finishedMeanwhile: string[] = []; + let slowDone = false; + strategy.processItem.mockImplementation(async (item: QueueItem) => { + if (item.depth === 0) return page(item.url, children); + if (item.url.endsWith("/slow")) { + await slow.closed; + slowDone = true; + } else if (!slowDone) { + finishedMeanwhile.push(item.url); + } + return page(item.url); + }); + const cb = vi.fn>(); + + const crawl = strategy.scrape(opts(), cb); + await vi.waitFor(() => expect(finishedMeanwhile).toHaveLength(4)); + slow.open(); + await crawl; + + expect(cb.mock.calls.at(-1)?.[0]).toMatchObject({ + pagesScraped: 6, + totalPages: 6, + pagesIndexed: 6, + }); + }); + + it("offers each URL at its shortest depth when items finish out of order", async () => { + // A links straight to X; B reaches X only through C. B finishing first must + // not let C claim X one level deeper than A offers it. + const strategy = new TestScraperStrategy(loadConfig()); + const a = gate(); + const links: Record = { + "https://example.com": ["https://example.com/A", "https://example.com/B"], + "https://example.com/A": ["https://example.com/X"], + "https://example.com/B": ["https://example.com/C"], + "https://example.com/C": ["https://example.com/X"], + }; + strategy.processItem.mockImplementation(async (item: QueueItem) => { + if (item.url === "https://example.com/A") await a.closed; + return page(item.url, links[item.url] ?? []); + }); + const cb = vi.fn>(); + + const crawl = strategy.scrape(opts(), cb); + await vi.waitFor(() => + expect(cb.mock.calls.map(([e]) => e.currentUrl)).toContain("https://example.com/B"), + ); + a.open(); + await crawl; + + const processed: QueueItem[] = strategy.processItem.mock.calls.map(([item]) => item); + expect(processed.map((item) => item.url)).toEqual([ + "https://example.com", + "https://example.com/A", + "https://example.com/B", + "https://example.com/X", + "https://example.com/C", + ]); + expect(processed.find((item) => item.url.endsWith("/X"))?.depth).toBe(2); + }); + + it("refills a freed slot without indexing past maxPages", async () => { + // A slot freed by an item that stored nothing may be refilled, but only + // while the running items could not already reach the limit between them. + const strategy = new TestScraperStrategy(loadConfig()); + const links = ["skip", "p1", "p2", "p3", "p4"].map((p) => `https://example.com/${p}`); + strategy.processItem.mockImplementation(async (item: QueueItem) => { + if (item.depth === 0) return page(item.url, links); + if (item.url.endsWith("/skip")) { + await new Promise((resolve) => setTimeout(resolve, 20)); + return { url: item.url, links: [], status: FetchStatus.SKIPPED }; + } + return page(item.url); + }); + const cb = vi.fn>(); + + await strategy.scrape(opts({ maxPages: 3 }), cb); + + expect(Math.max(...cb.mock.calls.map(([e]) => e.pagesIndexed ?? 0))).toBe(3); + expect(strategy.processItem).toHaveBeenCalledTimes(4); + expect(cb.mock.calls.at(-1)?.[0]).toMatchObject({ + pagesScraped: 4, + totalPages: 4, + pagesIndexed: 3, + }); + }); +}); + describe("BaseScraperStrategy empty-page reporting", () => { const opts = (): ScraperOptions => ({ url: "https://example.com/page", diff --git a/src/scraper/strategies/BaseScraperStrategy.ts b/src/scraper/strategies/BaseScraperStrategy.ts index 911332fc..84d7ecf3 100644 --- a/src/scraper/strategies/BaseScraperStrategy.ts +++ b/src/scraper/strategies/BaseScraperStrategy.ts @@ -76,6 +76,12 @@ export interface ProcessItemResult { class FailureThresholdExceededError extends ScraperError {} +/** Abort state shared by the items of one crawl. */ +interface CrawlAbortState { + /** Set once the child-page failure rate trips; running items stop on it. */ + abortError: FailureThresholdExceededError | null; +} + export abstract class BaseScraperStrategy implements ScraperStrategy { private static readonly FAILURE_RATE_MIN_SAMPLE = 10; @@ -88,7 +94,7 @@ export abstract class BaseScraperStrategy implements ScraperStrategy { * Usage flow: * 1. Initial queue setup: Root URL and initialQueue items are added to visited * 2. During processing: When a page returns links, each link is checked against visited - * 3. In processBatch deduplication: Only links NOT in visited are added to the queue AND to visited + * 3. In admitToQueue deduplication: Only links NOT in visited are added to the queue AND to visited * * This approach ensures: * - No URL is processed more than once @@ -322,388 +328,404 @@ export abstract class BaseScraperStrategy implements ScraperStrategy { } } - protected async processBatch( - batch: QueueItem[], + /** + * Processes one dequeued item to its outcome and returns the queue items it + * offers. + * + * The offered items have passed every discovered-link filter but are not yet + * deduplicated: admission belongs to {@link admitToQueue}, which the crawl loop + * calls in dequeue order so the queue stays breadth-first however the + * concurrent items finish. + * + * @param item The dequeued item. + * @param baseUrl Base for resolving links when the item has no final URL. + * @param options Scraper options for this crawl. + * @param progressCallback Receives the item's outcome. + * @param crawl Abort state shared by every item of this crawl. + * @param signal Cancels the crawl. + * @returns The queue items discovered by this item. + */ + protected async processQueueItem( + item: QueueItem, baseUrl: URL, options: ScraperOptions, progressCallback: ProgressCallback, - signal?: AbortSignal, // Add signal + crawl: CrawlAbortState, + signal?: AbortSignal, ): Promise { - let batchAbortError: FailureThresholdExceededError | null = null; - const ensureFailureRateWithinThreshold = (): void => { try { this.ensureFailureRateWithinThreshold(); } catch (error) { if (error instanceof FailureThresholdExceededError) { - batchAbortError ??= error; + crawl.abortError ??= error; } throw error; } }; - const throwIfBatchAborted = (): void => { - if (batchAbortError) { - throw batchAbortError; + // Items still running when the failure threshold trips stop at their next + // checkpoint instead of reporting outcomes for a crawl that has ended. + const throwIfCrawlAborted = (): void => { + if (crawl.abortError) { + throw crawl.abortError; } if (signal?.aborted) { - throw new CancellationError("Scraping cancelled during batch processing"); + throw new CancellationError("Scraping cancelled while processing"); } }; - const results = await Promise.all( - batch.map(async (item) => { - // Check signal before processing each item in the batch - throwIfBatchAborted(); - // Resolved for the progress log and the enqueue-time depth filter below. - // Items are no longer dropped here: every queued item reaches an outcome, - // which is what lets the processed count converge on the queued total. - const maxDepth = options.maxDepth ?? this.config.scraper.maxDepth; - - // Every dequeued item reaches exactly one outcome and advances the - // processed count once, which is what makes it converge on the queued - // total. Only `Stored` additionally advances the indexed count. - // - // Declared outside the try so the failure path reports through it too: - // the counter protocol has one implementation, not two. - const report = async ( - outcome: PageOutcome, - extra: Partial< - Pick - > = {}, - ): Promise => { - const pagesScraped = ++this.pageCount; - - // Two routes can reach one document — an llms.txt `.md` entry and the - // crawled HTML page — and both are sent to the store, which decides - // which representation to keep. Only the first of them is a new page: - // counting the second would let one document consume two units of the - // page budget and leave the crawl short of what the user asked for. - const identity = extra.result?.url ?? extra.emptyPage?.url; - const alreadyStored = - identity !== undefined && this.storedIdentities.has(identity); - const isNewPage = outcome === PageOutcome.Stored && !alreadyStored; - if (isNewPage) { - if (identity !== undefined) this.storedIdentities.add(identity); - this.pagesIndexed++; - } else { - this.retuneEffectiveTotal(options); - } - - logger.info( - `🌐 Scraping page ${pagesScraped}/${this.effectiveTotal} (depth ${item.depth}/${maxDepth}): ${item.url}`, - ); + // Check signal before processing the item + throwIfCrawlAborted(); + // Resolved for the progress log and the enqueue-time depth filter below. + // Items are no longer dropped here: every queued item reaches an outcome, + // which is what lets the processed count converge on the queued total. + const maxDepth = options.maxDepth ?? this.config.scraper.maxDepth; + + // Every dequeued item reaches exactly one outcome and advances the + // processed count once, which is what makes it converge on the queued + // total. Only `Stored` additionally advances the indexed count. + // + // Declared outside the try so the failure path reports through it too: + // the counter protocol has one implementation, not two. + const report = async ( + outcome: PageOutcome, + extra: Partial< + Pick + > = {}, + ): Promise => { + const pagesScraped = ++this.pageCount; + + // Two routes can reach one document — an llms.txt `.md` entry and the + // crawled HTML page — and both are sent to the store, which decides + // which representation to keep. Only the first of them is a new page: + // counting the second would let one document consume two units of the + // page budget and leave the crawl short of what the user asked for. + const identity = extra.result?.url ?? extra.emptyPage?.url; + const alreadyStored = identity !== undefined && this.storedIdentities.has(identity); + const isNewPage = outcome === PageOutcome.Stored && !alreadyStored; + if (isNewPage) { + if (identity !== undefined) this.storedIdentities.add(identity); + this.pagesIndexed++; + } else { + this.retuneEffectiveTotal(options); + } - await progressCallback({ - pagesScraped, - totalPages: this.effectiveTotal, - totalDiscovered: this.totalDiscovered, - pagesIndexed: this.pagesIndexed, - currentUrl: item.url, - depth: item.depth, - maxDepth, - outcome, - result: null, - pageId: item.pageId, - ...extra, - // Told to the store rather than derived there: only the crawl knows - // whether it has already stored this identity during this run, and - // that is what separates "a competing representation" from "the new - // state of the page", which are resolved in opposite directions. - ...(extra.result - ? { result: { ...extra.result, isAdditionalRepresentation: alreadyStored } } - : {}), - ...(extra.emptyPage - ? { - emptyPage: { - ...extra.emptyPage, - isAdditionalRepresentation: alreadyStored, - }, - } - : {}), - }); - }; - - try { - // Pass signal to processItem - const result = await this.processItem(item, options, signal); - throwIfBatchAborted(); - - if (result.status === FetchStatus.NOT_MODIFIED) { - // File/page hasn't changed, skip processing but count as processed - logger.debug(`Page unchanged (304): ${item.url}`); - await report(PageOutcome.Unchanged); - this.recordChildPageCompletion(item, options, result); - ensureFailureRateWithinThreshold(); - throwIfBatchAborted(); - return result.queueItems ?? []; - } + logger.info( + `🌐 Scraping page ${pagesScraped}/${this.effectiveTotal} (depth ${item.depth}/${maxDepth}): ${item.url}`, + ); - if (result.status === FetchStatus.NOT_FOUND) { - // A 404 at a representation is not a statement about the page. A - // refresh asks the location the content came from, which for a - // published Markdown file is not the page's own URL; a site that - // withdraws those files still serves the pages they described. - // Ask the identity before concluding anything, and send no stored - // validator with it — that one was issued by the resource that is - // now gone. If the identity 404s too, this item carries no further - // fallback and the page is deleted then. - const identityFallback: QueueItem[] = - item.identityUrl !== undefined && item.identityUrl !== item.url - ? // Spread so the item keeps whatever else governs how it is - // fetched and scoped; only the address and the validator - // change. `etag` is dropped because it was issued by the - // resource that just answered 404, and `identityUrl` because - // this IS the identity — a second 404 here is the page. - [{ ...item, url: item.identityUrl, etag: null, identityUrl: undefined }] - : []; - // Deliberately not named `isRefreshDeletion`: the method of that - // name answers "is this a refresh 404?" and is still consulted by - // `shouldCountTowardFailureThreshold`, which must keep ignoring - // these. This narrower question is whether the page is actually - // removed, which a pending fallback defers. - const deletesStoredPage = - this.isRefreshDeletion(item, result) && identityFallback.length === 0; - const fallbackQueueItems = [ - ...(result.queueItems ?? []), - ...identityFallback, - ]; - const hasNewFallbackQueueItem = fallbackQueueItems.some( - (queueItem) => - !this.visited.has( - normalizeUrl(queueItem.url, this.getUrlNormalizerOptions(options)), - ), - ); - - // Only the user's actual requested root should be fatal on 404 — an - // llms.txt-seeded depth-0 item is just one of several discovery seeds - // and a dead entry there shouldn't abort a scrape whose real root URL - // resolved fine (see llmstxt-discovery spec: llms.txt link failures - // are not supposed to fail the overall scrape). - if ( - this.isRequestedRoot(item, options) && - !deletesStoredPage && - !hasNewFallbackQueueItem - ) { - throw new ScraperError(`Root page not found: ${item.url}`, false); + await progressCallback({ + pagesScraped, + totalPages: this.effectiveTotal, + totalDiscovered: this.totalDiscovered, + pagesIndexed: this.pagesIndexed, + currentUrl: item.url, + depth: item.depth, + maxDepth, + outcome, + result: null, + pageId: item.pageId, + ...extra, + // Told to the store rather than derived there: only the crawl knows + // whether it has already stored this identity during this run, and + // that is what separates "a competing representation" from "the new + // state of the page", which are resolved in opposite directions. + ...(extra.result + ? { result: { ...extra.result, isAdditionalRepresentation: alreadyStored } } + : {}), + ...(extra.emptyPage + ? { + emptyPage: { + ...extra.emptyPage, + isAdditionalRepresentation: alreadyStored, + }, } + : {}), + }); + }; - if (!deletesStoredPage) { - this.recordChildPageFailure(item, options); - ensureFailureRateWithinThreshold(); - } + try { + // Pass signal to processItem + const result = await this.processItem(item, options, signal); + throwIfCrawlAborted(); + + if (result.status === FetchStatus.NOT_MODIFIED) { + // File/page hasn't changed, skip processing but count as processed + logger.debug(`Page unchanged (304): ${item.url}`); + await report(PageOutcome.Unchanged); + this.recordChildPageCompletion(item, options, result); + ensureFailureRateWithinThreshold(); + throwIfCrawlAborted(); + return result.queueItems ?? []; + } - throwIfBatchAborted(); + if (result.status === FetchStatus.NOT_FOUND) { + // A 404 at a representation is not a statement about the page. A + // refresh asks the location the content came from, which for a + // published Markdown file is not the page's own URL; a site that + // withdraws those files still serves the pages they described. + // Ask the identity before concluding anything, and send no stored + // validator with it — that one was issued by the resource that is + // now gone. If the identity 404s too, this item carries no further + // fallback and the page is deleted then. + const identityFallback: QueueItem[] = + item.identityUrl !== undefined && item.identityUrl !== item.url + ? // Spread so the item keeps whatever else governs how it is + // fetched and scoped; only the address and the validator + // change. `etag` is dropped because it was issued by the + // resource that just answered 404, and `identityUrl` because + // this IS the identity — a second 404 here is the page. + [{ ...item, url: item.identityUrl, etag: null, identityUrl: undefined }] + : []; + // Deliberately not named `isRefreshDeletion`: the method of that + // name answers "is this a refresh 404?" and is still consulted by + // `shouldCountTowardFailureThreshold`, which must keep ignoring + // these. This narrower question is whether the page is actually + // removed, which a pending fallback defers. + const deletesStoredPage = + this.isRefreshDeletion(item, result) && identityFallback.length === 0; + const fallbackQueueItems = [...(result.queueItems ?? []), ...identityFallback]; + const hasNewFallbackQueueItem = fallbackQueueItems.some( + (queueItem) => + !this.visited.has( + normalizeUrl(queueItem.url, this.getUrlNormalizerOptions(options)), + ), + ); - // File/page was deleted, count as processed - logger.debug(`Page deleted (404): ${item.url}`); - await report(PageOutcome.Absent, deletesStoredPage ? { deleted: true } : {}); - return fallbackQueueItems; - } + // Only the user's actual requested root should be fatal on 404 — an + // llms.txt-seeded depth-0 item is just one of several discovery seeds + // and a dead entry there shouldn't abort a scrape whose real root URL + // resolved fine (see llmstxt-discovery spec: llms.txt link failures + // are not supposed to fail the overall scrape). + if ( + this.isRequestedRoot(item, options) && + !deletesStoredPage && + !hasNewFallbackQueueItem + ) { + throw new ScraperError(`Root page not found: ${item.url}`, false); + } - if (result.status === FetchStatus.SKIPPED) { - // The resource was fetched but nothing can read its content type, so - // the body was abandoned at the headers. This is neither a success nor - // a failure: no page is stored, and the child-page failure rate is left - // untouched so an asset-heavy site cannot trip abortOnFailureRate. - // Only the user's actual requested root is fatal, matching the 404 - // branch above. An llms.txt seed is one of several discovery seeds, - // and a refresh replays stored pages at their stored depth — neither - // should abort a scrape whose real root resolved fine. - if (item.depth === 0 && !item.fromLlmsTxt && item.pageId === undefined) { - // Name the type: the user picked this URL, and "some content type" - // does not tell them whether they mistyped it or asked for a format - // this server cannot read. - const contentType = result.sourceContentType ?? "unknown content type"; - throw new ScraperError( - `Cannot process ${contentType} at ${item.url}: no pipeline can read it`, - false, - ); - } - logger.debug(`Skipped (unprocessable content): ${item.url}`); - await report(PageOutcome.Skipped); - return result.queueItems ?? []; - } + if (!deletesStoredPage) { + this.recordChildPageFailure(item, options); + ensureFailureRateWithinThreshold(); + } - if (result.status !== FetchStatus.SUCCESS) { - // Unreachable while FetchStatus has only the four handled members, but - // reporting here keeps "every dequeued item reaches exactly one - // outcome" true by construction rather than by enumeration. - logger.error(`❌ Unknown fetch status: ${result.status}`); - await report(PageOutcome.Failed); - return []; - } + throwIfCrawlAborted(); + + // File/page was deleted, count as processed + logger.debug(`Page deleted (404): ${item.url}`); + await report(PageOutcome.Absent, deletesStoredPage ? { deleted: true } : {}); + return fallbackQueueItems; + } + + if (result.status === FetchStatus.SKIPPED) { + // The resource was fetched but nothing can read its content type, so + // the body was abandoned at the headers. This is neither a success nor + // a failure: no page is stored, and the child-page failure rate is left + // untouched so an asset-heavy site cannot trip abortOnFailureRate. + // Only the user's actual requested root is fatal, matching the 404 + // branch above. An llms.txt seed is one of several discovery seeds, + // and a refresh replays stored pages at their stored depth — neither + // should abort a scrape whose real root resolved fine. + if (item.depth === 0 && !item.fromLlmsTxt && item.pageId === undefined) { + // Name the type: the user picked this URL, and "some content type" + // does not tell them whether they mistyped it or asked for a format + // this server cannot read. + const contentType = result.sourceContentType ?? "unknown content type"; + throw new ScraperError( + `Cannot process ${contentType} at ${item.url}: no pipeline can read it`, + false, + ); + } + logger.debug(`Skipped (unprocessable content): ${item.url}`); + await report(PageOutcome.Skipped); + return result.queueItems ?? []; + } + + if (result.status !== FetchStatus.SUCCESS) { + // Unreachable while FetchStatus has only the four handled members, but + // reporting here keeps "every dequeued item reaches exactly one + // outcome" true by construction rather than by enumeration. + logger.error(`❌ Unknown fetch status: ${result.status}`); + await report(PageOutcome.Failed); + return []; + } - // Handle successful processing - report result with content - // Use the final URL from the result (which may differ due to redirects) - // - // Canonicalised so that spellings differing only by a trailing slash or - // a fragment resolve to one page. Two routes to the same document — - // `/config/` from a crawl and `/config.md` from an llms.txt index — - // otherwise land as separate rows and split a page in two. - const finalUrl = this.canonicalizeStoredUrl(result.url || item.url, options); - - // Register the resolved identity so the other route to this page is - // recognised as already seen. The identity is only known after the - // response, which is why it cannot be settled when the URL is queued. - this.visited.add(normalizeUrl(finalUrl, this.getUrlNormalizerOptions(options))); - - // A result carrying no text is not a stored page. `WebScraperStrategy` - // already gates on this, but the local-file and GitHub processors pass - // their pipeline result through unconditionally, so an empty file would - // otherwise report Stored, inflate the indexed count and consume the - // page budget for a document the store then drops for having no chunks. - const producedContent = !!result.content?.textContent?.trim(); - if (result.content && producedContent) { - await report(PageOutcome.Stored, { - currentUrl: finalUrl, - result: { + // Handle successful processing - report result with content + // Use the final URL from the result (which may differ due to redirects) + // + // Canonicalised so that spellings differing only by a trailing slash or + // a fragment resolve to one page. Two routes to the same document — + // `/config/` from a crawl and `/config.md` from an llms.txt index — + // otherwise land as separate rows and split a page in two. + const finalUrl = this.canonicalizeStoredUrl(result.url || item.url, options); + + // Register the resolved identity so the other route to this page is + // recognised as already seen. The identity is only known after the + // response, which is why it cannot be settled when the URL is queued. + this.visited.add(normalizeUrl(finalUrl, this.getUrlNormalizerOptions(options))); + + // A result carrying no text is not a stored page. `WebScraperStrategy` + // already gates on this, but the local-file and GitHub processors pass + // their pipeline result through unconditionally, so an empty file would + // otherwise report Stored, inflate the indexed count and consume the + // page budget for a document the store then drops for having no chunks. + const producedContent = !!result.content?.textContent?.trim(); + if (result.content && producedContent) { + await report(PageOutcome.Stored, { + currentUrl: finalUrl, + result: { + url: finalUrl, + // Canonicalisation can move the identity too (`/docs/` to + // `/docs`), and then the bytes came from somewhere the identity + // no longer names. Record whichever URL actually served them. + contentUrl: + result.contentUrl ?? + (result.url && result.url !== finalUrl ? result.url : undefined), + title: result.content.title?.trim() || result.title?.trim() || "", + sourceContentType: result.sourceContentType || result.contentType || "", + contentType: result.contentType || "", + textContent: result.content.textContent || "", + links: result.content.links || [], + errors: result.content.errors || [], + chunks: result.content.chunks || [], + etag: result.etag || null, + lastModified: result.lastModified || null, + } satisfies ScrapeResult, + }); + throwIfCrawlAborted(); + } else { + // Fetched successfully but produced nothing to store: a directory + // listing, or a page whose pipeline extracted no text. Either way the + // item was processed, so it is reported rather than advancing the + // counter silently. + // A container that yields links — a directory listing, an archive — is + // processed and produces nothing, but it is not a page and must not be + // recorded as one. + const wasPage = result.isContainer !== true; + await report(PageOutcome.Empty, { + currentUrl: finalUrl, + emptyPage: !wasPage + ? undefined + : { url: finalUrl, - // Canonicalisation can move the identity too (`/docs/` to - // `/docs`), and then the bytes came from somewhere the identity - // no longer names. Record whichever URL actually served them. + // An empty page still has a retrieval location: `/guide` can + // be empty and have been read from `/guide.md`. Dropping it + // would send the next refresh to the identity carrying a + // validator the identity never issued. contentUrl: result.contentUrl ?? (result.url && result.url !== finalUrl ? result.url : undefined), - title: result.content.title?.trim() || result.title?.trim() || "", - sourceContentType: result.sourceContentType || result.contentType || "", - contentType: result.contentType || "", - textContent: result.content.textContent || "", - links: result.content.links || [], - errors: result.content.errors || [], - chunks: result.content.chunks || [], - etag: result.etag || null, - lastModified: result.lastModified || null, - } satisfies ScrapeResult, - }); - throwIfBatchAborted(); - } else { - // Fetched successfully but produced nothing to store: a directory - // listing, or a page whose pipeline extracted no text. Either way the - // item was processed, so it is reported rather than advancing the - // counter silently. - // A container that yields links — a directory listing, an archive — is - // processed and produces nothing, but it is not a page and must not be - // recorded as one. - const wasPage = result.isContainer !== true; - await report(PageOutcome.Empty, { - currentUrl: finalUrl, - emptyPage: !wasPage - ? undefined - : { - url: finalUrl, - // An empty page still has a retrieval location: `/guide` can - // be empty and have been read from `/guide.md`. Dropping it - // would send the next refresh to the identity carrying a - // validator the identity never issued. - contentUrl: - result.contentUrl ?? - (result.url && result.url !== finalUrl ? result.url : undefined), - title: result.title?.trim() || "", - sourceContentType: result.sourceContentType ?? null, - contentType: result.contentType ?? null, - // Withheld when the pipeline errored: storing the validator against - // a failure we do not understand would make the next refresh answer - // 304 and never retry, turning a transient fault permanent. - etag: result.pipelineFailed ? null : (result.etag ?? null), - lastModified: result.pipelineFailed - ? null - : (result.lastModified ?? null), - pipelineFailed: result.pipelineFailed === true, - }, - }); - throwIfBatchAborted(); - } - - // Extract discovered links - use the final URL as the base for resolving relative links - const nextItems = result.links || []; - const linkBaseUrl = finalUrl ? new URL(finalUrl) : baseUrl; - const internalAllowedFileRoots = - result.internalAllowedFileRoots ?? item.internalAllowedFileRoots; + title: result.title?.trim() || "", + sourceContentType: result.sourceContentType ?? null, + contentType: result.contentType ?? null, + // Withheld when the pipeline errored: storing the validator against + // a failure we do not understand would make the next refresh answer + // 304 and never retry, turning a transient fault permanent. + etag: result.pipelineFailed ? null : (result.etag ?? null), + lastModified: result.pipelineFailed + ? null + : (result.lastModified ?? null), + pipelineFailed: result.pipelineFailed === true, + }, + }); + throwIfCrawlAborted(); + } - this.recordChildPageCompletion(item, options, result); - ensureFailureRateWithinThreshold(); - throwIfBatchAborted(); - - // Depth is invariant across these links, so reject the whole set at - // once rather than parsing each URL and discarding it. - // - // Rejecting here rather than at dequeue is what keeps the progress - // denominator honest, and the rejection deliberately does NOT add the - // URL to `visited`. The queue is breadth-first in a normal crawl, so a - // URL is first offered at its minimum reachable depth and this cannot - // lose anything — but refresh mode builds its initial queue in database - // order, which is not sorted by depth, so the same URL can legitimately - // be offered again at a shallower depth later. Consuming a dedup slot - // here would discard it. - const childDepth = item.depth + 1; - const linkQueueItems = (childDepth > maxDepth ? [] : nextItems) - .map((value) => { - try { - const targetUrl = new URL(value, linkBaseUrl); - // Filter using shouldProcessUrl - if ( - !this.shouldProcessUrl(targetUrl.href, options, { - internalAllowedFileRoots, - }) - ) { - return null; - } - return { - url: targetUrl.href, - depth: childDepth, - ...(internalAllowedFileRoots ? { internalAllowedFileRoots } : {}), - } satisfies QueueItem; - } catch (_error) { - // Invalid URL or path - logger.warn(`❌ Invalid URL: ${value}`); - } + // Extract discovered links - use the final URL as the base for resolving relative links + const nextItems = result.links || []; + const linkBaseUrl = finalUrl ? new URL(finalUrl) : baseUrl; + const internalAllowedFileRoots = + result.internalAllowedFileRoots ?? item.internalAllowedFileRoots; + + this.recordChildPageCompletion(item, options, result); + ensureFailureRateWithinThreshold(); + throwIfCrawlAborted(); + + // Depth is invariant across these links, so reject the whole set at + // once rather than parsing each URL and discarding it. + // + // Rejecting here rather than at dequeue is what keeps the progress + // denominator honest, and the rejection deliberately does NOT add the + // URL to `visited`. The queue is breadth-first in a normal crawl, so a + // URL is first offered at its minimum reachable depth and this cannot + // lose anything — but refresh mode builds its initial queue in database + // order, which is not sorted by depth, so the same URL can legitimately + // be offered again at a shallower depth later. Consuming a dedup slot + // here would discard it. + const childDepth = item.depth + 1; + const linkQueueItems = (childDepth > maxDepth ? [] : nextItems) + .map((value) => { + try { + const targetUrl = new URL(value, linkBaseUrl); + // Filter using shouldProcessUrl + if ( + !this.shouldProcessUrl(targetUrl.href, options, { + internalAllowedFileRoots, + }) + ) { return null; - }) - .filter((item): item is QueueItem => item !== null); - - return [...(result.queueItems ?? []), ...linkQueueItems]; - } catch (error) { - if ( - error instanceof FailureThresholdExceededError || - error instanceof CancellationError - ) { - throw error; + } + return { + url: targetUrl.href, + depth: childDepth, + ...(internalAllowedFileRoots ? { internalAllowedFileRoots } : {}), + } satisfies QueueItem; + } catch (_error) { + // Invalid URL or path + logger.warn(`❌ Invalid URL: ${value}`); } + return null; + }) + .filter((item): item is QueueItem => item !== null); + + return [...(result.queueItems ?? []), ...linkQueueItems]; + } catch (error) { + if ( + error instanceof FailureThresholdExceededError || + error instanceof CancellationError + ) { + throw error; + } - // Never ignore errors for the root URL (depth 0) - if it fails, the job should fail - // There's no point in "successfully" completing with 0 documents - if (this.isRequestedRoot(item, options)) { - throw error; - } + // Never ignore errors for the root URL (depth 0) - if it fails, the job should fail + // There's no point in "successfully" completing with 0 documents + if (this.isRequestedRoot(item, options)) { + throw error; + } - if (batchAbortError) { - throw batchAbortError; - } + if (crawl.abortError) { + throw crawl.abortError; + } - this.recordChildPageFailure(item, options); - ensureFailureRateWithinThreshold(); + this.recordChildPageFailure(item, options); + ensureFailureRateWithinThreshold(); - if (options.ignoreErrors) { - logger.error(`❌ Failed to process ${item.url}: ${error}`); - await report(PageOutcome.Failed); - return []; - } - throw error; - } - }), - ); + if (options.ignoreErrors) { + logger.error(`❌ Failed to process ${item.url}: ${error}`); + await report(PageOutcome.Failed); + return []; + } + throw error; + } + } - // After all concurrent processing is done, deduplicate the results - const allLinks = results.flat().filter((item): item is QueueItem => item !== null); - const uniqueLinks: QueueItem[] = []; + /** + * Admits the offered items not yet seen and advances the queue totals. + * + * @param offered Items offered by one processed item, in discovery order. + * @param options Scraper options supplying the effective page limit. + * @returns The items admitted to the queue. + */ + private admitToQueue(offered: QueueItem[], options: ScraperOptions): QueueItem[] { + const admitted: QueueItem[] = []; - // Now perform deduplication once, after all parallel processing is complete - for (const item of allLinks) { + for (const item of offered) { const normalizedUrl = normalizeUrl(item.url, this.getUrlNormalizerOptions(options)); if (!this.visited.has(normalizedUrl)) { this.visited.add(normalizedUrl); - uniqueLinks.push(item); + admitted.push(item); // Always increment the unlimited counter this.totalDiscovered++; @@ -716,7 +738,7 @@ export abstract class BaseScraperStrategy implements ScraperStrategy { } } - return uniqueLinks; + return admitted; } async scrape( @@ -788,7 +810,18 @@ export abstract class BaseScraperStrategy implements ScraperStrategy { // `maxPages` bounds pages that produce content, not items processed: asking // for 100 pages should yield 100 pages, not stop at 100 attempts of which // some produced nothing. - while (queue.length > 0 && this.pagesIndexed < maxPages) { + // + // Items run in a sliding window: a slot is refilled as soon as its item + // finishes. Fixed batches waited for their slowest member, so one URL + // spending seconds in retry backoff idled every other slot and the crawl + // looked frozen on whichever page happened to finish last. + const crawl: CrawlAbortState = { abortError: null }; + const running = new Map>(); + const finished = new Map(); + let nextSeq = 0; + let nextToAdmit = 0; + + while (this.pagesIndexed < maxPages) { // Check for cancellation at the start of each loop iteration if (signal?.aborted) { logger.debug(`${isRefreshMode ? "Refresh" : "Scraping"} cancelled by signal.`); @@ -797,28 +830,52 @@ export abstract class BaseScraperStrategy implements ScraperStrategy { ); } - const remainingPages = maxPages - this.pagesIndexed; - if (remainingPages <= 0) { - break; + // Counting running items against the remaining budget is what stops + // concurrent work overshooting the limit: at most `maxPages - pagesIndexed` + // items run, so at most that many of them can index. + while ( + running.size < maxConcurrency && + this.pagesIndexed + running.size < maxPages + ) { + const item = queue.shift(); + if (!item) break; + const seq = nextSeq++; + // Always use latest canonical base (may have been updated after first fetch) + baseUrl = this.canonicalBaseUrl ?? baseUrl; + running.set( + seq, + this.processQueueItem( + item, + baseUrl, + options, + progressCallback, + crawl, + signal, + ).then((offered) => { + finished.set(seq, offered); + return seq; + }), + ); } - // Bounding the batch by the remaining budget is what stops a concurrent - // batch overshooting the limit: at most `remainingPages` items run, so at - // most `remainingPages` of them can index. - const batchSize = Math.min(maxConcurrency, remainingPages, queue.length); - const batch = queue.splice(0, batchSize); - - // Always use latest canonical base (may have been updated after first fetch) - baseUrl = this.canonicalBaseUrl ?? baseUrl; - const newUrls = await this.processBatch( - batch, - baseUrl, - options, - progressCallback, - signal, - ); + if (running.size === 0) { + break; + } - queue.push(...newUrls); + running.delete(await Promise.race(running.values())); + + // Admit in dequeue order, not completion order, so the queue stays + // breadth-first: a URL is first offered at its shortest depth however the + // running items finish, and the crawl order does not depend on timing. + for ( + let offered = finished.get(nextToAdmit); + offered !== undefined; + offered = finished.get(nextToAdmit) + ) { + finished.delete(nextToAdmit); + nextToAdmit++; + queue.push(...this.admitToQueue(offered, options)); + } } } diff --git a/src/scraper/strategies/WebScraperStrategy.test.ts b/src/scraper/strategies/WebScraperStrategy.test.ts index d2606f19..38ee6737 100644 --- a/src/scraper/strategies/WebScraperStrategy.test.ts +++ b/src/scraper/strategies/WebScraperStrategy.test.ts @@ -76,11 +76,12 @@ describe("WebScraperStrategy", () => { it("should carry internal archive roots to discovered temp archive members", async () => { const testStrategy = strategy as unknown as { - processBatch( - batch: QueueItem[], + processQueueItem( + item: QueueItem, baseUrl: URL, options: ScraperOptions, progressCallback: ProgressCallback, + crawl: { abortError: null }, signal?: AbortSignal, ): Promise; processItem( @@ -101,11 +102,12 @@ describe("WebScraperStrategy", () => { status: FetchStatus.SUCCESS, }); - const nextItems = await testStrategy.processBatch( - [{ url: "https://example.com/archive.zip", depth: 0 }], + const nextItems = await testStrategy.processQueueItem( + { url: "https://example.com/archive.zip", depth: 0 }, new URL("https://example.com/archive.zip"), { ...options, url: "https://example.com/archive.zip" }, vi.fn(), + { abortError: null }, ); expect(nextItems).toEqual([ @@ -121,11 +123,12 @@ describe("WebScraperStrategy", () => { it("should reject file:// links that escape the archive's internal roots", async () => { const testStrategy = strategy as unknown as { - processBatch( - batch: QueueItem[], + processQueueItem( + item: QueueItem, baseUrl: URL, options: ScraperOptions, progressCallback: ProgressCallback, + crawl: { abortError: null }, signal?: AbortSignal, ): Promise; processItem( @@ -148,11 +151,12 @@ describe("WebScraperStrategy", () => { status: FetchStatus.SUCCESS, }); - const nextItems = await testStrategy.processBatch( - [{ url: "https://example.com/archive.zip", depth: 0 }], + const nextItems = await testStrategy.processQueueItem( + { url: "https://example.com/archive.zip", depth: 0 }, new URL("https://example.com/archive.zip"), { ...options, url: "https://example.com/archive.zip", scope: "subpages" }, vi.fn(), + { abortError: null }, ); expect(nextItems).toEqual([ From 491a3162e8972eb22d7042e11ddab29328a6c273 Mon Sep 17 00:00:00 2001 From: Andre Date: Sun, 27 Sep 2026 16:12:07 -0700 Subject: [PATCH 2/2] fix(scraper): wait for running items before ending the crawl An item counts its page before its progress callback stores it. When the last pages that fit under maxPages ran concurrently, one could reach the limit while another was still storing, and the loop exited on the count alone, so scrape() resolved with a write still pending. Fixed batches never had this problem because Promise.all waited for every member. The loop now runs until nothing is left running. Reaching maxPages only stops new items from starting. --- .../strategies/BaseScraperStrategy.test.ts | 31 +++++++++++++++++++ src/scraper/strategies/BaseScraperStrategy.ts | 5 ++- 2 files changed, 35 insertions(+), 1 deletion(-) diff --git a/src/scraper/strategies/BaseScraperStrategy.test.ts b/src/scraper/strategies/BaseScraperStrategy.test.ts index 2370d31d..cb9028ae 100644 --- a/src/scraper/strategies/BaseScraperStrategy.test.ts +++ b/src/scraper/strategies/BaseScraperStrategy.test.ts @@ -2079,6 +2079,37 @@ describe("BaseScraperStrategy concurrent processing", () => { expect(processed.find((item) => item.url.endsWith("/X"))?.depth).toBe(2); }); + it("does not finish while a page at the limit is still being stored", async () => { + // Reaching maxPages stops new items from starting, but a running item may + // have counted its page and still be storing it. The crawl must not resolve + // before that write settles, or the job completes with a write pending. + const strategy = new TestScraperStrategy(loadConfig()); + const store = gate(); + strategy.processItem.mockImplementation(async (item: QueueItem) => + page( + item.url, + item.depth === 0 ? ["https://example.com/a", "https://example.com/b"] : [], + ), + ); + const cb = vi.fn(async (event: ScraperProgressEvent) => { + if (event.currentUrl === "https://example.com/a") await store.closed; + }); + + let settled = false; + const crawl = strategy.scrape(opts({ maxPages: 3 }), cb).finally(() => { + settled = true; + }); + await vi.waitFor(() => + expect(cb.mock.calls.map(([e]) => e.currentUrl)).toContain("https://example.com/b"), + ); + await new Promise((resolve) => setTimeout(resolve, 20)); + + expect(settled).toBe(false); + store.open(); + await crawl; + expect(cb.mock.calls.at(-1)?.[0]).toMatchObject({ pagesIndexed: 3 }); + }); + it("refills a freed slot without indexing past maxPages", async () => { // A slot freed by an item that stored nothing may be refilled, but only // while the running items could not already reach the limit between them. diff --git a/src/scraper/strategies/BaseScraperStrategy.ts b/src/scraper/strategies/BaseScraperStrategy.ts index 84d7ecf3..143b8c57 100644 --- a/src/scraper/strategies/BaseScraperStrategy.ts +++ b/src/scraper/strategies/BaseScraperStrategy.ts @@ -821,7 +821,10 @@ export abstract class BaseScraperStrategy implements ScraperStrategy { let nextSeq = 0; let nextToAdmit = 0; - while (this.pagesIndexed < maxPages) { + // Runs until nothing is left running. Reaching `maxPages` only stops new + // items from starting: an item counts its page before storing it, so exiting + // on the count alone would resolve the crawl with that write still pending. + while (true) { // Check for cancellation at the start of each loop iteration if (signal?.aborted) { logger.debug(`${isRefreshMode ? "Refresh" : "Scraping"} cancelled by signal.`);