Browse documentation
On this page

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

Queue send and ack
send, then ack or retry. PRODUCER BROKER CONSUMER producer send memory broker process-local queue handler ack retry

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;
  • sendBatch accepts at most 100 messages;
  • the sum of message bodies in one sendBatch may be at most 256000 bytes;
  • delaySeconds must 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.