From 3dee612c323e4d928c60de01e911b742d91d0b29 Mon Sep 17 00:00:00 2001 From: 7ttp <117663341+7ttp@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:19:44 +0530 Subject: [PATCH] refactor(cli): cover functions `serve-main-bundler` with effect lint (CLI-2592) --- .oxlintrc.effect.json | 3 + apps/cli/scripts/build-binary.ts | 5 +- apps/cli/scripts/build.ts | 5 +- .../serve-main-bundler.integration.test.ts | 881 ++++++++++-------- .../shared/functions/serve-main-bundler.ts | 31 +- .../functions/serve-main-bundler.unit.test.ts | 37 +- .../functions/serve-main-offline.e2e.test.ts | 6 +- apps/cli/src/shared/functions/serve.ts | 26 +- .../src/shared/functions/serve.unit.test.ts | 6 +- 9 files changed, 569 insertions(+), 431 deletions(-) diff --git a/.oxlintrc.effect.json b/.oxlintrc.effect.json index a605cd0500..0834ad9fba 100644 --- a/.oxlintrc.effect.json +++ b/.oxlintrc.effect.json @@ -28,6 +28,9 @@ "!apps/cli/src/shared/functions/functions-docker.unit.test.ts", "!apps/cli/src/shared/functions/functions.shared.ts", "!apps/cli/src/shared/functions/functions.shared.unit.test.ts", + "!apps/cli/src/shared/functions/serve-main-bundler.integration.test.ts", + "!apps/cli/src/shared/functions/serve-main-bundler.ts", + "!apps/cli/src/shared/functions/serve-main-bundler.unit.test.ts", "!apps/cli/src/shared/functions/serve.ts", "!apps/cli/src/shared/functions/serve.unit.test.ts", "!apps/cli/src/shared/git/**", diff --git a/apps/cli/scripts/build-binary.ts b/apps/cli/scripts/build-binary.ts index f2dcb02b00..28ed04fee9 100644 --- a/apps/cli/scripts/build-binary.ts +++ b/apps/cli/scripts/build-binary.ts @@ -1,3 +1,4 @@ +import { Effect } from "effect"; import { bundleServeMainTemplate } from "../src/shared/functions/serve-main-bundler.ts"; import { compileOptions, stackReleaseDefine } from "./compile-options.ts"; @@ -23,7 +24,9 @@ const result = await Bun.build({ define: { SUPABASE_CLI_VERSION: JSON.stringify(packageJson.version), ...(await stackReleaseDefine()), - SUPABASE_FUNCTIONS_SERVE_MAIN_TEMPLATE: JSON.stringify(await bundleServeMainTemplate()), + SUPABASE_FUNCTIONS_SERVE_MAIN_TEMPLATE: JSON.stringify( + await Effect.runPromise(bundleServeMainTemplate()), + ), // Skips msgpackr's native addon probe at the build host's path, which can hang macOS startup. "process.env.MSGPACKR_NATIVE_ACCELERATION_DISABLED": JSON.stringify("true"), }, diff --git a/apps/cli/scripts/build.ts b/apps/cli/scripts/build.ts index 87131b408e..52ee98efff 100644 --- a/apps/cli/scripts/build.ts +++ b/apps/cli/scripts/build.ts @@ -4,6 +4,7 @@ import { copyFile, mkdir, readFile, rm, writeFile } from "node:fs/promises"; import path from "node:path"; import process from "node:process"; import { parseArgs } from "node:util"; +import { Effect } from "effect"; import { bundleServeMainTemplate } from "../src/shared/functions/serve-main-bundler.ts"; import { compileOptions, stackReleaseDefine } from "./compile-options.ts"; import { darwinBinaries, MACOS_IDENTIFIERS } from "./macos-signing.ts"; @@ -89,7 +90,9 @@ const distDir = path.join(root, "dist"); const goSource = path.resolve(root, "apps/cli-go"); const buildDefines = { ...(await stackReleaseDefine()), - SUPABASE_FUNCTIONS_SERVE_MAIN_TEMPLATE: JSON.stringify(await bundleServeMainTemplate()), + SUPABASE_FUNCTIONS_SERVE_MAIN_TEMPLATE: JSON.stringify( + await Effect.runPromise(bundleServeMainTemplate()), + ), SUPABASE_CLI_POSTHOG_KEY: JSON.stringify(process.env.POSTHOG_API_KEY ?? ""), SUPABASE_CLI_POSTHOG_HOST: JSON.stringify(process.env.POSTHOG_ENDPOINT ?? ""), // Skips msgpackr's startup probe for its native addon at the build host's store path, which diff --git a/apps/cli/src/shared/functions/serve-main-bundler.integration.test.ts b/apps/cli/src/shared/functions/serve-main-bundler.integration.test.ts index bef3c5392a..ec046c2e75 100644 --- a/apps/cli/src/shared/functions/serve-main-bundler.integration.test.ts +++ b/apps/cli/src/shared/functions/serve-main-bundler.integration.test.ts @@ -1,5 +1,6 @@ import { createContext, SourceTextModule } from "node:vm"; import { describe, expect, it } from "@effect/vitest"; +import { Cause, Data, Deferred, Effect, Exit, Predicate, Schema, Scope } from "effect"; import { exportJWK, generateKeyPair, SignJWT } from "jose"; import { bundleServeMainTemplate } from "./serve-main-bundler.ts"; @@ -9,11 +10,14 @@ type ServeOptions = { }; type FetchImplementation = (input: string | URL | Request, init?: RequestInit) => Promise; -class InvalidWorkerCreation extends Error {} -class InvalidWorkerResponse extends Error {} -class WorkerRequestCancelled extends Error {} -class NotFoundError extends Error {} -class WorkerAlreadyRetired extends Error {} +class InvalidWorkerCreation extends Data.Error {} +class InvalidWorkerResponse extends Data.Error {} +class WorkerRequestCancelled extends Data.Error {} +class NotFoundError extends Data.Error {} +class WorkerAlreadyRetired extends Data.Error {} +class BootstrapLoadError extends Data.TaggedError("BootstrapLoadError")<{ + readonly cause: unknown; +}> {} type LoadOptions = { readonly errors?: Record Error>; @@ -26,7 +30,9 @@ type LoadOptions = { readonly metricError?: unknown; }; -const load = async ( +const encodeJson = Schema.encodeEffect(Schema.fromJsonString(Schema.Unknown)); + +const load = Effect.fnUntraced(function* ( bundle: string, env: Record, worker: Record, @@ -40,7 +46,7 @@ const load = async ( lstatError, metricError, }: LoadOptions = {}, -) => { +) { let options: ServeOptions | undefined; const state: { createOptions?: Record } = {}; const envRecord = { ...env }; @@ -48,10 +54,10 @@ const load = async ( Deno: { env: { get: (name: string) => envRecord[name], toObject: () => envRecord }, cwd: () => "/functions", - lstat: async () => { - if (lstatError !== undefined) throw lstatError; - return { isFile: true, isDirectory: false, isSymlink: false }; - }, + lstat: () => + lstatError === undefined + ? Promise.resolve({ isFile: true, isDirectory: false, isSymlink: false }) + : Promise.reject(lstatError), makeTempDirSync: () => "/tmp/worker", errors, version: { deno: "test" }, @@ -101,13 +107,22 @@ const load = async ( context: createContext(sandbox), identifier: "cli-serve.main.bundle.js", }); - await module.link(() => { - throw new Error("unexpected external import"); + yield* Effect.tryPromise({ + try: () => + module.link(() => { + throw new Error("unexpected external import"); + }), + catch: (cause) => new BootstrapLoadError({ cause }), + }); + yield* Effect.tryPromise({ + try: () => module.evaluate(), + catch: (cause) => new BootstrapLoadError({ cause }), }); - await module.evaluate(); - if (options === undefined) throw new Error("bootstrap did not register server"); + if (options === undefined) { + return yield* new BootstrapLoadError({ cause: "bootstrap did not register server" }); + } return { options, envRecord, state }; -}; +}); const baseEnv = (config: string) => ({ SUPABASE_INTERNAL_HOST_PORT: "8081", @@ -117,163 +132,196 @@ const baseEnv = (config: string) => ({ SUPABASE_INTERNAL_FUNCTIONS_ROOT: "/functions", }); +const okWorker = { fetch: () => Promise.resolve(new Response("ok")) }; + describe("CLI functions bootstrap bundle", () => { - it("fails startup for missing URL and malformed required config", async () => { - const bundle = await bundleServeMainTemplate(); - const config = JSON.stringify({ - hello: { - entrypointPath: "hello/index.ts", - importMapPath: "", - staticFiles: [], - verifyJWT: false, - }, - }); - await expect( - load( - bundle, - { ...baseEnv(config), SUPABASE_URL: "" }, - { fetch: async () => new Response("ok") }, - ), - ).rejects.toThrow(); - await expect( - load( - bundle, - { ...baseEnv("{"), SUPABASE_INTERNAL_FUNCTIONS_CONFIG: "{" }, - { fetch: async () => new Response("ok") }, - ), - ).rejects.toThrow(); - }); + it.effect("fails startup for missing URL and malformed required config", () => + Effect.gen(function* () { + const bundle = yield* bundleServeMainTemplate(); + const config = yield* encodeJson({ + hello: { + entrypointPath: "hello/index.ts", + importMapPath: "", + staticFiles: [], + verifyJWT: false, + }, + }); + const missingUrl = yield* Effect.exit( + load(bundle, { ...baseEnv(config), SUPABASE_URL: "" }, okWorker), + ); + expect(Exit.isFailure(missingUrl)).toBe(true); + const malformedConfig = yield* Effect.exit( + load(bundle, { ...baseEnv("{"), SUPABASE_INTERNAL_FUNCTIONS_CONFIG: "{" }, okWorker), + ); + expect(Exit.isFailure(malformedConfig)).toBe(true); + }), + ); - it("authenticates and forwards request, environment, and worker options", async () => { - const bundle = await bundleServeMainTemplate(); - const config = JSON.stringify({ - hello: { - entrypointPath: "hello/index.ts", - importMapPath: "", - staticFiles: [], - verifyJWT: true, - env: { FUNCTION_ONLY: "yes", SUPABASE_BLOCKED: "no" }, - }, - }); - let received: Request | undefined; - const loaded = await load( - bundle, - { ...baseEnv(config), KEEP: "yes", HOME: "hidden", SUPABASE_INTERNAL_SECRET_KEY: "internal" }, - { - fetch: async (request: Request) => { - received = request; - return new Response(request.body); + it.effect("authenticates and forwards request, environment, and worker options", () => + Effect.gen(function* () { + const bundle = yield* bundleServeMainTemplate(); + const config = yield* encodeJson({ + hello: { + entrypointPath: "hello/index.ts", + importMapPath: "", + staticFiles: [], + verifyJWT: true, + env: { FUNCTION_ONLY: "yes", SUPABASE_BLOCKED: "no" }, }, - }, - ); - const token = await new SignJWT({ sub: "test" }) - .setProtectedHeader({ alg: "HS256" }) - .sign(new TextEncoder().encode("secret")); - const response = await loaded.options.handler( - new Request("http://localhost/hello", { - method: "POST", - body: "body", - headers: { Authorization: `Bearer ${token}`, "sb-api-key": "remove", "x-tag": "tag" }, - }), - ); - expect(response.status).toBe(200); - expect(received?.headers.get("sb-api-key")).toBeNull(); - expect(received?.headers.get("x-tag")).toBe("tag"); - expect(await response.text()).toBe("body"); - expect(loaded.state.createOptions?.envVars).toEqual( - expect.arrayContaining([ - ["KEEP", "yes"], - ["FUNCTION_ONLY", "yes"], - ]), - ); - }); + }); + let received: Request | undefined; + const loaded = yield* load( + bundle, + { + ...baseEnv(config), + KEEP: "yes", + HOME: "hidden", + SUPABASE_INTERNAL_SECRET_KEY: "internal", + }, + { + fetch: (request: Request) => { + received = request; + return Promise.resolve(new Response(request.body)); + }, + }, + ); + const token = yield* Effect.promise(() => + new SignJWT({ sub: "test" }) + .setProtectedHeader({ alg: "HS256" }) + .sign(new TextEncoder().encode("secret")), + ); + const response = yield* Effect.promise(() => + loaded.options.handler( + new Request("http://localhost/hello", { + method: "POST", + body: "body", + headers: { Authorization: `Bearer ${token}`, "sb-api-key": "remove", "x-tag": "tag" }, + }), + ), + ); + expect(response.status).toBe(200); + expect(received?.headers.get("sb-api-key")).toBeNull(); + expect(received?.headers.get("x-tag")).toBe("tag"); + expect(yield* Effect.promise(() => response.text())).toBe("body"); + expect(loaded.state.createOptions?.envVars).toEqual( + expect.arrayContaining([ + ["KEEP", "yes"], + ["FUNCTION_ONLY", "yes"], + ]), + ); + }), + ); - it("rejects an invalid token", async () => { - const bundle = await bundleServeMainTemplate(); - const config = JSON.stringify({ - hello: { - entrypointPath: "hello/index.ts", - importMapPath: "", - staticFiles: [], - verifyJWT: true, - }, - }); - const loaded = await load(bundle, baseEnv(config), { - fetch: async () => new Response("should not run"), - }); - const response = await loaded.options.handler( - new Request("http://localhost/hello", { headers: { Authorization: "Bearer invalid" } }), - ); - expect(response.status).toBe(401); - expect(await response.json()).toMatchObject({ code: "UNAUTHORIZED_INVALID_JWT_FORMAT" }); - }); + it.effect("rejects an invalid token", () => + Effect.gen(function* () { + const bundle = yield* bundleServeMainTemplate(); + const config = yield* encodeJson({ + hello: { + entrypointPath: "hello/index.ts", + importMapPath: "", + staticFiles: [], + verifyJWT: true, + }, + }); + const loaded = yield* load(bundle, baseEnv(config), { + fetch: () => Promise.resolve(new Response("should not run")), + }); + const response = yield* Effect.promise(() => + loaded.options.handler( + new Request("http://localhost/hello", { headers: { Authorization: "Bearer invalid" } }), + ), + ); + expect(response.status).toBe(401); + expect(yield* Effect.promise(() => response.json())).toMatchObject({ + code: "UNAUTHORIZED_INVALID_JWT_FORMAT", + }); + }), + ); - it("retains non-abort handler failures", async () => { - const bundle = await bundleServeMainTemplate(); - const metricError = new Error("metrics unavailable"); - const loaded = await load( - bundle, - baseEnv("{}"), - { fetch: async () => new Response("ok") }, - { + it.effect("retains non-abort handler failures", () => + Effect.gen(function* () { + const bundle = yield* bundleServeMainTemplate(); + const metricError = new Error("metrics unavailable"); + const loaded = yield* load(bundle, baseEnv("{}"), okWorker, { metricError, - }, - ); - const failure = await loaded.options - .handler(new Request("http://localhost/_internal/metric")) - .catch((error) => error); - expect(failure).toMatchObject({ _tag: "BootstrapOperationError", cause: metricError }); - }); + }); + const exit = yield* Effect.exit( + Effect.promise(() => + loaded.options.handler(new Request("http://localhost/_internal/metric")), + ), + ); + const failure = Exit.isFailure(exit) ? Cause.squash(exit.cause) : undefined; + expect(Predicate.isTagged(failure, "BootstrapOperationError")).toBe(true); + expect(failure).toMatchObject({ cause: metricError }); + }), + ); - it.each([ - [InvalidWorkerCreation, 503, "BOOT_ERROR"], - [InvalidWorkerResponse, 500, "WORKER_ERROR"], - [WorkerRequestCancelled, 546, "WORKER_LIMIT"], - ] as const)("maps %s worker failure to the runtime response", async (ErrorType, status, code) => { - const bundle = await bundleServeMainTemplate(); - const config = JSON.stringify({ - hello: { - entrypointPath: "hello/index.ts", - importMapPath: "", - staticFiles: ["public/**"], - verifyJWT: false, - }, - }); - const failure = new ErrorType(); - const loaded = await load( - bundle, - baseEnv(config), - { - fetch: async () => { - throw failure; + it.effect.each([ + { + name: "InvalidWorkerCreation", + ErrorType: InvalidWorkerCreation, + status: 503, + code: "BOOT_ERROR", + }, + { + name: "InvalidWorkerResponse", + ErrorType: InvalidWorkerResponse, + status: 500, + code: "WORKER_ERROR", + }, + { + name: "WorkerRequestCancelled", + ErrorType: WorkerRequestCancelled, + status: 546, + code: "WORKER_LIMIT", + }, + ])("maps $name worker failure to the runtime response", ({ ErrorType, status, code }) => + Effect.gen(function* () { + const bundle = yield* bundleServeMainTemplate(); + const config = yield* encodeJson({ + hello: { + entrypointPath: "hello/index.ts", + importMapPath: "", + staticFiles: ["public/**"], + verifyJWT: false, }, - }, - { - errors: { InvalidWorkerCreation, InvalidWorkerResponse, WorkerRequestCancelled }, - creationError: ErrorType === InvalidWorkerCreation ? failure : undefined, - }, - ); - const response = await loaded.options.handler(new Request("http://localhost/hello")); - expect(response.status).toBe(status); - expect(await response.json()).toMatchObject({ code }); - }); + }); + const failure = new ErrorType(); + const loaded = yield* load( + bundle, + baseEnv(config), + { + fetch: () => Promise.reject(failure), + }, + { + errors: { InvalidWorkerCreation, InvalidWorkerResponse, WorkerRequestCancelled }, + creationError: ErrorType === InvalidWorkerCreation ? failure : undefined, + }, + ); + const response = yield* Effect.promise(() => + loaded.options.handler(new Request("http://localhost/hello")), + ); + expect(response.status).toBe(status); + expect(yield* Effect.promise(() => response.json())).toMatchObject({ code }); + }), + ); describe("retired worker dispatch", () => { - const config = JSON.stringify({ + const config = { hello: { entrypointPath: "hello/index.ts", importMapPath: "", staticFiles: [], verifyJWT: false, }, - }); + }; // Every create() hands out a distinct worker: worker n always rejects with failures[n - 1] when // one is given, otherwise it answers with its own number so the response names the worker. - const serve = async (failures: ReadonlyArray) => { + const serve = Effect.fnUntraced(function* (failures: ReadonlyArray) { let creates = 0; - const loaded = await load( - await bundleServeMainTemplate(), - baseEnv(config), + const loaded = yield* load( + yield* bundleServeMainTemplate(), + baseEnv(yield* encodeJson(config)), {}, { errors: { WorkerAlreadyRetired, InvalidWorkerResponse }, @@ -281,54 +329,76 @@ describe("CLI functions bootstrap bundle", () => { const worker = ++creates; const failure = failures[worker - 1]; return { - fetch: async (request: Request) => { - if (failure !== undefined) throw failure; - return new Response( - `fn-ok worker-${worker} ${request.method} ${await request.text()}`, - ); - }, + fetch: (request: Request) => + failure !== undefined + ? Promise.reject(failure) + : request + .text() + .then( + (text) => new Response(`fn-ok worker-${worker} ${request.method} ${text}`), + ), }; }, }, ); return { loaded, creates: () => creates }; - }; - - it("serves a bodyless request with a fresh worker after WorkerAlreadyRetired", async () => { - const { loaded, creates } = await serve([new WorkerAlreadyRetired()]); - const response = await loaded.options.handler(new Request("http://localhost/hello")); - expect(response.status).toBe(200); - expect(await response.text()).toBe("fn-ok worker-2 GET "); - expect(creates()).toBe(2); }); - it("does not retry a second consecutive WorkerAlreadyRetired", async () => { - const { loaded, creates } = await serve([ - new WorkerAlreadyRetired(), - new WorkerAlreadyRetired(), - ]); - const response = await loaded.options.handler(new Request("http://localhost/hello")); - expect(response.status).toBe(500); - expect(creates()).toBe(2); - }); + it.effect("serves a bodyless request with a fresh worker after WorkerAlreadyRetired", () => + Effect.gen(function* () { + const { loaded, creates } = yield* serve([new WorkerAlreadyRetired()]); + const response = yield* Effect.promise(() => + loaded.options.handler(new Request("http://localhost/hello")), + ); + expect(response.status).toBe(200); + expect(yield* Effect.promise(() => response.text())).toBe("fn-ok worker-2 GET "); + expect(creates()).toBe(2); + }), + ); - it("does not retry other worker failures", async () => { - const { loaded, creates } = await serve([new InvalidWorkerResponse()]); - const response = await loaded.options.handler(new Request("http://localhost/hello")); - expect(response.status).toBe(500); - expect(await response.json()).toMatchObject({ code: "WORKER_ERROR" }); - expect(creates()).toBe(1); - }); + it.effect("does not retry a second consecutive WorkerAlreadyRetired", () => + Effect.gen(function* () { + const { loaded, creates } = yield* serve([ + new WorkerAlreadyRetired(), + new WorkerAlreadyRetired(), + ]); + const response = yield* Effect.promise(() => + loaded.options.handler(new Request("http://localhost/hello")), + ); + expect(response.status).toBe(500); + expect(creates()).toBe(2); + }), + ); - it("does not replay a request whose body was already forwarded", async () => { - const { loaded, creates } = await serve([new WorkerAlreadyRetired()]); - const response = await loaded.options.handler( - new Request("http://localhost/hello", { method: "POST", body: "payload" }), - ); - expect(response.status).toBe(500); - expect(await response.json()).toMatchObject({ code: "Internal Server Error" }); - expect(creates()).toBe(1); - }); + it.effect("does not retry other worker failures", () => + Effect.gen(function* () { + const { loaded, creates } = yield* serve([new InvalidWorkerResponse()]); + const response = yield* Effect.promise(() => + loaded.options.handler(new Request("http://localhost/hello")), + ); + expect(response.status).toBe(500); + expect(yield* Effect.promise(() => response.json())).toMatchObject({ + code: "WORKER_ERROR", + }); + expect(creates()).toBe(1); + }), + ); + + it.effect("does not replay a request whose body was already forwarded", () => + Effect.gen(function* () { + const { loaded, creates } = yield* serve([new WorkerAlreadyRetired()]); + const response = yield* Effect.promise(() => + loaded.options.handler( + new Request("http://localhost/hello", { method: "POST", body: "payload" }), + ), + ); + expect(response.status).toBe(500); + expect(yield* Effect.promise(() => response.json())).toMatchObject({ + code: "Internal Server Error", + }); + expect(creates()).toBe(1); + }), + ); }); describe("request body ownership", () => { @@ -352,222 +422,267 @@ describe("CLI functions bootstrap bundle", () => { return { body, progress }; }; const functionConfig = (verifyJWT: boolean) => - JSON.stringify({ + encodeJson({ hello: { entrypointPath: "hello/index.ts", importMapPath: "", staticFiles: [], verifyJWT }, }); - it("reads the rest of a request body the worker abandons", async () => { - const loaded = await load(await bundleServeMainTemplate(), baseEnv(functionConfig(false)), { - fetch: async (request: Request) => { - const reader = request.body?.getReader(); - await reader?.read(); - // Not awaited: if the body were shared with the incoming request again, Bun - // would never settle this cancel and the test would hang instead of failing. - void reader?.cancel(); - return new Response("rejected", { status: 400 }); - }, - }); - const { body, progress } = upload(); + it.effect("reads the rest of a request body the worker abandons", () => + Effect.gen(function* () { + const loaded = yield* load( + yield* bundleServeMainTemplate(), + baseEnv(yield* functionConfig(false)), + { + fetch: (request: Request) => { + const reader = request.body?.getReader(); + return Promise.resolve(reader?.read()).then(() => { + // Not awaited: if the body were shared with the incoming request again, Bun + // would never settle this cancel and the test would hang instead of failing. + void reader?.cancel(); + return new Response("rejected", { status: 400 }); + }); + }, + }, + ); + const { body, progress } = upload(); - const response = await loaded.options.handler( - new Request("http://localhost/hello", { method: "POST", body, duplex: "half" }), - ); + const response = yield* Effect.promise(() => + loaded.options.handler( + new Request("http://localhost/hello", { method: "POST", body, duplex: "half" }), + ), + ); - expect(progress).toEqual({ readToEnd: true, cancelled: false }); - expect(response.status).toBe(400); - expect(await response.text()).toBe("rejected"); - }); + expect(progress).toEqual({ readToEnd: true, cancelled: false }); + expect(response.status).toBe(400); + expect(yield* Effect.promise(() => response.text())).toBe("rejected"); + }), + ); - it("settles an aborted request while the abandoned body read is pending", async () => { - const { promise: readPending, resolve: markReadPending } = Promise.withResolvers(); - const pendingRead = Promise.withResolvers().promise; - let pulls = 0; - const body = new ReadableStream({ - pull: (streamController) => { - pulls += 1; - if (pulls === 1) { - streamController.enqueue(new Uint8Array([1])); - return; - } - markReadPending(); - return pendingRead; - }, - }); - const loaded = await load(await bundleServeMainTemplate(), baseEnv(functionConfig(false)), { - fetch: async () => new Response("worker should not run"), - }); - const controller = new AbortController(); - const pending = loaded.options.handler( - new Request("http://localhost/missing", { - method: "POST", - body, - duplex: "half", - signal: controller.signal, - }), - ); + it.effect("settles an aborted request while the abandoned body read is pending", () => + Effect.gen(function* () { + const readPending = yield* Deferred.make(); + const hungRead = yield* Deferred.make(); + yield* Effect.addFinalizer(() => Deferred.interrupt(hungRead)); + const pendingRead = Effect.runPromiseWith(yield* Effect.context())( + Deferred.await(hungRead), + ); + let pulls = 0; + const body = new ReadableStream({ + pull: (streamController) => { + pulls += 1; + if (pulls === 1) { + streamController.enqueue(new Uint8Array([1])); + return; + } + Deferred.doneUnsafe(readPending, Effect.void); + return pendingRead; + }, + }); + const loaded = yield* load( + yield* bundleServeMainTemplate(), + baseEnv(yield* functionConfig(false)), + { + fetch: () => Promise.resolve(new Response("worker should not run")), + }, + ); + const requestScope = yield* Scope.fork(yield* Effect.scope); + const signal = yield* Scope.provide(Effect.abortSignal, requestScope); + const pending = loaded.options.handler( + new Request("http://localhost/missing", { + method: "POST", + body, + duplex: "half", + signal, + }), + ); - await readPending; - controller.abort(); + yield* Deferred.await(readPending); + yield* Scope.close(requestScope, Exit.void); - await expect(pending).resolves.toMatchObject({ status: 499 }); - }); + expect(yield* Effect.promise(() => pending)).toMatchObject({ status: 499 }); + }), + ); - it.each([ + it.effect.each([ ["an invalid token", "/hello", 401], ["an unknown function", "/missing", 404], - ] as const)("reads the whole upload before rejecting %s", async (_, path, status) => { - const loaded = await load(await bundleServeMainTemplate(), baseEnv(functionConfig(true)), { - fetch: async () => new Response("should not run"), - }); - const { body, progress } = upload(); - - const response = await loaded.options.handler( - new Request(`http://localhost${path}`, { - method: "POST", - body, - duplex: "half", - headers: { Authorization: "Bearer invalid" }, - }), - ); + ] as const)("reads the whole upload before rejecting %s", ([, path, status]) => + Effect.gen(function* () { + const loaded = yield* load( + yield* bundleServeMainTemplate(), + baseEnv(yield* functionConfig(true)), + { + fetch: () => Promise.resolve(new Response("should not run")), + }, + ); + const { body, progress } = upload(); - expect(progress).toEqual({ readToEnd: true, cancelled: false }); - expect(response.status).toBe(status); - }); - }); + const response = yield* Effect.promise(() => + loaded.options.handler( + new Request(`http://localhost${path}`, { + method: "POST", + body, + duplex: "half", + headers: { Authorization: "Bearer invalid" }, + }), + ), + ); - it("does not fetch after an aborted pending worker creation", async () => { - const bundle = await bundleServeMainTemplate(); - const config = JSON.stringify({ - hello: { - entrypointPath: "hello/index.ts", - importMapPath: "", - staticFiles: [], - verifyJWT: false, - }, - }); - let fetchCalls = 0; - let resolveCreation!: (worker: { fetch(request: Request): Promise }) => void; - const workerReady = new Promise<{ fetch(request: Request): Promise }>((resolve) => { - resolveCreation = resolve; - }); - let createInvoked!: () => void; - const createCalled = new Promise((resolve) => { - createInvoked = resolve; - }); - const loaded = await load( - bundle, - baseEnv(config), - { - fetch: async () => { - fetchCalls += 1; - return new Response("unreachable"); - }, - }, - { onCreate: createInvoked, creation: workerReady }, - ); - const controller = new AbortController(); - const pending = loaded.options.handler( - new Request("http://localhost/hello", { signal: controller.signal }), + expect(progress).toEqual({ readToEnd: true, cancelled: false }); + expect(response.status).toBe(status); + }), ); - await createCalled; - controller.abort(); - await expect(pending).resolves.toMatchObject({ status: 499 }); - resolveCreation({ - fetch: async () => { - fetchCalls += 1; - return new Response("unreachable"); - }, - }); - await workerReady; - expect(fetchCalls).toBe(0); }); - it("authenticates ES256 with injected keys and owned remote fallback", async () => { - const { publicKey, privateKey } = await generateKeyPair("ES256"); - const publicJwk = await exportJWK(publicKey); - const token = await new SignJWT({ sub: "test" }) - .setProtectedHeader({ alg: "ES256", kid: "test-key" }) - .sign(privateKey); - const config = JSON.stringify({ - hello: { - entrypointPath: "hello/index.ts", - importMapPath: "", - staticFiles: [], - verifyJWT: true, - }, - }); - const bundle = await bundleServeMainTemplate(); - const worker = { fetch: async () => new Response("ok") }; - const injected = await load( - bundle, - { - ...baseEnv(config), - SUPABASE_JWKS: JSON.stringify({ - keys: [{ ...publicJwk, kid: "test-key", alg: "ES256", use: "sig" }], - }), - }, - worker, - ); - const injectedResponse = await injected.options.handler( - new Request("http://localhost/hello", { headers: { Authorization: `Bearer ${token}` } }), - ); - expect(injectedResponse.status).toBe(200); - - for (const jwks of [undefined, "{"] as const) { + it.effect("does not fetch after an aborted pending worker creation", () => + Effect.gen(function* () { + const bundle = yield* bundleServeMainTemplate(); + const config = yield* encodeJson({ + hello: { + entrypointPath: "hello/index.ts", + importMapPath: "", + staticFiles: [], + verifyJWT: false, + }, + }); + const runPromise = Effect.runPromiseWith(yield* Effect.context()); let fetchCalls = 0; - const remote = await load( + const workerCreation = yield* Deferred.make<{ fetch(request: Request): Promise }>(); + yield* Effect.addFinalizer(() => Deferred.interrupt(workerCreation)); + const createCalled = yield* Deferred.make(); + const workerReady = runPromise(Deferred.await(workerCreation)); + const loaded = yield* load( bundle, + baseEnv(config), { - ...baseEnv(config), - ...(jwks === undefined ? {} : { SUPABASE_JWKS: jwks }), + fetch: () => { + fetchCalls += 1; + return Promise.resolve(new Response("unreachable")); + }, }, - worker, { - fetchImpl: async () => { - fetchCalls += 1; - return new Response( - JSON.stringify({ - keys: [{ ...publicJwk, kid: "test-key", alg: "ES256", use: "sig" }], - }), - { headers: { "content-type": "application/json" } }, - ); + onCreate: () => { + Deferred.doneUnsafe(createCalled, Effect.void); }, + creation: workerReady, }, ); - const remoteResponse = await remote.options.handler( - new Request("http://localhost/hello", { headers: { Authorization: `Bearer ${token}` } }), + const requestScope = yield* Scope.fork(yield* Effect.scope); + const signal = yield* Scope.provide(Effect.abortSignal, requestScope); + const pending = loaded.options.handler(new Request("http://localhost/hello", { signal })); + yield* Deferred.await(createCalled); + yield* Scope.close(requestScope, Exit.void); + expect(yield* Effect.promise(() => pending)).toMatchObject({ status: 499 }); + yield* Deferred.succeed(workerCreation, { + fetch: () => { + fetchCalls += 1; + return Promise.resolve(new Response("unreachable")); + }, + }); + yield* Effect.promise(() => workerReady); + expect(fetchCalls).toBe(0); + }), + ); + + it.effect("authenticates ES256 with injected keys and owned remote fallback", () => + Effect.gen(function* () { + const { publicKey, privateKey } = yield* Effect.promise(() => generateKeyPair("ES256")); + const publicJwk = yield* Effect.promise(() => exportJWK(publicKey)); + const token = yield* Effect.promise(() => + new SignJWT({ sub: "test" }) + .setProtectedHeader({ alg: "ES256", kid: "test-key" }) + .sign(privateKey), ); - expect(remoteResponse.status).toBe(200); - const secondResponse = await remote.options.handler( - new Request("http://localhost/hello", { headers: { Authorization: `Bearer ${token}` } }), + const config = yield* encodeJson({ + hello: { + entrypointPath: "hello/index.ts", + importMapPath: "", + staticFiles: [], + verifyJWT: true, + }, + }); + const keys = yield* encodeJson({ + keys: [{ ...publicJwk, kid: "test-key", alg: "ES256", use: "sig" }], + }); + const bundle = yield* bundleServeMainTemplate(); + const injected = yield* load( + bundle, + { + ...baseEnv(config), + SUPABASE_JWKS: keys, + }, + okWorker, ); - expect(secondResponse.status).toBe(200); - expect(fetchCalls).toBe(1); - } - }); + const injectedResponse = yield* Effect.promise(() => + injected.options.handler( + new Request("http://localhost/hello", { headers: { Authorization: `Bearer ${token}` } }), + ), + ); + expect(injectedResponse.status).toBe(200); - it("uses package discovery when package.json is present or lstat reports NotFound", async () => { - const bundle = await bundleServeMainTemplate(); - const config = JSON.stringify({ - hello: { - entrypointPath: "hello/index.ts", - importMapPath: "", - staticFiles: [], - verifyJWT: false, - }, - }); - const worker = { fetch: async () => new Response("ok") }; - const permissionFailure = await load(bundle, baseEnv(config), worker, { - lstatError: new Error("permission denied"), - }); - await permissionFailure.options.handler(new Request("http://localhost/hello")); - expect(permissionFailure.state.createOptions?.noNpm).toBe(false); + for (const jwks of [undefined, "{"] as const) { + let fetchCalls = 0; + const remote = yield* load( + bundle, + { + ...baseEnv(config), + ...(jwks === undefined ? {} : { SUPABASE_JWKS: jwks }), + }, + okWorker, + { + fetchImpl: () => { + fetchCalls += 1; + return Promise.resolve( + new Response(keys, { headers: { "content-type": "application/json" } }), + ); + }, + }, + ); + const remoteResponse = yield* Effect.promise(() => + remote.options.handler( + new Request("http://localhost/hello", { + headers: { Authorization: `Bearer ${token}` }, + }), + ), + ); + expect(remoteResponse.status).toBe(200); + const secondResponse = yield* Effect.promise(() => + remote.options.handler( + new Request("http://localhost/hello", { + headers: { Authorization: `Bearer ${token}` }, + }), + ), + ); + expect(secondResponse.status).toBe(200); + expect(fetchCalls).toBe(1); + } + }), + ); - const missing = await load(bundle, baseEnv(config), worker, { - errors: { NotFound: NotFoundError }, - lstatError: new NotFoundError(), - }); - await missing.options.handler(new Request("http://localhost/hello")); - expect(missing.state.createOptions?.noNpm).toBe(true); - }); + it.effect("uses package discovery when package.json is present or lstat reports NotFound", () => + Effect.gen(function* () { + const bundle = yield* bundleServeMainTemplate(); + const config = yield* encodeJson({ + hello: { + entrypointPath: "hello/index.ts", + importMapPath: "", + staticFiles: [], + verifyJWT: false, + }, + }); + const permissionFailure = yield* load(bundle, baseEnv(config), okWorker, { + lstatError: new Error("permission denied"), + }); + yield* Effect.promise(() => + permissionFailure.options.handler(new Request("http://localhost/hello")), + ); + expect(permissionFailure.state.createOptions?.noNpm).toBe(false); + + const missing = yield* load(bundle, baseEnv(config), okWorker, { + errors: { NotFound: NotFoundError }, + lstatError: new NotFoundError(), + }); + yield* Effect.promise(() => missing.options.handler(new Request("http://localhost/hello"))); + expect(missing.state.createOptions?.noNpm).toBe(true); + }), + ); }); diff --git a/apps/cli/src/shared/functions/serve-main-bundler.ts b/apps/cli/src/shared/functions/serve-main-bundler.ts index 4572e7b206..930515c70c 100644 --- a/apps/cli/src/shared/functions/serve-main-bundler.ts +++ b/apps/cli/src/shared/functions/serve-main-bundler.ts @@ -1,5 +1,6 @@ import { fileURLToPath } from "node:url"; +import { Effect } from "effect"; import { build } from "esbuild"; /** @@ -18,21 +19,25 @@ const serveMainEntrypoint = fileURLToPath(new URL("./serve.main.ts", import.meta * `platform: "browser"` selects `jose`'s Web Crypto build for the * edge-runtime's Deno; `Deno` and `EdgeRuntime` are left as free globals. */ -export async function bundleServeMainTemplate(): Promise { - const result = await build({ - entryPoints: [serveMainEntrypoint], - bundle: true, - format: "esm", - platform: "browser", - minify: true, - write: false, - legalComments: "none", - logLevel: "silent", - }); +export const bundleServeMainTemplate = Effect.fnUntraced(function* () { + const result = yield* Effect.promise(() => + build({ + entryPoints: [serveMainEntrypoint], + bundle: true, + format: "esm", + platform: "browser", + minify: true, + write: false, + legalComments: "none", + logLevel: "silent", + }), + ); const output = result.outputFiles[0]?.text; if (output === undefined) { - throw new Error("esbuild produced no output for the functions serve runtime template"); + return yield* Effect.die( + new Error("esbuild produced no output for the functions serve runtime template"), + ); } return output; -} +}); diff --git a/apps/cli/src/shared/functions/serve-main-bundler.unit.test.ts b/apps/cli/src/shared/functions/serve-main-bundler.unit.test.ts index 2779e430b3..6e76c8c09d 100644 --- a/apps/cli/src/shared/functions/serve-main-bundler.unit.test.ts +++ b/apps/cli/src/shared/functions/serve-main-bundler.unit.test.ts @@ -1,24 +1,29 @@ -import { describe, expect, it } from "vitest"; +import { describe, expect, it } from "@effect/vitest"; +import { Effect } from "effect"; import { bundleServeMainTemplate } from "./serve-main-bundler.ts"; describe("bundleServeMainTemplate", () => { - it("produces a self-contained runtime template with no remote import specifiers", async () => { - const bundled = await bundleServeMainTemplate(); + it.effect("produces a self-contained runtime template with no remote import specifiers", () => + Effect.gen(function* () { + const bundled = yield* bundleServeMainTemplate(); - // The offline failure (#45570) was caused by these being resolved over the - // network on every container start. They must be inlined into the bundle. - expect(bundled).not.toMatch(/(?:from|import)\s*["']https?:/); - expect(bundled).not.toContain("jsr:"); - expect(bundled).not.toMatch(/from\s*["']jose["']/); - }); + // The offline failure (#45570) was caused by these being resolved over the + // network on every container start. They must be inlined into the bundle. + expect(bundled).not.toMatch(/(?:from|import)\s*["']https?:/); + expect(bundled).not.toContain("jsr:"); + expect(bundled).not.toMatch(/from\s*["']jose["']/); + }), + ); - it("preserves the template's Deno.serve entrypoint and inlines jose", async () => { - const bundled = await bundleServeMainTemplate(); + it.effect("preserves the template's Deno.serve entrypoint and inlines jose", () => + Effect.gen(function* () { + const bundled = yield* bundleServeMainTemplate(); - // Template body survives bundling (Deno global left as a free reference). - expect(bundled).toContain("Deno.serve"); - // jose is inlined, so the bundle is materially larger than the ~12KB template. - expect(bundled.length).toBeGreaterThan(20_000); - }); + // Template body survives bundling (Deno global left as a free reference). + expect(bundled).toContain("Deno.serve"); + // jose is inlined, so the bundle is materially larger than the ~12KB template. + expect(bundled.length).toBeGreaterThan(20_000); + }), + ); }); diff --git a/apps/cli/src/shared/functions/serve-main-offline.e2e.test.ts b/apps/cli/src/shared/functions/serve-main-offline.e2e.test.ts index b83d3ae5e0..34952659dc 100644 --- a/apps/cli/src/shared/functions/serve-main-offline.e2e.test.ts +++ b/apps/cli/src/shared/functions/serve-main-offline.e2e.test.ts @@ -279,7 +279,7 @@ describe("functions serve runtime template (offline)", () => { const dir = await mkdtemp(join(tmpdir(), "supabase-serve-offline-e2e-")); const container = `supabase-serve-offline-e2e-${process.pid.toString()}`; try { - await writeFile(join(dir, "index.ts"), await bundleServeMainTemplate()); + await writeFile(join(dir, "index.ts"), await Effect.runPromise(bundleServeMainTemplate())); const run = spawnSync( "docker", @@ -341,7 +341,7 @@ describe("functions serve runtime template (offline)", () => { const dir = await mkdtemp(join(tmpdir(), "supabase-serve-auth-e2e-")); const container = `supabase-serve-auth-e2e-${process.pid.toString()}`; try { - await writeFile(join(dir, "index.ts"), await bundleServeMainTemplate()); + await writeFile(join(dir, "index.ts"), await Effect.runPromise(bundleServeMainTemplate())); const run = spawnSync( "docker", @@ -440,7 +440,7 @@ describe("functions serve runtime template (offline)", () => { const runtimeContainer = `${network}-runtime`; const kongContainer = `${network}-kong`; try { - await writeFile(join(dir, "index.ts"), await bundleServeMainTemplate()); + await writeFile(join(dir, "index.ts"), await Effect.runPromise(bundleServeMainTemplate())); await mkdir(join(dir, "functions", "custom"), { recursive: true }); await mkdir(join(dir, "functions", "_shared"), { recursive: true }); await mkdir(join(dir, "functions", "custom", ".supabase-worker", "custom"), { diff --git a/apps/cli/src/shared/functions/serve.ts b/apps/cli/src/shared/functions/serve.ts index 72dbcff0c8..fe1c997214 100644 --- a/apps/cli/src/shared/functions/serve.ts +++ b/apps/cli/src/shared/functions/serve.ts @@ -34,6 +34,7 @@ import { Duration, Effect, Exit, + Fiber, Option, Predicate, Redacted, @@ -360,24 +361,31 @@ declare const SUPABASE_FUNCTIONS_SERVE_MAIN_TEMPLATE: string | undefined; * so the shipped binary never bundles at runtime. Running from source * bundles on demand. */ -function getFunctionsServeMainTemplate(): Promise { +const getFunctionsServeMainTemplate = Effect.fnUntraced(function* () { if (cachedFunctionsServeMainTemplate !== undefined) { - return Promise.resolve(cachedFunctionsServeMainTemplate); + return cachedFunctionsServeMainTemplate; } if (typeof SUPABASE_FUNCTIONS_SERVE_MAIN_TEMPLATE === "string") { cachedFunctionsServeMainTemplate = SUPABASE_FUNCTIONS_SERVE_MAIN_TEMPLATE; - return Promise.resolve(cachedFunctionsServeMainTemplate); + return cachedFunctionsServeMainTemplate; } // Bundler (and its esbuild dependency) is imported lazily and only here, // so it's never loaded by shipped binaries, which always take the define // branch above. - return import("./serve-main-bundler.ts") - .then(({ bundleServeMainTemplate }) => bundleServeMainTemplate()) - .then((bundled) => { + // Detached so an interrupted bring-up still finishes and fills the cache; only this join observes it. + const bundling = yield* Effect.forkDetach( + Effect.gen(function* () { + const { bundleServeMainTemplate } = yield* Effect.promise( + () => import("./serve-main-bundler.ts"), + ); + const bundled = yield* bundleServeMainTemplate(); cachedFunctionsServeMainTemplate = bundled; return bundled; - }); -} + }), + { startImmediately: true }, + ); + return yield* Fiber.join(bundling); +}); function reveal(value: string | Redacted.Redacted | undefined): string | undefined { if (value === undefined) { @@ -1896,7 +1904,7 @@ export const startEdgeRuntimeContainer = Effect.fn("functions.startEdgeRuntimeCo ...buildFunctionsServeInspectArgs(input.inspectMode, input.inspectMain), ...(input.debug ? ["--verbose"] : []), ]; - const serveMainTemplate = yield* Effect.promise(() => getFunctionsServeMainTemplate()).pipe( + const serveMainTemplate = yield* getFunctionsServeMainTemplate().pipe( Effect.withSpan("functions.serve.bundleMainTemplate"), ); // Streamed in via `docker cp` between create and start: embedding the template in the diff --git a/apps/cli/src/shared/functions/serve.unit.test.ts b/apps/cli/src/shared/functions/serve.unit.test.ts index 4895274c63..f8b8ff4268 100644 --- a/apps/cli/src/shared/functions/serve.unit.test.ts +++ b/apps/cli/src/shared/functions/serve.unit.test.ts @@ -1,4 +1,3 @@ -import { nativeFailure } from "./functions-docker.ts"; import { it } from "@effect/vitest"; import { Effect } from "effect"; import { describe, expect } from "vitest"; @@ -20,10 +19,7 @@ describe("buildServeEntrypointCommand", () => { it.effect("keeps the spawned command short even with the real bundled template", () => { return Effect.gen(function* () { - const bundled = yield* Effect.tryPromise({ - try: () => bundleServeMainTemplate(), - catch: nativeFailure, - }); + const bundled = yield* bundleServeMainTemplate(); const script = buildServeEntrypointCommand(["edge-runtime", "start"]); expect(bundled.length).toBeGreaterThan(20_000); expect(script.length).toBeLessThan(128);