From beaf123a639361b58b0939e80045cb6df10fbe5b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jind=C5=99ich=20B=C3=A4r?= Date: Thu, 24 Sep 2026 13:18:58 +0200 Subject: [PATCH 1/5] fix(http-client): do not reuse cookies of the previous hop on same-origin redirects (#4154) Sets the jar cookies on a per-hop copy of the request, like `fetch` does, so a redirect recomputes them for its own URL. Closes #4151 --- packages/http-client/src/base-http-client.ts | 5 ++--- test/core/base-http-client.test.ts | 20 ++++++++++++++++++++ 2 files changed, 22 insertions(+), 3 deletions(-) diff --git a/packages/http-client/src/base-http-client.ts b/packages/http-client/src/base-http-client.ts index eb98b85df532..921285985de9 100644 --- a/packages/http-client/src/base-http-client.ts +++ b/packages/http-client/src/base-http-client.ts @@ -200,9 +200,8 @@ export abstract class BaseHttpClient implements BaseHttpClientInterface { currentRequest = initialRequest.clone(); while (true) { - await this.applyCookies(currentRequest, cookieJar); - - const response = await this.fetch(currentRequest, { + // Like `fetch`, set the jar cookies on a copy, so that the next redirect hop does not inherit them + const response = await this.fetch(await this.applyCookies(new Request(currentRequest), cookieJar), { signal, proxyUrl, cookieJar, diff --git a/test/core/base-http-client.test.ts b/test/core/base-http-client.test.ts index 0e12b9171226..562cc7c99fea 100644 --- a/test/core/base-http-client.test.ts +++ b/test/core/base-http-client.test.ts @@ -180,6 +180,8 @@ describe('BaseHttpClient credentials on redirects', () => { res.writeHead(302, { location: `${targetUrl}/echo` }).end(); } else if (pathname === '/same-origin') { res.writeHead(302, { location: '/echo' }).end(); + } else if (pathname === '/set-cookie') { + res.writeHead(302, { 'location': '/echo', 'set-cookie': 'session=new' }).end(); } else { echoCredentials(req, res); } @@ -217,6 +219,24 @@ describe('BaseHttpClient credentials on redirects', () => { expect(await response.json()).toEqual(credentials); }); + + test('sends cookies updated by a same-origin redirect', async () => { + const cookieJar = new CookieJar(); + await cookieJar.setCookie('session=old', redirectorUrl); + + const response = await httpClient.sendRequest(new Request(`${redirectorUrl}/set-cookie`), { cookieJar }); + + expect(await response.json()).toMatchObject({ cookie: 'session=new' }); + }); + + test('does not send cookies scoped to the path of the previous request on a same-origin redirect', async () => { + const cookieJar = new CookieJar(); + await cookieJar.setCookie('session=secret; Path=/same-origin', redirectorUrl); + + const response = await httpClient.sendRequest(new Request(`${redirectorUrl}/same-origin`), { cookieJar }); + + expect(await response.json()).toMatchObject({ cookie: null }); + }); }); describe('BaseHttpClient TLS error handling', () => { From d9c57223413b38394ff5a4427091d70db9fb7636 Mon Sep 17 00:00:00 2001 From: Jan Buchar Date: Thu, 24 Sep 2026 13:36:45 +0200 Subject: [PATCH 2/5] feat: Fully-fledged support for custom (de)serializers in RecoverableState (#4131) --- docs/public-api/crawlee-core.api.md | 8 ++- packages/core/src/recoverable_state.ts | 73 +++++++++++++++++++++++--- test/core/recoverable_state.test.ts | 47 +++++++++++++++++ 3 files changed, 120 insertions(+), 8 deletions(-) diff --git a/docs/public-api/crawlee-core.api.md b/docs/public-api/crawlee-core.api.md index 1e7295c883dc..e44dfa09d2ef 100644 --- a/docs/public-api/crawlee-core.api.md +++ b/docs/public-api/crawlee-core.api.md @@ -680,14 +680,20 @@ export class RecoverableState, TPersistedS } // @public -export interface RecoverableStateOptions, TPersistedState = TStateModel> extends RecoverableStatePersistenceOptions { +export interface RecoverableStateBaseOptions, TPersistedState = TStateModel> extends RecoverableStatePersistenceOptions { configuration?: Configuration; + contentType?: string; defaultState: TStateModel | (() => TStateModel); deserialize?: StateConversion; logger?: CrawleeLogger; serialize?: StateConversion; } +// @public +export type RecoverableStateOptions, TPersistedState = TStateModel> = RecoverableStateBaseOptions & ({ + contentType?: undefined; +} | Required, 'serialize' | 'deserialize' | 'contentType'>>); + // @public (undocumented) export interface RecoverableStatePersistenceOptions { keyValueStore?: KeyValueStore | PromiseLike; diff --git a/packages/core/src/recoverable_state.ts b/packages/core/src/recoverable_state.ts index b8f2f5b19bd2..f177f30fae47 100644 --- a/packages/core/src/recoverable_state.ts +++ b/packages/core/src/recoverable_state.ts @@ -1,3 +1,5 @@ +import { Readable } from 'node:stream'; + import { addTimeoutToPromise, storage as timeoutStorage } from '@apify/timeout'; import type { Configuration } from './configuration.js'; import { StateValidationError } from './errors.js'; @@ -84,9 +86,9 @@ export interface RecoverableStatePersistenceOptions { } /** - * Options for configuring the RecoverableState + * The fields of {@apilink RecoverableStateOptions}, without the constraint tying `contentType` to the conversions. */ -export interface RecoverableStateOptions< +export interface RecoverableStateBaseOptions< TStateModel = Record, TPersistedState = TStateModel, > extends RecoverableStatePersistenceOptions { @@ -111,16 +113,48 @@ export interface RecoverableStateOptions< /** * Optional conversion of the state to a plain JSON-serializable value before it is persisted. * If not provided, the state is persisted as is. + * + * With {@apilink RecoverableStateBaseOptions.contentType} set, it has to produce what + * {@apilink KeyValueStore.setValue} accepts alongside an explicit content type - a `string`, a `Buffer` or a + * stream. */ serialize?: StateConversion; /** * Optional conversion of a persisted value back to the state model, and the place to validate a record before * trusting it. If not provided, the persisted value is used as is. + * + * With {@apilink RecoverableStateBaseOptions.contentType} set, it receives a `Readable` of the record bytes + * instead of a parsed value. */ deserialize?: StateConversion; + + /** + * Content type of the persisted record. Setting it hands the record encoding over to + * {@apilink RecoverableStateBaseOptions.serialize} and {@apilink RecoverableStateBaseOptions.deserialize}, both of + * which are then required - the default JSON codec is bypassed in both directions. Meant for a state too large + * for `JSON.stringify`, which `serialize` can then stream out instead. + */ + contentType?: string; } +/** + * Options for configuring the RecoverableState + */ +export type RecoverableStateOptions< + TStateModel = Record, + TPersistedState = TStateModel, +> = RecoverableStateBaseOptions & + ( + | { contentType?: undefined } + | Required< + Pick< + RecoverableStateBaseOptions, + 'serialize' | 'deserialize' | 'contentType' + > + > + ); + /** * A class for managing persistent recoverable state using a plain JavaScript object. * @@ -145,6 +179,7 @@ export class RecoverableState, TPersistedS readonly #log: CrawleeLogger; readonly #serialize: (state: TStateModel) => Promise; readonly #deserialize: (persistedState: TPersistedState) => Promise; + readonly #contentType?: string; readonly #persistStateQuietly: (eventData?: Record) => Promise; /** @@ -165,6 +200,14 @@ export class RecoverableState, TPersistedS this.#configuration = options.configuration; this.#keyValueStore = options.keyValueStore ?? null; this.#log = options.logger ?? serviceLocator.getLogger().child({ prefix: 'RecoverableState' }); + this.#contentType = options.contentType; + + if (this.#contentType !== undefined && (options.serialize === undefined || options.deserialize === undefined)) { + throw new Error( + `A 'contentType' for the state persisted under key '${this.#persistStateKey}' requires both 'serialize' and 'deserialize' - the record is no longer JSON the default codec can handle.`, + ); + } + this.#serialize = this.#toConversion(options.serialize); this.#deserialize = this.#toConversion(options.deserialize); @@ -340,7 +383,12 @@ export class RecoverableState, TPersistedS const serializedState = await this.#serialize(this.currentValue); await this.#withTimeout( - async () => keyValueStore.setValue(this.#persistStateKey, serializedState), + async () => + this.#contentType === undefined + ? keyValueStore.setValue(this.#persistStateKey, serializedState) + : keyValueStore.setValue(this.#persistStateKey, serializedState, { + contentType: this.#contentType, + }), 'Persisting the state', ); } @@ -379,10 +427,21 @@ export class RecoverableState, TPersistedS return; } - const storedState = await this.#withTimeout( - async () => keyValueStore.getValue(this.#persistStateKey), - 'Loading the persisted state', - ); + // With a content type, the record is whatever `serialize` produced, so the codec's parse is skipped and + // `deserialize` gets the bytes as a stream - the shape a streaming parser wants. + // TODO: read the record as a stream instead of buffering it and wrapping (https://github.com/apify/crawlee/issues/2929). + const storedState = await this.#withTimeout(async () => { + if (this.#contentType === undefined) { + return keyValueStore.getValue(this.#persistStateKey); + } + + const record = await keyValueStore.getRecord(this.#persistStateKey); + if (record === null) { + return null; + } + + return Readable.from(Buffer.isBuffer(record.value) ? record.value : Buffer.from(record.value)); + }, 'Loading the persisted state'); if (storedState === null || storedState === undefined) { return; diff --git a/test/core/recoverable_state.test.ts b/test/core/recoverable_state.test.ts index 4d165ab75c1b..dae1938118cc 100644 --- a/test/core/recoverable_state.test.ts +++ b/test/core/recoverable_state.test.ts @@ -1,3 +1,5 @@ +import { Readable } from 'node:stream'; +import { text } from 'node:stream/consumers'; import { setTimeout as sleep } from 'node:timers/promises'; import { beforeEach, describe, expect, test, vi } from 'vitest'; @@ -344,6 +346,51 @@ describe('RecoverableState', () => { expect(restored.currentValue.data.value).toBe('updated'); }); + test('should hand the record encoding over to the conversions when a contentType is set', async () => { + // Streams both ways: the point of the option is a state too large for a single string. + const serialize = vi.fn((state: TestState) => Readable.from([Buffer.from(JSON.stringify(state))])); + const deserialize = vi.fn(async (bytes: Readable): Promise => JSON.parse(await text(bytes))); + + const build = () => + new RecoverableState({ + defaultState, + persistStateKey: 'test-key', + persistenceEnabled: true, + contentType: 'application/json; charset=utf-8', + serialize, + deserialize, + }); + + const recoverableState = build(); + await recoverableState.initialize(); + recoverableState.currentValue.counter = 42; + await recoverableState.persistState(); + + const record = await (await KeyValueStore.open()).getRecord('test-key'); + expect(record?.contentType).toBe('application/json; charset=utf-8'); + + const restored = build(); + await restored.initialize(); + + expect(deserialize).toHaveBeenCalledWith(expect.any(Readable)); + expect(restored.currentValue).toEqual({ ...defaultState, counter: 42 }); + }); + + test('should refuse a contentType without both conversions', () => { + const options = { + defaultState, + persistStateKey: 'test-key', + persistenceEnabled: true, + contentType: 'text/plain', + }; + + // The types reject both; the runtime check is for callers that get past them. + // @ts-expect-error + expect(() => new RecoverableState({ ...options, serialize: JSON.stringify })).toThrow(/'deserialize'/); + // @ts-expect-error + expect(() => new RecoverableState({ ...options, deserialize: JSON.parse })).toThrow(/'serialize'/); + }); + test('should call a defaultState factory afresh for every reset', async () => { const defaultStateFactory = vi.fn(() => ({ items: new Map([['a', 1]]) })); From 91291b2b443ee3663783788607d000476c58f7e4 Mon Sep 17 00:00:00 2001 From: Suliman Abdulrazzaq <144490671+SulimanAbdulrazzaq@users.noreply.github.com> Date: Thu, 24 Sep 2026 15:13:14 +0300 Subject: [PATCH 3/5] fix(utils): drop query string and fragment from guessed sitemap urls (#4132) `discoverValidSitemaps` builds the well-known sitemap candidates by reassigning `pathname` on a `URL` parsed from the domain's first input URL. `search` and `hash` survive that reassignment, so a seed like `https://example.com/products?page=2#top` is probed as: ``` https://example.com/sitemap.xml?page=2 https://example.com/sitemap.txt?page=2 https://example.com/sitemap_index.xml?page=2 ``` A server that keys on the full path-and-query answers those with a 404, and discovery silently finds nothing for a domain that does have a sitemap. A server that ignores the query answers 200, and the query-carrying URL is then what gets yielded and loaded as the sitemap. Clearing `search` and `hash` before setting the pathname fixes both. Userinfo and port are left alone, so sites behind basic auth or on a non-default port keep working as before. --- packages/utils/src/internals/sitemap.ts | 5 +++++ packages/utils/test/sitemap.test.ts | 19 +++++++++++++++++++ 2 files changed, 24 insertions(+) diff --git a/packages/utils/src/internals/sitemap.ts b/packages/utils/src/internals/sitemap.ts index b63349470ade..6117c657db06 100644 --- a/packages/utils/src/internals/sitemap.ts +++ b/packages/utils/src/internals/sitemap.ts @@ -653,6 +653,11 @@ export async function* discoverValidSitemaps( } } else { const firstUrl = new URL(domainUrls[0]); + // The guessed candidates are well-known paths on the same origin, so the query string + // and fragment of the input URL are not part of them - carrying them over probes (and + // yields) URLs like `/sitemap.xml?page=2`, which servers may well answer with a 404. + firstUrl.search = ''; + firstUrl.hash = ''; const possibleSitemapPathnames = ['/sitemap.xml', '/sitemap.txt', '/sitemap_index.xml']; const candidateSitemapUrls = possibleSitemapPathnames.map((pathname) => { firstUrl.pathname = pathname; diff --git a/packages/utils/test/sitemap.test.ts b/packages/utils/test/sitemap.test.ts index fbdd006e5a43..70cef14f2187 100644 --- a/packages/utils/test/sitemap.test.ts +++ b/packages/utils/test/sitemap.test.ts @@ -767,6 +767,25 @@ describe('discoverValidSitemaps', () => { expect(urls).toEqual(['http://sitemap-discovery.com/sitemap.xml']); }); + it('probes well-known paths without the query string and fragment of the input url', async () => { + nock('http://sitemap-discovery.com') + .get('/robots.txt') + .reply(404) + .head('/sitemap.xml') + .reply(200, '') + .head('/sitemap.txt') + .reply(404, '') + .head('/sitemap_index.xml') + .reply(404, ''); + + const urls = []; + for await (const url of discoverValidSitemaps(['http://sitemap-discovery.com/products?page=2#top'])) { + urls.push(url); + } + + expect(urls).toEqual(['http://sitemap-discovery.com/sitemap.xml']); + }); + it('extracts sitemaps from multiple domains with mixed order', async () => { nock('http://domain-a.com') .get('/robots.txt') From b4e6f2158b6f0d7f581fe3afc6eecedec64bcf0e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jind=C5=99ich=20B=C3=A4r?= Date: Thu, 24 Sep 2026 14:32:23 +0200 Subject: [PATCH 4/5] fix(http-crawler): forward the context proxy URL to the HTTP client (#4152) The `proxyUrl` from `crawlingContext.proxyInfo` got lost in #3295 and never reached `sendRequest`, so a proxy swapped in a pre-navigation hook was silently ignored (related to #4140). --- .../src/internals/http-crawler.ts | 1 + test/core/crawlers/http_crawler.test.ts | 30 +++++++++++++++++++ 2 files changed, 31 insertions(+) diff --git a/packages/http-crawler/src/internals/http-crawler.ts b/packages/http-crawler/src/internals/http-crawler.ts index 73578a3eca75..1b4f09832a53 100644 --- a/packages/http-crawler/src/internals/http-crawler.ts +++ b/packages/http-crawler/src/internals/http-crawler.ts @@ -920,6 +920,7 @@ export class HttpCrawler< { session, cookieJar, + proxyUrl: opts.proxyUrl, signal: cancelSignal, timeoutMillis: cancelSignal ? undefined : opts.timeout, ignoreTlsErrors: this.#ignoreTlsErrors, diff --git a/test/core/crawlers/http_crawler.test.ts b/test/core/crawlers/http_crawler.test.ts index cd8aaedef2a4..909285a0e0de 100644 --- a/test/core/crawlers/http_crawler.test.ts +++ b/test/core/crawlers/http_crawler.test.ts @@ -14,6 +14,7 @@ import { ThrottlingRequestManager, } from '@crawlee/http'; import { BaseHttpClient, ResponseWithUrl } from '@crawlee/http-client'; +import type { ProxyInfo, SendRequestOptions } from '@crawlee/types'; import { sleep } from '@crawlee/utils'; import iconv from 'iconv-lite'; @@ -643,6 +644,35 @@ test('works with a custom HttpClient', async () => { expect(results[1].includes('Schmexample Domain')).toBeTruthy(); }); +test('forwards the context proxyInfo url to the HttpClient', async () => { + const proxyUrls: (string | undefined)[] = []; + + const crawler = new HttpCrawler({ + maxRequestRetries: 0, + preNavigationHooks: [ + async (context) => { + // You probably don't want this in real code - let `ProxyConfiguration` assign proxies via sessions. + context.proxyInfo = { url: 'http://proxy.example.com:8000' } as ProxyInfo; + }, + ], + requestHandler: async () => {}, + httpClient: Object.assign(Object.create(BaseHttpClient.prototype) as BaseHttpClient, { + async sendRequest(request: Request, options?: SendRequestOptions) { + proxyUrls.push(options?.proxyUrl); + return new ResponseWithUrl('', { + url: request.url.toString(), + status: 200, + headers: { 'content-type': 'text/html; charset=utf-8' }, + }); + }, + }), + }); + + await crawler.run([url]); + + expect(proxyUrls).toEqual(['http://proxy.example.com:8000']); +}); + test('a 429 on a throttled domain paces the retry without spending it or the session', async () => { const hits: number[] = []; router.set('/429-then-ok', (req, res) => { From 1b31982541e4cf02e07457812e40accc261dc4f5 Mon Sep 17 00:00:00 2001 From: Atirna Date: Thu, 24 Sep 2026 20:32:55 +0530 Subject: [PATCH 5/5] feat: extend the request lock from `extendTimeout` (#4041) Co-authored-by: Jan Buchar --- docs/public-api/crawlee-basic.api.md | 1 + docs/public-api/crawlee-core.api.md | 4 + docs/public-api/crawlee-types.api.md | 1 + docs/upgrading/upgrading_v4.md | 4 +- .../src/internals/basic-crawler.ts | 5 ++ .../src/internals/crawlers/crawler_commons.ts | 5 +- .../internals/throttling_request_manager.ts | 15 ++++ .../test/throttling_request_manager.test.ts | 52 +++++++++++++ packages/core/src/storages/request_manager.ts | 11 +++ .../src/storages/request_manager_tandem.ts | 4 + packages/core/src/storages/request_queue.ts | 13 ++++ packages/fs-storage/package.json | 2 +- .../src/resource-clients/request-queue.ts | 8 ++ .../prolong-request-lock.test.ts | 30 ++++++++ packages/types/src/storages.ts | 10 +++ pnpm-lock.yaml | 74 +++++++++---------- test/core/crawlers/basic_crawler.test.ts | 38 ++++++++++ test/core/request_manager_tandem.test.ts | 16 ++++ test/core/storages/request_queue.test.ts | 35 +++++++++ 19 files changed, 286 insertions(+), 42 deletions(-) create mode 100644 packages/fs-storage/test/request-queue/prolong-request-lock.test.ts diff --git a/docs/public-api/crawlee-basic.api.md b/docs/public-api/crawlee-basic.api.md index 60fb0fa96ddb..c74f9e01ac1a 100644 --- a/docs/public-api/crawlee-basic.api.md +++ b/docs/public-api/crawlee-basic.api.md @@ -1140,6 +1140,7 @@ export class ThrottlingRequestManager; // (undocumented) drop(): Promise; + extendRequestProcessingTimeSecs(request: Request_2, secs: number): Promise; fetchNextRequest(): Promise | null>; // (undocumented) getHandledCount(): Promise; diff --git a/docs/public-api/crawlee-core.api.md b/docs/public-api/crawlee-core.api.md index e44dfa09d2ef..04f28552054b 100644 --- a/docs/public-api/crawlee-core.api.md +++ b/docs/public-api/crawlee-core.api.md @@ -427,6 +427,7 @@ export interface IRequestManager extends IRequestLoader { addRequest(requestLike: Source, options?: RequestQueueOperationOptions): Promise; // (undocumented) addRequestsBatched(requests: RequestsLike, options?: AddRequestsBatchedOptions): Promise; + extendRequestProcessingTimeSecs?(request: Request_2, secs: number): Promise; purge?(): Promise; reclaimRequest(request: Request_2, options?: RequestQueueOperationOptions): Promise; recordPacingSignal(signal: PacingSignal): boolean; @@ -787,6 +788,8 @@ export class RequestManagerTandem implements IRequestManager { // (undocumented) addRequestsBatched(requests: RequestsLike, options?: AddRequestsBatchedOptions): Promise; checkReadiness(): Promise; + // (undocumented) + extendRequestProcessingTimeSecs(request: Request_2, secs: number): Promise; fetchNextRequest(): Promise | null>; // (undocumented) getHandledCount(): Promise; @@ -833,6 +836,7 @@ export class RequestQueue implements IStorage, IRequestManager { readonly backend: RequestQueueBackend; checkReadiness(): Promise; drop(): Promise; + extendRequestProcessingTimeSecs(request: Request_2, secs: number): Promise; fetchNextRequest(): Promise | null>; getHandledCount(): Promise; getInfo(): Promise; diff --git a/docs/public-api/crawlee-types.api.md b/docs/public-api/crawlee-types.api.md index 8b1b7ad43f75..e5fd04ede5b3 100644 --- a/docs/public-api/crawlee-types.api.md +++ b/docs/public-api/crawlee-types.api.md @@ -374,6 +374,7 @@ export type RedirectHandler = (redirectResponse: Response, updatedRequest: { export interface RequestQueueBackend { addBatchOfRequests(requests: RequestSchema[], options?: RequestQueueOperationOptions): Promise; drop(): Promise; + extendRequestProcessingTimeSecs?(requestId: string, secs: number): Promise; fetchNextRequest(): Promise; getMetadata(): Promise; getRequest(uniqueKey: string): Promise; diff --git a/docs/upgrading/upgrading_v4.md b/docs/upgrading/upgrading_v4.md index dfd08c622f50..be633dc3675d 100644 --- a/docs/upgrading/upgrading_v4.md +++ b/docs/upgrading/upgrading_v4.md @@ -1717,7 +1717,7 @@ The sub-backend interfaces (`DatasetBackend`, `KeyValueStoreBackend`, `RequestQu **`RequestQueueBackend`:** -The request queue backend was reduced from 12 methods to 10. The distributed-locking protocol (`listAndLockHead` → `prolongRequestLock` → `deleteRequestLock`) and the queue-head/consistency bookkeeping that used to live in the `RequestQueue` frontend have been removed from the interface; coordinating multiple clients accessing the same queue (e.g. request locking on the Apify platform) is now an internal concern of the backend implementation. +The request queue backend's surface was reshaped. The frontend-owned distributed-locking protocol (`listAndLockHead` → `prolongRequestLock` → `deleteRequestLock`) was removed. Queue-head and consistency bookkeeping are now internal concerns of the backend implementation. | Before (v3) | After (v4) | |---|---| @@ -1727,7 +1727,7 @@ The request queue backend was reduced from 12 methods to 10. The distributed-loc | `updateRequest(request, opts?)` | `markRequestAsHandled(request)` / `reclaimRequest(request, opts?)` | | `listHead(opts?)` | `fetchNextRequest()` (returns a single request, marks it in progress) | | `listAndLockHead(opts)` | Removed (locking is internal to the client) | -| `prolongRequestLock(id, opts)` | Removed | +| `prolongRequestLock(id, opts)` | `extendRequestProcessingTimeSecs(requestId, secs)` (optional — per-request lock extension for locking backends; wired to `context.extendTimeout`) | | `deleteRequestLock(id, opts?)` | Removed | | `deleteRequest(id)` | Removed | | _(n/a)_ | `isEmpty()` (new — `true` when no pending requests are left to fetch) | diff --git a/packages/basic-crawler/src/internals/basic-crawler.ts b/packages/basic-crawler/src/internals/basic-crawler.ts index 8d515afbca6d..3960b3b500c9 100644 --- a/packages/basic-crawler/src/internals/basic-crawler.ts +++ b/packages/basic-crawler/src/internals/basic-crawler.ts @@ -1499,6 +1499,11 @@ export class BasicCrawler< if (context[navigationDeadlineKey] !== undefined) { context[navigationDeadlineKey] += extraMillis; } + + // Extension failure must not fail the request that asked for more time. + this.requestManager?.extendRequestProcessingTimeSecs?.(context.request, secs)?.catch((error) => { + this.log.debug('Extending the request processing time failed', { url: context.request.url, error }); + }); }, }; } diff --git a/packages/basic-crawler/src/internals/crawlers/crawler_commons.ts b/packages/basic-crawler/src/internals/crawlers/crawler_commons.ts index 8843397b2fae..024436fba0ec 100644 --- a/packages/basic-crawler/src/internals/crawlers/crawler_commons.ts +++ b/packages/basic-crawler/src/internals/crawlers/crawler_commons.ts @@ -251,8 +251,9 @@ export interface CrawlingContext exten * ``` * * Extends the request handler's own timeout and the crawler's internal one together, so the extension - * is not immediately undone by the latter. Calling it from a handler that has already timed out does - * nothing. + * is not immediately undone by the latter. On a locking storage backend the request's lock is + * extended by the same amount, so the extra time is not spent while the queue considers the request + * free to hand out again. Calling it from a handler that has already timed out does nothing. */ extendTimeout(secs: number): void; } diff --git a/packages/basic-crawler/src/internals/throttling_request_manager.ts b/packages/basic-crawler/src/internals/throttling_request_manager.ts index a88f54bb3e61..17fa5d8eb8c4 100644 --- a/packages/basic-crawler/src/internals/throttling_request_manager.ts +++ b/packages/basic-crawler/src/internals/throttling_request_manager.ts @@ -1042,6 +1042,21 @@ export class ThrottlingRequestManager manager.setExpectedRequestProcessingTimeSecs?.(secs)); } + /** + * @inheritdoc + * Extending a lock must not consume the in-flight marker that later routes + * {@link ThrottlingRequestManager.markRequestAsHandled} / {@link ThrottlingRequestManager.reclaimRequest}. + */ + async extendRequestProcessingTimeSecs(request: Request, secs: number): Promise { + const key = request.id ?? request.uniqueKey; + + const manager = this.#inFlightFromInner.has(key) + ? this.#resolvedInner! + : await this.#selectManagerOrThrow(request.url); + + return (await manager.extendRequestProcessingTimeSecs?.(request, secs)) ?? false; + } + /** * Runs `fn` over the sub-queues and, if it has been resolved, the wrapped manager - bookkeeping never forces * a lazily-opened `inner`, since there is no point opening a queue purely to tell it something. diff --git a/packages/basic-crawler/test/throttling_request_manager.test.ts b/packages/basic-crawler/test/throttling_request_manager.test.ts index 09bed0c79f84..92494a6e7cce 100644 --- a/packages/basic-crawler/test/throttling_request_manager.test.ts +++ b/packages/basic-crawler/test/throttling_request_manager.test.ts @@ -1,6 +1,11 @@ +import { rm } from 'node:fs/promises'; +import { resolve } from 'node:path'; + import type { ThrottlingRequestManagerOptions } from '@crawlee/basic'; import { ThrottlingRequestManager } from '@crawlee/basic'; import type { AddRequestsBatchedResult, StorageIdentifier } from '@crawlee/core'; + +import { vi } from 'vitest'; import { KeyValueStore, MemoryStorageBackend, @@ -9,6 +14,7 @@ import { withStorageTransaction, } from '@crawlee/core'; import { sleep } from '@crawlee/utils'; +import { FileSystemStorageBackend } from '@crawlee/fs-storage'; describe('ThrottlingRequestManager', () => { beforeEach(() => { @@ -499,6 +505,52 @@ describe('ThrottlingRequestManager', () => { }); }); + test('extendRequestProcessingTimeSecs reaches the manager holding the request without disturbing its routing', async () => { + // The file-system backend actually locks fetched requests, so the spies observe + // the real forwarding without replacing it. + const tmpLocation = resolve(import.meta.dirname, './tmp/extend-routing'); + const storageBackend = new FileSystemStorageBackend({ localDataDirectory: tmpLocation }); + try { + await withLockingRoutingTest(); + } finally { + await rm(tmpLocation, { force: true, recursive: true }); + } + + async function withLockingRoutingTest() { + const inner = await RequestQueue.open({ name: 'inner-queue' }, { storageBackend }); + const manager = new ThrottlingRequestManager({ + inner, + domains: ['example.com'], + requestManagerOpener: async (identifier, options) => + RequestQueue.open(identifier, { ...options, storageBackend }), + }); + + await manager.addRequest({ url: 'https://example.com/routed' }); + await inner.addRequest({ url: 'https://example.com/inner' }); + + const routed = (await manager.fetchNextRequest())!; + expect(routed.url).toBe('https://example.com/routed'); + const fromInner = (await manager.fetchNextRequest())!; + expect(fromInner.url).toBe('https://example.com/inner'); + + // Re-opening the deterministic alias yields the same cached frontend the manager routes to. + const subQueue = await RequestQueue.open({ alias: 'throttled-example.com' }, { storageBackend }); + const innerSpy = vi.spyOn(inner.backend, 'extendRequestProcessingTimeSecs'); + const subSpy = vi.spyOn(subQueue.backend, 'extendRequestProcessingTimeSecs'); + + await expect(manager.extendRequestProcessingTimeSecs(routed, 30)).resolves.toBe(true); + expect(subSpy).toHaveBeenCalledExactlyOnceWith(routed.id, 30); + expect(innerSpy).not.toHaveBeenCalled(); + + await expect(manager.extendRequestProcessingTimeSecs(fromInner, 30)).resolves.toBe(true); + expect(innerSpy).toHaveBeenCalledExactlyOnceWith(fromInner.id, 30); + expect(subSpy).toHaveBeenCalledTimes(1); + + await manager.markRequestAsHandled(fromInner); + expect(await inner.getPendingCount()).toBe(0); + } + }); + test('recordPacingSignal enforces throttling and fair scheduling', async () => { const inner = await createQueue(); const manager = new ThrottlingRequestManager({ diff --git a/packages/core/src/storages/request_manager.ts b/packages/core/src/storages/request_manager.ts index cfa20dcfd77a..6c3c80f7253a 100644 --- a/packages/core/src/storages/request_manager.ts +++ b/packages/core/src/storages/request_manager.ts @@ -54,6 +54,17 @@ export interface IRequestManager extends IRequestLoader { * @returns `true` if anything in the composition took responsibility for the signal. */ recordPacingSignal(signal: PacingSignal): boolean; + + /** + * Extends the lock on a request previously handed out by `fetchNextRequest()` and still being + * processed, on storage backends that reserve requests via locking (e.g. via + * {@apilink CrawlingContext.extendTimeout|`context.extendTimeout`}). + * + * @returns `true` when the lock was prolonged, `false` when this manager does not lock + * requests or no longer holds this one; non-locking implementations may leave it `undefined`, + * which callers treat as `false`. + */ + extendRequestProcessingTimeSecs?(request: Request, secs: number): Promise; } /** diff --git a/packages/core/src/storages/request_manager_tandem.ts b/packages/core/src/storages/request_manager_tandem.ts index 1a5c4c607cb5..c101a5ddb879 100644 --- a/packages/core/src/storages/request_manager_tandem.ts +++ b/packages/core/src/storages/request_manager_tandem.ts @@ -245,4 +245,8 @@ export class RequestManagerTandem implements IRequestManager { recordPacingSignal(signal: PacingSignal): boolean { return this.#resolvedRequestManager?.recordPacingSignal(signal) ?? false; } + + async extendRequestProcessingTimeSecs(request: Request, secs: number): Promise { + return (await this.getRequestManager()).extendRequestProcessingTimeSecs?.(request, secs) ?? false; + } } diff --git a/packages/core/src/storages/request_queue.ts b/packages/core/src/storages/request_queue.ts index 2dbcb5ea9e2e..f213fd0720cb 100644 --- a/packages/core/src/storages/request_queue.ts +++ b/packages/core/src/storages/request_queue.ts @@ -860,6 +860,19 @@ export class RequestQueue implements IStorage, IRequestManager { await this.backend.setExpectedRequestProcessingTimeSecs?.(secs); } + /** + * @inheritdoc + * Unlike {@link RequestQueue.setExpectedRequestProcessingTimeSecs}, which sizes every future lock, + * this only touches the one request it is given. + */ + async extendRequestProcessingTimeSecs(request: Request, secs: number): Promise { + if (!request.id) { + return false; + } + + return (await this.backend.extendRequestProcessingTimeSecs?.(request.id, secs)) ?? false; + } + /** * Caches information about request to beware of unneeded addRequest() calls. */ diff --git a/packages/fs-storage/package.json b/packages/fs-storage/package.json index f15178fba773..17a5ca27c81b 100644 --- a/packages/fs-storage/package.json +++ b/packages/fs-storage/package.json @@ -42,7 +42,7 @@ "access": "public" }, "dependencies": { - "@crawlee/fs-storage-native": ">=0.2.1 <0.3", + "@crawlee/fs-storage-native": ">=0.2.2 <0.3", "@crawlee/types": "workspace:*", "@crawlee/utils": "workspace:*", "zod": "catalog:" diff --git a/packages/fs-storage/src/resource-clients/request-queue.ts b/packages/fs-storage/src/resource-clients/request-queue.ts index afcd97ca8737..b72ec36b0d82 100644 --- a/packages/fs-storage/src/resource-clients/request-queue.ts +++ b/packages/fs-storage/src/resource-clients/request-queue.ts @@ -75,6 +75,14 @@ export class RequestQueueBackend extends CachedIdClient implements storage.Reque await this.#nativeBackend.setExpectedRequestProcessingTime(secs); } + /** + * @inheritdoc + * The lock is owned and extended by the native client; `false` means it no longer holds the request. + */ + async extendRequestProcessingTimeSecs(requestId: string, secs: number): Promise { + return this.#nativeBackend.prolongRequestLock(requestId, secs); + } + async getMetadata(): Promise { return this.#nativeBackend.getMetadata(); } diff --git a/packages/fs-storage/test/request-queue/prolong-request-lock.test.ts b/packages/fs-storage/test/request-queue/prolong-request-lock.test.ts new file mode 100644 index 000000000000..44ef0717c678 --- /dev/null +++ b/packages/fs-storage/test/request-queue/prolong-request-lock.test.ts @@ -0,0 +1,30 @@ +import { rm } from 'node:fs/promises'; +import { resolve } from 'node:path'; + +import { FileSystemStorageBackend } from '@crawlee/fs-storage'; + +describe('FileSystemStorageBackend extendRequestProcessingTimeSecs', () => { + const tmpLocation = resolve(import.meta.dirname, './tmp/prolong-request-lock'); + + afterEach(async () => { + await rm(tmpLocation, { force: true, recursive: true }); + }); + + test('extends the lock on a fetched request, false for one never locked', async () => { + const storage = new FileSystemStorageBackend({ + localDataDirectory: tmpLocation, + }); + const queue = await storage.createRequestQueueBackend({ + name: 'default', + }); + await queue.addBatchOfRequests([{ url: 'http://example.com/1', uniqueKey: '1' }]); + + const locked = await queue.fetchNextRequest(); + expect(locked).toBeDefined(); + + expect(await queue.extendRequestProcessingTimeSecs(locked!.id!, 30)).toBe(true); + expect(await queue.extendRequestProcessingTimeSecs('no-such-request', 30)).toBe(false); + + await queue.drop(); + }); +}); diff --git a/packages/types/src/storages.ts b/packages/types/src/storages.ts index d57ff19efb48..49c337fa59f5 100644 --- a/packages/types/src/storages.ts +++ b/packages/types/src/storages.ts @@ -393,6 +393,16 @@ export interface RequestQueueBackend { * still being processed. Clients that do not lock may ignore it. */ setExpectedRequestProcessingTimeSecs?(secs: number): Promise; + + /** + * Extends the lock on a request previously handed out by {@link fetchNextRequest}, for a consumer + * that needs more time than the sizing hint reserved. + * + * @returns `true` when prolonged, `false` when this client does not hold the request locked (or + * does not lock at all). `false` is information, not an error: the request may be processed twice. + * Non-locking backends leave this unimplemented. + */ + extendRequestProcessingTimeSecs?(requestId: string, secs: number): Promise; } /** diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 1bf31d3b68e2..d5e7d80903e2 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -835,8 +835,8 @@ importers: packages/fs-storage: dependencies: '@crawlee/fs-storage-native': - specifier: '>=0.2.1 <0.3' - version: 0.2.1 + specifier: '>=0.2.2 <0.3' + version: 0.2.2 '@crawlee/types': specifier: workspace:* version: link:../types @@ -2233,60 +2233,60 @@ packages: resolution: {integrity: sha512-TzlTVpKPjaqW6qOYjQcYUDuGsLCNsvFHVBXkYGTAnf5V37jCWrE5haKNXzz0WZUtVHjrpV76L1buANjwXMfT8w==} engines: {node: '>=22'} - '@crawlee/fs-storage-native-darwin-arm64@0.2.1': - resolution: {integrity: sha512-vpjIyzhv2XB8oz5UbcXVMm6RwyfJgAE7xEiXG3pgxYjjtWOP2xGZX17dJiCw9fcOjH3bawfNZTZiKSDvxF0Q1A==} + '@crawlee/fs-storage-native-darwin-arm64@0.2.2': + resolution: {integrity: sha512-8R9fm6iyUVSxPdRA/m7XTgDawm6bjZIaDzmLmtQPwEIDTE97z6mk2MOyrlT53rr7B5aObN/nYJeCDENIGfcH1g==} engines: {node: '>= 20'} cpu: [arm64] os: [darwin] - '@crawlee/fs-storage-native-darwin-x64@0.2.1': - resolution: {integrity: sha512-os5MWPkUoNKM+JSLN5MxeVoRxlwqVlef2gMqXc+a359kSS3iqGjP+tUZQ8MTIU5s370igxIx6dnJl+xr6FplJA==} + '@crawlee/fs-storage-native-darwin-x64@0.2.2': + resolution: {integrity: sha512-ETYUt8T9cNEpq4qfaDdIGlWrdXiDgYL7BzGkNRTYb+NCc3mN3SvY6kEKLSrezn8S6Hwmgud5tfPPe4HjlcxRoQ==} engines: {node: '>= 20'} cpu: [x64] os: [darwin] - '@crawlee/fs-storage-native-linux-arm64-gnu@0.2.1': - resolution: {integrity: sha512-7rO5NC6HKjl+rh8zK/bayKBP3MMQs/Ib5eyzZJi4U8Jb2Vke1pU9yta+XjKzVU2rQZB8ErntISTxczbUfZfhFQ==} + '@crawlee/fs-storage-native-linux-arm64-gnu@0.2.2': + resolution: {integrity: sha512-rRyd22FOwlvuh+YVwwazlvxWRX2eUh8v2Wpp2So5oZjCEbfoY4Wp4R2Lbd8kjsMYWPBC9MIwv1RkiICpxVfqCA==} engines: {node: '>= 20'} cpu: [arm64] os: [linux] libc: [glibc] - '@crawlee/fs-storage-native-linux-arm64-musl@0.2.1': - resolution: {integrity: sha512-+oNToqNyXkcaT7oi8WBcKaQxJNZD4oXv1QWyUq+z69h5H6GfzcQe/nqslXALY5SNx9yTCzTtVS5AgrgqvGzyTg==} + '@crawlee/fs-storage-native-linux-arm64-musl@0.2.2': + resolution: {integrity: sha512-B0436STNN9aLMdlEeqEsfysarhcggQ0dS6p/9iffLUnqFC2OWGcXxHOVl+2pH597tPVyv0gsVfAC8BxJ5a1c4Q==} engines: {node: '>= 20'} cpu: [arm64] os: [linux] libc: [musl] - '@crawlee/fs-storage-native-linux-x64-gnu@0.2.1': - resolution: {integrity: sha512-qHQTBNFlBHajKfcwvzN8u68mLBWsEbaZQb7yNBmocUcBSYXm0U6JAFn1BcpNJyFb+CfpERR3QaCIUxBVdbRBDw==} + '@crawlee/fs-storage-native-linux-x64-gnu@0.2.2': + resolution: {integrity: sha512-OXB+0qlAzcZCbdncZ1koITkHotrzYCwm8r7p4jhcv9MUbhtx7h7EUJ/90aMARRvX0T00HxWyCzvGTiJQxNUidA==} engines: {node: '>= 20'} cpu: [x64] os: [linux] libc: [glibc] - '@crawlee/fs-storage-native-linux-x64-musl@0.2.1': - resolution: {integrity: sha512-IjrfVP7mnyQqAgoRKTLxWmjYRyktVI/cBfXg+RIEtjss9a9OIt4jrHhjLhUo8cQr5IT6wyEXRto0AOxZth96SA==} + '@crawlee/fs-storage-native-linux-x64-musl@0.2.2': + resolution: {integrity: sha512-uF92wHy7IExtblv3Jsi55z+05qtqnigucXtnsXmol0i1JFjwLGxiDCJC9cskZVph8bHVGyLQFitIlI4w5U1sUQ==} engines: {node: '>= 20'} cpu: [x64] os: [linux] libc: [musl] - '@crawlee/fs-storage-native-win32-arm64-msvc@0.2.1': - resolution: {integrity: sha512-GsVD+XbiL5tTQlmIPtiThyGFp7B3vg1l8ERXyCHR6Gsz+85ZeMQ9jg3/ILFBKFY15I8v5FoDV0cQtdH7D5+FKg==} + '@crawlee/fs-storage-native-win32-arm64-msvc@0.2.2': + resolution: {integrity: sha512-SKx5GIZJ1jPfnOmJCsHHwMwrBM+ZQUB+c1xcw+YijcD4un06MxuKRYS79nFSM2SvLzvMROBhsN+3LoIXOZjNow==} engines: {node: '>= 20'} cpu: [arm64] os: [win32] - '@crawlee/fs-storage-native-win32-x64-msvc@0.2.1': - resolution: {integrity: sha512-H8uigGqYChXySUDpnDhwl9JrDbp/rZMFpjuoSDSssF1gr+83d7qRjISYKVC0TUlKiYAnT7SqYCBqNrsuSbJBUw==} + '@crawlee/fs-storage-native-win32-x64-msvc@0.2.2': + resolution: {integrity: sha512-d7oDXVLAnmpmDhQ9mkN/iTJTg+XJ60mwLyANxzPtyGqiapxWPT0Y2Yw7N+dUloSUFyuWDVl1R4A45eORuwVCYg==} engines: {node: '>= 20'} cpu: [x64] os: [win32] - '@crawlee/fs-storage-native@0.2.1': - resolution: {integrity: sha512-x6umw+kVjuA6k+Qhx9RwJSjl/DRyeo1jlCiQH04ZKz5DP54NhQjMlZSQD1Y0v4SYuN0u+wqDVINQQQsarsFkeg==} + '@crawlee/fs-storage-native@0.2.2': + resolution: {integrity: sha512-WGKNyPoRq1LKqndeJ9LdEWigGNm9Z+3U8jJaSlEVgPq8CkmcqHANXxpZXkgGXl9jfEiWzFlJmWm4JUMqT/TxfA==} engines: {node: '>= 20'} '@crawlee/types@3.16.0': @@ -14590,40 +14590,40 @@ snapshots: '@conventional-changelog/template@1.2.1': {} - '@crawlee/fs-storage-native-darwin-arm64@0.2.1': + '@crawlee/fs-storage-native-darwin-arm64@0.2.2': optional: true - '@crawlee/fs-storage-native-darwin-x64@0.2.1': + '@crawlee/fs-storage-native-darwin-x64@0.2.2': optional: true - '@crawlee/fs-storage-native-linux-arm64-gnu@0.2.1': + '@crawlee/fs-storage-native-linux-arm64-gnu@0.2.2': optional: true - '@crawlee/fs-storage-native-linux-arm64-musl@0.2.1': + '@crawlee/fs-storage-native-linux-arm64-musl@0.2.2': optional: true - '@crawlee/fs-storage-native-linux-x64-gnu@0.2.1': + '@crawlee/fs-storage-native-linux-x64-gnu@0.2.2': optional: true - '@crawlee/fs-storage-native-linux-x64-musl@0.2.1': + '@crawlee/fs-storage-native-linux-x64-musl@0.2.2': optional: true - '@crawlee/fs-storage-native-win32-arm64-msvc@0.2.1': + '@crawlee/fs-storage-native-win32-arm64-msvc@0.2.2': optional: true - '@crawlee/fs-storage-native-win32-x64-msvc@0.2.1': + '@crawlee/fs-storage-native-win32-x64-msvc@0.2.2': optional: true - '@crawlee/fs-storage-native@0.2.1': + '@crawlee/fs-storage-native@0.2.2': optionalDependencies: - '@crawlee/fs-storage-native-darwin-arm64': 0.2.1 - '@crawlee/fs-storage-native-darwin-x64': 0.2.1 - '@crawlee/fs-storage-native-linux-arm64-gnu': 0.2.1 - '@crawlee/fs-storage-native-linux-arm64-musl': 0.2.1 - '@crawlee/fs-storage-native-linux-x64-gnu': 0.2.1 - '@crawlee/fs-storage-native-linux-x64-musl': 0.2.1 - '@crawlee/fs-storage-native-win32-arm64-msvc': 0.2.1 - '@crawlee/fs-storage-native-win32-x64-msvc': 0.2.1 + '@crawlee/fs-storage-native-darwin-arm64': 0.2.2 + '@crawlee/fs-storage-native-darwin-x64': 0.2.2 + '@crawlee/fs-storage-native-linux-arm64-gnu': 0.2.2 + '@crawlee/fs-storage-native-linux-arm64-musl': 0.2.2 + '@crawlee/fs-storage-native-linux-x64-gnu': 0.2.2 + '@crawlee/fs-storage-native-linux-x64-musl': 0.2.2 + '@crawlee/fs-storage-native-win32-arm64-msvc': 0.2.2 + '@crawlee/fs-storage-native-win32-x64-msvc': 0.2.2 '@crawlee/types@3.16.0': dependencies: diff --git a/test/core/crawlers/basic_crawler.test.ts b/test/core/crawlers/basic_crawler.test.ts index 2b0d51a8ba91..35be620d3245 100644 --- a/test/core/crawlers/basic_crawler.test.ts +++ b/test/core/crawlers/basic_crawler.test.ts @@ -1,4 +1,5 @@ import { readFile, rm } from 'node:fs/promises'; +import { resolve } from 'node:path'; import type { Server } from 'node:http'; import http from 'node:http'; import type { AddressInfo } from 'node:net'; @@ -45,6 +46,7 @@ import { import type { StorageTransaction } from '@crawlee/core'; import { currentStorageTransaction, MemoryStorageBackend } from '@crawlee/core'; import { BaseHttpClient } from '@crawlee/http-client'; +import { FileSystemStorageBackend } from '@crawlee/fs-storage'; import type { Dictionary, ISession, ProxyInfo } from '@crawlee/types'; import { RobotsTxtFile, sleep } from '@crawlee/utils'; import express from 'express'; @@ -1914,6 +1916,42 @@ describe('BasicCrawler', () => { expect(failed).toHaveLength(0); }); + test('context.extendTimeout prolongs the request lock on a locking request manager', async () => { + // The file-system backend is the real locking one: its queue takes actual locks, + // so the whole crawler -> queue -> backend -> native lock handoff runs. + const tmpLocation = resolve(import.meta.dirname, './tmp/extend-timeout-lock'); + const requestQueue = await RequestQueue.open(null, { + storageBackend: new FileSystemStorageBackend({ localDataDirectory: tmpLocation }), + }); + const extendSpy = vitest.spyOn(requestQueue.backend, 'extendRequestProcessingTimeSecs'); + + try { + let sawFailure = false; + + const crawler = new BasicCrawler({ + requestManager: requestQueue, + requestHandlerTimeoutSecs: 60, + maxRequestRetries: 0, + requestHandler: async ({ request, extendTimeout }) => { + await sleep(50); + extendTimeout(42); + expect(request.id).toBeTruthy(); + }, + failedRequestHandler: async () => { + sawFailure = true; + }, + }); + + await crawler.run(['https://example.com']); + + expect(sawFailure).toBe(false); + expect(extendSpy).toHaveBeenCalledExactlyOnceWith(expect.any(String), 42); + await expect(extendSpy.mock.results[0]!.value).resolves.toBe(true); + } finally { + await rm(tmpLocation, { force: true, recursive: true }); + } + }); + test('a route can override requestHandlerTimeoutSecs, other routes keep the default', async () => { const requestList = await RequestList.open({ sources: [ diff --git a/test/core/request_manager_tandem.test.ts b/test/core/request_manager_tandem.test.ts index d414f597a0dc..547c91536be8 100644 --- a/test/core/request_manager_tandem.test.ts +++ b/test/core/request_manager_tandem.test.ts @@ -361,4 +361,20 @@ describe('RequestManagerTandem', () => { await tandem.fetchNextRequest(); expect(hintSpy).toHaveBeenCalledWith(600); }); + + test('extendRequestProcessingTimeSecs forwards to the writable request manager', async () => { + const requestList = await RequestList.open(null, []); + const requestQueue = await RequestQueue.open(); + await requestQueue.addRequest({ url: 'https://example.com/1' }); + const fetched = await requestQueue.fetchNextRequest(); + expect(fetched).not.toBeNull(); + + const stub = vi.fn(async () => true); + (requestQueue.backend as any).extendRequestProcessingTimeSecs = stub; + + const tandem = new RequestManagerTandem(requestList, requestQueue); + + await expect(tandem.extendRequestProcessingTimeSecs!(fetched!, 30)).resolves.toBe(true); + expect(stub).toHaveBeenCalledExactlyOnceWith(fetched!.id, 30); + }); }); diff --git a/test/core/storages/request_queue.test.ts b/test/core/storages/request_queue.test.ts index a75f0b07f5e4..6ad7202ae2d0 100644 --- a/test/core/storages/request_queue.test.ts +++ b/test/core/storages/request_queue.test.ts @@ -1,8 +1,11 @@ +import { rm } from 'node:fs/promises'; +import { resolve } from 'node:path'; /* eslint-disable dot-notation */ import { CrawlingRequest } from '@crawlee/basic'; import { MemoryStorageBackend, ProxyConfiguration, Request, RequestQueue, serviceLocator } from '@crawlee/core'; import { BaseHttpClient } from '@crawlee/http-client'; +import { FileSystemStorageBackend } from '@crawlee/fs-storage'; import { sleep } from '@crawlee/utils'; // `vitest.mockObject` clones the object and drops its prototype, so build the mock manually to @@ -396,6 +399,38 @@ describe('RequestQueue remote', () => { }); }); + describe('extendRequestProcessingTimeSecs', () => { + test('forwards the request id and extension to the backend', async () => { + // The file-system backend actually locks fetched requests; the spy observes + // the real forwarding and the real lock extension. + const tmpLocation = resolve(import.meta.dirname, './tmp/extend-forwarding'); + try { + const queue = await RequestQueue.open(null, { + storageBackend: new FileSystemStorageBackend({ localDataDirectory: tmpLocation }), + }); + await queue.addRequest({ url: 'http://example.com/a' }); + const fetched = await queue.fetchNextRequest(); + expect(fetched).not.toBeNull(); + + const stub = vitest.spyOn(queue.backend, 'extendRequestProcessingTimeSecs'); + + await expect(queue.extendRequestProcessingTimeSecs(fetched!, 30)).resolves.toBe(true); + expect(stub).toHaveBeenCalledExactlyOnceWith(fetched!.id, 30); + } finally { + await rm(tmpLocation, { force: true, recursive: true }); + } + }); + + test('resolves to false when the backend does not implement per-request locking', async () => { + const queue = await RequestQueue.open(); + await queue.addRequest({ url: 'http://example.com/a' }); + const fetched = await queue.fetchNextRequest(); + expect(fetched).not.toBeNull(); + + await expect(queue.extendRequestProcessingTimeSecs(fetched!, 30)).resolves.toBe(false); + }); + }); + describe('stats', () => { test('start at zero', async () => { const queue = await RequestQueue.open();