SDKs: TypeScript, Python, PHP, Go, Rust, Elixir
Client SDKs for every runtime.
Official client SDKs for Node.js, Deno, Bun, Python, PHP, Go, Rust, Elixir and Cloudflare Workers. Each speaks the native TCP protocol and MessagePack through an idiomatic Queue and Worker API.
The design is simple: the server owns every queue semantic, retries with backoff, priorities, scheduling, stall detection, the dead letter queue. Your applications only add jobs and process them. A Next.js API written in TypeScript can enqueue work that a Python service consumes, both against the same queue, with no shared runtime and no translation layer.
Supported platforms
Section titled “Supported platforms”| Platform | Package | Distribution |
|---|---|---|
| Node.js ≥ 20, Bun, Deno ≥ 2, Cloudflare Workers | bunqueue-client | npm · source |
| Python ≥ 3.9 | bunqueue-client | PyPI · source |
| PHP ≥ 8.1 | bunqueue/client | Packagist · source |
| Go ≥ 1.26.5 | github.com/egeominotti/bunqueue/sdk/go | go get, source |
| Rust ≥ 1.85 | bunqueue-client | crates.io · source |
| Elixir ≥ 1.15 | bunqueue_client | Hex upcoming, source |
| Bun, embedded in process, no server | bunqueue | Quick Start |
Every SDK is versioned independently from the server and follows semantic
versioning; each changelog is linked in Resources. The wire
protocol is a formal, public, versioned contract (protocolVersion: 2,
negotiated via Hello), so a client written today keeps working across
server upgrades, and services in different languages interoperate on the
same queues out of the box.
Getting started
Section titled “Getting started”1. Run the server
Section titled “1. Run the server”The server is the only component that requires Bun, and
only if you run it with bunx. Docker and the prebuilt binary need nothing
at all.
bunx bunqueue startOne command, no install. Without --data-path (or BUNQUEUE_DATA_PATH)
the queue is in-memory: pass it, e.g. --data-path ./data/bunq.db, to
persist jobs to SQLite.
docker run -d --name bunqueue \ -p 6789:6789 -p 6790:6790 \ -v bunqueue-data:/app/data \ ghcr.io/egeominotti/bunqueue:latestThe named volume keeps the SQLite file across container restarts and upgrades.
services: bunqueue: image: ghcr.io/egeominotti/bunqueue:latest ports: - "6789:6789" # TCP protocol (SDKs) - "6790:6790" # HTTP API (/health, /metrics) volumes: - bunqueue-data:/app/data # environment: # AUTH_TOKENS: "your-secret-token"
volumes: bunqueue-data:docker compose up -d# download the binary for your platform from GitHub Releasescurl -LO https://github.com/egeominotti/bunqueue/releases/latest/download/bunqueue-darwin-arm64chmod +x bunqueue-darwin-arm64./bunqueue-darwin-arm64 startPrebuilt binaries for macOS and Linux are attached to every release; no runtime required.
Port 6789 serves the TCP protocol used by the SDKs, port 6790 serves the
HTTP API with /health and /metrics. Additional deployment options are
covered in Running the Server.
2. Install the client
Section titled “2. Install the client”npm install bunqueue-clientbun add bunqueue-clientdeno add npm:bunqueue-clientpip install bunqueue-clientcomposer require bunqueue/clientgo get github.com/egeominotti/bunqueue/sdk/gocargo add bunqueue-client# Hex release upcoming; use sdk/elixir as a path dependency today{:bunqueue_client, path: "../bunqueue/sdk/elixir"}3. Produce jobs
Section titled “3. Produce jobs”import { Queue } from 'bunqueue-client';
const queue = new Queue('emails');
await queue.add('welcome', { to: 'user@example.com' });from bunqueue import Queue
queue = Queue("emails")
queue.add("welcome", {"to": "user@example.com"})use Bunqueue\Queue;
$queue = new Queue('emails');
$queue->add('welcome', ['to' => 'user@example.com']);queue := bunqueue.NewQueue("emails", bunqueue.Options{})defer queue.Close()
queue.Add("welcome", map[string]any{"to": "user@example.com"}, nil)use bunqueue_client::{ConnectionOptions, JobOptions, Queue, Value};
let queue = Queue::new("emails", ConnectionOptions::default());let data = Value::Map(vec![(Value::from("to"), Value::from("user@example.com"))]);queue.add("welcome", data, JobOptions::default())?;queue = Bunqueue.queue("emails"){:ok, _job} = Bunqueue.Queue.add(queue, "welcome", %{to: "user@example.com"})4. Process jobs
Section titled “4. Process jobs”import { Worker } from 'bunqueue-client';
const worker = new Worker('emails', async (job) => { await sendEmail(job.data.to); return { sent: true };});
worker.on('completed', (job) => console.log('done:', job.id));worker.on('error', (err) => console.error(err)); // always attach (see Worker semantics)from bunqueue import Worker
def process(job): send_email(job.data["to"]) return {"sent": True}
Worker("emails", process, concurrency=10).run()use Bunqueue\Worker;
$worker = new Worker('emails', function (Bunqueue\Job $job) { sendEmail($job->data()['to']); return ['sent' => true];});
$worker->installSignalHandlers();$worker->run(); // blocking loop; or $worker->runOnce() from a cron tickworker := bunqueue.NewWorker("emails", func(job *bunqueue.Job) (any, error) { return sendEmail(job.Data()["to"].(string))}, bunqueue.WorkerOptions{Concurrency: 8})
worker.Run()use bunqueue_client::{ProcessError, Value, Worker, WorkerOptions};
let worker = Worker::new( "emails", |job| { deliver(job.data()) .map(|_| Value::from(true)) .map_err(|error| ProcessError::retryable(error.to_string())) }, WorkerOptions::default(),);worker.run()?;worker = Bunqueue.Worker.new("emails", fn job -> send_email(job.data) {:ok, %{sent: true}} end, concurrency: 8)
Bunqueue.Worker.run(worker)Run either file with the runtime you already use:
node --experimental-strip-types app.ts # Node 22 or laterbun app.ts # Bundeno run -A app.ts # Deno 2 or laterpython app.py # Pythonphp worker.php # PHPgo run . # Gocargo run # Rustmix run app.exs # ElixirProducer and worker are usually separate services, often in different
languages: a Next.js API adds jobs, a Python service processes them,
against the same queue and the same protocol. Constructors default to
host: 'localhost' and port: 6789, so no options are needed for a local
setup.
Protocol and architecture
Section titled “Protocol and architecture”Understanding four facts about the transport explains most SDK behavior:
- Framing: every message is a 4-byte big-endian length prefix followed by a standard msgpack map. Maximum frame size is 64 MB; maximum job payload is 10 MB.
- Pipelining: every request carries a
reqIdthe server echoes back, so many commands are in flight on one socket concurrently. A single connection is usually all a service needs. - Authentication-first: when a token is configured,
Authis guaranteed to be the first frame on every (re)connection, in every SDK: TypeScript through synchronous write ordering, Python through a connection lock (safe under free-threaded concurrency), and PHP, Go, Rust, and Elixir inside the connect sequence itself. - The job name travels inside
data:add('welcome', {to})is encoded asdata: { name: 'welcome', to }. Non-object payloads are wrapped as{ name, payload }.
Connection options
Section titled “Connection options”const queue = new Queue('emails', { host: 'queue.example.com', port: 6789, token: process.env.BUNQUEUE_TOKEN, tls: true, // or { caFile } or { rejectUnauthorized: false } commandTimeoutMs: 10000, // default maxInFlight: 10_000, // backpressure bound; 0 = unbounded (default) poolSize: 4, // producer-side connection pool; default 1});queue = Queue( "emails", host="queue.example.com", port=6789, token=os.environ["BUNQUEUE_TOKEN"], tls=True, # or {"ca_file": "./ca.pem"} or an ssl.SSLContext command_timeout=10.0, # default)$queue = new Queue('emails', [ 'host' => 'queue.example.com', 'port' => 6789, 'token' => getenv('BUNQUEUE_TOKEN'), 'tls' => true, // or ['caFile' => './ca.pem'] or ['verifyPeer' => false] 'connectTimeout' => 10.0, // seconds, default 'commandTimeout' => 30.0, // seconds, default]);queue := bunqueue.NewQueue("emails", bunqueue.Options{ Host: "queue.example.com", Port: 6789, Token: os.Getenv("BUNQUEUE_TOKEN"), TLS: &bunqueue.TLSOptions{CAFile: "./ca.pem"}, // or &TLSOptions{} for system CAs ConnectTimeout: 10 * time.Second, // default CommandTimeout: 30 * time.Second, // default})use std::{path::PathBuf, time::Duration};use bunqueue_client::{ConnectionOptions, Queue, TlsOptions};
let queue = Queue::new("emails", ConnectionOptions { host: "queue.example.com".into(), port: 6789, token: std::env::var("BUNQUEUE_TOKEN").ok(), tls: Some(TlsOptions { ca_file: Some(PathBuf::from("./ca.pem")) }), connect_timeout: Duration::from_secs(10), command_timeout: Duration::from_secs(30), ..Default::default()});queue = Bunqueue.queue("emails", host: "queue.example.com", port: 6789, token: System.fetch_env!("BUNQUEUE_TOKEN"), tls: true, ca_file: "./ca.pem", timeout: 30_000 )The command timeout governs how long each in-flight command waits for a response. TypeScript and Python keep the TCP connect timeout internal; PHP, Go, and Rust expose it separately. Elixir applies its timeout to connect and command exchange.
Connection resilience
Section titled “Connection resilience”Every SDK preserves the same recovery invariants; implementation details follow the runtime:
| Mechanism | Availability and behavior |
|---|---|
| Lazy reconnect | All six: a lost connection reconnects on the next call; workers re-register after every connection generation |
| Worker retry backoff | All six retry a failed pull loop without spinning; TypeScript, Python, PHP, Go, and Rust use bounded backoff, while Elixir uses a short fixed delay |
| Producer fast-fail window | TypeScript and Python throttle repeated failed connection attempts (500 ms → 5 s) |
| TCP keepalive | TypeScript and Python enable an approximately 15 s idle probe where the OS exposes the controls |
| Timeout-driven teardown | All six discard a stream whose framing state is ambiguous; TypeScript/Python tolerate a configurable consecutive-timeout threshold, PHP/Go/Rust/Elixir tear down immediately |
| Backpressure | TypeScript optionally parks callers at maxInFlight; the synchronous clients naturally serialize/bound calls |
| Auth ordering | All six prevent any command racing ahead of Auth after a reconnect |
Producing jobs
Section titled “Producing jobs”Job options
Section titled “Job options”Every option is validated server-side; identical semantics in every SDK.
Naming follows each language’s idiom: TypeScript and PHP use camelCase keys
(attempts, jobId, removeOnComplete), Python uses snake_case
(attempts, job_id, remove_on_complete), Go takes a
bunqueue.JobOptions map with the camelCase keys, Rust uses typed
JobOptions fields, and Elixir accepts keyword or map options. Unknown options
raise an error in every SDK instead of being silently dropped.
| Option | Type | Default | Notes |
|---|---|---|---|
priority | number | 0 | Higher runs sooner; −1 000 000 … 1 000 000 |
delay | ms | 0 | Up to 1 year |
attempts | number | 3 | Max attempts including the first; up to 1000 |
backoff | number | { type, delay } | 1000 | type: 'fixed' or 'exponential'; delay up to 1 day |
ttl | ms | - | Expires the job if not processed in time |
timeout | ms | - | Per-job processing timeout, up to 1 day |
jobId | string | - | Custom id; idempotent, re-adding an unfinished id is a no-op |
deduplication | { id, ttl?, extend?, replace? } | - | Dedup window keyed on id |
debounce | { id, ttl? } | - | Collapses bursts to the last job |
dependsOn | string[] | - | Job ids that must complete first |
parentId / childrenIds | string / string[] | - | Flow relationships (usually set via FlowProducer) |
tags / groupId | string[] / string | - | Metadata; groupId also scopes group rate limits |
lifo | boolean | false | Orders only among LIFO jobs; mixed queues stay FIFO |
removeOnComplete / removeOnFail | boolean | false | Drop the job record at the terminal state |
durable | boolean | false | Bypass the server write buffer, fsync before the ACK returns |
repeat | { every } or { pattern, tz? } + limit? | - | Repeatable jobs (see Cron) |
stallTimeout | ms | - | Per-job stall detection override |
stackTraceLimit / keepLogs / sizeLimit | number | - | Failure stack cap, retained log lines, payload cap |
await queue.add('report', data, { priority: 10, delay: 5000, attempts: 5 });await queue.add('charge', payment, { jobId: `order-${orderId}`, durable: true });queue.add("report", data, priority=10, delay=5000, attempts=5)queue.add("charge", payment, job_id=f"order-{order_id}", durable=True)$queue->add('report', $data, ['priority' => 10, 'delay' => 5000, 'attempts' => 5]);$queue->add('charge', $payment, ['jobId' => "order-{$orderId}", 'durable' => true]);queue.Add("report", data, bunqueue.JobOptions{"priority": 10, "delay": 5000, "attempts": 5})queue.Add("charge", payment, bunqueue.JobOptions{"jobId": "order-" + orderID, "durable": true})Idempotency and bulk
Section titled “Idempotency and bulk”jobId makes add idempotent: re-adding an id whose job is still
unfinished (waiting, active, waiting-children) returns the existing job
instead of creating a duplicate. This holds for addBulk too, each bulk
entry’s jobId is preserved on the wire, so an idempotent batch ingest can
be re-run safely after a crash:
await queue.addBulk( orders.map((o) => ({ name: 'ingest', data: o, opts: { jobId: `order-${o.id}` } })));queue.add_bulk([ {"name": "ingest", "data": o, "job_id": f"order-{o['id']}"} for o in orders])$queue->addBulk(array_map( fn ($o) => ['name' => 'ingest', 'data' => $o, 'jobId' => "order-{$o['id']}"], $orders));entries := make([]bunqueue.BulkEntry, 0, len(orders))for _, o := range orders { entries = append(entries, bunqueue.BulkEntry{ Name: "ingest", Data: o, Opts: bunqueue.JobOptions{"jobId": "order-" + o.ID}, })}ids, err := queue.AddBulk(entries)Producer throughput
Section titled “Producer throughput”For high-volume producers the TypeScript SDK can fan commands across a round-robin connection pool, one line, no API change (TypeScript only):
const queue = new Queue('ingest', { poolSize: 4 });The pool is producer-side only: workers intentionally keep a single
connection, because job leases (lock tokens) and worker registration are
per-connection server state. In every SDK, addBulk is the first tool for
throughput: one round-trip for the whole batch.
Processing jobs
Section titled “Processing jobs”Worker options
Section titled “Worker options”| Option | Default | Notes |
|---|---|---|
concurrency | 4 in TypeScript/Python/Go/Rust; 1 in Elixir/PHP | Jobs processed in parallel except PHP, which is sequential by design |
batchSize | usually 10; Elixir defaults to concurrency | Jobs fetched per PULLB, capped by free slots and the server max (1000) |
pollTimeoutMs | usually 5000; Elixir 1000 | Server-side long-poll; max 30 000 |
lockTtlMs | 30 000 | Job lease TTL |
heartbeatIntervalS | 10; Go defaults to disabled | Worker + per-job lock heartbeats; 0 disables (Go: negative or DisableHeartbeat: true) |
ackBatch | off | Opt-in ACK batching (below; TypeScript and Python) |
autorun | true where exposed | TypeScript/Python can start at construction; PHP, Go, Rust, and Elixir start explicitly |
Names follow each language: Python, Rust, and Elixir use snake_case; PHP uses
camelCase array keys; Go uses a WorkerOptions struct
(PollTimeoutMs, LockTtlMs, …).
Lease model
Section titled “Lease model”A pulled job carries a lock token. The worker heartbeats every active
job’s lock on the heartbeat interval, so a job that legitimately runs
longer than the lock TTL survives. If the worker dies, the lease expires
and the server requeues the job (or moves it to the DLQ once maxStalls
is exceeded), at-least-once delivery, so make handlers idempotent.
PHP is the deliberate sequential exception: it can heartbeat between jobs,
but it cannot interrupt a running user callback. A PHP handler that can exceed
lockTtlMs must call $job->extendLock(...) from the callback or split the
work into shorter jobs.
Failures
Section titled “Failures”import { UnrecoverableError } from 'bunqueue-client';
const worker = new Worker('emails', async (job) => { if (!isValid(job.data)) { throw new UnrecoverableError('malformed payload'); // skip retries → DLQ } return await send(job.data);});from bunqueue import UnrecoverableError, Worker
def process(job): if not is_valid(job.data): raise UnrecoverableError("malformed payload") # skip retries -> DLQ return send(job.data)
worker = Worker("emails", process)use Bunqueue\UnrecoverableError;use Bunqueue\Worker;
$worker = new Worker('emails', function (Bunqueue\Job $job) { if (!isValid($job->data())) { throw new UnrecoverableError('malformed payload'); // skip retries -> DLQ } return send($job->data());});worker := bunqueue.NewWorker("emails", func(job *bunqueue.Job) (any, error) { if !isValid(job.Data()) { return nil, bunqueue.NewUnrecoverableError("malformed payload") // skip retries -> DLQ } return send(job.Data())}, bunqueue.WorkerOptions{})Panics are recovered, failed with their real stack, and never kill the worker.
A thrown error fails the job with its message and the leading stack lines
(persisted server-side, capped by stackTraceLimit); the server applies
the retry/backoff policy and eventually the dead letter queue.
UnrecoverableError bypasses retries entirely.
Worker events
Section titled “Worker events”TypeScript, Python, PHP, and Go expose worker lifecycle listeners. Rust and Elixir use their structured telemetry callback for transport/command lifecycle and normal language control flow for per-job handler outcomes.
worker.on('completed', (job, result) => log.info('done', job.id));worker.on('failed', (job, err) => log.warn('failed', job.id, err.message));worker.on('error', (err) => log.error(err)); // always attachworker.on("completed", lambda job, result: log.info("done %s", job.id))worker.on("failed", lambda job, err: log.warning("failed %s: %s", job.id, err))worker.on("error", lambda err: log.error(err))$worker->on('completed', fn ($job, $result) => $log->info("done {$job->id()}"));$worker->on('failed', fn ($job, $err) => $log->warning("failed {$job->id()}"));$worker->on('error', fn ($err) => $log->error($err->getMessage()));worker.On("completed", func(args ...any) { job := args[0].(*bunqueue.Job) log.Printf("done %s", job.ID())})worker.On("error", func(args ...any) { log.Println(args[0]) })In TypeScript, always attach an error listener: per Node
EventEmitter semantics an unhandled error event throws. The other SDKs
swallow listener exceptions. Every worker frees each job’s concurrency
slot before emitting, so a throwing listener cannot leak a slot or
degrade throughput in any SDK. The error itself is yours to observe.
ACK batching (high volume)
Section titled “ACK batching (high volume)”Opt-in: coalesce completed-job acknowledgements into ACKB round-trips.
Available in TypeScript and Python; the PHP and Go workers acknowledge
each job individually.
const worker = new Worker('ingest', process, { concurrency: 32, ackBatch: { enabled: true, maxSize: 50, maxDelayMs: 5 },});worker = Worker( "ingest", process, concurrency=32, ack_batch={"max_size": 50, "max_delay_ms": 5},)Semantics are strict: a job stays active, its lock still heartbeated,
until the server confirms the batch; the batch is flushed on close();
every job settles exactly once even if an event listener throws. Defaults
are off, so nothing changes unless you enable it.
Graceful shutdown
Section titled “Graceful shutdown”await worker.close(); // stop pulling, flush batched ACKs, drain in-flight jobsawait worker.close(true); // force: skip the in-flight drainworker.close() # stop pulling, wait for in-flight jobs to drainworker.close(timeout=5) # bound the wait; the drain continues in the backgroundPython has no force flag: close(timeout=...) bounds how long the call
waits, and close(timeout=0) detaches immediately while in-flight jobs
finish in the background.
$worker->installSignalHandlers(); // SIGTERM / SIGINT -> graceful stop$worker->stop(); // finish the in-flight job, then return from run()$worker->close(); // unregister and close the connectionThe PHP worker is sequential, so stopping waits for at most one in-flight job; the unprocessed rest of the batch is re-leased by the server.
worker.Stop() // stop pulling; in-flight jobs finishworker.Close() // unregister and close the connectionQuery, control and operations
Section titled “Query, control and operations”All SDKs implement the core producer, lookup, wait, queue-control, DLQ,
scheduler, rate-limit, and global-concurrency operations with identical server
semantics. TypeScript and PHP use camelCase, Python/Rust/Elixir use snake_case,
and Go uses exported Go style (GetJobCounts, RetryDlq, …). The broader
webhook/stats/metrics convenience surface remains richest in TypeScript,
Python, PHP, and Go; Rust and Elixir can always use the documented protocol or
HTTP operations until typed helpers are added.
| Area | Methods |
|---|---|
| Query | getJob, getJobByCustomId, getJobs + per-state helpers (getWaiting, getActive, getCompleted, getFailed, getDelayed, getPrioritized, getWaitingChildren), getJobState (Python name: get_state), getResult, getProgress, counts + counts-per-priority, getChildrenValues, job logs |
| Waiting | waitForJob(id, ttlMs), resolves with the result; throws on timeout (CommandTimeoutError) and on a failed job (CommandError), never silently returns undefined |
| Control | pause / resume / isPaused, drain, obliterate, clean, remove, discard, promoteJob(s), retryJob (failed → waiting), updateJobData, updateJobProgress, changeJobPriority, changeJobDelay, moveJobToWait/Delayed/Completed/Failed, extendJobLock |
| DLQ | getDlq, retryDlq(id?), purgeDlq, retryJobs (BullMQ semantics: retries the whole DLQ by default, or pass count to retry only the first N entries, honored by server 2.8.29 and later; Python retry_dlq accepts count= too), DLQ configuration |
| Schedulers | upsertJobScheduler (cron pattern or fixed interval, timezone, skipIfNoWorker, preventOverlap, per-scheduler job options), removeJobScheduler, getJobScheduler(s), addCron, every |
| Admin | rate limiting, global concurrency, stall configuration, webhooks, getStats, getMetrics, listQueues, getWorkers |
A representative slice, in every language:
const job = await queue.getJob(id); // null when missingconst state = await queue.getJobState(id); // 'waiting' | 'active' | ...const result = await queue.waitForJob(id, 30_000);const counts = await queue.getJobCounts(); // { waiting, active, completed, ... }await queue.pause();const dropped = await queue.drain(); // number of removed jobsawait queue.retryDlq(); // re-queue dead-lettered jobsjob = queue.get_job(id) # None when missingstate = queue.get_state(id) # "waiting" | "active" | ...result = queue.wait_for_job(id, timeout_ms=30000)counts = queue.get_job_counts() # {"waiting": ..., "active": ..., ...}queue.pause()dropped = queue.drain() # number of removed jobsqueue.retry_dlq() # re-queue dead-lettered jobs$job = $queue->getJob($id); // null when missing$state = $queue->getState($id); // 'waiting' | 'active' | ...$result = $queue->waitForJob($id, 30000);$counts = $queue->getJobCounts(); // ['waiting' => ..., 'active' => ..., ...]$queue->pause();$dropped = $queue->drain(); // number of removed jobs$queue->retryDlq(); // re-queue dead-lettered jobsjob, _ := queue.GetJob(id) // nil when missingstate, _ := queue.GetState(id) // "waiting" | "active" | ...result, _ := queue.WaitForJob(id, 30000)counts, _ := queue.GetJobCounts() // map[string]int{"waiting": ..., ...}_ = queue.Pause()dropped, _ := queue.Drain() // number of removed jobs_, _ = queue.RetryDlq("", 0) // re-queue dead-lettered jobslet job = queue.get_job(&id)?; // None when missinglet state = queue.get_state(&id)?;let result = queue.wait_for_job(&id, 30_000)?;let counts = queue.get_job_counts()?;queue.pause()?;let dropped = queue.drain()?;queue.retry_dlq(None, None)?;{:ok, job} = Bunqueue.Queue.get_job(queue, id) # nil when missing{:ok, state} = Bunqueue.Queue.get_state(queue, id){:ok, result} = Bunqueue.Queue.wait_for_job(queue, id, 30_000){:ok, counts} = Bunqueue.Queue.get_job_counts(queue):ok = Bunqueue.Queue.pause(queue){:ok, dropped} = Bunqueue.Queue.drain(queue){:ok, _count} = Bunqueue.Queue.retry_dlq(queue)Two behaviors worth knowing:
- Not-found is
null/None/nil, never an exception,getJob,getJobByCustomIdandgetJobSchedulermap the server’s not-found response for you. - Progress updates require an active job; the server rejects progress updates on waiting jobs by design.
Pipelines and parent/child trees with automatic ordering, children complete before their parent, results are readable from the parent:
import { FlowProducer } from 'bunqueue-client';
const flow = new FlowProducer();
// Sequential pipelineawait flow.addChain([ { name: 'extract', queueName: 'etl' }, { name: 'transform', queueName: 'etl' }, { name: 'load', queueName: 'etl' },]);
// Tree: parent waits for its childrenconst node = await flow.add({ name: 'assemble', queueName: 'orders', children: [ { name: 'reserve-stock', queueName: 'orders' }, { name: 'charge-card', queueName: 'orders' }, ],});
// Fan-in: N parallel jobs converge into oneawait flow.addBulkThen(parts, { name: 'merge', queueName: 'orders' });from bunqueue import FlowProducer
flow = FlowProducer()
# Sequential pipelineflow.add_chain([ {"name": "extract", "queueName": "etl"}, {"name": "transform", "queueName": "etl"}, {"name": "load", "queueName": "etl"},])
# Tree: parent waits for its childrennode = flow.add({ "name": "assemble", "queueName": "orders", "children": [ {"name": "reserve-stock", "queueName": "orders"}, {"name": "charge-card", "queueName": "orders"}, ],})
# Fan-in: N parallel jobs converge into oneflow.add_bulk_then(parts, {"name": "merge", "queueName": "orders"})use Bunqueue\FlowProducer;
$flow = new FlowProducer();
// Sequential pipeline$flow->addChain([ ['name' => 'extract', 'queueName' => 'etl'], ['name' => 'transform', 'queueName' => 'etl'], ['name' => 'load', 'queueName' => 'etl'],]);
// Tree: parent waits for its children$node = $flow->add([ 'name' => 'assemble', 'queueName' => 'orders', 'children' => [ ['name' => 'reserve-stock', 'queueName' => 'orders'], ['name' => 'charge-card', 'queueName' => 'orders'], ],]);flow := bunqueue.NewFlowProducer(bunqueue.Options{})defer flow.Close()
// Sequential pipelineids, err := flow.AddChain([]bunqueue.ChainStep{ {Name: "extract", QueueName: "etl"}, {Name: "transform", QueueName: "etl"}, {Name: "load", QueueName: "etl"},})
// Tree: parent waits for its childrennode, err := flow.Add(bunqueue.FlowJob{ Name: "assemble", QueueName: "orders", Children: []bunqueue.FlowJob{ {Name: "reserve-stock", QueueName: "orders"}, {Name: "charge-card", QueueName: "orders"}, },})use bunqueue_client::{ ChainStep, ConnectionOptions, FlowProducer, JobOptions, Value,};
let flow = FlowProducer::new(ConnectionOptions::default());let step = |name: &str| ChainStep { name: name.into(), queue_name: "etl".into(), data: Value::Nil, options: JobOptions::default(),};let ids = flow.add_chain(vec![ step("extract"), step("transform"), step("load"),])?;flow = Bunqueue.FlowProducer.new()
{:ok, ids} = Bunqueue.FlowProducer.add_chain(flow, [ %{name: "extract", queue: "etl"}, %{name: "transform", queue: "etl"}, %{name: "load", queue: "etl"} ])Fan-in (addBulkThen, N parallel jobs converging into one final job) is
available in TypeScript and Python.
Creation is transactional per call: if any push fails, every job already
created by that call is cancelled (rollback). getFlow(id) reconstructs a
tree from a root id, it returns null for a missing root, skips children
that were since removed (removeOnComplete), and is cycle-safe, so it
always returns the surviving partial tree instead of throwing.
Observability
Section titled “Observability”All six SDKs expose opt-in, dependency-free structured telemetry. Bring
OpenTelemetry, Prometheus, Logger, tracing, or your own collector; SDKs stay
silent by default. Consumer callback exceptions and Rust callback panics are
isolated so an observer can never break transport or worker correctness.
import { Queue, consoleLogger, type TelemetryEvent } from 'bunqueue-client';
const queue = new Queue('emails', { logger: consoleLogger('info'), // or any { debug, info, warn, error } onTelemetry: (e: TelemetryEvent) => { if (e.type === 'command') histogram.observe({ cmd: e.cmd }, e.durationMs); if (e.type === 'reconnect_scheduled') reconnects.inc(); },});
// Connection is an EventEmitter for imperative lifecycle hooks:queue.connection.on('connect', (i) => log.info('link up', i));queue.connection.on('disconnect', (i) => log.warn('link down', i));queue.connection.on('reconnect_scheduled', (i) => log.warn('retrying', i));queue = Queue("emails", on_telemetry=lambda event: metrics.observe(event))$queue = new Queue('emails', [ 'onEvent' => fn (array $event) => $metrics->observe($event),]);queue := bunqueue.NewQueue("emails", bunqueue.Options{ OnEvent: func(event bunqueue.TelemetryEvent) { metrics.Observe(event) },})use std::sync::Arc;use bunqueue_client::{ConnectionOptions, Queue, TelemetryCallback};
let telemetry: TelemetryCallback = Arc::new(|event| tracing::debug!(?event));let queue = Queue::new("emails", ConnectionOptions { telemetry: Some(telemetry), ..Default::default()});queue = Bunqueue.queue("emails", event_handler: fn event -> Logger.info("bunqueue", bunqueue: event) end )The idiomatic event sets cover connection/reconnection, authentication,
command latency and outcome, timeout, transport error, and close as applicable
to each runtime; worker retry events are added where the client has a retry
loop. Tokens, job payloads, results, private keys, and CA contents are never
recorded. Scrape the server’s HTTP /metrics endpoint for authoritative queue
and broker metrics.
Simple Mode
Section titled “Simple Mode”Bunqueue bundles a Queue and a Worker into a single object, a 1:1 port
of the official client’s Simple Mode. It brings
routes, onion middleware, in-process retry strategies, a circuit breaker,
batch accumulation, event triggers, job TTL, priority aging, cooperative
cancellation, and deduplication or debounce defaults. Available in
TypeScript and Python; in PHP and Go, compose Queue and Worker
directly.
import { Bunqueue } from 'bunqueue-client';
const app = new Bunqueue('notifications', { routes: { 'send-email': async (job) => ({ sent: true }), 'send-sms': async (job) => ({ sent: true }), }, concurrency: 10, retry: { maxAttempts: 5, strategy: 'jitter' }, circuitBreaker: { threshold: 5, resetTimeout: 30_000 },});
app.use(async (job, next) => { console.time(job.name); const result = await next(); console.timeEnd(job.name); return result;});
await app.add('send-email', { to: 'alice@example.com' });await app.cron('daily-digest', '0 9 * * *', { to: 'all' });from bunqueue import Bunqueue
app = Bunqueue( "notifications", routes={ "send-email": lambda job: {"sent": True}, "send-sms": lambda job: {"sent": True}, }, concurrency=10, retry={"max_attempts": 5, "strategy": "jitter"}, circuit_breaker={"threshold": 5, "reset_timeout": 30000},)
def timing(job, next_fn): result = next_fn() print(f"{job.name} done") return result
app.use(timing)
app.add("send-email", {"to": "alice@example.com"})app.cron("daily-digest", "0 9 * * *", {"to": "all"})Scheduler job options are honored end to end: the template’s attempts,
backoff, timeout, delay, remove_on_complete / remove_on_fail,
priority and deduplication are applied to every cron-spawned job, in
every SDK (scheduler support included: upsertJobScheduler exists in PHP
and Go too).
One difference from the Bun original: embedded: true is not available,
it requires the in-process Bun runtime. Connect to a server instead; the
option raises a clear error if you try.
Cloudflare Workers
Section titled “Cloudflare Workers”The same client runs inside Workers. Enable Node.js compatibility and add jobs directly from your fetch handlers:
compatibility_flags = ["nodejs_compat"]compatibility_date = "2025-01-01"import { Queue } from 'bunqueue-client';
export default { async fetch(req: Request, env: Env): Promise<Response> { const queue = new Queue('signups', { host: env.BQ_HOST, port: 6789, token: env.BQ_TOKEN, tls: true }); const job = await queue.add('welcome', await req.json()); queue.close(); return Response.json({ queued: job.id }); },};Workers are request-scoped, so there is no long-lived worker loop. Instead: produce from fetch handlers, and consume in batches from a Cron Trigger, pull, process, acknowledge, return. Both patterns work with Simple Mode too. Two requirements apply: the server must be reachable from the internet, and TLS needs a publicly trusted certificate.
Security
Section titled “Security”Authentication uses server-side tokens; transport security uses native TLS.
const queue = new Queue('emails', { host: 'queue.example.com', port: 6789, token: process.env.BUNQUEUE_TOKEN, // server started with AUTH_TOKENS=... tls: true, // or { caFile: './ca.pem' } for a custom CA});queue = Queue( "emails", host="queue.example.com", port=6789, token=os.environ["BUNQUEUE_TOKEN"], tls={"ca_file": "./ca.pem"}, # or True for system CAs, or an ssl.SSLContext)$queue = new Queue('emails', [ 'host' => 'queue.example.com', 'port' => 6789, 'token' => getenv('BUNQUEUE_TOKEN'), // server started with AUTH_TOKENS=... 'tls' => ['caFile' => './ca.pem'], // or true for system CAs]);queue := bunqueue.NewQueue("emails", bunqueue.Options{ Host: "queue.example.com", Port: 6789, Token: os.Getenv("BUNQUEUE_TOKEN"), // server started with AUTH_TOKENS=... TLS: &bunqueue.TLSOptions{CAFile: "./ca.pem"}, // or &TLSOptions{} for system CAs})use std::path::PathBuf;use bunqueue_client::{ConnectionOptions, Queue, TlsOptions};
let queue = Queue::new("emails", ConnectionOptions { host: "queue.example.com".into(), token: std::env::var("BUNQUEUE_TOKEN").ok(), tls: Some(TlsOptions { ca_file: Some(PathBuf::from("./ca.pem")) }), ..Default::default()});queue = Bunqueue.queue("emails", host: "queue.example.com", token: System.fetch_env!("BUNQUEUE_TOKEN"), tls: true, ca_file: "./ca.pem" )Certificate verification is on by default for every TLS connection in
every SDK, a wrong or missing CA rejects the connection instead of
silently connecting; only an explicit opt-out
(rejectUnauthorized: false in TypeScript, {"verify": False} in Python,
['verifyPeer' => false] in PHP, InsecureSkipVerify: true in Go,
verify: false in Elixir) switches to encryption-only mode for development.
Rust deliberately exposes no insecure TLS mode. Start the server with
AUTH_TOKENS to require authentication, and with TLS_CERT_FILE plus
TLS_KEY_FILE for encrypted transport. Full hardening guidance lives in the
deployment guide.
Errors and guarantees
Section titled “Errors and guarantees”Every SDK exposes typed equivalents of the same error categories, so retry
logic can branch precisely (PHP uses exception classes under
Bunqueue\Exception; Go uses errors.As; Rust uses the Error enum; Elixir
returns {:error, exception} or raises through bang APIs):
| Error | Meaning |
|---|---|
| Connection | Link lost or server unreachable (in-flight commands reject; the next call reconnects) |
| Timeout | No response within the command timeout, including waitForJob timeout |
| Command | The server answered ok: false (validation, failed wait, or a write-side not-found) |
| Authentication | Token rejected during the connection handshake |
| Protocol / serialization | Invalid map, frame, MessagePack extension, or oversized outgoing body |
| Unrecoverable processing | Raised/returned by your processor to skip retries and fail terminally |
Delivery is at-least-once: a worker crash after processing but before
the ACK means the job runs again, design handlers to be idempotent (the
jobId option helps on the producing side). Two numeric-precision rules:
- JavaScript numbers are IEEE 754 doubles, exact up to 2⁵³. Pass larger
64-bit identifiers (snowflake ids) as strings; never put
BigIntin job data. - Python, PHP, Go, Rust, and Elixir integers outside the int32 range are recursively encoded as float64 on the wire (exact up to 2⁵³, safe for millisecond timestamps); the same string rule applies to larger identifiers.
Write a client in any language
Section titled “Write a client in any language”The wire protocol is small and fully documented: length-prefixed msgpack frames, one command per message, plain request/response. If your language is not covered yet, you can build a client against the wire protocol specification and certify it with the conformance suite: point its runner at your client and it tells you exactly what to fix. The official SDKs are built and certified the same way.
Two environments are out of scope today. Browsers have no raw-TCP API, so route through your own backend instead of connecting directly. WebAssembly is feasible but not one portable runtime: a WASM target becomes official only when it can pass the same conformance suite without weakening authentication, TLS, or framing.
Resources
Section titled “Resources”bunqueue-clienton npm, the full step by step README, from server start to production- Changelogs: TypeScript · Python · PHP · Go · Rust · Elixir
- Queue API, every job option explained
- Worker, concurrency, events, graceful shutdown
- Simple Mode, everything in one object
- Flows, pipelines, parent and child jobs, fan in
- Cron Jobs, schedules and repeatable jobs
- Dead Letter Queue, what happens when jobs fail
- Deployment guide, Docker, TLS, authentication, monitoring
- Wire protocol specification, the normative contract every SDK implements
- Conformance suite, certify a client in any language
- SDK sources:
sdk/typescript·sdk/python·sdk/php·sdk/go·sdk/rust·sdk/elixir