Peren documentation
Queues
Produce and consume asynchronous work through a queue binding backed by memory, file, cell, or an external broker.
A queue binding lets a Worker send messages and lets Peren deliver them to a queue handler in batches. You declare the binding with type = "queue". The broker that stores and leases messages lives under [queues].broker.
Queue send and ack

Produce messages
env.JOBS is a producer. send enqueues one body. sendBatch enqueues many.
export default {
async fetch(request, env) {
await env.JOBS.send({ path: new URL(request.url).pathname });
await env.JOBS.sendBatch([
{ body: { kind: "welcome", user: "user_123" } },
{ body: { kind: "digest", user: "user_123" }, delaySeconds: 60 },
]);
return new Response("queued", { status: 202 });
},
};
Bodies may be objects, strings, Uint8Array, or ArrayBuffer. Objects are JSON-encoded. Optional fields on send and on each sendBatch entry are contentType, delaySeconds, and dedupId.
Peren refuses oversized work in the isolate before the broker sees it:
- one message body may be at most 128000 bytes;
sendBatchaccepts at most 100 messages;- the sum of message bodies in one
sendBatchmay be at most 256000 bytes; delaySecondsmust be between 0 and 86400 inclusive.
A body over the message limit throws RangeError with the text Queue message body exceeds 128000 bytes.
Consume batches
List the queue under consumes_queues on the same service. Peren leases a batch and calls queue on the Worker.
export default {
async queue(batch) {
for (const message of batch.messages) {
const job = JSON.parse(message.body);
if (job.path === "/fail") {
message.retry({ delaySeconds: 30 });
continue;
}
message.ack();
}
},
};
Each message exposes id, body, attempts, timestamp, ack(), and retry(options). body is a UTF-8 string. batch.ackAll() and batch.retryAll(options) apply the same disposition to every message in the batch. batch.metrics reports ready, delayed, leased, and oldestReadyTimestamp.
If a leased message has no disposition after the handler returns, Peren acknowledges it.
Consumer batching and retry policy come from consumes_queues and optional [queues.consumer_defaults]. Defaults are max_batch_size = 10, max_batch_timeout_secs = 5, max_retries = 3, max_concurrency = 1, and retry_delay_secs = 0. A consumer may set dead_letter_queue to move messages that exhaust retries.
Smallest working example
Use the memory broker for a single process. Messages exist only in that process and disappear when it exits.
fleet.toml:
[node]
node_id = "00000000-0000-0000-0000-000000000001"
advertise_addr = "127.0.0.1:7000"
listen = "127.0.0.1:7000"
[bucket]
kind = "memory"
[mtls]
ca_cert_path = "./certs/ca.pem"
leaf_cert_path = "./certs/leaf-cert.pem"
leaf_key_path = "./certs/leaf-key.pem"
[queues]
broker = "memory"
[[services]]
name = "api"
worker_bundle_path = "worker.js"
compatibility_date = "2026-01-01"
consumes_queues = [{ queue = "jobs", max_batch_size = 10, max_retries = 3, dead_letter_queue = "dead" }]
[services.bindings.JOBS]
type = "queue"
queue_name = "jobs"
[[sockets]]
name = "public"
listen = "127.0.0.1:8080"
service = "api"
worker.js:
export default {
async fetch(request, env) {
await env.JOBS.send({ path: new URL(request.url).pathname });
return Response.json({ queued: true, provider: env.JOBS.provider.kind });
},
async queue(batch) {
batch.ackAll();
},
};
Create development certificates with peren devcert ./certs, then run peren serve fleet.toml. Send work with curl -i http://127.0.0.1:8080/demo. The Worker returns JSON that includes "queued": true and "provider": "memory". The consumer leases the message and acknowledges the batch.
Broker scope
| Broker | Scope | Required fields |
|---|---|---|
memory |
Process-local. Lost on exit. | None beyond broker = "memory". |
file |
Node-local file on this host. | file_path |
cell |
Shared through cell storage. Survives process restart when the cell store remains. | Optional cell_path. When omitted, Peren uses data/queues/cell.sqlite under the node data directory. |
nats |
External NATS JetStream broker. | nats_url |
rabbitmq |
External AMQP broker. | amqp_url |
kafka |
External Kafka cluster. | kafka_bootstrap_servers |
Config validation refuses a selected broker when its connection field is missing. The problem names queues connection and states that the field is required for the selected provider.
File broker
Persist the queue on one host:
[queues]
broker = "file"
file_path = "./data/queues.db"
Messages remain across process restarts on that host. They are not shared with other nodes.
Cell broker
Store the queue in cell storage:
[queues]
broker = "cell"
cell_path = "./data/queues/cell.sqlite"
Omit cell_path to use the default path under the node data directory. Cell-backed messages are shared through that store rather than remaining in one process heap.
NATS
[queues]
broker = "nats"
nats_url = "nats://127.0.0.1:4222"
The binding and consumer configuration stay the same. Peren opens the NATS URL from [queues].
RabbitMQ
[queues]
broker = "rabbitmq"
amqp_url = "amqp://guest:[email protected]:5672/%2f"
Supply the AMQP URL your broker accepts. Keep credentials out of committed files when they are not local development defaults.
Kafka
[queues]
broker = "kafka"
kafka_bootstrap_servers = "127.0.0.1:9092"
kafka_bootstrap_servers is the bootstrap address list Peren passes to the Kafka client.
Failure: message too large
await env.JOBS.send("x".repeat(128001));
The isolate throws RangeError: Queue message body exceeds 128000 bytes. The broker does not receive the message. Shrink the body or split work across multiple messages before retrying.
Related: Operate queues, Bindings overview.