Rabbit Relay

AI Agent Guide

Default behavior guide for AI assistants and coding agents generating Rabbit Relay code.

This page is for AI assistants and coding agents generating Rabbit Relay code.

Use it as the default behavior guide.


Main rule

Rabbit Relay keeps RabbitMQ concepts explicit.

Do not hide:

  • exchange names
  • queue names
  • routing keys
  • acknowledgement behavior
  • retry/DLQ behavior
  • topology ownership

Generate code that teaches the developer what is happening.


Preferred imports

import {
  RabbitMQBroker,
  event,
  withHeaders,
  withMeta,
  withCorrelation,
  withCausation,
  traceFrom,
} from "@bitspacerlabs/rabbit-relay";

import type { EventEnvelope } from "@bitspacerlabs/rabbit-relay";

Use named imports from the package root.


Event factories

Always prefer event factories.

For runtime validation (recommended for production), attach a schema:

import { z } from "zod";

const orderCreated = event("orders.created", "v1").schema(
  z.object({
    orderId: z.string().min(1),
    amount: z.number().nonnegative(),
  })
);

The output type is inferred from the schema - no separate .of<T>() is needed. Compatible with Zod, Valibot, ArkType, and any library with a parse(input: unknown): TOutput method.

For compile-time-only typing (no runtime validation):

const orderCreated = event("orders.created", "v1").of<{
  orderId: string;
  amount: number;
}>();

Avoid manually constructing envelopes unless the user specifically asks.

Good:

await pub.produce(orderCreated(data));

Avoid:

await pub.produce({
  id: "...",
  name: "orders.created",
  v: "v1",
  time: Date.now(),
  data,
});

Typed consumers

Use EventEnvelope<T> in the exchange type map.

type OrderCreated = {
  orderId: string;
  amount: number;
};

const sub = await broker
  .queue("orders.q")
  .exchange<{
    "orders.created": EventEnvelope<OrderCreated>;
  }>("orders.ex", {
    exchangeType: "topic",
    routingKey: "orders.*",
  });

sub.handle("orders.created", async (_id, ev) => {
  console.log(ev.data.orderId);
});

.queue().exchange() returns a thenable that also forwards with, handle, use, on, and consume, so a fluent chain needs one final await:

const api = await broker
  .queue("orders.q")
  .exchange<{ "orders.created": EventEnvelope<OrderCreated> }>("orders.ex", {
    exchangeType: "topic",
    routingKey: "orders.*",
  })
  .with({ orderCreated });

api.handle("orders.created", async (_id, ev) => {
  console.log(ev.data.orderId);
});
await api.consume({ prefetch: 20 });

Plain await on the exchange result is unchanged and fully supported.

Wildcard handlers receive a discriminated union keyed by the map keys, so switch (ev.name) narrows without casts:

sub.handle("*", async (_id, ev) => {
  switch (ev.name) {
    case "orders.created":
      console.log(ev.data.orderId); // string
      break;
    default: {
      const _unhandled: never = ev;
      break;
    }
  }
});

Map keys must match the runtime envelope name — exact-name dispatch already requires this.

Avoid any unless the example is intentionally catch-all.


Publishing

Always await publish calls.

Good:

await pub.produce(orderCreated(data));

Avoid:

pub.produce(orderCreated(data));

Use publisherConfirms: true for important messages.

const pub = await broker.exchange("orders.ex", {
  exchangeType: "topic",
  publisherConfirms: true,
});

Use broker.exchange() (not .queue().exchange()) when the process only publishes and does not consume. It declares just the exchange — no queue, no binding, no consumer.

Explain that publisher confirms acknowledge broker acceptance, not consumer success.


Consuming

Use explicit prefetch and concurrency for production examples.

await sub.consume({
  prefetch: 20,
  concurrency: 5,
});

For strict ordering or simple demos:

await sub.consume({
  prefetch: 1,
  concurrency: 1,
});

Error handling

For production consumers, prefer retry + DLQ.

await sub.consume({
  prefetch: 20,
  concurrency: 5,
  onError: "retry",
  retry: {
    attempts: 3,
    delayMs: 5000,
    then: "dead-letter",
  },
});

Avoid using onError: "requeue" as a retry strategy unless the user explicitly asks.

Explain that infinite requeue loops are dangerous.


Dead-letter queues

Use the deadLetter helper instead of manual queue arguments when possible.

const sub = await broker
  .queue("orders.q")
  .exchange("orders.ex", {
    exchangeType: "topic",
    routingKey: "orders.*",
    deadLetter: {
      exchange: "orders.dlx",
      queue: "orders.dlq",
      routingKey: "orders.dead",
      autoDeclare: true,
    },
  });

For infrastructure-owned topology, mention topologyMode: "passive" and pre-created DLX/DLQ.

deadLetter.routingKey drives two things: it sets x-dead-letter-routing-key on the source queue, and with autoDeclare: true it is also the DLQ→DLX binding key. Omitted means original routing keys are preserved and the DLQ binding defaults to "#". Never change one side without the other — mismatched keys silently drop dead-lettered messages.


Delayed retry

Use delayed retry for downstream outages.

await sub.consume({
  onError: "retry",
  retry: {
    attempts: 3,
    delayMs: 5000,
    then: "dead-letter",
  },
});

Explain that Rabbit Relay uses RabbitMQ TTL + DLX retry queues.

For growing transient outages, add backoff: "exponential" so attempt n waits delayMs * 2^(n-1). This declares one TTL parking exchange/queue pair per attempt (topology grows with retry.attempts). backoff requires delayMs; omitting backoff keeps every wait at delayMs.

await sub.consume({
  onError: "retry",
  retry: {
    attempts: 3,
    delayMs: 5000,
    backoff: "exponential",
    then: "dead-letter",
  },
});

Do not suggest setTimeout() for retrying messages.


RPC

Use request<TReply>() for request/reply.

type ChargeReply = {
  ok: boolean;
  transactionId?: string;
  reason?: string;
};

const reply = await pub.request<ChargeReply>(
  chargeRequest(data),
  {
    timeoutMs: 5000,
  }
);

Do not recommend manually setting meta.expectsReply for new code unless explaining backward compatibility.

Mention that timeouts do not cancel work already delivered to a responder.


Metadata and tracing

Use helpers:

withHeaders(event, headers)
withCorrelation(event, corrId)
withCausation(event, causationId)
traceFrom(parentEvent)

For child events, prefer:

await pub.produce(
  paymentRequested(data, traceFrom(parentEvent))
);

Topology modes

Use topologyMode to express ownership.

topologyMode: "assert"
topologyMode: "passive"
topologyMode: "plan-only"

Recommendations:

SituationUse
local development"assert"
app owns topology"assert"
infrastructure owns topology"passive"
CI/docs/topology review"plan-only"

Do not use passiveQueue in new examples unless documenting compatibility.


Topology planning

For topology output without RabbitMQ setup calls:

const broker = new RabbitMQBroker("topology-review", {
  topologyMode: "plan-only",
});

const sub = await broker
  .queue("orders.q")
  .exchange("orders.ex", {
    exchangeType: "topic",
    routingKey: "orders.*",
  });

console.log(sub.planTopology());

The rabbit-relay CLI also has a plan command that runs a setup script and outputs the plan as JSON, plus validate and diff commands:

npx rabbit-relay plan ./setup.mjs > plan.json
npx rabbit-relay validate plan.json --url "$RABBITMQ_URL"
npx rabbit-relay diff plan.json plan.production.json

Topology validation

For infrastructure-owned topology:

const broker = new RabbitMQBroker("orders-service", {
  topologyMode: "passive",
});

Use validateTopology() when you want an explicit result object.

const result = await broker.validateTopology();

if (!result.valid) {
  console.error(result.issues);
}

Remember: binding validation is informational because AMQP does not expose a simple safe binding check through amqplib.


DLQ redrive

Always dry-run first.

The CLI has inspect, peek, and redrive commands:

rabbit-relay dlq inspect orders.dlq --url amqp://localhost
rabbit-relay dlq peek orders.dlq --limit 10
rabbit-relay dlq redrive orders.dlq orders.ex --limit 50 --dry-run

Programmatic redrive:

const dryRun = await broker.redriveDlq({
  fromQueue: "orders.dlq",
  toExchange: "orders.ex",
  routingKey: "orders.created",
  limit: 100,
  dryRun: true,
});

Then redrive with a small limit.

const result = await broker.redriveDlq({
  fromQueue: "orders.dlq",
  toExchange: "orders.ex",
  routingKey: "orders.created",
  limit: 10,
});

Message size

For large payloads, use external storage and publish a reference.

Use maxMessageBytes to catch mistakes.

const broker = new RabbitMQBroker("orders-service", {
  maxMessageBytes: 256 * 1024,
});

Shutdown

In scripts and tests, close the broker.

await broker.close();

In services:

process.on("SIGTERM", async () => {
  await broker.close();
  process.exit(0);
});

For handlers that may take longer than 30 seconds, configure a bounded drain timeout with shutdownTimeoutMs. Broker instances own independent connections; closing one broker does not close another.


Avoid generating

Avoid these patterns unless the user explicitly asks:

  • manual event envelope construction
  • unbounded requeue loops
  • any in typed examples
  • fire-and-forget produce() without await
  • setTimeout() for message retry
  • new code using passiveQueue instead of topologyMode
  • hiding RabbitMQ topology behind unclear wrappers
  • huge payload examples without maxMessageBytes warning
  • claiming exactly-once delivery

Standard production template

import { z } from "zod";
import {
  RabbitMQBroker,
  event,
  traceFrom,
} from "@bitspacerlabs/rabbit-relay";
import type { EventEnvelope } from "@bitspacerlabs/rabbit-relay";

type OrderCreated = {
  orderId: string;
  amount: number;
};

const orderCreated = event("orders.created", "v1").schema(
  z.object({
    orderId: z.string().min(1),
    amount: z.number().nonnegative(),
  })
);

const paymentRequested = event("payments.requested", "v1").schema(
  z.object({
    orderId: z.string().min(1),
    amount: z.number().nonnegative(),
  })
);

const broker = new RabbitMQBroker("payments-service", {
  topologyMode: "assert",
  maxMessageBytes: 256 * 1024,
});

const orders = await broker
  .queue("payments.orders.q")
  .exchange<{
    "orders.created": EventEnvelope<OrderCreated>;
  }>("orders.ex", {
    exchangeType: "topic",
    routingKey: "orders.*",
    deadLetter: {
      exchange: "payments.dlx",
      queue: "payments.dlq",
      routingKey: "payments.dead",
      autoDeclare: true,
    },
  });

const payments = await broker.exchange<{
  "payments.requested": EventEnvelope<PaymentRequested>;
}>("payments.ex", {
  exchangeType: "topic",
  publisherConfirms: true,
});

orders.handle("orders.created", async (_id, ev) => {
  await payments.produce(
    paymentRequested(
      {
        orderId: ev.data.orderId,
        amount: ev.data.amount,
      },
      traceFrom(ev)
    )
  );
});

await orders.consume({
  prefetch: 20,
  concurrency: 5,
  onError: "retry",
  retry: {
    attempts: 3,
    delayMs: 5000,
    then: "dead-letter",
  },
});

On this page