@plinthjs/horizon
v0.1.0
Published
Mason queue supervisor and metrics (the Laravel Horizon equivalent, headless).
Maintainers
Readme
@plinthjs/horizon
Mason's queue supervisor and metrics, a headless Laravel Horizon equivalent built on
@plinthjs/queue (there is no dashboard UI). A Supervisor manages a pool of worker slots draining one
or more queues, a MasterSupervisor coordinates several supervisors, and a MetricsRecorder reports
throughput, runtimes, wait times and processed/failed counts.
Mason models a Horizon "process" as a logical concurrency slot over an injected work loop, so
scaling, load balancing and the running / paused / terminated lifecycle are deterministic and
testable without spawning OS processes.
Install
npm install @plinthjs/horizonUsage
import { ArrayQueue } from '@plinthjs/queue'
import { pendingSizeFor, Supervisor, workLoopFor } from '@plinthjs/horizon'
const queue = new ArrayQueue()
const supervisor = new Supervisor(
workLoopFor(queue), // each slot drains the queue once per tick via queue.work()
{ queues: ['default', 'emails'], maxProcesses: 4, balance: 'auto' },
pendingSizeFor(queue), // load probe used by 'auto' balancing
)
supervisor.start() // supervisors start 'paused'
await supervisor.tick() // { processed, failed } for one round across every slot
// Or loop until paused/terminated, or until the predicate says stop.
let rounds = 0
await supervisor.run(() => ++rounds < 10)
supervisor.terminate() // a terminated supervisor cannot be restartedScaling and balancing
supervisor.processCount() // 4
supervisor.processes() // { default: 2, emails: 2 }
supervisor.processesFor('emails') // 2
supervisor.scale(2) // clamped to 1..maxProcesses
supervisor.rebalance() // force an 'auto' rebalance now; returns the allocation
supervisor.watching() // ['default', 'emails']The balance option picks the strategy:
'off'(default): every slot works the first queue.'simple': slots are split evenly across the queues.'auto': slots follow each queue's pending size, moving at mostbalanceMaxShiftslots per rebalance, keeping at leastminProcessesper queue, and waitingbalanceCooldownMsbetween rebalances (measured with the injectablenowclock).
Master supervisor
import { MasterSupervisor } from '@plinthjs/horizon'
const master = new MasterSupervisor({ web: webSupervisor }).add('mail', mailSupervisor)
master.start()
await master.tick() // summed { processed, failed } across children
master.state() // 'running' if any child runs
master.states() // { web: 'running', mail: 'running' }
master.processes() // slots summed per queue
master.pause().continue().terminate()Metrics
import { MetricsRecorder } from '@plinthjs/horizon'
let now = 0
const metrics = new MetricsRecorder({ now: () => now, throughputWindowMs: 60_000 })
metrics
.recordJob({ queue: 'emails', runtimeMs: 120, status: 'processed', waitMs: 30 })
.recordJob({ queue: 'emails', runtimeMs: 80, status: 'failed' })
metrics.processed('emails') // 1
metrics.failed() // 1
metrics.averageRuntime('emails') // 100
metrics.throughput('emails') // jobs per minute within the window
metrics.forQueue('emails') // { queue, processed, failed, averageRuntimeMs, averageWaitMs, throughputPerMinute }
metrics.snapshot() // { at, overall, queues }