Skip to content
Files SDK
Esc
↑↓navigate↵open⌘Jpreview
On this page

Process R2 uploads with Queues and idempotent handlers

A Queue consumer Worker that runs work for every file R2 stores, survives duplicate and out-of-order notifications, and replays what failed from a dead-letter queue.

R2 can send a message to a Cloudflare Queue whenever an object is created or deleted, so work runs for every file that lands in the bucket, even when the browser uploaded it directly and closed the tab before telling your server. A consumer Worker passes each message to files.events.dispatch(), which normalizes it into a FileEvent and runs your handlers. The consumer calls message.ack() only after the work has committed, and message.retry() with a growing delay when it hasn’t.

Queues delivers at least once and in no guaranteed order. Correctness therefore comes from the handler, not the queue. This guide’s handler checks D1 for work it has already done. It re-reads the object to skip events that a later write has overtaken, and it makes its one write conditional on the event’s time. Messages that exhaust their retries go to a dead-letter queue, which the same Worker drains into a D1 table you can replay from.

Before you start

  • A Cloudflare account with an R2 bucket (this guide calls it uploads) and access to Queues and D1. Uploads can reach the bucket any way you like. Use R2 in Cloudflare Workers and the Next.js uploader both store files under users/<id>/, and this guide processes that prefix.
  • A Worker project with Wrangler. The examples type env by hand with @cloudflare/workers-types; npx wrangler types can generate the same Env interface.
  • Written against files-sdk 3.0, Wrangler 4.148, and the 2026-10-01 compatibility date. The behavior in What happened under wrangler dev comes from running this Worker locally against Wrangler’s Queues, R2, and D1 simulation, with notifications sent by hand. It was not run against a deployed bucket.
npm install files-sdk
pnpm add files-sdk
yarn add files-sdk
bun add files-sdk
nub add files-sdk
aube add files-sdk

How a notification turns into work

  1. An object lands under users/, from the gateway, a presigned PUT, the dashboard, or another service.
  2. R2 matches a notification rule and sends a message to the r2-events queue. The body names the action, the bucket, the key, the size and ETag (creates only), and the time.
  3. Queues hands the consumer Worker a batch. For each message, dispatch() parses the body, drops events for other buckets, and runs the handler that matches the key.
  4. A message whose handler threw is retried with a delay. After max_retries it moves to r2-events-dlq.
  5. The same Worker consumes r2-events-dlq and parks each message in a D1 table. After you fix the cause, a replay sends the parked bodies back to r2-events.

R2 notifications carry no event ID. files-sdk/events builds event.id from the action type, bucket, key, event time, and ETag, for example created:uploads/users/42/report.pdf@2024-05-24T19:36:44.379Z#c846ff7a18f28c2e262116d6e8719ef0. A delete has no ETag, so its ID ends at the time. The same body always produces the same ID, but two writes of identical bytes at different times produce different ones.

Create the queues, database, and notification rule

npx wrangler queues create r2-events
npx wrangler queues create r2-events-dlq
npx wrangler d1 create uploads-db

npx wrangler r2 bucket notification create uploads \
  --event-type object-create --event-type object-delete \
  --prefix users/ \
  --queue r2-events \
  --description "users/ uploads to the processing consumer"

object-create fires for PutObject, CopyObject, and CompleteMultipartUpload. object-delete fires for DeleteObject and LifecycleDeletion. All five arrive at your handlers as created or deleted.

The prefix does two jobs. It limits the rule to user uploads, and it keeps anything the consumer writes back into the bucket, such as thumbnails under thumbnails/, from triggering the consumer again. --suffix narrows further, for example --suffix .pdf, and when you pass both, R2 applies both. Neither accepts a regular expression.

R2’s documentation caps a bucket at 100 rules and rejects rules that overlap, so one write can’t produce two notifications through two rules. One queue takes up to 5,000 messages per second; split rules across queues above that.

Configure the consumer

{
  "name": "uploads-consumer",
  "main": "src/index.ts",
  "compatibility_date": "2026-10-01",
  "r2_buckets": [{ "binding": "UPLOADS", "bucket_name": "uploads" }],
  "d1_databases": [
    {
      "binding": "DB",
      "database_name": "uploads-db",
      "database_id": "<id from wrangler d1 create>",
    },
  ],
  "queues": {
    // Only replays send to the queue. R2 sends the notifications itself.
    "producers": [{ "binding": "R2_EVENTS", "queue": "r2-events" }],
    "consumers": [
      {
        "queue": "r2-events",
        "max_batch_size": 10,
        "max_retries": 3,
        "dead_letter_queue": "r2-events-dlq",
      },
      { "queue": "r2-events-dlq", "max_retries": 10 },
    ],
  },
}

One Worker consumes both queues and tells them apart by batch.queue. Set the admin token for the replay route as a secret with npx wrangler secret put ADMIN_TOKEN.

A message that fails max_retries times without a dead-letter queue is deleted, so keep the dead_letter_queue line. The dead-letter consumer gets its own, higher max_retries because nothing catches its failures: a dead-letter queue has no dead-letter queue of its own. Messages in a dead-letter queue with no consumer persist for four days, and on the Workers Free plan queues keep messages for 24 hours. Parking them in D1 keeps them until you decide.

Record what’s been processed

CREATE TABLE uploads (
  key TEXT PRIMARY KEY,
  etag TEXT,
  size INTEGER,
  sha256 TEXT,
  state TEXT NOT NULL,          -- 'processed' or 'deleted'
  event_time INTEGER NOT NULL,  -- the newest event applied to this row, ms
  event_id TEXT NOT NULL
);

CREATE TABLE failed_events (
  message_id TEXT PRIMARY KEY,
  body TEXT NOT NULL,
  reason TEXT NOT NULL,
  failed_at INTEGER NOT NULL
);
npx wrangler d1 execute uploads-db --remote --file schema.sql

The processing state lives in D1 rather than a Durable Object for three reasons:

  • The app reads it. A file list can join on uploads to show which files are processed and their checksums. Rows in a Durable Object would need a second read path.
  • One conditional write is enough. Each handler finishes with a single statement whose WHERE clause refuses to overwrite a newer event’s result. That holds however many consumer invocations run at once, and Queues runs several by default.
  • Nothing has to be serialized. The work is cheap to repeat. If processing a file took minutes or cost money, a Durable Object per key would give you one place to queue that key’s events and run them one at a time, at the cost of an extra hop and more code.

Write the handlers

import {
  createFiles,
  type Files,
  FilesError,
  type StoredFile,
} from "files-sdk";
import { events, type FileEvent } from "files-sdk/events";
import { r2 } from "files-sdk/r2";

export interface Env {
  UPLOADS: R2Bucket;
  DB: D1Database;
  R2_EVENTS: Queue;
  ADMIN_TOKEN: string;
}

export function createStorage(env: Env) {
  const files = createFiles({
    adapter: r2({ binding: env.UPLOADS }),
    // A binding doesn't know its bucket's name. Name it, so events for any
    // other bucket that shares the queue are dropped.
    plugins: [events({ bucket: "uploads" })],
  });

  files.events.on("created", "users/**", (event) =>
    processUpload(env, files, event)
  );
  files.events.on("deleted", "users/**", (event) =>
    recordDelete(env, files, event)
  );
  return files;
}

async function processUpload(env: Env, files: Files, event: FileEvent) {
  // A redelivery, or a later write of the same bytes: nothing left to do.
  const done = await env.DB.prepare(
    "SELECT 1 FROM uploads WHERE key = ? AND etag = ? AND state = 'processed'"
  )
    .bind(event.key, event.etag)
    .first();
  if (done) {
    return;
  }

  let file: StoredFile;
  try {
    file = await files.download(event.key);
  } catch (error) {
    if (error instanceof FilesError && error.code === "NotFound") {
      return; // Deleted since. The delete has its own event.
    }
    throw error;
  }
  if (file.etag !== event.etag) {
    await file.stream().cancel();
    return; // Overwritten since. The newer write has its own event.
  }

  // Hash as the bytes stream in, so the file never sits in memory whole.
  const hasher = new crypto.DigestStream("SHA-256");
  await file.stream().pipeTo(hasher);
  const sha256 = [...new Uint8Array(await hasher.digest)]
    .map((byte) => byte.toString(16).padStart(2, "0"))
    .join("");

  // The only write. A newer event that committed first wins.
  await env.DB.prepare(
    `INSERT INTO uploads (key, etag, size, sha256, state, event_time, event_id)
     VALUES (?1, ?2, ?3, ?4, 'processed', ?5, ?6)
     ON CONFLICT (key) DO UPDATE SET
       etag = excluded.etag, size = excluded.size, sha256 = excluded.sha256,
       state = 'processed', event_time = excluded.event_time,
       event_id = excluded.event_id
     WHERE excluded.event_time >= uploads.event_time`
  )
    .bind(event.key, event.etag, file.size, sha256, event.time, event.id)
    .run();
}

async function recordDelete(env: Env, files: Files, event: FileEvent) {
  // Re-created since. That write's event decides what the row says.
  if (await files.exists(event.key)) {
    return;
  }
  await env.DB.prepare(
    `UPDATE uploads
     SET state = 'deleted', etag = NULL, sha256 = NULL,
         event_time = ?2, event_id = ?3
     WHERE key = ?1 AND event_time <= ?2`
  )
    .bind(event.key, event.time, event.id)
    .run();
}

Hashing stands in for whatever your app does with a new file: extract text, scan it, or tell its owner. Each handler makes three checks:

  • Has this generation already been processed? The lookup is by key and ETag, not by event.id. A redelivered message, a second write of identical bytes, and a gateway event for the same upload all carry the same ETag, so they all stop here. If the work must run once per write rather than once per content, as with a “you uploaded a file” email, key it on event.id instead.
  • Is the event still current? download() returns the object as it is now. A different ETag means a later write replaced it, and NotFound means it was deleted, and either one has its own notification coming. The ETag check reads the stored object, which closes the gap a separate head() would leave between checking and reading.
  • Is a newer result already stored? Two consumer invocations can work on the same key at once. The WHERE clause on each write refuses to replace a row that a later event (by R2’s eventTime) already wrote, so the newest event wins whichever commits last.

events({ bucket: "uploads" }) matters with a binding. The plugin drops other buckets’ events by default, but only when the adapter knows its bucket’s name, and r2({ binding }) doesn’t. Without the option, a queue that also carries another bucket’s notifications would run these handlers for that bucket’s keys.

Acknowledge after the work commits

import { FilesError } from "files-sdk";
import { createStorage, type Env } from "./storage";

const DEAD_LETTER_QUEUE = "r2-events-dlq";

// 30 s, 60 s, 120 s, ... capped at an hour. `attempts` starts at 1.
const backoff = (attempts: number) => Math.min(30 * 2 ** (attempts - 1), 3600);

async function park(env: Env, messages: readonly Message[], reason: string) {
  await env.DB.batch(
    messages.map((message) =>
      env.DB.prepare(
        `INSERT OR IGNORE INTO failed_events (message_id, body, reason, failed_at)
         VALUES (?, ?, ?, ?)`
      ).bind(message.id, JSON.stringify(message.body), reason, Date.now())
    )
  );
}

export default {
  async queue(batch, env) {
    if (batch.queue === DEAD_LETTER_QUEUE) {
      await park(env, batch.messages, "retries exhausted");
      batch.ackAll();
      return;
    }

    const files = createStorage(env);
    for (const message of batch.messages) {
      try {
        await files.events.dispatch(message.body);
        message.ack();
      } catch (error) {
        if (error instanceof FilesError && error.code === "Invalid") {
          // Not an R2 notification. Redelivering it won't change that.
          await park(env, [message], error.message);
          message.ack();
        } else {
          message.retry({ delaySeconds: backoff(message.attempts) });
        }
      }
    }
  },

  // `fetch` (the replay route) follows in the next section.
} satisfies ExportedHandler<Env>;

Each message is acknowledged on its own. If the Worker threw out of queue() instead, Queues would retry the whole batch except the messages already acknowledged. Per-message ack() and retry() keep one bad file from holding back the other nine.

dispatch() runs the handlers and resolves once they finish. When a handler throws, it rejects with a Provider error after trying every event in the message, and it reports the handler’s error to the plugin’s onError first, which logs it with console.error by default. A body that isn’t an R2 notification at all is different: dispatch() rejects with Invalid before any handler runs, and no retry can fix it, so the consumer parks it straight away.

The backoff waits 30, 60, then 120 seconds between attempts. delaySeconds can go up to 24 hours, and retry_delay on the consumer sets a default when you call retry() without one.

Replay what failed

Add a fetch handler to the same default export. It sends the parked messages back to r2-events and removes them from the table:

  async fetch(request, env) {
    const url = new URL(request.url);
    if (request.method !== "POST" || url.pathname !== "/admin/replay") {
      return new Response("Not found", { status: 404 });
    }
    if (!(await isAdmin(request, env))) {
      return new Response("Unauthorized", { status: 401 });
    }

    const { results } = await env.DB.prepare(
      `SELECT message_id, body FROM failed_events
       WHERE reason = 'retries exhausted' ORDER BY failed_at LIMIT 100`
    ).all<{ message_id: string; body: string }>();
    if (results.length > 0) {
      await env.R2_EVENTS.sendBatch(
        results.map((row) => ({ body: JSON.parse(row.body) }))
      );
      await env.DB.batch(
        results.map((row) =>
          env.DB.prepare("DELETE FROM failed_events WHERE message_id = ?").bind(
            row.message_id
          )
        )
      );
    }
    return Response.json({ replayed: results.length });
  },
async function isAdmin(request: Request, env: Env) {
  const encoder = new TextEncoder();
  const given = encoder.encode(request.headers.get("authorization") ?? "");
  const expected = encoder.encode(`Bearer ${env.ADMIN_TOKEN}`);
  return (
    given.byteLength === expected.byteLength &&
    crypto.subtle.timingSafeEqual(given, expected)
  );
}
curl -X POST -H "Authorization: Bearer $ADMIN_TOKEN" \
  https://uploads-consumer.<your-subdomain>.workers.dev/admin/replay

A replayed body is the original notification byte for byte, so it gets the same event.id and goes through the same three checks. If the object has changed since it failed, the replay finds that out and skips it. The route sends before it deletes, so a crash between the two sends a message twice at worst, which the handler absorbs. Parked malformed bodies stay out of the replay: the first version of this route replayed everything, and the {"hello":"world"} test body went straight back into the table under a new message ID.

What happened under wrangler dev

Notification rules live on the real bucket, not in wrangler.jsonc. Under wrangler dev, writes to the local R2 simulation sent nothing to the local queue. To exercise the consumer locally, send notification-shaped bodies to the queue yourself from a route that only exists in development:

// Development only: stands in for R2's notification.
if (url.pathname === "/dev/notify") {
  const key = url.searchParams.get("key") ?? "";
  const object = await env.UPLOADS.head(key);
  await env.R2_EVENTS.send({
    action: object ? "PutObject" : "DeleteObject",
    bucket: "uploads",
    eventTime: new Date().toISOString(),
    object: object ? { key, size: object.size, eTag: object.etag } : { key },
  });
  return new Response(null, { status: 202 });
}

The local R2 binding reports ETags as bare hex, the same form as the eTag in R2’s documented notification body. The table below comes from running the code above with logging around its bindings, sending bodies by hand in the order shown, and injecting failures into the uploads write:

Scenario What arrived What the consumer did uploads row after
Duplicate The same body twice First: lookup missed, read the object, hashed it, inserted, acked. Second: lookup hit, acked without reading R2 Processed once
Overwrite, newest event first v2’s event, then v1’s v2 processed. v1: lookup missed, download() returned v2’s ETag, acked without writing v2
Create then delete, delete first The delete, then the create Delete: exists() false, the update matched no row. Create: download() got NotFound, acked None
Delete, then a late copy of the create The delete, then the original create again Delete marked the row deleted. Create: lookup missed (state is deleted), NotFound, acked deleted
Deleted and re-created, events reversed v2’s create, the delete, v1’s create v2 processed. Delete: the object exists again, acked. v1: stale ETag, acked v2
Write committed, then the call threw One body Attempt 1: insert committed, the call threw, retry with 30 s. Attempt 2, 31 s later: lookup hit, acked Processed once
Write fails every time One body Attempts 1 to 4 at 0, 31, 92, and 213 s, each retry. One second after the fourth, the message arrived on r2-events-dlq with attempts back at 1 and the same message ID, and was parked failed_events row, then processed after replay
Another bucket A body for other-bucket dispatch() dropped it, no handler ran, acked None
Not a notification {"hello":"world"} Invalid, parked, acked failed_events row
Outside users/ A body for thumbnails/7/a.png No handler matched, acked without reading R2 None

With max_retries: 3, the message was attempted four times: the first delivery plus three retries. The retry() call on the fourth attempt asked for 240 seconds, and Queues sent the message to the dead-letter queue instead.

files events parse from the CLI prints the event a body produces, which helps when one doesn’t match the handler you expect:

npx -p files-sdk files events parse delivery.json --format r2

Gateway completions, SDK writes, and notifications

The same upload can reach on() handlers from three sources, and their IDs never match:

source Fires when id Retried on failure
"provider" R2 stores or deletes any object the rule matches Built from the body (above) Yes, by Queues
"gateway" A gateway upload completes, if events() is installed on the gateway’s files gateway:<uploadId> No. Failures go to onError and the upload still succeeds
"sdk" A write through this files instance, with events({ sdk: true }) sdk:<random UUID> No

onUploadComplete is a separate hook on the gateway, not an event source, and it fires only for uploads that reach complete.

For R2, do the processing in the queue consumer only:

  • Notifications are the only source that sees every write. A browser that closes before complete, a file added in the dashboard, and a lifecycle deletion never reach the gateway.
  • Gateway event handlers run inside the complete request. Hashing or scanning there makes the user wait, and a failure is reported to onError but never retried.
  • Event times come from different clocks. A gateway event’s time is the object’s lastModified, while a notification’s comes from R2’s eventTime. The ordering guard above compares times, so feed it from one source. During testing, one notification sent with the current time and others with fixed times was enough to make a delete lose to the create it followed.
  • Leave sdk: true off. R2 sends notifications for the consumer’s own writes too, so every write would be handled twice, and the random sdk: IDs can’t be deduplicated against the provider’s.

Use onUploadComplete for what needs the user’s request: record the owner from authorize’s context, or return a row ID to the browser. The consumer then fills in the processing results for that key.

Limits and tradeoffs

  • No exactly-once delivery. Queues delivers at least once and doesn’t preserve order. The three checks make repeated and reordered events harmless for work stored in D1. A side effect outside D1, such as an email or a payment, needs its own idempotency key. Pass event.id, or a key and ETag, to an API that accepts one, or use the plugin’s dedupe store, which still lets two racing copies through.
  • Ordering relies on R2’s eventTime. Two writes to one key within the same millisecond have equal times, and the >= lets whichever commits last win. If that matters, compare the stored ETag against the live object when you read the row.
  • Deletes carry no ETag or size. The delete handler re-reads the bucket instead.
  • Keys arrive as R2 sends them. R2 doesn’t document an encoding for object.key, and the parser doesn’t decode it.
  • Consumer limits apply to the work. A consumer invocation has a 15-minute wall-clock limit and 128 MB of memory. CPU time is capped at 10 ms on Workers Free and can be raised to 5 minutes on Workers Paid with limits.cpu_ms. Hashing a stream keeps memory flat, but CPU grows with file size. Heavy processing belongs in a Workflow that the consumer starts.
  • The dead-letter consumer can fail too. If D1 is down, the dead-letter batch retries up to its own max_retries and is then deleted. Alert on errors from that queue.
  • Operations are billed per message. A delivery is about three operations, and every retry adds a read. A rule without --prefix or --suffix sends a message for every object in the bucket.

Troubleshooting

Nothing arrives at the consumer under wrangler dev. R2 notifications don’t reach the local queue. Send bodies yourself, as in What happened under wrangler dev, or check the deployed consumer’s logs with npx wrangler tail.

files-sdk/events: not a r2 notification (no action/object, message body, or messages[]). Something other than R2 published to the queue, or a body was wrapped twice. The consumer parks these with the message as the reason. Run the body through files events parse --format r2 to see what it produces.

files-sdk/events: the r2-binding adapter has no notification format on this instance: a plugin turned provider events off…. A plugin such as tiering({ fallback: true }) turns provider events off for the instance. Give the consumer its own createFiles() with only events().

Handlers never run, and dispatch() returns an empty array. The event was filtered out before any handler. Check the body’s bucket against events({ bucket }), the key against the on() pattern, and the instance prefix if you set one. Actions other than the five listed above also parse to nothing.

files-sdk/events: created handler failed for users/… Error: … in the logs, with the message retried. That’s onError reporting a handler that threw. The message is retried with backoff and lands in failed_events after the last attempt. Fix the cause, then call the replay route.

A file’s row still says processed after it was deleted. The delete’s event_time was older than the row’s, so the guard refused it. Within R2’s own notifications that means the file was re-created later. If you also write this table from the gateway or a script, their clocks differ from R2’s; see Gateway completions, SDK writes, and notifications.

Last updated on

Was this page helpful?