diff --git a/src/scraper/strategies/BaseScraperStrategy.test.ts b/src/scraper/strategies/BaseScraperStrategy.test.ts index e9278a0c..cb9028ae 100644 --- a/src/scraper/strategies/BaseScraperStrategy.test.ts +++ b/src/scraper/strategies/BaseScraperStrategy.test.ts @@ -1985,6 +1985,158 @@ 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("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. + 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..143b8c57 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,21 @@ 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; + + // 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.`); @@ -797,28 +833,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([