Skip to content

refactor!: Split off ConcurrencySystem from AutoscaledPool - #3917

Merged
janbuchar merged 55 commits into
v4from
injectable-autoscaled-pool
Jul 30, 2026
Merged

janbuchar merged 55 commits into
v4from
injectable-autoscaled-pool

Conversation

@janbuchar

@janbuchar janbuchar commented Jul 27, 2026 •

Copy link
Copy Markdown
Contributor

Accepting a whole pre-built AutoscaledPool (as the issue suggested) can't work — the pool has a single task source, so a shared instance ignores the second crawler's work. Split out the only shareable part instead: a new ConcurrencySystem (snapshotter + system status + budget + autoscaling) that AutoscaledPool delegates to. Inject the same one into several crawlers to cap their combined concurrency.

  • Pool and crawlers depend on a minimal read-only IConcurrencySystem contract; task booking is atomic (including the per-minute cap), so sharing a budget can't oversubscribe it.
  • Inject-or-default ownership (Unify the inject-or-default ownership pattern #3888): crawlers build and drive their own default (fresh per run()), while an injected system is caller-owned — start()/stop() it yourself, and run() throws if you forget.
  • Runtime concurrency tuning moved from the pool to the ConcurrencySystem instance.
  • Load-signal configuration is now a single loadSignals bag (per-resource tuning + custom), and a built-in signal can be switched off with false. Both evaluation windows apply to every signal alike, and LoadSignal.start() hands signals the window they'll be sampled over so retention isn't guesswork.
  • Snapshotter/SystemStatus and the load-signal internals are now private to ConcurrencySystem: @crawlee/core drops 17 exported symbols and gains 11. Custom LoadSignals still supported.

Breaking changes are in the upgrading guide.

@janbuchar janbuchar added the t-tooling Issues with this label are in the ownership of the tooling team. label Jul 27, 2026
@janbuchar
janbuchar requested a review from barjin July 27, 2026 20:38
janbuchar added 15 commits July 27, 2026 23:17
With several pools borrowing one ConcurrencySystem, the capacity check,
the await of isTaskReadyFunction and the concurrency increment in
maybeRunTask formed a check-then-act race - two pools could both pass
the check while one slot remained and both book a task, exceeding the
shared budget by up to N-1. Replace registerTaskStart() with an atomic
tryRegisterTaskStart() that re-checks and books in a single synchronous
step; the pool's up-front hasCapacityForTask() call remains as a cheap
early-out before querying task readiness.
Any of the minConcurrency/maxConcurrency/maxRequestsPerMinute shortcuts
made HttpCrawler abandon its HTTP-optimized defaults entirely (higher
starting concurrency, relaxed event loop signal), silently degrading
e.g. CheerioCrawler({ maxConcurrency: 5 }) to desiredConcurrency 1 with
default snapshotter tuning - on master the shortcuts merged into the
optimized options instead.

Worse, the optimized system was passed to BasicCrawler as an injected
concurrencySystem, so the crawler treated it as borrowed and never
started or stopped it - no snapshots and no autoscaling for any default
HttpCrawler.

Replace the option juggling with a protected
createDefaultConcurrencySystem() hook: BasicCrawler calls it to build
the owned default (so the lifecycle handling stays correct), and
HttpCrawler overrides it to fold the HTTP-optimized tuning underneath
the user's shortcuts. Also fix the cheerio-max-requests e2e actor to
drive the lifecycle of the system it injects, per the documented
contract.
The class JSDoc claimed a shared instance is reference-counted, with
stop() only tearing things down once the last borrower leaves. The
implementation is a plain idempotency flag - the first stop() stops the
snapshotter for every borrower. Describe the actual single-owner
contract instead of promising bookkeeping that does not exist.
The crawler-owned governor was built once in the constructor and merely
restarted across repeated run() calls, so state leaked between runs:
snapshot stores are not cleared on stop() (the event loop signal would
even measure the whole idle gap as blockage on the next start, tripping
an immediate scale-down), and the autoscaled desired concurrency and
per-minute task counts carried over. Pre-split, every run got a fresh
pool with a fresh snapshotter.

Build the owned default lazily in _init() instead, once per run - this
also moves the createDefaultConcurrencySystem() call out of the base
constructor, so subclass overrides no longer run against a partially
constructed instance. An injected system is untouched: its state is
long-lived by design and its lifecycle belongs to the caller.
The pre-split pool skipped its autoscaling tick while paused. The
autoscaling loop now lives on the ConcurrencySystem, which deliberately
knows nothing about its borrowing pools' pause state (other pools
sharing the system may still be running), so it keeps evaluating and
logging while a pool is paused. Spell this out in the upgrading guide
and on pause() itself instead of leaving the behavior change silent.
_autoscale, _scaleUp, _scaleDown and _incrementTasksDonePerSecond were
declared protected, leaking into the public API report of a class that
has no subclasses - the very thing #3890 privatized on AutoscaledPool
before the split moved them here. Tests already reach them through
ts-expect-error suppressions, which work the same for private members.
…plit

The scaling guide (and its v4.0 snapshot) still steered users to pass
desiredConcurrency, scaling ratios, autoscaleIntervalSecs,
loggingIntervalSecs and maxTasksPerMinute through autoscaledPoolOptions,
none of which live there anymore - point it at an injected
ConcurrencySystem instead, including the caller-owned lifecycle and the
shared-budget use case. Also drop the maybeRunIntervalSecs section
(no longer exposed through crawler options), correct the documented
desiredConcurrencyRatio default (0.9, not 0.95), and refresh the
Snapshotter class JSDoc, which referenced the removed flat option names
and a nonexistent maxUsedCpuRatio snapshotter option.
The warning that concurrency shortcuts are ignored alongside an
injected concurrencySystem used a ??-chain coerced to a boolean, so a
falsy-but-present value (e.g. minConcurrency: 0) short-circuited the
chain and silently suppressed the warning even when other shortcuts
were set. Compare against undefined explicitly.
AutoscaledPool and BasicCrawler consumed the concrete ConcurrencySystem
class, coupling them to its whole surface (snapshotting, system status,
autoscaling internals) when all they use is the budget accessors and the
task-booking protocol. Extract exactly that contract into an
IConcurrencySystem interface - the pool and the crawler's
concurrencySystem option now depend on it, ConcurrencySystem stays the
canonical implementation the crawler builds by default, and alternate
governors can be substituted without subclassing. The interface also
spells out the atomicity requirement on tryRegisterTaskStart() that
shared budgets rely on.

getCurrentStatus() deliberately stays off the interface: SystemInfo
mandates all four built-in resource infos, which is telemetry of the
canonical implementation rather than part of the governor contract.
The pool checked isOverMaxRequestLimit separately before booking, which
had the same check-then-await weakness as the concurrency budget once
had: with several pools sharing a governor, two pools could both pass
the cap check and both book. Fold the cap into the atomic booking
instead - it still only applies once a task is known to be ready
(tryRegisterTaskStart is called after the readiness query, and
hasCapacityForTask deliberately ignores it), so an empty queue never
stalls the pool for an extra minute.

This also removes isOverMaxRequestLimit from IConcurrencySystem - a
per-minute cap is scaling policy of the canonical implementation, not
part of the governor contract - and makes it private on
ConcurrencySystem.
The systemStatusOptions passthrough carried three members: loadSignals
(a genuine extension point), currentHistorySecs, and a snapshotter
injection hook only tests ever used. Hoist the two useful ones directly
onto ConcurrencySystemOptions and drop the bag - with that, SystemStatus
and SystemStatusOptions have no public consumers left and become
internal implementation details of the ConcurrencySystem, per #3109.
With SystemStatus folded away, the Snapshotter class, the built-in
signal implementations (MemoryLoadSignal, the create*LoadSignal
factories), their per-class option interfaces, the concrete snapshot
types and evaluateLoadSignalSample have no public consumers - users
only ever touch the data-only SnapshotterOptions bags, and nothing
in the codebase references any of them from outside the autoscaling
module. Mark them all internal.

The custom-signal extension point stays public: the LoadSignal
interface, LoadSnapshot and the SnapshotStore composition helper,
now documented as wired in via ConcurrencySystemOptions.loadSignals.

@barjin barjin left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you @janbuchar , a few initial ideas I got reading this ⬇️

Comment on lines +455 to +456
this.concurrencySystem.registerTaskEnd();
this.ownConcurrency--;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Perhaps we should do this in the catch (or finally) branch as well?

import { BasicCrawler, CriticalError } from './packages/basic-crawler/src/index.js';

const c = new BasicCrawler({
    requestHandler: async () => {
        throw new CriticalError('Request failed');
    },
});

await c.run(['https://crawlee.dev/x']).catch(() => {}); // do not exit

console.log(c.autoscaledPool?.currentConcurrency);
// ⬆️  will log 1, as CriticalError goes straight to the AutoscaledPool, but the catch branch doesn't `registerTaskEnd`

Since this PR wants users to share the ConcurrencySystem, this might be a problem.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice catch! 3029248

* free slot at once, overshooting the shared budget.
*/
tryRegisterTaskStart(): boolean {
if (!this.hasCapacityForTask()) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We never check whether the ConcurrencySystem has actually been started. Since the .start() call is now users' responsibility, we should imo warn them, if they haven't called it.

migrating this

const c = new BasicCrawler({
    autoscaledPoolOptions: {
        maxConcurrency: 1,
    },
});

await c.run(['https://crawlee.dev/x']);

to this

const c = new BasicCrawler({
    concurrencySystem: new ConcurrencySystem({
        maxConcurrency: 1,
    }),
});

await c.run(['https://crawlee.dev/x']);

is easier than it seems (and currently logs no warning), and will leave the user with scaling hard-pinned to desiredConcurrency (iiuc).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

True, I made it crash in that case. 640a297

Comment thread docs/upgrading/upgrading_v4.md Outdated
Comment on lines +1242 to +1249
const crawler = new CheerioCrawler({
concurrencySystem: new ConcurrencySystem({
desiredConcurrency: 10,
maxTasksPerMinute: 120,
systemStatusOptions: { currentHistorySecs: 10 },
}),
requestHandler,
});

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

where did the start() call go? 😄

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Comment thread packages/core/src/autoscaling/autoscaled_pool.ts
getHistoricalStatus() asked each signal for everything it had retained,
so the window autoscaling reasons over was whatever retention each
signal happened to have: snapshotHistorySecs for the built-ins, but a
private, unsettable value for custom signals - a signal keeping five
minutes of snapshots let stale overloads hold desired concurrency down
long after the resource recovered.

Request the window explicitly instead, so snapshotHistorySecs bounds
the historical evaluation for every signal, custom ones included. This
is behavior-preserving at the defaults (built-ins retain exactly the
30s they are evaluated over, and SnapshotStore defaults to the same);
it only diverges where retention and window were deliberately set
apart, which is precisely where the old behavior was surprising.

Snapshot retention keeps its role as the memory bound - it still caps
what any window can see.
Snapshot retention was out-of-band knowledge: a custom signal had to
pass a millisecond value to new SnapshotStore(...) that happened to
match the windows of whichever ConcurrencySystem would drive it. Guess
too low and the signal silently contributed a narrower view of its
resource than every other signal, with nothing to warn about it - and
the value it needed was not reachable from anywhere in the interface.

LoadSignal.start() now receives a context with the widest window the
signal will be queried with (the wider of the gating and autoscaling
windows), so retention is derived rather than guessed. Signals built on
SnapshotStore apply it automatically via useSampleWindow(); hand-written
ones get the number explicitly. Implementations that declare no
parameter still satisfy the interface, so this breaks no custom signal.

Consequently retention stops being separately configurable: the
per-signal snapshotHistoryMillis options are gone, the Snapshotter no
longer threads a retention value into four constructors, and
snapshotHistorySecs now means one thing - the autoscaling window, which
signals size themselves to. This also removes the possibility of a
gating window wider than retention silently collapsing the two windows
into one, since retention covers the wider of them by construction.
janbuchar added 14 commits July 29, 2026 14:16
- IConcurrencySystem.isRunning is required: an implementation with no startup
  lifecycle reports `true`, instead of the member being optional with only an
  explicit `false` meaning anything.
- AutoscaledPoolTaskLoopOptions is derived from AutoscaledPoolOptions rather
  than being a third hand-written options interface.
- Document what AutoscaledPool.system and ConcurrencySystem.getCurrentStatus()
  are actually for; the former claimed a re-injection workflow nothing uses, the
  latter claimed a caller it does not have.
Three fields (the resolved dependency, the injected instance and a builder for
the default) collapse into a resolver closure plus the dependency it produces.
The constructor no longer builds a placeholder OwnedOrInjected that `_init()`
immediately discards - the field is simply absent until the first run, which is
what made a pre-`run()` `teardown()` a no-op anyway.
- The upgrading guide states the start()/stop() ownership rule once, in an
  admonition, instead of five times across the section.
- The two genuinely new capabilities it was documenting - switching a built-in
  load signal off, and wrapping one by constructing its class - move to the
  scaling guide, where features belong; the upgrading guide links to them and
  keeps only the behavioral change (a duplicate signal name now throws).
- The scaling guide gains a Load signals section covering those, custom signals
  and the two evaluation windows.
- Signal docblocks stop repeating the same paragraph four times over, and the
  contract-level explanation lives on LoadSignal itself. Drop the essay from the
  private assertUniqueSignalNames, and the @category on classes marked
  @internal.
The one snippet teaching the replacement for `autoscaledPool.maxConcurrency = 10`
was the one snippet that skipped `stop()`, and it dropped the HTTP-optimized
preset that the same section warns about a few paragraphs later - so the
canonical example walked into both traps it documents. It now shows the full
ownership cycle, keeps the preset, and tunes from an `errorHandler` where a
reader would actually want it.

Also spell out that this is the only way in: `pool.system` is typed as the
read-only `IConcurrencySystem`, so a crawler's own default system cannot be
retuned after the fact.
website/versioned_docs is generated content, synced by its own pass - the
ConcurrencySystem split has no business hand-editing the 4.0 snapshot. Reverts
this branch's edits to version-4.0/guides/scaling_crawlers.mdx and the example
file swap alongside it.
@janbuchar

Copy link
Copy Markdown
Contributor Author

@barjin this grew in scope by a fair bit, but I consider it done. Please re-check it.

@janbuchar
janbuchar requested a review from barjin July 29, 2026 22:26

@barjin barjin left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you @janbuchar , the changes are mostly fine with me 👍

My AI junior associate flagged the unfair scheduling with shared ConcurrencySystem. I managed to reproduce this and tried describing it in the comment.

On one hand, it's real and leads to consumer starvation, on the other, I cannot think of an easy fix. Let's perhaps make a separate issue and solve this in another PR?

Comment on lines +418 to 420
// Run task after the previous one finished. Only on success: a failed task rejects the pool, and
// nudging the loop afterwards could start work on an already destroyed pool.
setImmediate(this.maybeRunTask);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This leads to unfair task allocation. If one consumer (A) has a head start, it:

  1. A maxes out the shared ConcurrencySystem, e.g., running 10 tasks in parallel
  2. B is introduced, wants to share the ConcurrencySystem
  3. Every A task finished will immediately spawn another A task, leading to B starving.

Basically, the resource allocation is determined at initialization time, but won't change afterwards (until either consumer doesn't have any requests to process).

minimal example:

import { AutoscaledPool, ConcurrencySystem } from './packages/core/src/index.js';

const concurrencySystem = new ConcurrencySystem({
    minConcurrency: 10,
    maxConcurrency: 10,
    desiredConcurrency: 10,
    loggingIntervalSecs: null,
    loadSignals: { memory: false, eventLoop: false, cpu: false, client: false },
});

const done = { A: 0, B: 0 };
const pool = (name: 'A' | 'B') =>
    new AutoscaledPool({
        concurrencySystem,
        isFinishedFunction: async () => false,
        isTaskReadyFunction: async () => true,
        runTaskFunction: async () => {
            await new Promise((r) => setTimeout(r, 250));
            done[name]++;
        },
    });

const [a, b] = [pool('A'), pool('B')];
await concurrencySystem.start();

void a.run();
await new Promise((r) => setTimeout(r, 100)); // long enough for A to fill all 10 slots
void b.run();
await new Promise((r) => setTimeout(r, 2000));

console.log(done); // { A: 80, B: 0 }

await a.abort();
await b.abort();
await concurrencySystem.stop();

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yeah, I don't think we absolutely need to come up with a solution, but maybe we could make it possible for custom concurrency systems to resolve this? What if the interface accepted a crawler identity in the task allocation methods? 🤔

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

t-tooling Issues with this label are in the ownership of the tooling team.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Accept a pre-built AutoscaledPool instance in crawler options

4 participants