refactor!: Split off ConcurrencySystem from AutoscaledPool - #3917
Conversation
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
left a comment
There was a problem hiding this comment.
Thank you @janbuchar , a few initial ideas I got reading this ⬇️
| this.concurrencySystem.registerTaskEnd(); | ||
| this.ownConcurrency--; |
There was a problem hiding this comment.
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.
| * free slot at once, overshooting the shared budget. | ||
| */ | ||
| tryRegisterTaskStart(): boolean { | ||
| if (!this.hasCapacityForTask()) { |
There was a problem hiding this comment.
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).
There was a problem hiding this comment.
True, I made it crash in that case. 640a297
| const crawler = new CheerioCrawler({ | ||
| concurrencySystem: new ConcurrencySystem({ | ||
| desiredConcurrency: 10, | ||
| maxTasksPerMinute: 120, | ||
| systemStatusOptions: { currentHistorySecs: 10 }, | ||
| }), | ||
| requestHandler, | ||
| }); |
There was a problem hiding this comment.
where did the start() call go? 😄
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.
- 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.
|
@barjin this grew in scope by a fair bit, but I consider it done. Please re-check it. |
barjin
left a comment
There was a problem hiding this comment.
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?
| // 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); |
There was a problem hiding this comment.
This leads to unfair task allocation. If one consumer (A) has a head start, it:
Amaxes out the sharedConcurrencySystem, e.g., running 10 tasks in parallelBis introduced, wants to share theConcurrencySystem- Every
Atask finished will immediately spawn anotherAtask, leading toBstarving.
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();There was a problem hiding this comment.
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? 🤔


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 newConcurrencySystem(snapshotter + system status + budget + autoscaling) thatAutoscaledPooldelegates to. Inject the same one into several crawlers to cap their combined concurrency.IConcurrencySystemcontract; task booking is atomic (including the per-minute cap), so sharing a budget can't oversubscribe it.run()), while an injected system is caller-owned —start()/stop()it yourself, andrun()throws if you forget.ConcurrencySysteminstance.loadSignalsbag (per-resource tuning +custom), and a built-in signal can be switched off withfalse. Both evaluation windows apply to every signal alike, andLoadSignal.start()hands signals the window they'll be sampled over so retention isn't guesswork.Snapshotter/SystemStatusand the load-signal internals are now private toConcurrencySystem:@crawlee/coredrops 17 exported symbols and gains 11. CustomLoadSignals still supported.Breaking changes are in the upgrading guide.