RabbitMQ prefetch: set it to your Sume concurrency_limit

Set channel.prefetch to the workspace concurrency_limit so each RabbitMQ message holds one Sume video job until it is terminal. Node sample.

5 min readSume
All posts

Set channel.prefetch(n) to the concurrency_limit that Sume reports for your workspace (4 on Pro, 8 on Startup, 20 on Scale). Keep each message unacked until its video job is terminal. Then RabbitMQ never hands a consumer more jobs than your workspace can process at once, and you do not rely on Sume's queue to absorb the excess.

Sume would accept the extra work anyway. Generation admission is queue-first: a valid submit becomes queued when every processing seat is busy, and only a full queue returns 429 queue_full. Prefetch is therefore pacing on your side. It keeps a backlog of 10,000 broker messages from turning into 10,000 paid, reserved jobs.

Where the number comes from

The default processing concurrency comes from the plan, and top-ups do not raise it. Prefer the live generation_limits.concurrency_limit field on a submit response, because an admin override can change it.

Default processing concurrency and queue capacity by plan (Sume docs, read 2026-10-04)
PlanProcessing concurrencyQueue capacityAccepted jobs
Free156
Pro42024
Startup84048
Scale20100120

A consumer that holds one message per job

The consumer below submits to /v1/videos with the message id as Idempotency-Key, then polls the polling_url every 30 seconds, which is the interval the video docs use. It acks only after a terminal status. Set SUME_CONCURRENCY from your plan.

import amqp from "amqplib";
const key = process.env.SUME_API_KEY;
if (!key) throw new Error("SUME_API_KEY is empty");
const h = { Authorization: `Bearer ${key}`, "Content-Type": "application/json" };
const sleep = (s) => new Promise((r) => setTimeout(r, s * 1000));
const conn = await amqp.connect(process.env.AMQP_URL ?? "amqp://localhost");
const ch = await conn.createChannel();
await ch.assertQueue("renders", { durable: true });
await ch.prefetch(Number(process.env.SUME_CONCURRENCY ?? 4));
await ch.consume("renders", async (msg) => {
  const job = JSON.parse(msg.content.toString());
  const res = await fetch("https://api.sume.com/v1/videos", {
    method: "POST",
    headers: { ...h, "Idempotency-Key": job.id },
    body: JSON.stringify({ model: job.model, prompt: job.prompt, duration: job.duration }),
  });
  if (res.status === 429 || res.status >= 500) {
    await sleep(Number(res.headers.get("retry-after") ?? 10));
    return ch.nack(msg, false, true);
  }
  if (!res.ok) return ch.nack(msg, false, false);
  const { polling_url } = await res.json();
  let s = "pending";
  while (!["completed", "failed", "cancelled"].includes(s)) {
    await sleep(30);
    s = (await (await fetch(polling_url, { headers: h })).json()).status;
  }
  ch.ack(msg);
});

What each response does to the message

Decide per outcome whether the message returns to the queue. A requeue without a delay turns a 429 into a hot loop, which is why the sample sleeps first.

Submit outcomes and message handling (Sume docs, read 2026-10-04)
ResponseMeaningMessage
202 with polling_urlJob acceptedKeep unacked, poll, ack at a terminal status
429 queue_fullAccepted capacity is used upRequeue after a delay; same key is safe
429 rate_limitedRequest volume limitWait for retry-after, requeue
402 insufficient_creditsBalance cannot cover the reservationDead-letter; a retry cannot succeed
409 idempotency_conflictSame key, different payloadDead-letter; this is a bug in your key scheme
503Sume could not dispatch safelyRequeue with the same key

Caveats

  • The limit is per workspace, not per key. Two services sharing one workspace share the seats, so divide prefetch between them.
  • Several consumers each hold their own prefetch window. Total in flight is consumers times prefetch.
  • RabbitMQ enforces a consumer acknowledgement timeout. Keep your own client deadline for a render (the jobs docs suggest 20 minutes for video) well inside it, or ack after submit and finish with a webhook.
  • A failed job is still a terminal job. Ack it, record the error, and decide separately whether to enqueue a new request with a new idempotency key.

Sources

Related posts

More in Developers

All Developers posts

Written by Sume