shrk-neo
v1.2.0
Published
Pass JS objects between Node.js worker processes straight from SharedArrayBuffer — no JSON, no IPC/network, zero-copy buffers. powered by vexify
Maintainers
Readme
shrk-neo
Pass JS objects between Node.js worker processes straight from a
SharedArrayBuffer— no JSON, no IPC, no sockets. Buffer / typed-array values travel zero-copy, with no serialization at all; arbitrary JS objects use V8's native serializer.powered by vexify · Apache-2.0
shrk-neo gives Node.js "multi-process" (worker thread) programs a real
shared heap. One process creates the region, every other process attaches to
the very same memory. Writers and readers touch the memory directly with
Atomics-guarded lock-free operations — there is no pipe, no socket, no
postMessage involved in moving the data.
It ships seven primitives:
| primitive | purpose |
| --- | --- |
| SharedMemory | key/value object store (set/get/wait/getMany) |
| SharedBus | lock-free message queue / mailbox, topic-filterable (send/receive) |
| SharedCounter | atomic 32-bit counter (inc/dec/fetchAdd/exchange) |
| SharedLock | reentrant cross-worker mutex (lock/unlock/withLock) |
| SharedTicketLock | FAIR (FIFO) cross-worker mutex — no starvation |
| SharedCond | condition variable (wait/signal/broadcast) |
| SharedArray | typed array living directly in shared memory |
Why "no serialization"?
- Buffer / typed-array / ArrayBuffer values are copied once into the shared region and read back as live views of that memory. Zero copies, zero encoding — a true zero-copy, no-serialization path.
- Arbitrary JS objects cannot physically live outside a V8 heap, so they
are encoded with
v8.serialize()— the fastest native encoding Node has. It beatsJSON.stringifyby a wide margin and natively handlesBuffer,Map,Set,Date, typed arrays, and circular references.
Install
npm install shrk-neoQuick start
// master.js — create the shared region, spawn workers
const { Worker } = require('node:worker_threads')
const { SharedMemory } = require('shrk-neo')
const mem = new SharedMemory({ size: 16 * 1024 * 1024, slots: 256 })
new Worker('./worker.js', { workerData: { sab: mem.sab } })
// Write objects from the main process…
mem.set('config', { retries: 3, backoff: [100, 200, 500], labels: new Map([['env', 'prod']]) })
mem.set('image', Buffer.alloc(1024 * 1024, 7)) // zero-copy blob// worker.js — attach, read straight from shared memory
const { workerData } = require('node:worker_threads')
const { SharedMemory } = require('shrk-neo')
const mem = SharedMemory.attach(workerData.sab)
const cfg = mem.get('config') // -> { retries: 3, backoff: [100,200,500], labels: Map{env->prod} }
const img = mem.get('image') // -> Buffer that views the shared memory (zero-copy)
// …and write back the other way
mem.set('result', { ok: true, checksum: img.length })That's it. The data never crosses an IPC channel — worker.js reads the
object bytes directly out of the region the master created.
API
new SharedMemory(options?)
Creates a new shared region. Only the process that spawns the workers calls
this. Exposes .sab (the SharedArrayBuffer) to hand to workers.
| option | default | description |
| --- | --- | --- |
| size | 4 * 1024 * 1024 | total size in bytes of the region (header + data) |
| slots | 64 | number of key slots — the max number of live keys |
SharedMemory.attach(sab)
Attaches to a region created by another process. No allocation — reads/writes
go directly to the shared memory. Returns a new SharedMemory instance.
set(key, value) → this
Writes value under key. Replacing a key never leaves readers seeing a
half-written value (new slot is published first, old one freed after).
get(key, options?) → value | undefined
Reads key directly from the shared region.
- Buffer values return a
Bufferthat views the SharedArrayBuffer (zero-copy, no serialization). Pass{ copy: true }for an independent snapshot. - Other values are returned via
v8.deserialize().
has(key) → boolean
True if the key currently exists.
keys() → string[]
All currently stored keys.
getMany(keys) → object
Reads many keys in one call. Returns { key: value }; missing keys are
omitted.
entries() → [key, value][] and iteration
entries() returns all [key, value] pairs; SharedMemory is also
iterable (for (const [k, v] of mem)). size is the number of live keys.
delete(key) → boolean
Removes the key; returns whether it existed.
clear()
Removes every key and resets the allocator (reclaims all space).
wait(key, timeoutMs?) → value | undefined
Blocks until key has a value, then returns it (as get). Sleeps on the
change counter with Atomics.wait — no busy loop. Blocks the calling
thread's event loop, so prefer calling it inside a worker.
waitAsync(key, timeoutMs?) → Promise<value | undefined>
Non-blocking wait for the main thread, built on Atomics.waitAsync
(Node.js ≥ 16.17).
stats() → { slots, liveKeys, freeSlots, bytesUsed, bytesCapacity }
Usage stats for the region. bytesUsed is the allocator's high-water mark
since the last clear().
SharedBus — message queue
const bus = new SharedBus({ size: 4 * 1024 * 1024, slots: 128 })
// main → workers
bus.send({ job: 'render', frames: [1, 2, 3] })
bus.send(Buffer.alloc(64 * 1024)) // raw bytes, zero-copy on receive
// …and from a worker:
bus.send({ result: 'ok' })// worker.js
const bus = SharedBus.attach(workerData.busSab)
const job = bus.receive(5000) // blocking receive (use inside a worker)
const next = bus.tryReceive() // never blocks, undefined if empty
const also = await bus.receiveAsync(5000) // non-blocking, for the main threadnew SharedBus(options?)
| option | default | description |
| --- | --- | --- |
| size | 1 * 1024 * 1024 | total bytes of the region |
| slots | 64 | number of message slots |
Each slot owns a fixed capacity = (size - header) / slots bytes, so space is
reclaimed the instant a consumer claims a message. A payload bigger than
capacity throws (message too large).
send(value, options?) → this
Publishes a message. Objects are v8-encoded; Buffer / typed-array / ArrayBuffer values are stored raw. Throws when every slot is busy (bounded queue).
Options: { topic } — an integer topic tag (default 0). Consumers
that receive with a specific topic only see messages tagged with it; a
receiver with no topic sees everything. Use it to multiplex several logical
channels on one bus:
bus.send({ job: 'render' }, { topic: 1 }) // only topic-1 consumers see it
const job = bus.receive(5000, { topic: 1 }) // filter by topictryReceive(options?) → value | undefined
Non-blocking receive; undefined when no matching message. Options:
{ topic, copy }.
receive(timeoutMs?, options?) → value | undefined
Blocking receive (Atomics.wait, no busy loop). Blocks the calling thread's
event loop — use it in workers, receiveAsync on the main thread. Options:
{ topic, copy }.
receiveAsync(timeoutMs?, options?) → Promise<value | undefined>
Main-thread-friendly blocking receive via Atomics.waitAsync. Options:
{ topic, copy }.
pending / empty / clear()
pending counts queued + in-flight messages; empty is pending === 0;
clear() drops everything.
Like SharedMemory.get, Buffer values are returned as zero-copy views by
default ({ copy: true } for a snapshot). Because a slot is freed as soon as
it is consumed, a zero-copy message's bytes may be reused by a later send —
process it promptly or use copy: true.
SharedCounter — atomic counter
const counter = new SharedCounter({ initial: 0 }) // main
// in a worker: SharedCounter.attach(workerData.counterSab)
counter.inc() // +1, returns new value
counter.inc(5) // +5, returns new value
counter.dec(2) // -2, returns new value
counter.value // current value (atomic load)
counter.reset(0) // store a new value
counter.fetchAdd(3) // +3, returns the OLD value
counter.exchange(10) // store 10, returns the OLD value
counter.compareAndSet(10, 20) // true only if it was 10Useful for progress meters, per-worker credits, or aggregate statistics
across processes. inc/dec/fetchAdd are Atomics.add — safe from many
workers at once; exchange and compareAndSet give you atomic
store/CAS-style semantics.
SharedLock — reentrant mutex
const lock = new SharedLock() // main
// in a worker: SharedLock.attach(workerData.lockSab)
lock.lock()
try {
// critical section, exclusive across all workers
lock.lock() // reentrant: the owner may lock again (depth grows)
lock.unlock() // each lock() needs a matching unlock()
} finally {
lock.unlock()
}
lock.withLock(() => { /* same, with auto-release */ })
lock.tryLock() // true/false, never blocks (reentrant for the owner)
lock.locked // is anyone holding it?
lock.forceUnlock() // escape hatch: release regardless of ownerlock() blocks the calling thread via Atomics.wait (use inside workers).
The lock is reentrant and owner-checked: only the thread that acquired it
may unlock(), and a different process calling unlock() throws. If a
holder dies while holding the lock, any process can forceUnlock() to break
the deadlock.
SharedTicketLock — fair mutex
const lock = new SharedTicketLock() // main
// in a worker: SharedTicketLock.attach(workerData.lockSab)
lock.lock()
try { /* critical section — granted in FIFO arrival order */ }
finally { lock.unlock() }
lock.tryLock() // true/false, never blocks
lock.queued // waiters behind the current holder
lock.lockedA ticket lock grants the mutex in arrival order, so no waiter can be starved.
Not owner-checked — the holder is expected to unlock().
SharedCond — condition variable
const lock = new SharedLock()
const cond = new SharedCond() // main
// in a worker: SharedCond.attach(workerData.condSab)
lock.lock()
while (!predicate()) cond.wait(1000) // blocks the thread; use in a worker
lock.unlock()
cond.signal() // wake one waiter
cond.broadcast() // wake all waiters
await cond.waitAsync(1000) // main-thread friendly version
cond.waiting // how many processes are waiting right nowwait() sleeps via Atomics.wait (no busy loop) until signal()/broadcast()
or a timeout, returning 'ok' or 'timed-out'. Like every condition
variable, pair it with a SharedLock + predicate and re-check the predicate
after waking.
SharedArray — shared typed array
const arr = new SharedArray({ type: 'float64', length: 1024 }) // main
// in a worker: SharedArray.attach(workerData.arraySab)
arr.view[i] = 3.14 // every process sees the SAME live view — no IPC
arr.set(i, 2.5) // convenience setter
arr.at(i) // convenience getter
arr.fill(0)
const snapshot = arr.toArray() // defensive copy (safe to transfer)
arr.bytes // data-region bytesSupported types: float64, float32, int32, uint32, int16, uint16,
int8, uint8, bigint64, biguint64 (also pass size instead of
length). .view is a plain TypedArray into the shared buffer — writes are
instantly visible to all attached processes. Element access is not atomic;
coordinate with a SharedLock/SharedTicketLock when readers must observe
consistent snapshots.
Semantics & caveats
- Scope —
SharedArrayBuffercan only be shared between worker threads of the same process, which is Node's native "multi-process" model. Data does not cross OS-process boundaries (that would require IPC, which this package deliberately avoids). - Concurrency — the store is lock-free for reads. The safe pattern is:
writers own their keys, readers read. Concurrent
set/deleteracing withgeton the same key may return the previous value; coordinate such writes yourself (e.g. one producer per key). - Memory reuse — deleted values free their slot immediately; their
bytes are reclaimed by the next
clear(). If the region fills upsetthrowsshared memory exhausted— callclear(), delete keys, or create a larger region. - Zero-copy buffers are live views of shared memory. If a writer later
overwrites the same key, a long-held zero-copy buffer may see its content
change; use
{ copy: true }or snapshot the bytes when that matters.
Examples
node examples/main.jsLicense
Apache-2.0 · powered by vexify
