From 3a93d7aa425773bce44707a5f12de9424ea4dfea Mon Sep 17 00:00:00 2001 From: Arman Charan Date: Sun, 27 Sep 2026 11:42:15 +1000 Subject: [PATCH 1/5] feat: expose ctx.waitUntil on the Next.js App Router adapter onUploadComplete has to return before the client receives serverData. Hand extra work to Next.js after so that return is not held up. Co-authored-by: Cursor --- .changeset/next-wait-until.md | 9 ++++ docs/src/app/(docs)/file-routes/page.mdx | 25 +++++++++ packages/uploadthing/src/next.ts | 45 +++++++++++++++- .../uploadthing/test/node/adapters.test.ts | 52 ++++++++++++++++++- 4 files changed, 128 insertions(+), 3 deletions(-) create mode 100644 .changeset/next-wait-until.md diff --git a/.changeset/next-wait-until.md b/.changeset/next-wait-until.md new file mode 100644 index 0000000000..12b07db50d --- /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 for the task to outlive the +response. diff --git a/docs/src/app/(docs)/file-routes/page.mdx b/docs/src/app/(docs)/file-routes/page.mdx index 4eac2854fe..1dff825fa3 100644 --- a/docs/src/app/(docs)/file-routes/page.mdx +++ b/docs/src/app/(docs)/file-routes/page.mdx @@ -299,4 +299,29 @@ 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. + Requires Next.js 15. 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 { deletedPrevious: true }; + }); +``` diff --git a/packages/uploadthing/src/next.ts b/packages/uploadthing/src/next.ts index 7b48ba7655..c7033ae5b5 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,46 @@ 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. + * + * The task is handed to Next.js `after`. `onUploadComplete` can return + * `serverData` without awaiting the task, and Next.js keeps the invocation + * alive until the task settles. Requires Next.js 15. + */ + waitUntil: (task: AfterTask) => void; +}; + +let didWarnMissingAfter = false; + +const scheduleAfterResponse = (task: AfterTask): void => { + const after = (NextServer as { after?: (task: AfterTask) => void }).after; + if (typeof after === "function") { + after(task); + return; + } + + if (!didWarnMissingAfter) { + didWarnMissingAfter = true; + // eslint-disable-next-line no-console + console.warn( + "[uploadthing] ctx.waitUntil needs Next.js 15. The task was started, but a serverless function may freeze it when the response ends.", + ); + } + + void Promise.resolve(typeof task === "function" ? task() : task); +}; + type AdapterArgs = { req: NextRequest; + ctx: RequestContext; }; export const createUploadthing = ( @@ -30,7 +69,11 @@ export const createRouteHandler = ( opts: RouteHandlerOptions, ) => { const handler = makeAdapterHandler<[NextRequest], AdapterArgs>( - (req) => Effect.succeed({ req }), + (req) => + Effect.succeed({ + req, + ctx: { waitUntil: scheduleAfterResponse }, + }), (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..4ca078c349 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"; @@ -33,6 +34,14 @@ import { uploadCompleteMock, } from "../__test-helpers"; +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" })); afterAll(() => server.close()); @@ -311,6 +320,9 @@ describe("adapters:next", async () => { middlewareMock(opts); expectTypeOf<{ req: NextRequest; + ctx: { + waitUntil: (task: Promise | (() => unknown)) => void; + }; }>(opts); return {}; }) @@ -361,7 +373,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 +411,39 @@ describe("adapters:next", async () => { method: "POST", }); }); + + it("hands ctx.waitUntil to Next.js after", async () => { + const task = vi.fn(() => Promise.resolve("deleted")); + const handlers = createRouteHandler({ + router: { + background: f({ blob: {} }) + .middleware((opts) => { + opts.ctx.waitUntil(task); + return {}; + }) + .onUploadComplete(uploadCompleteMock), + }, + config: { token: testToken.encoded }, + }); + + const res = await handlers.POST( + 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" }], + }), + }), + ); + + expect(res.status).toBe(200); + expect(after).toHaveBeenCalledWith(task); + expect(task).not.toHaveBeenCalled(); + }); }); describe("adapters:next-legacy", async () => { From 0fa01204b7d11ccefb78f350cd86e0425460665a Mon Sep 17 00:00:00 2001 From: Arman Charan Date: Sun, 27 Sep 2026 11:52:27 +1000 Subject: [PATCH 2/5] fix: keep ctx.waitUntil from failing the upload hook Next.js after throws outside the request scope. Development can hit that because the callback fiber outlives the response. Run the task anyway, and keep its failures off the hook. Co-authored-by: Cursor --- packages/uploadthing/src/next.ts | 45 ++++-- .../uploadthing/test/node/adapters.test.ts | 140 ++++++++++++++++++ 2 files changed, 174 insertions(+), 11 deletions(-) diff --git a/packages/uploadthing/src/next.ts b/packages/uploadthing/src/next.ts index c7033ae5b5..824eecf22a 100644 --- a/packages/uploadthing/src/next.ts +++ b/packages/uploadthing/src/next.ts @@ -38,22 +38,45 @@ export type RequestContext = { 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.", + ); +}; + +/** Run a task outside Next.js `after`. Failures stay in the background. */ +const runTask = (task: AfterTask): void => { + try { + const pending = typeof task === "function" ? task() : task; + void Promise.resolve(pending).catch((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); + } +}; + const scheduleAfterResponse = (task: AfterTask): void => { const after = (NextServer as { after?: (task: AfterTask) => void }).after; if (typeof after === "function") { - after(task); - return; - } - - if (!didWarnMissingAfter) { - didWarnMissingAfter = true; - // eslint-disable-next-line no-console - console.warn( - "[uploadthing] ctx.waitUntil needs Next.js 15. The task was started, but a serverless function may freeze it when the response ends.", - ); + try { + after(task); + return; + } catch { + // `after` throws outside the request scope. In development the callback + // fiber can outlive that scope, so the task still has to run. + warnTaskMayFreeze(); + } + } else { + warnTaskMayFreeze(); } - void Promise.resolve(typeof task === "function" ? task() : task); + runTask(task); }; type AdapterArgs = { diff --git a/packages/uploadthing/test/node/adapters.test.ts b/packages/uploadthing/test/node/adapters.test.ts index 4ca078c349..4f9853e946 100644 --- a/packages/uploadthing/test/node/adapters.test.ts +++ b/packages/uploadthing/test/node/adapters.test.ts @@ -22,6 +22,8 @@ import { vi, } from "vitest"; +import { signPayload } from "@uploadthing/shared"; + import { baseHeaders, createApiUrl, @@ -31,8 +33,11 @@ import { requestSpy, requestsToDomain, testToken, + UFS_HOST, uploadCompleteMock, + UTFS_URL, } from "../__test-helpers"; +import { UploadedFileData } from "../../src/_internal/shared-schemas"; vi.mock("next/server", async () => { const actual = (await vi.importActual("next/server")) as typeof NextServer; @@ -444,6 +449,141 @@ describe("adapters:next", async () => { expect(after).toHaveBeenCalledWith(task); expect(task).not.toHaveBeenCalled(); }); + + it("runs the task when Next.js after throws", async () => { + vi.mocked(after).mockImplementation(() => { + throw new Error("after was called outside a request scope"); + }); + const errorLog = vi.spyOn(console, "error").mockImplementation(() => undefined); + const warnLog = vi.spyOn(console, "warn").mockImplementation(() => undefined); + const task = vi.fn(() => { + throw new Error("boom"); + }); + const handlers = createRouteHandler({ + router: { + background: f({ blob: {} }) + .middleware((opts) => { + opts.ctx.waitUntil(task); + return {}; + }) + .onUploadComplete(uploadCompleteMock), + }, + config: { token: testToken.encoded }, + }); + + const res = await handlers.POST( + 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" }], + }), + }), + ); + + expect(res.status).toBe(200); + expect(task).toHaveBeenCalledOnce(); + expect(errorLog).toHaveBeenCalled(); + errorLog.mockRestore(); + warnLog.mockRestore(); + }); + + it("handles a rejected task when Next.js after throws", async () => { + vi.mocked(after).mockImplementation(() => { + throw new Error("after was called outside a request scope"); + }); + const errorLog = vi.spyOn(console, "error").mockImplementation(() => undefined); + const warnLog = vi.spyOn(console, "warn").mockImplementation(() => undefined); + const unhandled: Array = []; + const onUnhandled = (reason: unknown) => { + unhandled.push(reason); + }; + process.on("unhandledRejection", onUnhandled); + const handlers = createRouteHandler({ + router: { + background: f({ blob: {} }) + .middleware((opts) => { + opts.ctx.waitUntil(Promise.reject(new Error("delete failed"))); + return {}; + }) + .onUploadComplete(uploadCompleteMock), + }, + config: { token: testToken.encoded }, + }); + + const res = await handlers.POST( + 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" }], + }), + }), + ); + await new Promise((resolve) => setTimeout(resolve, 0)); + process.off("unhandledRejection", onUnhandled); + + expect(res.status).toBe(200); + expect(unhandled).toEqual([]); + errorLog.mockRestore(); + warnLog.mockRestore(); + }); + + it("hands ctx.waitUntil to Next.js after from onUploadComplete", async () => { + const task = vi.fn(() => Promise.resolve("deleted")); + const handlers = createRouteHandler({ + router: { + background: f({ blob: {} }) + .middleware(() => ({})) + .onUploadComplete((opts) => { + opts.ctx.waitUntil(task); + }), + }, + config: { token: testToken.encoded }, + }); + const payload = 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 signature = await Effect.runPromise( + signPayload(payload, testToken.decoded.apiKey), + ); + + const res = await handlers.POST( + new NextRequest(createApiUrl("background"), { + method: "POST", + headers: { + "uploadthing-hook": "callback", + "x-uploadthing-signature": signature, + }, + body: payload, + }), + ); + + expect(res.status).toBe(200); + expect(after).toHaveBeenCalledWith(task); + expect(task).not.toHaveBeenCalled(); + }); }); describe("adapters:next-legacy", async () => { From f222248b13a2c3cd535b3772e89477dd1aba557b Mon Sep 17 00:00:00 2001 From: Arman Charan Date: Mon, 28 Sep 2026 14:33:42 +1000 Subject: [PATCH 3/5] fix: register ctx.waitUntil with Next.js after inside the request after reads the request store at call time. Callback hooks can run after that store is gone, so register one flush while the request is open and let the hooks enqueue. The task outlives the response on Next.js 15.1. Co-authored-by: Cursor --- .changeset/next-wait-until.md | 4 +- docs/src/app/(docs)/file-routes/page.mdx | 8 +- packages/uploadthing/src/next.ts | 111 ++++++-- .../uploadthing/test/node/adapters.test.ts | 254 ++++++++++++++---- 4 files changed, 293 insertions(+), 84 deletions(-) diff --git a/.changeset/next-wait-until.md b/.changeset/next-wait-until.md index 12b07db50d..3429d0d459 100644 --- a/.changeset/next-wait-until.md +++ b/.changeset/next-wait-until.md @@ -5,5 +5,5 @@ 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 for the task to outlive the -response. +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 1dff825fa3..d659b290fe 100644 --- a/docs/src/app/(docs)/file-routes/page.mdx +++ b/docs/src/app/(docs)/file-routes/page.mdx @@ -304,8 +304,10 @@ properties: `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. - Requires Next.js 15. The same `ctx` is passed to `middleware` and - `onUploadError`. + 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`. @@ -322,6 +324,6 @@ f(["image"]) }) .onUploadComplete(async ({ metadata, ctx }) => { ctx.waitUntil(utapi.deleteFiles(metadata.oldFileKey)); - return { deletedPrevious: true }; + return { deletionScheduled: true }; }); ``` diff --git a/packages/uploadthing/src/next.ts b/packages/uploadthing/src/next.ts index 824eecf22a..8955067de3 100644 --- a/packages/uploadthing/src/next.ts +++ b/packages/uploadthing/src/next.ts @@ -29,9 +29,11 @@ export type RequestContext = { /** * Schedule work that should not delay the client callback. * - * The task is handed to Next.js `after`. `onUploadComplete` can return - * `serverData` without awaiting the task, and Next.js keeps the invocation - * alive until the task settles. Requires Next.js 15. + * 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; }; @@ -43,42 +45,84 @@ const warnTaskMayFreeze = (): void => { 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.", + "[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.", ); }; -/** Run a task outside Next.js `after`. Failures stay in the background. */ -const runTask = (task: AfterTask): void => { +/** 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; - void Promise.resolve(pending).catch((error: unknown) => { - // eslint-disable-next-line no-console - console.error("[uploadthing] ctx.waitUntil task failed.", error); - }); + 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(); } }; -const scheduleAfterResponse = (task: AfterTask): void => { - const after = (NextServer as { after?: (task: AfterTask) => void }).after; - if (typeof after === "function") { +type RequestScheduler = { + enqueue: (task: AfterTask) => void; +}; + +/** + * `after` reads the request store at call time. Hooks run later, sometimes + * after that store is gone, so the queue is opened here and hooks only enqueue. + */ +const openScheduler = (): RequestScheduler => { + const tasks: Array = []; + let sealed = false; + let failed = false; + + const flush = (): Promise => { + sealed = true; + const batch = tasks.splice(0, tasks.length); + return Promise.all(batch.map(settleTask)).then(() => undefined); + }; + + if (typeof NextServer.after === "function") { try { - after(task); - return; + NextServer.after(flush); } catch { - // `after` throws outside the request scope. In development the callback - // fiber can outlive that scope, so the task still has to run. - warnTaskMayFreeze(); + // Thrown before registration. Leave `flush` uncalled so the task runs once. + failed = true; + sealed = true; } } else { - warnTaskMayFreeze(); + failed = true; + sealed = true; } - runTask(task); + return { + enqueue(task) { + if (failed) { + warnTaskMayFreeze(); + void settleTask(task); + return; + } + if (sealed) { + // Development detaches callback hooks after the response. `after` + // already ran, and the dev server still finishes the task. + void settleTask(task); + return; + } + tasks.push(task); + }, + }; }; +/** + * Set for the synchronous start of `POST`, then captured into adapter args. + * Cleared before the handler awaits so concurrent requests do not share it. + */ +let requestScheduler: RequestScheduler | undefined; + type AdapterArgs = { req: NextRequest; ctx: RequestContext; @@ -92,14 +136,31 @@ export const createRouteHandler = ( opts: RouteHandlerOptions, ) => { const handler = makeAdapterHandler<[NextRequest], AdapterArgs>( - (req) => - Effect.succeed({ + (req) => { + const scheduler = requestScheduler; + return Effect.succeed({ req, - ctx: { waitUntil: scheduleAfterResponse }, - }), + ctx: { + waitUntil: (task) => { + if (scheduler) scheduler.enqueue(task); + else void settleTask(task); + }, + }, + }); + }, (req) => Effect.succeed(req), opts, "nextjs-app", ); - return { POST: handler, GET: handler }; + + const POST = (req: NextRequest) => { + requestScheduler = openScheduler(); + try { + return handler(req); + } finally { + requestScheduler = undefined; + } + }; + + return { POST, GET: handler }; }; diff --git a/packages/uploadthing/test/node/adapters.test.ts b/packages/uploadthing/test/node/adapters.test.ts index 4f9853e946..2980feb8a7 100644 --- a/packages/uploadthing/test/node/adapters.test.ts +++ b/packages/uploadthing/test/node/adapters.test.ts @@ -15,6 +15,7 @@ import { setupServer } from "msw/node"; import { afterAll, beforeAll, + beforeEach, describe, expect, expectTypeOf, @@ -319,6 +320,10 @@ describe("adapters:next", async () => { ); const f = createUploadthing(); + beforeEach(() => { + vi.mocked(after).mockReset(); + }); + const router = { middleware: f({ blob: {} }) .middleware((opts) => { @@ -346,6 +351,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([ @@ -446,98 +452,119 @@ describe("adapters:next", async () => { ); expect(res.status).toBe(200); - expect(after).toHaveBeenCalledWith(task); + expect(after).toHaveBeenCalledOnce(); expect(task).not.toHaveBeenCalled(); + const flush = vi.mocked(after).mock.calls[0]?.[0]; + expect(typeof flush).toBe("function"); + if (typeof flush !== "function") return; + await flush(); + expect(task).toHaveBeenCalledOnce(); + await flush(); + expect(task).toHaveBeenCalledOnce(); }); - it("runs the task when Next.js after throws", async () => { - vi.mocked(after).mockImplementation(() => { - throw new Error("after was called outside a request scope"); - }); - const errorLog = vi.spyOn(console, "error").mockImplementation(() => undefined); - const warnLog = vi.spyOn(console, "warn").mockImplementation(() => undefined); - const task = vi.fn(() => { - throw new Error("boom"); - }); + it("hands ctx.waitUntil to Next.js after from onUploadComplete", async () => { + const task = vi.fn(() => Promise.resolve("deleted")); const handlers = createRouteHandler({ router: { background: f({ blob: {} }) - .middleware((opts) => { + .middleware(() => ({})) + .onUploadComplete((opts) => { opts.ctx.waitUntil(task); - return {}; - }) - .onUploadComplete(uploadCompleteMock), + }), }, config: { token: testToken.encoded }, }); + const payload = 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 signature = await Effect.runPromise( + signPayload(payload, testToken.decoded.apiKey), + ); const res = await handlers.POST( - new NextRequest(createApiUrl("background", "upload"), { + new NextRequest(createApiUrl("background"), { method: "POST", headers: { - ...baseHeaders, - host: "localhost:3000", - "x-forwarded-proto": "http", + "uploadthing-hook": "callback", + "x-uploadthing-signature": signature, }, - body: JSON.stringify({ - files: [{ name: "foo.txt", size: 48, type: "text/plain" }], - }), + body: payload, }), ); expect(res.status).toBe(200); + expect(after).toHaveBeenCalledOnce(); + expect(task).not.toHaveBeenCalled(); + const flush = vi.mocked(after).mock.calls[0]?.[0]; + expect(typeof flush).toBe("function"); + if (typeof flush !== "function") return; + await flush(); expect(task).toHaveBeenCalledOnce(); - expect(errorLog).toHaveBeenCalled(); - errorLog.mockRestore(); - warnLog.mockRestore(); }); - it("handles a rejected task when Next.js after throws", async () => { - vi.mocked(after).mockImplementation(() => { - throw new Error("after was called outside a request scope"); - }); - const errorLog = vi.spyOn(console, "error").mockImplementation(() => undefined); - const warnLog = vi.spyOn(console, "warn").mockImplementation(() => undefined); - const unhandled: Array = []; - const onUnhandled = (reason: unknown) => { - unhandled.push(reason); - }; - process.on("unhandledRejection", onUnhandled); + it("hands ctx.waitUntil to Next.js after from onUploadError", async () => { + const task = vi.fn(() => Promise.resolve("deleted")); const handlers = createRouteHandler({ router: { background: f({ blob: {} }) - .middleware((opts) => { - opts.ctx.waitUntil(Promise.reject(new Error("delete failed"))); - return {}; + .middleware(() => ({})) + .onUploadError((opts) => { + opts.ctx.waitUntil(task); }) .onUploadComplete(uploadCompleteMock), }, config: { token: testToken.encoded }, }); + const payload = JSON.stringify({ + fileKey: "some-random-key.png", + error: "network", + }); + const signature = await Effect.runPromise( + signPayload(payload, testToken.decoded.apiKey), + ); const res = await handlers.POST( - new NextRequest(createApiUrl("background", "upload"), { + new NextRequest(createApiUrl("background"), { method: "POST", headers: { - ...baseHeaders, - host: "localhost:3000", - "x-forwarded-proto": "http", + "uploadthing-hook": "error", + "x-uploadthing-signature": signature, }, - body: JSON.stringify({ - files: [{ name: "foo.txt", size: 48, type: "text/plain" }], - }), + body: payload, }), ); - await new Promise((resolve) => setTimeout(resolve, 0)); - process.off("unhandledRejection", onUnhandled); expect(res.status).toBe(200); - expect(unhandled).toEqual([]); - errorLog.mockRestore(); - warnLog.mockRestore(); + expect(after).toHaveBeenCalledOnce(); + expect(task).not.toHaveBeenCalled(); + const flush = vi.mocked(after).mock.calls[0]?.[0]; + expect(typeof flush).toBe("function"); + if (typeof flush !== "function") return; + await flush(); + expect(task).toHaveBeenCalledOnce(); }); - it("hands ctx.waitUntil to Next.js after from onUploadComplete", async () => { + it("runs a detached dev hook once when after already flushed", async () => { + vi.mocked(after).mockImplementation((flush) => { + if (typeof flush === "function") void flush(); + }); + const warnLog = vi + .spyOn(console, "warn") + .mockImplementation(() => undefined); const task = vi.fn(() => Promise.resolve("deleted")); const handlers = createRouteHandler({ router: { @@ -547,7 +574,7 @@ describe("adapters:next", async () => { opts.ctx.waitUntil(task); }), }, - config: { token: testToken.encoded }, + config: { token: testToken.encoded, isDev: true }, }); const payload = JSON.stringify({ status: "uploaded", @@ -581,8 +608,127 @@ describe("adapters:next", async () => { ); expect(res.status).toBe(200); - expect(after).toHaveBeenCalledWith(task); - expect(task).not.toHaveBeenCalled(); + await vi.waitUntil(() => task.mock.calls.length > 0); + const flush = vi.mocked(after).mock.calls[0]?.[0]; + expect(typeof flush).toBe("function"); + if (typeof flush === "function") await flush(); + expect(task).toHaveBeenCalledOnce(); + expect(warnLog).not.toHaveBeenCalled(); + warnLog.mockRestore(); + }); + + it("runs the task once when Next.js after is missing", async () => { + const nextServer = await import("next/server"); + const originalAfter = nextServer.after; + Object.defineProperty(nextServer, "after", { + configurable: true, + value: undefined, + }); + const errorLog = vi + .spyOn(console, "error") + .mockImplementation(() => undefined); + const warnLog = vi + .spyOn(console, "warn") + .mockImplementation(() => undefined); + const unhandled: Array = []; + const onUnhandled = (reason: unknown) => { + unhandled.push(reason); + }; + process.on("unhandledRejection", onUnhandled); + const task = vi.fn(() => { + throw new Error("boom"); + }); + const handlers = createRouteHandler({ + router: { + background: f({ blob: {} }) + .middleware((opts) => { + opts.ctx.waitUntil(task); + opts.ctx.waitUntil(Promise.reject(new Error("delete failed"))); + return {}; + }) + .onUploadComplete(uploadCompleteMock), + }, + config: { token: testToken.encoded }, + }); + + const post = () => + handlers.POST( + 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" }], + }), + }), + ); + + try { + const first = await post(); + await new Promise((resolve) => setTimeout(resolve, 0)); + const second = await post(); + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(first.status).toBe(200); + expect(second.status).toBe(200); + expect(task).toHaveBeenCalledTimes(2); + expect(unhandled).toEqual([]); + expect(errorLog).toHaveBeenCalled(); + expect(warnLog).toHaveBeenCalledOnce(); + } finally { + Object.defineProperty(nextServer, "after", { + configurable: true, + value: originalAfter, + }); + process.off("unhandledRejection", onUnhandled); + errorLog.mockRestore(); + warnLog.mockRestore(); + } + }); + + it("runs the task once when Next.js after throws", async () => { + vi.mocked(after).mockImplementation(() => { + throw new Error("after was called outside a request scope"); + }); + const errorLog = vi + .spyOn(console, "error") + .mockImplementation(() => undefined); + const task = vi.fn(() => { + throw new Error("boom"); + }); + const handlers = createRouteHandler({ + router: { + background: f({ blob: {} }) + .middleware((opts) => { + opts.ctx.waitUntil(task); + return {}; + }) + .onUploadComplete(uploadCompleteMock), + }, + config: { token: testToken.encoded }, + }); + + const res = await handlers.POST( + 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" }], + }), + }), + ); + + expect(res.status).toBe(200); + expect(task).toHaveBeenCalledOnce(); + expect(errorLog).toHaveBeenCalled(); + errorLog.mockRestore(); }); }); From d7baf101a507a86e939e035fca5ebe519541ec5e Mon Sep 17 00:00:00 2001 From: Arman Charan Date: Mon, 28 Sep 2026 15:07:28 +1000 Subject: [PATCH 4/5] refactor: open the ctx.waitUntil queue in the adapter args makeAdapterHandler builds the adapter args before its first await, so the queue can be opened there without a module-level handoff. Co-authored-by: Cursor --- packages/uploadthing/src/next.ts | 21 +++------------------ 1 file changed, 3 insertions(+), 18 deletions(-) diff --git a/packages/uploadthing/src/next.ts b/packages/uploadthing/src/next.ts index 8955067de3..fa6b65c814 100644 --- a/packages/uploadthing/src/next.ts +++ b/packages/uploadthing/src/next.ts @@ -117,12 +117,6 @@ const openScheduler = (): RequestScheduler => { }; }; -/** - * Set for the synchronous start of `POST`, then captured into adapter args. - * Cleared before the handler awaits so concurrent requests do not share it. - */ -let requestScheduler: RequestScheduler | undefined; - type AdapterArgs = { req: NextRequest; ctx: RequestContext; @@ -137,7 +131,8 @@ export const createRouteHandler = ( ) => { const handler = makeAdapterHandler<[NextRequest], AdapterArgs>( (req) => { - const scheduler = requestScheduler; + // Built eagerly, while the route handler still holds the request scope. + const scheduler = req.method === "POST" ? openScheduler() : undefined; return Effect.succeed({ req, ctx: { @@ -152,15 +147,5 @@ export const createRouteHandler = ( opts, "nextjs-app", ); - - const POST = (req: NextRequest) => { - requestScheduler = openScheduler(); - try { - return handler(req); - } finally { - requestScheduler = undefined; - } - }; - - return { POST, GET: handler }; + return { POST: handler, GET: handler }; }; From d9656925eec5d6548bcfe260bcd55052418ec9e9 Mon Sep 17 00:00:00 2001 From: Arman Charan Date: Mon, 28 Sep 2026 16:17:00 +1000 Subject: [PATCH 5/5] fix: handle queued ctx.waitUntil rejections before the flush A promise handed to ctx.waitUntil is already running. Attach its rejection handler when it is queued, not when after runs. The flush also waits for tasks queued while it runs. Tests cover each hook, the rejection, the fallbacks, and a flush-time enqueue. Tests that assert on the one-time warning load a fresh adapter, so they pass in any order. Co-authored-by: Cursor --- packages/uploadthing/src/next.ts | 80 +-- .../uploadthing/test/node/adapters.test.ts | 474 ++++++++---------- 2 files changed, 253 insertions(+), 301 deletions(-) diff --git a/packages/uploadthing/src/next.ts b/packages/uploadthing/src/next.ts index fa6b65c814..60b2d7a7e8 100644 --- a/packages/uploadthing/src/next.ts +++ b/packages/uploadthing/src/next.ts @@ -67,54 +67,64 @@ const settleTask = (task: AfterTask): Promise => { } }; -type RequestScheduler = { - enqueue: (task: AfterTask) => void; +/** + * 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 the queue is opened here and hooks only enqueue. + * after that store is gone, so `after` is registered here and hooks enqueue. */ -const openScheduler = (): RequestScheduler => { - const tasks: Array = []; - let sealed = false; - let failed = false; +const openScheduler = (): Scheduler => { + const queue: Array<() => Promise> = []; + let state: SchedulerState = "queueing"; const flush = (): Promise => { - sealed = true; - const batch = tasks.splice(0, tasks.length); - return Promise.all(batch.map(settleTask)).then(() => undefined); + 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. Leave `flush` uncalled so the task runs once. - failed = true; - sealed = true; + // Thrown before registration, so `flush` never runs. + state = "unregistered"; } } else { - failed = true; - sealed = true; + state = "unregistered"; } - return { - enqueue(task) { - if (failed) { - warnTaskMayFreeze(); - void settleTask(task); - return; - } - if (sealed) { - // Development detaches callback hooks after the response. `after` - // already ran, and the dev server still finishes the task. - void settleTask(task); - return; - } - tasks.push(task); + 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 = { @@ -132,16 +142,8 @@ export const createRouteHandler = ( const handler = makeAdapterHandler<[NextRequest], AdapterArgs>( (req) => { // Built eagerly, while the route handler still holds the request scope. - const scheduler = req.method === "POST" ? openScheduler() : undefined; - return Effect.succeed({ - req, - ctx: { - waitUntil: (task) => { - if (scheduler) scheduler.enqueue(task); - else void settleTask(task); - }, - }, - }); + const scheduler = req.method === "POST" ? openScheduler() : inProcess; + return Effect.succeed({ req, ctx: { waitUntil: scheduler.enqueue } }); }, (req) => Effect.succeed(req), opts, diff --git a/packages/uploadthing/test/node/adapters.test.ts b/packages/uploadthing/test/node/adapters.test.ts index 2980feb8a7..8117366cf1 100644 --- a/packages/uploadthing/test/node/adapters.test.ts +++ b/packages/uploadthing/test/node/adapters.test.ts @@ -14,12 +14,14 @@ import { createApp, H3Event, toWebHandler } from "h3"; import { setupServer } from "msw/node"; import { afterAll, + afterEach, beforeAll, beforeEach, describe, expect, expectTypeOf, it, + onTestFinished, vi, } from "vitest"; @@ -39,6 +41,7 @@ import { 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; @@ -315,9 +318,8 @@ 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(() => { @@ -423,59 +425,11 @@ describe("adapters:next", async () => { }); }); - it("hands ctx.waitUntil to Next.js after", async () => { - const task = vi.fn(() => Promise.resolve("deleted")); - const handlers = createRouteHandler({ - router: { - background: f({ blob: {} }) - .middleware((opts) => { - opts.ctx.waitUntil(task); - return {}; - }) - .onUploadComplete(uploadCompleteMock), - }, - config: { token: testToken.encoded }, - }); - - const res = await handlers.POST( - 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" }], - }), - }), - ); + describe("ctx.waitUntil", () => { + const outsideRequestScope = "after was called outside a request scope"; + const taskFailed = "[uploadthing] ctx.waitUntil task failed."; - expect(res.status).toBe(200); - expect(after).toHaveBeenCalledOnce(); - expect(task).not.toHaveBeenCalled(); - const flush = vi.mocked(after).mock.calls[0]?.[0]; - expect(typeof flush).toBe("function"); - if (typeof flush !== "function") return; - await flush(); - expect(task).toHaveBeenCalledOnce(); - await flush(); - expect(task).toHaveBeenCalledOnce(); - }); - - it("hands ctx.waitUntil to Next.js after from onUploadComplete", async () => { - const task = vi.fn(() => Promise.resolve("deleted")); - const handlers = createRouteHandler({ - router: { - background: f({ blob: {} }) - .middleware(() => ({})) - .onUploadComplete((opts) => { - opts.ctx.waitUntil(task); - }), - }, - config: { token: testToken.encoded }, - }); - const payload = JSON.stringify({ + const uploadedPayload = JSON.stringify({ status: "uploaded", metadata: {}, origin: "https://example.com", @@ -491,244 +445,240 @@ describe("adapters:next", async () => { fileHash: "some-md5-hash", }), }); - const signature = await Effect.runPromise( - signPayload(payload, testToken.decoded.apiKey), - ); + const failedPayload = JSON.stringify({ + fileKey: "some-random-key.png", + error: "network", + }); - const res = await handlers.POST( - new NextRequest(createApiUrl("background"), { + const uploadRequest = () => + new NextRequest(createApiUrl("background", "upload"), { method: "POST", headers: { - "uploadthing-hook": "callback", - "x-uploadthing-signature": signature, + ...baseHeaders, + host: "localhost:3000", + "x-forwarded-proto": "http", }, - body: payload, - }), - ); - - expect(res.status).toBe(200); - expect(after).toHaveBeenCalledOnce(); - expect(task).not.toHaveBeenCalled(); - const flush = vi.mocked(after).mock.calls[0]?.[0]; - expect(typeof flush).toBe("function"); - if (typeof flush !== "function") return; - await flush(); - expect(task).toHaveBeenCalledOnce(); - }); - - it("hands ctx.waitUntil to Next.js after from onUploadError", async () => { - const task = vi.fn(() => Promise.resolve("deleted")); - const handlers = createRouteHandler({ - router: { - background: f({ blob: {} }) - .middleware(() => ({})) - .onUploadError((opts) => { - opts.ctx.waitUntil(task); - }) - .onUploadComplete(uploadCompleteMock), - }, - config: { token: testToken.encoded }, - }); - const payload = JSON.stringify({ - fileKey: "some-random-key.png", - error: "network", - }); - const signature = await Effect.runPromise( - signPayload(payload, testToken.decoded.apiKey), - ); + body: JSON.stringify({ + files: [{ name: "foo.txt", size: 48, type: "text/plain" }], + }), + }); - const res = await handlers.POST( + const hookRequest = async (hook: "callback" | "error", payload: string) => new NextRequest(createApiUrl("background"), { method: "POST", headers: { - "uploadthing-hook": "error", - "x-uploadthing-signature": signature, + "uploadthing-hook": hook, + "x-uploadthing-signature": await Effect.runPromise( + signPayload(payload, testToken.decoded.apiKey), + ), }, body: payload, - }), - ); + }); - expect(res.status).toBe(200); - expect(after).toHaveBeenCalledOnce(); - expect(task).not.toHaveBeenCalled(); - const flush = vi.mocked(after).mock.calls[0]?.[0]; - expect(typeof flush).toBe("function"); - if (typeof flush !== "function") return; - await flush(); - expect(task).toHaveBeenCalledOnce(); - }); + /** 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 }, + }); + }; - it("runs a detached dev hook once when after already flushed", async () => { - vi.mocked(after).mockImplementation((flush) => { - if (typeof flush === "function") void flush(); - }); - const warnLog = vi - .spyOn(console, "warn") - .mockImplementation(() => undefined); - const task = vi.fn(() => Promise.resolve("deleted")); - const handlers = createRouteHandler({ - router: { - background: f({ blob: {} }) - .middleware(() => ({})) - .onUploadComplete((opts) => { - opts.ctx.waitUntil(task); - }), - }, - config: { token: testToken.encoded, isDev: true }, + /** 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(); }); - const payload = 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", - }), + + 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(); }); - const signature = await Effect.runPromise( - signPayload(payload, testToken.decoded.apiKey), - ); - const res = await handlers.POST( - new NextRequest(createApiUrl("background"), { - method: "POST", - headers: { - "uploadthing-hook": "callback", - "x-uploadthing-signature": signature, - }, - body: payload, - }), + 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(); + }, ); - expect(res.status).toBe(200); - await vi.waitUntil(() => task.mock.calls.length > 0); - const flush = vi.mocked(after).mock.calls[0]?.[0]; - expect(typeof flush).toBe("function"); - if (typeof flush === "function") await flush(); - expect(task).toHaveBeenCalledOnce(); - expect(warnLog).not.toHaveBeenCalled(); - warnLog.mockRestore(); - }); + 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"); + }); + }); - it("runs the task once when Next.js after is missing", async () => { - const nextServer = await import("next/server"); - const originalAfter = nextServer.after; - Object.defineProperty(nextServer, "after", { - configurable: true, - value: undefined, - }); - const errorLog = vi - .spyOn(console, "error") - .mockImplementation(() => undefined); - const warnLog = vi - .spyOn(console, "warn") - .mockImplementation(() => undefined); - const unhandled: Array = []; - const onUnhandled = (reason: unknown) => { - unhandled.push(reason); - }; - process.on("unhandledRejection", onUnhandled); - const task = vi.fn(() => { - 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" }), + ); }); - const handlers = createRouteHandler({ - router: { - background: f({ blob: {} }) - .middleware((opts) => { - opts.ctx.waitUntil(task); - opts.ctx.waitUntil(Promise.reject(new Error("delete failed"))); - return {}; - }) - .onUploadComplete(uploadCompleteMock), - }, - config: { token: testToken.encoded }, + + 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(); }); - const post = () => - handlers.POST( - 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" }], - }), - }), + 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(); - try { - const first = await post(); - await new Promise((resolve) => setTimeout(resolve, 0)); - const second = await post(); - await new Promise((resolve) => setTimeout(resolve, 0)); + expect(res.status).toBe(200); + expect(task).toHaveBeenCalledOnce(); + expect(warnLog).not.toHaveBeenCalled(); + }); - expect(first.status).toBe(200); - expect(second.status).toBe(200); - expect(task).toHaveBeenCalledTimes(2); - expect(unhandled).toEqual([]); - expect(errorLog).toHaveBeenCalled(); - expect(warnLog).toHaveBeenCalledOnce(); - } finally { + 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: originalAfter, + 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, }); - process.off("unhandledRejection", onUnhandled); - errorLog.mockRestore(); - warnLog.mockRestore(); - } - }); - it("runs the task once when Next.js after throws", async () => { - vi.mocked(after).mockImplementation(() => { - throw new Error("after was called outside a request scope"); - }); - const errorLog = vi - .spyOn(console, "error") - .mockImplementation(() => undefined); - const task = vi.fn(() => { - throw new Error("boom"); - }); - const handlers = createRouteHandler({ - router: { - background: f({ blob: {} }) - .middleware((opts) => { - opts.ctx.waitUntil(task); - return {}; - }) - .onUploadComplete(uploadCompleteMock), - }, - config: { token: testToken.encoded }, + 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(); }); - const res = await handlers.POST( - 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" }], - }), - }), - ); + 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, + }); - expect(res.status).toBe(200); - expect(task).toHaveBeenCalledOnce(); - expect(errorLog).toHaveBeenCalled(); - errorLog.mockRestore(); + const res = await handlers.POST(uploadRequest()); + await endResponse(); + + expect(res.status).toBe(200); + expect(task).toHaveBeenCalledOnce(); + expect(warnLog).toHaveBeenCalledOnce(); + }); }); });