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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
152 changes: 152 additions & 0 deletions src/scraper/strategies/BaseScraperStrategy.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1985,6 +1985,158 @@ describe("BaseScraperStrategy progress counter semantics", () => {
});
});

describe("BaseScraperStrategy concurrent processing", () => {
const opts = (o: Partial<ScraperOptions> = {}): 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<void>((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<ProgressCallback<ScraperProgressEvent>>();

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<string, string[]> = {
"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<ProgressCallback<ScraperProgressEvent>>();

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<ProgressCallback<ScraperProgressEvent>>();

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",
Expand Down
Loading
Loading