@plinthjs/concurrency
v0.1.0
Published
Mason concurrency helpers (the Laravel Illuminate\Concurrency equivalent).
Maintainers
Readme
@plinthjs/concurrency
Mason's concurrency helpers, the Illuminate\Concurrency equivalent. Run a batch of tasks at once
and collect the results in the same shape you passed in (run), or schedule work to happen after the
current tick, fire-and-forget (defer). Dependency-free: it uses Node built-ins only.
Two drivers ship: AsyncDriver (the default, cooperative overlap on the event loop, ideal for I/O
fan-out) and WorkerDriver (true parallelism on node:worker_threads, for CPU-bound work).
Install
npm install @plinthjs/concurrencyUsage
import { Concurrency } from '@plinthjs/concurrency'
// An array batch resolves to an array, in the same order.
const [user, orders] = await Concurrency.run([() => fetchUser(1), () => fetchOrders(1)])
// A record batch resolves to a record, with the same keys.
const stats = await Concurrency.run({
users: () => countUsers(),
orders: () => countOrders(),
})
stats.users // number
// If any task throws, run() rejects with that error.Tasks that each await independent I/O genuinely overlap, so a batch finishes in about the time of
its slowest task rather than the sum.
Deferred work
defer schedules a batch to run after the current tick. Its results are discarded and its errors are
swallowed, so a failing deferred task never crashes the request. The returned handle lets you (or a
test) force the run and wait for it:
const handle = Concurrency.defer([() => recordAnalytics(event), () => warmCache()])
await handle.invoke() // runs now if not started yet; never rejects, idempotentWorker threads
For CPU-bound work, use the worker driver. A worker runs in a separate JS realm, so its tasks must
be self-contained: build each one with workerJob(task, payload), where task captures no outer
variables and payload (and the result) are structured-cloneable.
import { WorkerDriver, workerJob } from '@plinthjs/concurrency'
const workers = new WorkerDriver({ timeoutMs: 5_000 }) // optional per-job timeout
const results = await workers.run({
sum: workerJob((p: { a: number; b: number }) => p.a + p.b, { a: 2, b: 3 }),
product: workerJob((p: { a: number; b: number }) => p.a * p.b, { a: 2, b: 3 }),
})
// { sum: 5, product: 6 }Referencing a captured variable inside a worker task throws a ReferenceError in the worker, which
surfaces as a rejected run.
The manager
Concurrency is a process-wide ConcurrencyManager on the async driver. Build your own for a
different default, a worker timeout, or custom drivers:
import { AsyncDriver, ConcurrencyManager } from '@plinthjs/concurrency'
const manager = new ConcurrencyManager({ default: 'async', workerTimeoutMs: 10_000 })
manager.driver() // the default driver (cached after first use)
manager.driver('worker') // the WorkerDriver
manager.extend('manual', () => new AsyncDriver({ schedule: () => {} })) // returns the manager
manager.driver('missing') // throws: Concurrency driver [missing] is not defined.AsyncDriver accepts an injectable schedule function (defaulting to queueMicrotask), so tests can
control exactly when deferred work fires.
