diff --git a/.changeset/next-wait-until.md b/.changeset/next-wait-until.md
new file mode 100644
index 0000000000..3429d0d459
--- /dev/null
+++ b/.changeset/next-wait-until.md
@@ -0,0 +1,9 @@
+---
+"uploadthing": minor
+---
+
+The Next.js App Router adapter passes `ctx.waitUntil` into `middleware`,
+`onUploadComplete`, and `onUploadError`. Work scheduled there runs through
+Next.js `after`, so the client receives `onUploadComplete`'s return value
+without waiting for it. Requires Next.js 15.1 for the task to outlive the
+response. Earlier versions start the task and warn once.
diff --git a/docs/src/app/(docs)/file-routes/page.mdx b/docs/src/app/(docs)/file-routes/page.mdx
index 4eac2854fe..d659b290fe 100644
--- a/docs/src/app/(docs)/file-routes/page.mdx
+++ b/docs/src/app/(docs)/file-routes/page.mdx
@@ -299,4 +299,31 @@ properties:
An object with info for the file that was uploaded, such as the name, key,
size, url etc.
+
+ Present on the Next.js App Router adapter (`uploadthing/next`).
+ `ctx.waitUntil` schedules work with Next.js
+ [`after`](https://nextjs.org/docs/app/api-reference/functions/after), so the
+ client receives this function's return value without waiting for that work.
+ On Next.js 15.1 or later, that work stays alive after the response. Earlier
+ versions start the task and warn once. In development, callback hooks run
+ after the response, so those tasks start in-process. The same `ctx` is
+ passed to `middleware` and `onUploadError`.
+
+
+```ts
+import { createUploadthing } from "uploadthing/next";
+import { UTApi } from "uploadthing/server";
+
+const f = createUploadthing();
+const utapi = new UTApi();
+
+f(["image"])
+ .middleware(async () => {
+ return { oldFileKey: "previous-key" };
+ })
+ .onUploadComplete(async ({ metadata, ctx }) => {
+ ctx.waitUntil(utapi.deleteFiles(metadata.oldFileKey));
+ return { deletionScheduled: true };
+ });
+```
diff --git a/packages/uploadthing/src/next.ts b/packages/uploadthing/src/next.ts
index 7b48ba7655..60b2d7a7e8 100644
--- a/packages/uploadthing/src/next.ts
+++ b/packages/uploadthing/src/next.ts
@@ -1,3 +1,4 @@
+import * as NextServer from "next/server";
import type { NextRequest } from "next/server";
import * as Effect from "effect/Effect";
@@ -18,8 +19,117 @@ export {
UTRegion as experimental_UTRegion,
} from "./_internal/types";
+type AfterTask = Promise | (() => unknown);
+
+/**
+ * Request helpers for the Next.js App Router adapter.
+ * Passed to `middleware`, `onUploadComplete`, and `onUploadError`.
+ */
+export type RequestContext = {
+ /**
+ * Schedule work that should not delay the client callback.
+ *
+ * On Next.js 15.1 or later, the task is handed to `after` while the request
+ * is still open. `onUploadComplete` can return `serverData` without awaiting
+ * the task, and Next.js keeps the invocation alive until the task settles.
+ * Earlier versions start the task and warn once. In development, callback
+ * hooks run after the response, so those tasks start in-process.
+ */
+ waitUntil: (task: AfterTask) => void;
+};
+
+let didWarnMissingAfter = false;
+
+const warnTaskMayFreeze = (): void => {
+ if (didWarnMissingAfter) return;
+ didWarnMissingAfter = true;
+ // eslint-disable-next-line no-console
+ console.warn(
+ "[uploadthing] ctx.waitUntil could not register with Next.js after. The task was started, but a serverless function may freeze it when the response ends. Requires Next.js 15.1 or later for the task to outlive the response.",
+ );
+};
+
+/** Failures stay off the upload hook, including when `after` runs the task. */
+const settleTask = (task: AfterTask): Promise => {
+ try {
+ const pending = typeof task === "function" ? task() : task;
+ return Promise.resolve(pending).then(
+ () => undefined,
+ (error: unknown) => {
+ // eslint-disable-next-line no-console
+ console.error("[uploadthing] ctx.waitUntil task failed.", error);
+ },
+ );
+ } catch (error) {
+ // eslint-disable-next-line no-console
+ console.error("[uploadthing] ctx.waitUntil task failed.", error);
+ return Promise.resolve();
+ }
+};
+
+/**
+ * A promise is already running, so its rejection handler is attached now.
+ * A callback stays deferred until it is run.
+ */
+const deferTask = (task: AfterTask): (() => Promise) => {
+ if (typeof task === "function") return () => settleTask(task);
+ const settled = settleTask(task);
+ return () => settled;
+};
+
+type Scheduler = { enqueue: (task: AfterTask) => void };
+
+type SchedulerState = "queueing" | "flushed" | "unregistered";
+
+const inProcess: Scheduler = {
+ enqueue: (task) => void deferTask(task)(),
+};
+
+/**
+ * `after` reads the request store at call time. Hooks run later, sometimes
+ * after that store is gone, so `after` is registered here and hooks enqueue.
+ */
+const openScheduler = (): Scheduler => {
+ const queue: Array<() => Promise> = [];
+ let state: SchedulerState = "queueing";
+
+ const flush = (): Promise => {
+ const batch = queue.splice(0);
+ if (batch.length === 0) {
+ state = "flushed";
+ return Promise.resolve();
+ }
+ return Promise.all(batch.map((run) => run())).then(flush);
+ };
+
+ if (typeof NextServer.after === "function") {
+ try {
+ NextServer.after(flush);
+ } catch {
+ // Thrown before registration, so `flush` never runs.
+ state = "unregistered";
+ }
+ } else {
+ state = "unregistered";
+ }
+
+ const byState: Record Promise) => void> = {
+ queueing: (run) => void queue.push(run),
+ // Development detaches callback hooks past the response, and the dev
+ // server still finishes the task.
+ flushed: (run) => void run(),
+ unregistered: (run) => {
+ warnTaskMayFreeze();
+ void run();
+ },
+ };
+
+ return { enqueue: (task) => byState[state](deferTask(task)) };
+};
+
type AdapterArgs = {
req: NextRequest;
+ ctx: RequestContext;
};
export const createUploadthing = (
@@ -30,7 +140,11 @@ export const createRouteHandler = (
opts: RouteHandlerOptions,
) => {
const handler = makeAdapterHandler<[NextRequest], AdapterArgs>(
- (req) => Effect.succeed({ req }),
+ (req) => {
+ // Built eagerly, while the route handler still holds the request scope.
+ const scheduler = req.method === "POST" ? openScheduler() : inProcess;
+ return Effect.succeed({ req, ctx: { waitUntil: scheduler.enqueue } });
+ },
(req) => Effect.succeed(req),
opts,
"nextjs-app",
diff --git a/packages/uploadthing/test/node/adapters.test.ts b/packages/uploadthing/test/node/adapters.test.ts
index 6ba3fc822c..8117366cf1 100644
--- a/packages/uploadthing/test/node/adapters.test.ts
+++ b/packages/uploadthing/test/node/adapters.test.ts
@@ -1,6 +1,7 @@
/* eslint-disable no-restricted-globals */
import type { NextApiRequest, NextApiResponse } from "next";
-import { NextRequest } from "next/server";
+import { after, NextRequest } from "next/server";
+import type * as NextServer from "next/server";
import * as FetchHttpClient from "@effect/platform/FetchHttpClient";
import * as HttpServerRequest from "@effect/platform/HttpServerRequest";
import * as HttpServerResponse from "@effect/platform/HttpServerResponse";
@@ -13,14 +14,19 @@ import { createApp, H3Event, toWebHandler } from "h3";
import { setupServer } from "msw/node";
import {
afterAll,
+ afterEach,
beforeAll,
+ beforeEach,
describe,
expect,
expectTypeOf,
it,
+ onTestFinished,
vi,
} from "vitest";
+import { signPayload } from "@uploadthing/shared";
+
import {
baseHeaders,
createApiUrl,
@@ -30,8 +36,20 @@ import {
requestSpy,
requestsToDomain,
testToken,
+ UFS_HOST,
uploadCompleteMock,
+ UTFS_URL,
} from "../__test-helpers";
+import { UploadedFileData } from "../../src/_internal/shared-schemas";
+import type { RequestContext } from "../../src/next";
+
+vi.mock("next/server", async () => {
+ const actual = (await vi.importActual("next/server")) as typeof NextServer;
+ return {
+ ...actual,
+ after: vi.fn(),
+ };
+});
const server = setupServer(...handlers);
beforeAll(() => server.listen({ onUnhandledRequest: "bypass" }));
@@ -300,17 +318,23 @@ describe("adapters:server", async () => {
});
describe("adapters:next", async () => {
- const { createUploadthing, createRouteHandler } = await import(
- "../../src/next"
- );
+ const nextAdapter = await import("../../src/next");
+ const { createUploadthing, createRouteHandler } = nextAdapter;
const f = createUploadthing();
+ beforeEach(() => {
+ vi.mocked(after).mockReset();
+ });
+
const router = {
middleware: f({ blob: {} })
.middleware((opts) => {
middlewareMock(opts);
expectTypeOf<{
req: NextRequest;
+ ctx: {
+ waitUntil: (task: Promise | (() => unknown)) => void;
+ };
}>(opts);
return {};
})
@@ -329,6 +353,7 @@ describe("adapters:next", async () => {
expect(res.status).toBe(200);
expect(res.headers.get("content-type")).toBe("application/json");
+ expect(after).not.toHaveBeenCalled();
const json = await res.json();
expect(json).toEqual([
@@ -361,7 +386,10 @@ describe("adapters:next", async () => {
expect(middlewareMock).toHaveBeenCalledOnce();
expect(middlewareMock).toHaveBeenCalledWith(
- expect.objectContaining({ req }),
+ expect.objectContaining({
+ req,
+ ctx: { waitUntil: expect.any(Function) },
+ }),
);
// Should proceed to generate a signed URL
@@ -396,6 +424,262 @@ describe("adapters:next", async () => {
method: "POST",
});
});
+
+ describe("ctx.waitUntil", () => {
+ const outsideRequestScope = "after was called outside a request scope";
+ const taskFailed = "[uploadthing] ctx.waitUntil task failed.";
+
+ const uploadedPayload = JSON.stringify({
+ status: "uploaded",
+ metadata: {},
+ origin: "https://example.com",
+ file: new UploadedFileData({
+ url: `${UTFS_URL}/f/some-random-key.png`,
+ appUrl: `${UTFS_URL}/a/${testToken.decoded.appId}/f/some-random-key.png`,
+ ufsUrl: `https://${testToken.decoded.appId}.${UFS_HOST}/f/some-random-key.png`,
+ name: "foo.png",
+ key: "some-random-key.png",
+ size: 48,
+ type: "image/png",
+ customId: null,
+ fileHash: "some-md5-hash",
+ }),
+ });
+ const failedPayload = JSON.stringify({
+ fileKey: "some-random-key.png",
+ error: "network",
+ });
+
+ const uploadRequest = () =>
+ new NextRequest(createApiUrl("background", "upload"), {
+ method: "POST",
+ headers: {
+ ...baseHeaders,
+ host: "localhost:3000",
+ "x-forwarded-proto": "http",
+ },
+ body: JSON.stringify({
+ files: [{ name: "foo.txt", size: 48, type: "text/plain" }],
+ }),
+ });
+
+ const hookRequest = async (hook: "callback" | "error", payload: string) =>
+ new NextRequest(createApiUrl("background"), {
+ method: "POST",
+ headers: {
+ "uploadthing-hook": hook,
+ "x-uploadthing-signature": await Effect.runPromise(
+ signPayload(payload, testToken.decoded.apiKey),
+ ),
+ },
+ body: payload,
+ });
+
+ /** Every hook on the route schedules the same work. */
+ const handlersFor = (
+ schedule: (waitUntil: RequestContext["waitUntil"]) => void,
+ { adapter = nextAdapter, isDev = false } = {},
+ ) => {
+ const route = adapter.createUploadthing();
+ return adapter.createRouteHandler({
+ router: {
+ background: route({ blob: {} })
+ .middleware(({ ctx }) => {
+ schedule(ctx.waitUntil);
+ return {};
+ })
+ .onUploadError(({ ctx }) => schedule(ctx.waitUntil))
+ .onUploadComplete(({ ctx }) => schedule(ctx.waitUntil)),
+ },
+ config: { token: testToken.encoded, isDev },
+ });
+ };
+
+ /** Stands in for Next.js running `after` callbacks once the response ends. */
+ const endResponse = () =>
+ Promise.all(
+ vi
+ .mocked(after)
+ .mock.calls.map(([task]) =>
+ typeof task === "function" ? task() : task,
+ ),
+ );
+
+ /**
+ * The one-time warning is module state, so tests that assert on it load a
+ * fresh adapter. The `next/server` mock stays shared across the reset.
+ */
+ const loadFreshAdapter = () => {
+ vi.resetModules();
+ return import("../../src/next");
+ };
+
+ afterEach(() => {
+ vi.restoreAllMocks();
+ });
+
+ it("registers with after before POST first awaits", async () => {
+ const task = vi.fn();
+ const handlers = handlersFor((waitUntil) => waitUntil(task));
+ let inPostCall = false;
+ vi.mocked(after).mockImplementation(() => {
+ if (!inPostCall) throw new Error(outsideRequestScope);
+ });
+
+ inPostCall = true;
+ const pending = handlers.POST(uploadRequest());
+ inPostCall = false;
+
+ expect((await pending).status).toBe(200);
+ expect(task).not.toHaveBeenCalled();
+ await endResponse();
+ expect(task).toHaveBeenCalledOnce();
+ });
+
+ it.each([
+ { hook: "callback", payload: uploadedPayload },
+ { hook: "error", payload: failedPayload },
+ ] as const)(
+ "runs $hook hook tasks once the response ends",
+ async ({ hook, payload }) => {
+ const task = vi.fn();
+ const handlers = handlersFor((waitUntil) => waitUntil(task));
+
+ const res = await handlers.POST(await hookRequest(hook, payload));
+
+ expect(res.status).toBe(200);
+ expect(task).not.toHaveBeenCalled();
+ await endResponse();
+ expect(task).toHaveBeenCalledOnce();
+ },
+ );
+
+ it("contains failures, including promises that reject before the response ends", async () => {
+ const errorLog = vi
+ .spyOn(console, "error")
+ .mockImplementation(() => undefined);
+ const unhandled: Array = [];
+ const collect = (reason: unknown) => void unhandled.push(reason);
+ process.on("unhandledRejection", collect);
+ onTestFinished(() => void process.off("unhandledRejection", collect));
+ const handlers = handlersFor((waitUntil) => {
+ waitUntil(Promise.reject(new Error("delete failed")));
+ waitUntil(() => {
+ throw new Error("boom");
+ });
+ });
+
+ const res = await handlers.POST(uploadRequest());
+ await new Promise((resolve) => setTimeout(resolve, 0));
+ await endResponse();
+
+ expect(res.status).toBe(200);
+ expect(unhandled).toEqual([]);
+ expect(errorLog).toHaveBeenCalledWith(
+ taskFailed,
+ expect.objectContaining({ message: "delete failed" }),
+ );
+ expect(errorLog).toHaveBeenCalledWith(
+ taskFailed,
+ expect.objectContaining({ message: "boom" }),
+ );
+ });
+
+ it("keeps the flush open for tasks queued while it runs", async () => {
+ const finished = vi.fn();
+ const handlers = handlersFor((waitUntil) =>
+ waitUntil(() =>
+ waitUntil(() =>
+ new Promise((resolve) => setTimeout(resolve, 0)).then(finished),
+ ),
+ ),
+ );
+
+ const res = await handlers.POST(uploadRequest());
+ await endResponse();
+
+ expect(res.status).toBe(200);
+ expect(finished).toHaveBeenCalledOnce();
+ });
+
+ it("runs a development hook that outlives the response once, without warning", async () => {
+ const adapter = await loadFreshAdapter();
+ vi.mocked(after).mockImplementation((flush) => {
+ if (typeof flush === "function") void flush();
+ });
+ const warnLog = vi
+ .spyOn(console, "warn")
+ .mockImplementation(() => undefined);
+ const task = vi.fn();
+ const handlers = handlersFor((waitUntil) => waitUntil(task), {
+ adapter,
+ isDev: true,
+ });
+
+ const res = await handlers.POST(
+ await hookRequest("callback", uploadedPayload),
+ );
+ await vi.waitUntil(() => task.mock.calls.length > 0);
+ await endResponse();
+
+ expect(res.status).toBe(200);
+ expect(task).toHaveBeenCalledOnce();
+ expect(warnLog).not.toHaveBeenCalled();
+ });
+
+ it("warns once and runs each task in-process when after is missing", async () => {
+ const adapter = await loadFreshAdapter();
+ const nextServer = await import("next/server");
+ const mockedAfter = nextServer.after;
+ Object.defineProperty(nextServer, "after", {
+ configurable: true,
+ value: undefined,
+ });
+ onTestFinished(() => {
+ Object.defineProperty(nextServer, "after", {
+ configurable: true,
+ value: mockedAfter,
+ });
+ });
+ const warnLog = vi
+ .spyOn(console, "warn")
+ .mockImplementation(() => undefined);
+ const task = vi.fn();
+ const handlers = handlersFor((waitUntil) => waitUntil(task), {
+ adapter,
+ });
+
+ const responses = [
+ await handlers.POST(uploadRequest()),
+ await handlers.POST(uploadRequest()),
+ ];
+
+ expect(responses.map((res) => res.status)).toEqual([200, 200]);
+ expect(task).toHaveBeenCalledTimes(2);
+ expect(warnLog).toHaveBeenCalledOnce();
+ });
+
+ it("runs the task once when after throws, even if it kept the callback", async () => {
+ const adapter = await loadFreshAdapter();
+ vi.mocked(after).mockImplementation(() => {
+ throw new Error(outsideRequestScope);
+ });
+ const warnLog = vi
+ .spyOn(console, "warn")
+ .mockImplementation(() => undefined);
+ const task = vi.fn();
+ const handlers = handlersFor((waitUntil) => waitUntil(task), {
+ adapter,
+ });
+
+ const res = await handlers.POST(uploadRequest());
+ await endResponse();
+
+ expect(res.status).toBe(200);
+ expect(task).toHaveBeenCalledOnce();
+ expect(warnLog).toHaveBeenCalledOnce();
+ });
+ });
});
describe("adapters:next-legacy", async () => {