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 @@ -221,6 +221,7 @@ export class ConcurrencySystem implements IConcurrencySystem {
get isRunning(): boolean;
get maxConcurrency(): number;
set maxConcurrency(value: number);
readonly maxTasksPerMinute: number;
get minConcurrency(): number;
set minConcurrency(value: number);
registerTaskEnd(_consumer?: ConcurrencyConsumer): void;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -216,14 +216,15 @@
* @category Scaling
*/
export class ConcurrencySystem implements IConcurrencySystem {
private readonly log: CrawleeLogger;

Check warning on line 219 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.

private readonly desiredConcurrencyRatio: number;

Check warning on line 221 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.
private readonly scaleUpStepRatio: number;

Check warning on line 222 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.
private readonly scaleDownStepRatio: number;

Check warning on line 223 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.
readonly #loggingIntervalMillis: number;
readonly #autoscaleIntervalMillis: number;
private readonly maxTasksPerMinute: number;
/** The cap on tasks started per minute, or `Infinity` when uncapped. */
readonly maxTasksPerMinute: number;

#minConcurrency: number;
#maxConcurrency: number;
Expand All @@ -232,9 +233,9 @@
#lastLoggingTime?: number;
#tasksPerMinute: number[] = Array.from({ length: 60 }, () => 0);

private readonly snapshotter: Snapshotter;

Check warning on line 236 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.
readonly #loadSignals: LoadSignal[];
private readonly systemStatus: SystemStatus;

Check warning on line 238 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.

#autoscaleInterval?: BetterIntervalID;
#tasksDonePerSecondInterval?: BetterIntervalID;
Expand Down Expand Up @@ -359,7 +360,7 @@
* meaningless. A contradictory pair (`minConcurrency > maxConcurrency`) resolves in favour of the maximum, since
* that is the limit callers set in order to protect something.
*/
private clampDesiredConcurrency(): void {

Check warning on line 363 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.
const atLeastMin = Math.max(this.#desiredConcurrency, this.#minConcurrency);
this.#desiredConcurrency = Math.min(atLeastMin, this.#maxConcurrency);
}
Expand Down Expand Up @@ -390,7 +391,7 @@
await this.#startPromise;
}

private async boot(): Promise<void> {

Check warning on line 394 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.
// Per-session measurement state, reset so a restarted system isn't judged on the previous session. The
// per-minute window matters most: its ageing interval is cleared while we are down, so starts from before an
// arbitrarily long stop would otherwise still count against "this minute" and trip the cap immediately.
Expand Down Expand Up @@ -434,7 +435,7 @@
await this.shutDown();
}

private async shutDown(): Promise<void> {

Check warning on line 438 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.
if (this.#autoscaleInterval) betterClearInterval(this.#autoscaleInterval);
if (this.#tasksDonePerSecondInterval) betterClearInterval(this.#tasksDonePerSecondInterval);
await this.snapshotter.stop();
Expand All @@ -447,7 +448,7 @@
* {@apilink ConcurrencySystem.isRunning|`isRunning`} on the way in. Both the overload verdict and
* `desiredConcurrency` are frozen at that point, so the borrowing pool would otherwise just quietly mis-scale.
*/
private warnIfNotRunning(): void {

Check warning on line 451 in packages/basic-crawler/src/internals/autoscaling/concurrency_system.ts

View workflow job for this annotation

GitHub Actions / Lint

crawlee(prefer-private-fields)

Use a native `#private` field instead of the TypeScript `private` modifier.
if (this.#running || this.#warnedAboutQueryWhileStopped) {
return;
}
Expand Down
8 changes: 5 additions & 3 deletions packages/basic-crawler/src/internals/basic-crawler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -706,8 +706,7 @@ export class BasicCrawler<
* request queue; subsequent ones get their own queue via a unique alias so they don't
* collide.
*/
// kept as TS-private: tests reset the counter at runtime
private static instanceCount = 0;
static #instanceCount = 0;

/**
* Tracks crawler instances that accessed shared state without having an explicit id.
Expand Down Expand Up @@ -1066,7 +1065,7 @@ export class BasicCrawler<
// Initialize the Configuration instance to avoid lazy loading in the components
serviceLocator.getConfiguration();

const instanceIndex = BasicCrawler.instanceCount++;
const instanceIndex = BasicCrawler.#instanceCount++;
this.#identity = { instanceIndex, hasExplicitId: id !== undefined, id: id ?? String(instanceIndex) };

if (requestManager !== undefined && (requestList !== undefined || requestQueue !== undefined)) {
Expand Down Expand Up @@ -2637,6 +2636,7 @@ export class BasicCrawler<
}

/** Handles a single request - runs the request handler with retries, error handling, and lifecycle management. */
// oxlint-disable-next-line crawlee/prefer-private-fields -- patched by @crawlee/otel
private async handleRequest(
crawlingContext: ExtendedContext,
requestSource: IRequestManager,
Expand Down Expand Up @@ -2817,6 +2817,7 @@ export class BasicCrawler<
*
* @param request The request object, passed separately to circumvent potential dynamic logic in crawlingContext.request
*/
// oxlint-disable-next-line crawlee/prefer-private-fields -- patched by @crawlee/otel
private async requestFunctionErrorHandler(
error: Error,
crawlingContext: CrawlingContext,
Expand Down Expand Up @@ -2895,6 +2896,7 @@ export class BasicCrawler<
await this.handleFailedRequestHandler(crawlingContext, error); // This function prints an error message.
}

// oxlint-disable-next-line crawlee/prefer-private-fields -- patched by @crawlee/otel
private async handleFailedRequestHandler(crawlingContext: CrawlingContext, error: Error): Promise<void> {
// Always log the last error regardless if the user provided a failedRequestHandler
const { id, url, method, uniqueKey } = crawlingContext.request;
Expand Down
35 changes: 17 additions & 18 deletions packages/basic-crawler/src/internals/throttling_request_manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -279,9 +279,8 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage

readonly #subManagers = new Map<string, Promise<T>>();

// Not `#private`, unlike the rest: the tests reach for these two.
private readonly domainStates = new Map<string, DomainState>();
private readonly log: CrawleeLogger;
readonly #domainStates = new Map<string, DomainState>();
readonly #log: CrawleeLogger;

/** Domains from the `domains` option, which are throttled whether or not the crawl ever visits them. */
readonly #listedDomains = new Set<string>();
Expand Down Expand Up @@ -355,7 +354,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
this.#throttleBy = options.throttleBy ?? 'hostname';
this.#maxThrottledDomains = options.maxThrottledDomains ?? 100;
this.#persistStateKey = options.persistStateKey ?? DEFAULT_PERSIST_STATE_KEY;
this.log = serviceLocator.getLogger().child({ prefix: 'ThrottlingRequestManager' });
this.#log = serviceLocator.getLogger().child({ prefix: 'ThrottlingRequestManager' });

for (const domain of Array.isArray(options.domains) ? options.domains : []) {
let hostname: string;
Expand All @@ -371,7 +370,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage

const key = this.#domainKey(hostname);
this.#listedDomains.add(key);
this.domainStates.set(key, newDomainState(key));
this.#domainStates.set(key, newDomainState(key));
}
}

Expand Down Expand Up @@ -440,7 +439,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
return;
}
this.#warnedAbout.add(key);
this.log.warning(message);
this.#log.warning(message);
}

#extractDomain(url: string): string {
Expand All @@ -453,7 +452,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage

#getDomainState(url: string): DomainState | null {
const domain = this.#extractDomain(url);
return this.domainStates.get(domain) ?? null;
return this.#domainStates.get(domain) ?? null;
}

/**
Expand Down Expand Up @@ -529,11 +528,11 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
}

#ensureDomainState(domain: string): DomainState {
let state = this.domainStates.get(domain);
let state = this.#domainStates.get(domain);

if (!state) {
state = newDomainState(domain);
this.domainStates.set(domain, state);
this.#domainStates.set(domain, state);
}

return state;
Expand Down Expand Up @@ -589,7 +588,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
*/
#fetchableDomains(): string[] {
const now = Date.now();
return Array.from(this.domainStates.values())
return Array.from(this.#domainStates.values())
.filter((state) => now >= throttledUntil(state) && this.#subManagers.has(state.domain))
.sort((a, b) => throttledUntil(a) - throttledUntil(b))
.map((state) => state.domain);
Expand Down Expand Up @@ -642,7 +641,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
}

this.#minCrawlDelayMs = Math.max(this.#minCrawlDelayMs, intervalMs);
this.log.debug(`Crawl-delay floor for every domain set to ${(this.#minCrawlDelayMs / 1000).toFixed(1)}s`);
this.#log.debug(`Crawl-delay floor for every domain set to ${(this.#minCrawlDelayMs / 1000).toFixed(1)}s`);

return true;
}
Expand All @@ -664,7 +663,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage

if (scope === 'hostname' && this.#throttleBy === 'registrableDomain') {
// Only debug: grouping by registrable domain deliberately paces whole sites, subdomains included.
this.log.debug(
this.#log.debug(
`Applying a pacing signal scoped to "hostname" across the whole registrable domain, because that ` +
`is how this manager groups requests (\`throttleBy\`). Sibling subdomains are paced with it.`,
);
Expand Down Expand Up @@ -719,7 +718,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage

if (delayMs > this.#maxDelayMs) {
const source = waitGiven ? 'requested wait' : 'exponential backoff';
this.log.warning(
this.#log.warning(
`Capping ${source} delay of ${(delayMs / 1000).toFixed(1)}s for domain "${state.domain}" ` +
`to maxDelaySecs (${(this.#maxDelayMs / 1000).toFixed(1)}s); the domain may continue to rate-limit. ` +
`Consider increasing maxDelaySecs if this recurs.`,
Expand All @@ -730,7 +729,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
state.backoffUntil = now + delayMs;
state.backoffDecaysAt = state.backoffUntil + delayMs;

this.log.info(
this.#log.info(
`Rate limit (429) detected for domain "${state.domain}" ` +
`(consecutive: ${state.consecutive429Count}, delay: ${(delayMs / 1000).toFixed(1)}s)`,
);
Expand Down Expand Up @@ -759,7 +758,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage

if (state.declaredCrawlDelayMs === null) {
state.declaredCrawlDelayMs = intervalMs;
this.log.debug(`Set crawl-delay for domain "${state.domain}" to ${(intervalMs / 1000).toFixed(1)}s`);
this.#log.debug(`Set crawl-delay for domain "${state.domain}" to ${(intervalMs / 1000).toFixed(1)}s`);
}

return true;
Expand Down Expand Up @@ -926,7 +925,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
let readyAt: number | undefined;
let stallCandidates: DomainState[] | undefined;

for (const state of this.domainStates.values()) {
for (const state of this.#domainStates.values()) {
// A `Crawl-delay` can give a domain a clock before its first request gives it a queue - nothing
// to fetch from and nothing to wait for until then.
if (!this.#subManagers.has(state.domain)) {
Expand Down Expand Up @@ -1027,7 +1026,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage

this.#migratedFromInner = 0;

for (const state of this.domainStates.values()) {
for (const state of this.#domainStates.values()) {
state.consecutive429Count = 0;
state.backoffUntil = 0;
state.crawlDelayUntil = 0;
Expand Down Expand Up @@ -1077,7 +1076,7 @@ export class ThrottlingRequestManager<T extends IRequestManager = IRequestManage
await this.#ensureSubManagers();

for (const domain of this.#fetchableDomains()) {
const state = this.domainStates.get(domain)!;
const state = this.#domainStates.get(domain)!;

// Armed while the fetch below is still suspended, so that a concurrent `fetchNextRequest` cannot
// find the domain fetchable and dispatch into the same window - which would pace each task
Expand Down
Loading
Loading