Skip to content
Merged
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
1 change: 1 addition & 0 deletions docs/public-api/crawlee-basic.api.md
Original file line number Diff line number Diff line change
Expand Up @@ -1140,6 +1140,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
checkReadiness(): Promise<RequestSourceStatus>;
// (undocumented)
drop(): Promise<void>;
extendRequestProcessingTimeSecs(request: Request_2, secs: number): Promise<boolean>;
fetchNextRequest<R extends Dictionary = Dictionary>(): Promise<Request_2<R> | null>;
// (undocumented)
getHandledCount(): Promise<number>;
Expand Down
12 changes: 11 additions & 1 deletion docs/public-api/crawlee-core.api.md
Original file line number Diff line number Diff line change
Expand Up @@ -427,6 +427,7 @@ export interface IRequestManager extends IRequestLoader {
addRequest(requestLike: Source, options?: RequestQueueOperationOptions): Promise<RequestQueueOperationInfo>;
// (undocumented)
addRequestsBatched(requests: RequestsLike, options?: AddRequestsBatchedOptions): Promise<AddRequestsBatchedResult>;
extendRequestProcessingTimeSecs?(request: Request_2, secs: number): Promise<boolean>;
purge?(): Promise<void>;
reclaimRequest(request: Request_2, options?: RequestQueueOperationOptions): Promise<RequestQueueOperationInfo | null>;
recordPacingSignal(signal: PacingSignal): boolean;
Expand Down Expand Up @@ -680,14 +681,20 @@ export class RecoverableState<TStateModel = Record<string, unknown>, TPersistedS
}

// @public
export interface RecoverableStateOptions<TStateModel = Record<string, unknown>, TPersistedState = TStateModel> extends RecoverableStatePersistenceOptions {
export interface RecoverableStateBaseOptions<TStateModel = Record<string, unknown>, TPersistedState = TStateModel> extends RecoverableStatePersistenceOptions {
configuration?: Configuration;
contentType?: string;
defaultState: TStateModel | (() => TStateModel);
deserialize?: StateConversion<TPersistedState, TStateModel>;
logger?: CrawleeLogger;
serialize?: StateConversion<TStateModel, TPersistedState>;
}

// @public
export type RecoverableStateOptions<TStateModel = Record<string, unknown>, TPersistedState = TStateModel> = RecoverableStateBaseOptions<TStateModel, TPersistedState> & ({
contentType?: undefined;
} | Required<Pick<RecoverableStateBaseOptions<TStateModel, TPersistedState>, 'serialize' | 'deserialize' | 'contentType'>>);

// @public (undocumented)
export interface RecoverableStatePersistenceOptions {
keyValueStore?: KeyValueStore | PromiseLike<KeyValueStore>;
Expand Down Expand Up @@ -781,6 +788,8 @@ export class RequestManagerTandem implements IRequestManager {
// (undocumented)
addRequestsBatched(requests: RequestsLike, options?: AddRequestsBatchedOptions): Promise<AddRequestsBatchedResult>;
checkReadiness(): Promise<RequestSourceStatus>;
// (undocumented)
extendRequestProcessingTimeSecs(request: Request_2, secs: number): Promise<boolean>;
fetchNextRequest<T extends Dictionary = Dictionary>(): Promise<Request_2<T> | null>;
// (undocumented)
getHandledCount(): Promise<number>;
Expand Down Expand Up @@ -827,6 +836,7 @@ export class RequestQueue implements IStorage, IRequestManager {
readonly backend: RequestQueueBackend;
checkReadiness(): Promise<RequestLoaderStatus>;
drop(): Promise<void>;
extendRequestProcessingTimeSecs(request: Request_2, secs: number): Promise<boolean>;
fetchNextRequest<T extends Dictionary = Dictionary>(): Promise<Request_2<T> | null>;
getHandledCount(): Promise<number>;
getInfo(): Promise<RequestQueueInfo>;
Expand Down
1 change: 1 addition & 0 deletions docs/public-api/crawlee-types.api.md
Original file line number Diff line number Diff line change
Expand Up @@ -374,6 +374,7 @@ export type RedirectHandler = (redirectResponse: Response, updatedRequest: {
export interface RequestQueueBackend {
addBatchOfRequests(requests: RequestSchema[], options?: RequestQueueOperationOptions): Promise<BatchAddRequestsResult>;
drop(): Promise<void>;
extendRequestProcessingTimeSecs?(requestId: string, secs: number): Promise<boolean>;
fetchNextRequest(): Promise<UpdateRequestSchema | undefined>;
getMetadata(): Promise<RequestQueueInfo>;
getRequest(uniqueKey: string): Promise<UpdateRequestSchema | undefined>;
Expand Down
4 changes: 2 additions & 2 deletions docs/upgrading/upgrading_v4.md
Original file line number Diff line number Diff line change
Expand Up @@ -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) |
|---|---|
Expand All @@ -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) |
Expand Down
5 changes: 5 additions & 0 deletions packages/basic-crawler/src/internals/basic-crawler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 });
});
},
};
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -251,8 +251,9 @@ export interface CrawlingContext<UserData extends Dictionary = Dictionary> 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;
}
15 changes: 15 additions & 0 deletions packages/basic-crawler/src/internals/throttling_request_manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1042,6 +1042,21 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
await this.#forEachManager((manager) => 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<boolean> {
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.
Expand Down
52 changes: 52 additions & 0 deletions packages/basic-crawler/test/throttling_request_manager.test.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -9,6 +14,7 @@ import {
withStorageTransaction,
} from '@crawlee/core';
import { sleep } from '@crawlee/utils';
import { FileSystemStorageBackend } from '@crawlee/fs-storage';

describe('ThrottlingRequestManager', () => {
beforeEach(() => {
Expand Down Expand Up @@ -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({
Expand Down
73 changes: 66 additions & 7 deletions packages/core/src/recoverable_state.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand Down Expand Up @@ -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<string, unknown>,
TPersistedState = TStateModel,
> extends RecoverableStatePersistenceOptions {
Expand All @@ -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<TStateModel, TPersistedState>;

/**
* 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<TPersistedState, TStateModel>;

/**
* 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<string, unknown>,
TPersistedState = TStateModel,
> = RecoverableStateBaseOptions<TStateModel, TPersistedState> &
(
| { contentType?: undefined }
| Required<
Pick<
RecoverableStateBaseOptions<TStateModel, TPersistedState>,
'serialize' | 'deserialize' | 'contentType'
>
>
);

/**
* A class for managing persistent recoverable state using a plain JavaScript object.
*
Expand All @@ -145,6 +179,7 @@ export class RecoverableState<TStateModel = Record<string, unknown>, TPersistedS
readonly #log: CrawleeLogger;
readonly #serialize: (state: TStateModel) => Promise<TPersistedState>;
readonly #deserialize: (persistedState: TPersistedState) => Promise<TStateModel>;
readonly #contentType?: string;
readonly #persistStateQuietly: (eventData?: Record<string, unknown>) => Promise<void>;

/**
Expand All @@ -165,6 +200,14 @@ export class RecoverableState<TStateModel = Record<string, unknown>, 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);

Expand Down Expand Up @@ -340,7 +383,12 @@ export class RecoverableState<TStateModel = Record<string, unknown>, 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',
);
}
Expand Down Expand Up @@ -379,10 +427,21 @@ export class RecoverableState<TStateModel = Record<string, unknown>, 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;
Expand Down
11 changes: 11 additions & 0 deletions packages/core/src/storages/request_manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<boolean>;
}

/**
Expand Down
4 changes: 4 additions & 0 deletions packages/core/src/storages/request_manager_tandem.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<boolean> {
return (await this.getRequestManager()).extendRequestProcessingTimeSecs?.(request, secs) ?? false;
}
}
13 changes: 13 additions & 0 deletions packages/core/src/storages/request_queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<boolean> {
if (!request.id) {
return false;
}

return (await this.backend.extendRequestProcessingTimeSecs?.(request.id, secs)) ?? false;
}

/**
* Caches information about request to beware of unneeded addRequest() calls.
*/
Expand Down
2 changes: 1 addition & 1 deletion packages/fs-storage/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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:"
Expand Down
Loading
Loading