diff --git a/packages/basic-crawler/src/internals/basic-crawler.ts b/packages/basic-crawler/src/internals/basic-crawler.ts index 54fd7ef66f3a..b82a023e63a1 100644 --- a/packages/basic-crawler/src/internals/basic-crawler.ts +++ b/packages/basic-crawler/src/internals/basic-crawler.ts @@ -47,6 +47,7 @@ import { RequestQueue, RequestQueueV1, RequestState, + RateLimitError, RetryRequestError, Router, SessionError, @@ -55,7 +56,7 @@ import { validators, } from '@crawlee/core'; import type { Awaitable, BatchAddRequestsResult, Dictionary, SetStatusMessageOptions } from '@crawlee/types'; -import { getObjectType, isAsyncIterable, isIterable, RobotsTxtFile, ROTATE_PROXY_ERRORS } from '@crawlee/utils'; +import { getObjectType, isAsyncIterable, isIterable, RobotsTxtFile, ROTATE_PROXY_ERRORS, sleep } from '@crawlee/utils'; import { stringify } from 'csv-stringify/sync'; import { ensureDir, writeFile, writeJSON } from 'fs-extra'; import ow, { ArgumentError } from 'ow'; @@ -263,6 +264,13 @@ export interface BasicCrawlerOptions; protected maxSessionRotations: number; protected maxRequestsPerCrawl?: number; @@ -598,6 +607,7 @@ export class BasicCrawler { + const getOperationMode = (): { + mode: 'ERROR' | 'REGULAR'; + failedDelta: number; + } => { const { requestsFailed } = this.stats.state; const { requestsFailed: previousRequestsFailed } = previousState; @@ -1253,15 +1271,24 @@ export class BasicCrawler { - return this.handleSkippedRequest({ url, reason: 'robotsTxt' }); + return this.handleSkippedRequest({ + url, + reason: 'robotsTxt', + }); }) .concat( skippedBecauseOfLimit.map((request) => { const url = typeof request === 'string' ? request : request.url!; - return this.handleSkippedRequest({ url, reason: 'limit' }); + return this.handleSkippedRequest({ + url, + reason: 'limit', + }); }), [...skippedBecauseOfMaxCrawlDepth].map((url) => { - return this.handleSkippedRequest({ url, reason: 'depth' }); + return this.handleSkippedRequest({ + url, + reason: 'depth', + }); }), ), ); @@ -1533,7 +1560,9 @@ export class BasicCrawler 0) { + this.log.debug(`Waiting ${delayMillis}ms before next attempt due to rate limiting (Retry-After).`); + await sleep(delayMillis); + } + } } finally { await this._cleanupContext(crawlingContext); @@ -1850,7 +1889,18 @@ export class BasicCrawler= 500 || blockedStatusCodes.includes(statusCode!); + const isTransientContentType = + statusCode! >= 500 || statusCode === 429 || blockedStatusCodes.includes(statusCode!); if (!this.supportedMimeTypes.has(type) && !this.supportedMimeTypes.has('*/*') && !isTransientContentType) { request.noRetry = true; diff --git a/test/core/crawlers/browser_crawler.test.ts b/test/core/crawlers/browser_crawler.test.ts index 1b07bd84a28e..e889301512ce 100644 --- a/test/core/crawlers/browser_crawler.test.ts +++ b/test/core/crawlers/browser_crawler.test.ts @@ -615,7 +615,7 @@ describe('BrowserCrawler', () => { await crawler.run(); - expect(failedRequests.length).toBe(3); + expect(failedRequests.length).toBe(BLOCKED_STATUS_CODES.length); failedRequests.forEach((fr) => { const [msg] = fr.errorMessages; expect(msg).toContain(`Request blocked - received ${fr.userData.statusCode} status code.`); @@ -750,7 +750,7 @@ describe('BrowserCrawler', () => { await crawler.run(); - expect(failedRequests.length).toBe(3); + expect(failedRequests.length).toBe(BLOCKED_STATUS_CODES.length); failedRequests.forEach((fr) => { const [msg] = fr.errorMessages; expect(msg).toContain(`Request blocked - received ${fr.userData.statusCode} status code.`); @@ -807,6 +807,55 @@ describe('BrowserCrawler', () => { } }); + test.concurrent('should handle 429 Rate Limit with Retry-After header', async () => { + const localStorageEmulator = new MemoryStorageEmulator(); + await localStorageEmulator.init(); + const puppeteerPlugin = new PuppeteerPlugin(puppeteer); + + try { + const succeeded: Request[] = []; + const crawler = new BrowserCrawlerTest({ + browserPoolOptions: { + browserPlugins: [puppeteerPlugin], + }, + useSessionPool: true, + sessionPoolOptions: { + maxPoolSize: 1, + }, + maxConcurrency: 1, + maxRequestRetries: 1, + requestHandler: async ({ request }) => { + succeeded.push(request); + }, + }); + + // @ts-expect-error Overriding protected method + crawler._navigationHandler = async ({ request }) => { + if (request.retryCount === 0) { + return { + status: () => 429, + headers: () => ({ 'retry-after': '1' }), + }; + } + + return { + status: () => 200, + headers: () => ({}), + }; + }; + + const start = Date.now(); + await crawler.run([serverAddress]); + const end = Date.now(); + + expect(succeeded).toHaveLength(1); + expect(succeeded[0].retryCount).toBe(1); + expect(end - start).toBeGreaterThanOrEqual(1000); + } finally { + await localStorageEmulator.destroy(); + } + }); + test.concurrent('should increment session usage correctly', async () => { const localStorageEmulator = new MemoryStorageEmulator(); await localStorageEmulator.init(); diff --git a/test/core/crawlers/cheerio_crawler.test.ts b/test/core/crawlers/cheerio_crawler.test.ts index cfba7f7b8044..2fee9c531778 100644 --- a/test/core/crawlers/cheerio_crawler.test.ts +++ b/test/core/crawlers/cheerio_crawler.test.ts @@ -933,7 +933,7 @@ describe('CheerioCrawler', () => { }); test('should retire session on "blocked" status codes', async () => { - for (const code of [401, 403, 429]) { + for (const code of [401, 403]) { const failed: Request[] = []; const sessions: Session[] = []; const crawler = new CheerioCrawler({ diff --git a/test/core/crawlers/http_crawler.test.ts b/test/core/crawlers/http_crawler.test.ts index 0aa96976aca7..2b7a6e9c00bd 100644 --- a/test/core/crawlers/http_crawler.test.ts +++ b/test/core/crawlers/http_crawler.test.ts @@ -61,6 +61,17 @@ router.set('/403-with-octet-stream', (req, res) => { res.end(); }); +router.set('/429-rate-limit', (req, res) => { + res.statusCode = 429; + res.setHeader('Retry-After', '1'); // 1 second + res.end(); +}); + +router.set('/429-rate-limit-no-header', (req, res) => { + res.statusCode = 429; + res.end(); +}); + let server: http.Server; let url: string; @@ -397,6 +408,81 @@ describe.each( expect(succeeded[0].retryCount).toBe(1); }); + test('should handle 429 Rate Limit with Retry-After header', async () => { + const succeeded: any[] = []; + const sessionIds: string[] = []; + const crawler = new HttpCrawler({ + httpClient, + maxConcurrency: 1, + maxRequestRetries: 1, + sessionPoolOptions: { + maxPoolSize: 1, + }, + preNavigationHooks: [ + async ({ request, session }, gotOptions) => { + sessionIds.push(session!.id); + if (request.retryCount === 0) { + request.url = `${url}/429-rate-limit`; + } else { + request.url = url; + } + if (gotOptions) { + gotOptions.throwHttpErrors = false; + } + }, + ], + requestHandler: async ({ request }) => { + succeeded.push(request); + }, + failedRequestHandler: async ({ request, error }) => { + console.error('FAILED', request.retryCount, request.url, error.message); + }, + }); + + const start = Date.now(); + await crawler.run([url]); + const end = Date.now(); + + expect(succeeded).toHaveLength(1); + expect(succeeded[0].retryCount).toBe(1); + expect(sessionIds).toHaveLength(2); + expect(sessionIds[0]).toBe(sessionIds[1]); + expect(end - start).toBeGreaterThanOrEqual(1000); // Should delay for at least 1s + }); + + test('should handle 429 Rate Limit without Retry-After header via cooldown', async () => { + const succeeded: any[] = []; + const crawler = new HttpCrawler({ + httpClient, + maxConcurrency: 1, + maxRequestRetries: 1, + rateLimitCooldownSecs: 1, + preNavigationHooks: [ + async ({ request }, gotOptions) => { + if (request.retryCount === 0) { + request.url = `${url}/429-rate-limit-no-header`; + } else { + request.url = url; + } + if (gotOptions) { + gotOptions.throwHttpErrors = false; + } + }, + ], + requestHandler: async ({ request }) => { + succeeded.push(request); + }, + }); + + const start = Date.now(); + await crawler.run([url]); + const end = Date.now(); + + expect(succeeded).toHaveLength(1); + expect(succeeded[0].retryCount).toBe(1); + expect(end - start).toBeGreaterThanOrEqual(1000); + }); + test.skipIf(httpClient instanceof ImpitHttpClient)('should work with cacheable-request', async () => { const isFromCache: Record = {}; const cache = new Map();