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 () => {