From b8afb4bfc231d77e24d109450bced269b27347fd Mon Sep 17 00:00:00 2001 From: 7ttp <117663341+7ttp@users.noreply.github.com> Date: Fri, 2 Oct 2026 16:33:50 +0530 Subject: [PATCH 1/2] test(cli): cover functions offline serve e2e with effect lint (CLI-2592) --- .oxlintrc.effect.json | 22 +- .../functions/serve-main-offline.e2e.test.ts | 874 +++++++++--------- 2 files changed, 463 insertions(+), 433 deletions(-) diff --git a/.oxlintrc.effect.json b/.oxlintrc.effect.json index 0834ad9fba..23de31a024 100644 --- a/.oxlintrc.effect.json +++ b/.oxlintrc.effect.json @@ -12,27 +12,7 @@ "!apps/cli/src/shared/compute/**", "!apps/cli/src/shared/config/**", "!apps/cli/src/shared/feedback/**", - "!apps/cli/src/shared/functions/delete.errors.ts", - "!apps/cli/src/shared/functions/delete.ts", - "!apps/cli/src/shared/functions/deploy.errors.ts", - "!apps/cli/src/shared/functions/deploy.errors.unit.test.ts", - "!apps/cli/src/shared/functions/deploy.ts", - "!apps/cli/src/shared/functions/deploy.unit.test.ts", - "!apps/cli/src/shared/functions/download.errors.ts", - "!apps/cli/src/shared/functions/download.errors.unit.test.ts", - "!apps/cli/src/shared/functions/download.ts", - "!apps/cli/src/shared/functions/functions-api.errors.ts", - "!apps/cli/src/shared/functions/functions-api.errors.unit.test.ts", - "!apps/cli/src/shared/functions/functions-config.ts", - "!apps/cli/src/shared/functions/functions-docker.ts", - "!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/functions/**", "!apps/cli/src/shared/git/**", "!apps/cli/src/shared/init/**", "!apps/cli/src/shared/issue/**", 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 34952659dc..092a6db063 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 @@ -1,11 +1,21 @@ -import { execFile, execSync, spawnSync } from "node:child_process"; -import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises"; -import { tmpdir } from "node:os"; -import { join } from "node:path"; -import { promisify } from "node:util"; - -import { Effect } from "effect"; -import { describe, expect, test } from "vitest"; +import { BunServices } from "@effect/platform-bun"; +import { describe, expect, it } from "@effect/vitest"; +import { + Clock, + Data, + Duration, + Effect, + FileSystem, + Layer, + Path, + Predicate, + Result, + Schedule, + Stream, +} from "effect"; +import { FetchHttpClient, HttpBody, HttpClient } from "effect/unstable/http"; +import * as HttpClientRequest from "effect/unstable/http/HttpClientRequest"; +import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; import { START_KONG_YML_TEMPLATE } from "../../commands/start/templates/kong.yml.ts"; import { edgeRuntimeDockerfileImage } from "../../command-internal/edge-runtime-image.ts"; @@ -25,16 +35,41 @@ import { bundleServeMainTemplate } from "./serve-main-bundler.ts"; * isolating the network as the only variable. */ -function hasDocker(): boolean { - try { - execSync("docker info", { stdio: "ignore" }); - return true; - } catch { - return false; - } -} +class ServeOfflineE2eError extends Data.TaggedError("ServeOfflineE2eError")<{ + readonly message: string; + readonly cause?: unknown; +}> {} -const dockerAvailable = hasDocker(); +const docker = (args: ReadonlyArray, timeoutMs?: number) => + Effect.gen(function* () { + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const child = yield* spawner.spawn( + ChildProcess.make("docker", args, { stdin: "ignore", stdout: "pipe", stderr: "pipe" }), + ); + const [exitCode, stdout, stderr] = yield* Effect.all( + [ + child.exitCode, + Stream.mkString(Stream.decodeText(child.stdout)), + Stream.mkString(Stream.decodeText(child.stderr)), + ], + { concurrency: "unbounded" }, + ); + return { status: Number(exitCode), stdout, stderr }; + }).pipe( + Effect.scoped, + (effect) => (timeoutMs === undefined ? effect : Effect.timeout(effect, timeoutMs)), + Effect.mapError( + (cause) => new ServeOfflineE2eError({ message: `docker ${args.join(" ")} failed`, cause }), + ), + ); + +const dockerAvailable = await Effect.runPromise( + docker(["info"]).pipe( + Effect.map(({ status }) => status === 0), + Effect.orElseSucceed(() => false), + Effect.provide(BunServices.layer), + ), +); const SERVE_OFFLINE_STARTUP_TIMEOUT_MS = 60_000; const SERVE_OFFLINE_ATTEMPT_TIMEOUT_MS = 10_000; const DOCKER_COMMAND_TIMEOUT_MS = 5_000; @@ -148,489 +183,504 @@ const authFailureCases = [ }, ]; -function containerLogs(container: string): string { - const result = spawnSync("docker", ["logs", container], { - encoding: "utf8", - timeout: DOCKER_COMMAND_TIMEOUT_MS, - }); - const failure = result.error ? `\n` : ""; - return `${result.stdout ?? ""}\n${result.stderr ?? ""}${failure}`; -} - -const execFileAsync = promisify(execFile); +const containerLogs = (container: string) => + docker(["logs", container], DOCKER_COMMAND_TIMEOUT_MS).pipe( + Effect.map(({ stdout, stderr }) => `${stdout}\n${stderr}`), + Effect.catch((error) => + Effect.succeed(`\n\n`), + ), + ); -async function containerState(container: string): Promise { - try { - const { stdout } = await execFileAsync( - "docker", - [ - "inspect", - "--format", - '{{.State.Status}}{{if ne .State.Status "running"}} (exit code {{.State.ExitCode}}{{if .State.OOMKilled}}, OOM-killed{{end}}){{end}}', - container, - ], - { encoding: "utf8", timeout: DOCKER_COMMAND_TIMEOUT_MS }, - ); - return stdout.trim(); - } catch (error) { - const stderr = error instanceof Error && "stderr" in error ? String(error.stderr) : ""; - return `not inspectable (${stderr.trim() || String(error)})`; - } -} +const containerState = (container: string) => + docker( + [ + "inspect", + "--format", + '{{.State.Status}}{{if ne .State.Status "running"}} (exit code {{.State.ExitCode}}{{if .State.OOMKilled}}, OOM-killed{{end}}){{end}}', + container, + ], + DOCKER_COMMAND_TIMEOUT_MS, + ).pipe( + Effect.map(({ status, stdout, stderr }) => + status === 0 ? stdout.trim() : `not inspectable (${stderr.trim() || `exit ${status}`})`, + ), + Effect.catch((error) => + Effect.succeed(`not inspectable (${String(error.cause ?? error.message)})`), + ), + ); function isTerminalContainerState(state: string): boolean { return /^(exited|dead)\b/u.test(state); } -async function containerDiagnostics(containers: readonly string[]): Promise { - const blocks = await Promise.all( - containers.map( - async (container) => - `${container} (${await containerState(container)}) logs:\n${containerLogs(container)}`, +const containerDiagnostics = (containers: readonly string[]) => + Effect.forEach( + containers, + (container) => + Effect.all([containerState(container), containerLogs(container)]).pipe( + Effect.map(([state, logs]) => `${container} (${state}) logs:\n${logs}`), + ), + { concurrency: "unbounded" }, + ).pipe(Effect.map((blocks) => blocks.join("\n"))); + +const fetchFunctionWithDiagnostics = ( + request: HttpClientRequest.HttpClientRequest, + diagnosticContainers: readonly string[], + timeout: Duration.Input, +) => + HttpClient.execute(request).pipe( + Effect.timeout(timeout), + Effect.catch((cause) => + Effect.gen(function* () { + return yield* new ServeOfflineE2eError({ + message: `Function request to ${request.url} failed.\n${yield* containerDiagnostics(diagnosticContainers)}`, + cause, + }); + }), ), ); - return blocks.join("\n"); -} -async function fetchFunctionWithDiagnostics( - url: string, - diagnosticContainers: readonly string[], - init: RequestInit, -): Promise { - try { - return await fetch(url, init); - } catch (cause) { - throw new Error( - `Function request to ${url} failed.\n${await containerDiagnostics(diagnosticContainers)}`, - { cause }, - ); - } -} +class TransientFunctionResponse extends Data.TaggedError("TransientFunctionResponse")<{ + readonly cause: unknown; +}> {} + +class UnexpectedResponseStatus extends Data.TaggedError("UnexpectedResponseStatus")<{ + readonly message: string; +}> {} -async function fetchColdFunction( +const fetchColdFunction = ( url: string, diagnosticContainers: readonly string[], - init?: RequestInit, -): Promise { + headers: Record = {}, +) => { // Runtime health does not start user workers, and Edge Runtime exposes no // per-worker readiness signal. A cold worker can briefly disconnect or // return 502/503, so retry only those transient outcomes, each bounded so // one hung request cannot eat the budget. A container that exited or died // ends the wait: nothing will answer. - const deadline = Date.now() + SERVE_OFFLINE_STARTUP_TIMEOUT_MS; - let lastError: unknown; - - for (;;) { - try { - const response = await fetch(url, { - ...init, - signal: AbortSignal.timeout( - Math.max(1, Math.min(SERVE_OFFLINE_ATTEMPT_TIMEOUT_MS, deadline - Date.now())), - ), - }); - if (response.status !== 502 && response.status !== 503) { - return response; + return Effect.gen(function* () { + const deadline = (yield* Clock.currentTimeMillis) + SERVE_OFFLINE_STARTUP_TIMEOUT_MS; + const attempt = Effect.gen(function* () { + const attemptMs = Math.min( + SERVE_OFFLINE_ATTEMPT_TIMEOUT_MS, + deadline - (yield* Clock.currentTimeMillis), + ); + const result = yield* HttpClient.execute(HttpClientRequest.get(url, { headers })).pipe( + Effect.timeout(Duration.millis(Math.max(1, attemptMs))), + Effect.result, + ); + if ( + Result.isSuccess(result) && + result.success.status !== 502 && + result.success.status !== 503 + ) { + return result.success; + } + const cause = Result.isSuccess(result) + ? new UnexpectedResponseStatus({ message: `Received ${result.success.status} from ${url}` }) + : result.failure; + if (Result.isSuccess(result)) { + yield* Stream.runDrain(result.success.stream).pipe( + Effect.timeout(Duration.millis(Math.max(1, attemptMs))), + Effect.ignore, + ); + } + const states = yield* Effect.forEach( + diagnosticContainers, + (container) => Effect.map(containerState(container), (state) => ({ container, state })), + { concurrency: "unbounded" }, + ); + const dead = states.filter(({ state }) => isTerminalContainerState(state)); + const remainingMs = deadline - (yield* Clock.currentTimeMillis); + if (dead.length === 0 && remainingMs > 0) { + return yield* new TransientFunctionResponse({ cause }); } - lastError = new Error(`Received ${response.status} from ${url}`); - await response.body?.cancel(); - } catch (error) { - lastError = error; - } - - const states = await Promise.all( - diagnosticContainers.map(async (container) => ({ - container, - state: await containerState(container), - })), - ); - const dead = states.filter(({ state }) => isTerminalContainerState(state)); - const remainingMs = deadline - Date.now(); - if (dead.length > 0 || remainingMs <= 0) { const reason = dead.length > 0 ? `: ${dead.map(({ container, state }) => `${container} is ${state}`).join(", ")}` : ""; - throw new Error( - `Function at ${url} did not become ready${reason}.\n${await containerDiagnostics(diagnosticContainers)}`, - { cause: lastError }, - ); - } + return yield* new ServeOfflineE2eError({ + message: `Function at ${url} did not become ready${reason}.\n${yield* containerDiagnostics(diagnosticContainers)}`, + cause, + }); + }); - await Bun.sleep(Math.min(250, remainingMs)); - } -} + return yield* attempt.pipe( + Effect.retry({ + schedule: Schedule.spaced("250 millis"), + while: Predicate.isTagged("TransientFunctionResponse"), + }), + ); + }); +}; + +const awaitRuntimeReady = (url: string, diagnosticContainers: readonly string[]) => + Effect.gen(function* () { + const deadline = (yield* Clock.currentTimeMillis) + SERVE_OFFLINE_STARTUP_TIMEOUT_MS; + yield* HttpClient.execute(HttpClientRequest.get(url)).pipe( + Effect.timeout(Duration.millis(SERVE_OFFLINE_ATTEMPT_TIMEOUT_MS)), + Effect.filterOrFail( + (response) => response.status === 401, + (response) => + new UnexpectedResponseStatus({ message: `Received ${response.status} from ${url}` }), + ), + Effect.retry({ + schedule: Schedule.spaced("250 millis"), + while: () => Effect.map(Clock.currentTimeMillis, (now) => now < deadline), + }), + Effect.catch((cause) => + Effect.gen(function* () { + return yield* new ServeOfflineE2eError({ + message: `Runtime at ${url} did not become ready.\n${yield* containerDiagnostics(diagnosticContainers)}`, + cause, + }); + }), + ), + ); + }); -async function writeKongConfig(dir: string, edgeRuntimeContainer: string) { +const writeKongConfig = Effect.fnUntraced(function* (dir: string, edgeRuntimeContainer: string) { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; // Uses the TS transcription of the Kong template that `start`'s Kong // service already ports byte-for-byte. const config = START_KONG_YML_TEMPLATE.replaceAll("{{ .EdgeRuntimeId }}", edgeRuntimeContainer) .replaceAll("{{ .BearerToken }}", "$((headers.authorization or headers.apikey))") .replaceAll("{{ .QueryToken }}", "$((query_params.apikey))") .replace(/{{ \.[A-Za-z]+ }}/g, "unused"); - await writeFile(join(dir, "kong.yml"), config); -} + yield* fs.writeFileString(path.join(dir, "kong.yml"), config); +}); + +const resolveImage = (image: string, deadline?: number) => + Effect.tryPromise({ + try: () => ensureImage(image, deadline), + catch: (cause) => + new ServeOfflineE2eError({ + message: cause instanceof Error ? cause.message : String(cause), + cause, + }), + }); + +const removeOnClose = (containers: ReadonlyArray, network?: string) => + Effect.addFinalizer(() => + docker(["rm", "-f", ...containers]).pipe( + Effect.andThen(network === undefined ? Effect.void : docker(["network", "rm", network])), + Effect.ignore, + ), + ); + +const publishedPort = Effect.fnUntraced(function* (container: string, containerPort: string) { + const portResult = yield* docker(["port", container, containerPort]); + expect(portResult.status, portResult.stderr).toBe(0); + const port = Number(portResult.stdout.trim().split(":").at(-1)); + expect(port).toBeGreaterThan(0); + return port; +}); + +const testLayer = Layer.mergeAll( + BunServices.layer, + FetchHttpClient.layer, + Layer.succeed(HttpClient.TracerPropagationEnabled, false), +); describe("functions serve runtime template (offline)", () => { - test.skipIf(!dockerAvailable)( + it.live.skipIf(!dockerAvailable)( "boots under edge-runtime with networking disabled and fetches nothing remote", - { timeout: SERVE_OFFLINE_TEST_TIMEOUT_MS }, - async () => { - const runtimeImage = await ensureImage(await Effect.runPromise(edgeRuntimeDockerfileImage)); - 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 Effect.runPromise(bundleServeMainTemplate())); - - const run = spawnSync( - "docker", - [ - "run", - "-d", - "--name", - container, - "--network", - "none", - "-e", - "SUPABASE_INTERNAL_HOST_PORT=8081", - "-e", - "SUPABASE_INTERNAL_JWT_SECRET=offline-e2e", - "-e", - "SUPABASE_URL=http://127.0.0.1:54321", - "-e", - "SUPABASE_INTERNAL_FUNCTIONS_CONFIG={}", - "-e", - "SUPABASE_INTERNAL_WALLCLOCK_LIMIT_SEC=400", - "-v", - `${dir}:/app:ro`, - "--entrypoint", - "edge-runtime", - runtimeImage, - "start", - "--main-service=/app", - "--port=8081", - ], - { encoding: "utf8" }, - ); + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const runtimeImage = yield* resolveImage(yield* edgeRuntimeDockerfileImage); + const dir = yield* fs.makeTempDirectoryScoped({ prefix: "supabase-serve-offline-e2e-" }); + const container = `supabase-serve-offline-e2e-${process.pid.toString()}`; + yield* removeOnClose([container]); + yield* fs.writeFileString(path.join(dir, "index.ts"), yield* bundleServeMainTemplate()); + + const run = yield* docker([ + "run", + "-d", + "--name", + container, + "--network", + "none", + "-e", + "SUPABASE_INTERNAL_HOST_PORT=8081", + "-e", + "SUPABASE_INTERNAL_JWT_SECRET=offline-e2e", + "-e", + "SUPABASE_URL=http://127.0.0.1:54321", + "-e", + "SUPABASE_INTERNAL_FUNCTIONS_CONFIG={}", + "-e", + "SUPABASE_INTERNAL_WALLCLOCK_LIMIT_SEC=400", + "-v", + `${dir}:/app:ro`, + "--entrypoint", + "edge-runtime", + runtimeImage, + "start", + "--main-service=/app", + "--port=8081", + ]); expect(run.status, run.stderr).toBe(0); - const deadline = Date.now() + SERVE_OFFLINE_STARTUP_TIMEOUT_MS; - let logs = ""; - while (Date.now() < deadline) { - logs = containerLogs(container); - if (/Serving functions on/.test(logs) || /worker boot error/i.test(logs)) { - break; - } - await new Promise((resolve) => setTimeout(resolve, 250)); - } + const deadline = (yield* Clock.currentTimeMillis) + SERVE_OFFLINE_STARTUP_TIMEOUT_MS; + const logs = yield* containerLogs(container).pipe( + Effect.repeat({ + schedule: Schedule.spaced("250 millis"), + until: (current) => + /Serving functions on/.test(current) || /worker boot error/i.test(current) + ? Effect.succeed(true) + : Effect.map(Clock.currentTimeMillis, (now) => now >= deadline), + }), + ); expect(logs).toMatch(/Serving functions on/); expect(logs).not.toMatch(/deno\.land|jsr\.io/); expect(logs).not.toMatch(/dns error|name resolution|worker boot error/i); - } finally { - spawnSync("docker", ["rm", "-f", container], { stdio: "ignore" }); - await rm(dir, { recursive: true, force: true }); - } - }, + }).pipe(Effect.provide(testLayer)), + SERVE_OFFLINE_TEST_TIMEOUT_MS, ); - test.skipIf(!dockerAvailable)( + it.live.skipIf(!dockerAvailable)( "returns canonical JWT auth failures", - { timeout: SERVE_OFFLINE_TEST_TIMEOUT_MS }, - async () => { - const runtimeImage = await ensureImage(await Effect.runPromise(edgeRuntimeDockerfileImage)); - 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 Effect.runPromise(bundleServeMainTemplate())); - - const run = spawnSync( - "docker", - [ - "run", - "-d", - "--name", - container, - "-p", - "127.0.0.1::8081", - "-e", - "SUPABASE_INTERNAL_HOST_PORT=8081", - "-e", - "SUPABASE_INTERNAL_JWT_SECRET=auth-e2e", - "-e", - "SUPABASE_URL=http://127.0.0.1:54321", - "-e", - `SUPABASE_INTERNAL_FUNCTIONS_CONFIG=${AUTH_FUNCTIONS_CONFIG}`, - "-e", - "SUPABASE_INTERNAL_WALLCLOCK_LIMIT_SEC=400", - "-e", - 'SUPABASE_JWKS={"keys":[]}', - "-v", - `${dir}:/app:ro`, - "--entrypoint", - "edge-runtime", - runtimeImage, - "start", - "--main-service=/app", - "--port=8081", - ], - { encoding: "utf8" }, - ); + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const runtimeImage = yield* resolveImage(yield* edgeRuntimeDockerfileImage); + const dir = yield* fs.makeTempDirectoryScoped({ prefix: "supabase-serve-auth-e2e-" }); + const container = `supabase-serve-auth-e2e-${process.pid.toString()}`; + yield* removeOnClose([container]); + yield* fs.writeFileString(path.join(dir, "index.ts"), yield* bundleServeMainTemplate()); + + const run = yield* docker([ + "run", + "-d", + "--name", + container, + "-p", + "127.0.0.1::8081", + "-e", + "SUPABASE_INTERNAL_HOST_PORT=8081", + "-e", + "SUPABASE_INTERNAL_JWT_SECRET=auth-e2e", + "-e", + "SUPABASE_URL=http://127.0.0.1:54321", + "-e", + `SUPABASE_INTERNAL_FUNCTIONS_CONFIG=${AUTH_FUNCTIONS_CONFIG}`, + "-e", + "SUPABASE_INTERNAL_WALLCLOCK_LIMIT_SEC=400", + "-e", + 'SUPABASE_JWKS={"keys":[]}', + "-v", + `${dir}:/app:ro`, + "--entrypoint", + "edge-runtime", + runtimeImage, + "start", + "--main-service=/app", + "--port=8081", + ]); expect(run.status, run.stderr).toBe(0); - const portResult = spawnSync("docker", ["port", container, "8081/tcp"], { - encoding: "utf8", - }); - expect(portResult.status, portResult.stderr).toBe(0); - const port = Number(portResult.stdout.trim().split(":").at(-1)); - expect(port).toBeGreaterThan(0); + const port = yield* publishedPort(container, "8081/tcp"); const url = `http://127.0.0.1:${port}/test`; - const deadline = Date.now() + SERVE_OFFLINE_STARTUP_TIMEOUT_MS; - let ready = false; - let lastError: unknown; - while (Date.now() < deadline) { - try { - const response = await fetch(url, { - signal: AbortSignal.timeout(SERVE_OFFLINE_ATTEMPT_TIMEOUT_MS), - }); - if (response.status === 401) { - ready = true; - break; - } - lastError = new Error(`Received ${response.status} from ${url}`); - } catch (error) { - lastError = error; - } - await new Promise((resolve) => setTimeout(resolve, 250)); - } - if (!ready) { - throw new Error( - `Runtime at ${url} did not become ready.\n${await containerDiagnostics([container])}`, - { cause: lastError }, - ); - } + yield* awaitRuntimeReady(url, [container]); for (const { name, authorization, code, message } of authFailureCases) { - const response = await fetch(url, { - headers: authorization === undefined ? undefined : { authorization }, - }); + const response = yield* HttpClient.execute( + HttpClientRequest.get(url, { + headers: authorization === undefined ? {} : { authorization }, + }), + ); expect(response.status, name).toBe(401); - expect(response.headers.get("content-type"), name).toContain("application/json"); - expect(response.headers.get("sb-error-code"), name).toBe(code); - expect(await response.json(), name).toEqual({ code, message, msg: message }); + expect(response.headers["content-type"], name).toContain("application/json"); + expect(response.headers["sb-error-code"], name).toBe(code); + expect(yield* response.json, name).toEqual({ code, message, msg: message }); } - } finally { - spawnSync("docker", ["rm", "-f", container], { stdio: "ignore" }); - await rm(dir, { recursive: true, force: true }); - } - }, + }).pipe(Effect.provide(testLayer)), + SERVE_OFFLINE_TEST_TIMEOUT_MS, ); - test.skipIf(!dockerAvailable)( + it.live.skipIf(!dockerAvailable)( "preserves function env and CORS headers, exposes JWT errors, and returns early responses through Kong", - { timeout: SERVE_OFFLINE_TEST_TIMEOUT_MS }, - async () => { - const imageDeadline = resolveDeadline(); - const [runtimeImage, kongImage] = await Promise.all([ - ensureImage(await Effect.runPromise(edgeRuntimeDockerfileImage), imageDeadline), - ensureImage(dockerfileServiceImage("kong", false), imageDeadline), - ]); - const dir = await mkdtemp(join(tmpdir(), "supabase-serve-kong-e2e-")); - const network = `supabase-serve-kong-e2e-${process.pid.toString()}`; - const runtimeContainer = `${network}-runtime`; - const kongContainer = `${network}-kong`; - try { - 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"), { - recursive: true, - }); - await writeFile(join(dir, "functions", "custom", "index.ts"), CUSTOM_FUNCTION); - await writeFile( - join(dir, "functions", "custom", ".supabase-worker", "custom", "index.ts"), + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const imageDeadline = resolveDeadline(); + const [runtimeImage, kongImage] = yield* Effect.all( + [ + resolveImage(yield* edgeRuntimeDockerfileImage, imageDeadline), + resolveImage(dockerfileServiceImage("kong", false), imageDeadline), + ], + { concurrency: "unbounded" }, + ); + const dir = yield* fs.makeTempDirectoryScoped({ prefix: "supabase-serve-kong-e2e-" }); + const network = `supabase-serve-kong-e2e-${process.pid.toString()}`; + const runtimeContainer = `${network}-runtime`; + const kongContainer = `${network}-kong`; + yield* removeOnClose([kongContainer, runtimeContainer], network); + yield* fs.writeFileString(path.join(dir, "index.ts"), yield* bundleServeMainTemplate()); + yield* fs.makeDirectory(path.join(dir, "functions", "custom"), { recursive: true }); + yield* fs.makeDirectory(path.join(dir, "functions", "_shared"), { recursive: true }); + yield* fs.makeDirectory( + path.join(dir, "functions", "custom", ".supabase-worker", "custom"), + { + recursive: true, + }, + ); + yield* fs.writeFileString( + path.join(dir, "functions", "custom", "index.ts"), + CUSTOM_FUNCTION, + ); + yield* fs.writeFileString( + path.join(dir, "functions", "custom", ".supabase-worker", "custom", "index.ts"), NESTED_FUNCTION, ); - await writeFile( - join(dir, "functions", "_shared", "value.ts"), + yield* fs.writeFileString( + path.join(dir, "functions", "_shared", "value.ts"), 'export const sharedValue = "shared-import-ok";\n', ); - await writeKongConfig(dir, runtimeContainer); + yield* writeKongConfig(dir, runtimeContainer); - const createNetwork = spawnSync("docker", ["network", "create", network], { - encoding: "utf8", - }); + const createNetwork = yield* docker(["network", "create", network]); expect(createNetwork.status, createNetwork.stderr).toBe(0); - const runRuntime = spawnSync( - "docker", - [ - "run", - "-d", - "--name", - runtimeContainer, - "--network", - network, - "-e", - "SUPABASE_INTERNAL_HOST_PORT=8081", - "-e", - "SUPABASE_INTERNAL_JWT_SECRET=auth-e2e", - "-e", - `SUPABASE_URL=http://${kongContainer}:8000`, - "-e", - `SUPABASE_INTERNAL_FUNCTIONS_CONFIG=${KONG_FUNCTIONS_CONFIG}`, - "-e", - "SUPABASE_INTERNAL_DEBUG=true", - "-e", - "SHARED=shared", - "-e", - "GLOBAL_ONLY=global", - "-e", - "SUPABASE_INTERNAL_WALLCLOCK_LIMIT_SEC=400", - "-e", - 'SUPABASE_JWKS={"keys":[]}', - "-v", - `${dir}:/app:ro`, - "--entrypoint", - "edge-runtime", - runtimeImage, - "start", - "--main-service=/app", - "--port=8081", - ], - { encoding: "utf8" }, - ); + const runRuntime = yield* docker([ + "run", + "-d", + "--name", + runtimeContainer, + "--network", + network, + "-e", + "SUPABASE_INTERNAL_HOST_PORT=8081", + "-e", + "SUPABASE_INTERNAL_JWT_SECRET=auth-e2e", + "-e", + `SUPABASE_URL=http://${kongContainer}:8000`, + "-e", + `SUPABASE_INTERNAL_FUNCTIONS_CONFIG=${KONG_FUNCTIONS_CONFIG}`, + "-e", + "SUPABASE_INTERNAL_DEBUG=true", + "-e", + "SHARED=shared", + "-e", + "GLOBAL_ONLY=global", + "-e", + "SUPABASE_INTERNAL_WALLCLOCK_LIMIT_SEC=400", + "-e", + 'SUPABASE_JWKS={"keys":[]}', + "-v", + `${dir}:/app:ro`, + "--entrypoint", + "edge-runtime", + runtimeImage, + "start", + "--main-service=/app", + "--port=8081", + ]); expect(runRuntime.status, runRuntime.stderr).toBe(0); - const runKong = spawnSync( - "docker", - [ - "run", - "-d", - "--name", - kongContainer, - "--network", - network, - "-p", - "127.0.0.1::8000", - "-e", - "KONG_DATABASE=off", - "-e", - "KONG_DECLARATIVE_CONFIG=/home/kong/kong.yml", - "-e", - "KONG_PLUGINS=request-transformer,cors", - "-e", - "KONG_NGINX_WORKER_PROCESSES=1", - "-v", - `${join(dir, "kong.yml")}:/home/kong/kong.yml:ro`, - kongImage, - "kong", - "docker-start", - ], - { encoding: "utf8" }, - ); + const runKong = yield* docker([ + "run", + "-d", + "--name", + kongContainer, + "--network", + network, + "-p", + "127.0.0.1::8000", + "-e", + "KONG_DATABASE=off", + "-e", + "KONG_DECLARATIVE_CONFIG=/home/kong/kong.yml", + "-e", + "KONG_PLUGINS=request-transformer,cors", + "-e", + "KONG_NGINX_WORKER_PROCESSES=1", + "-v", + `${path.join(dir, "kong.yml")}:/home/kong/kong.yml:ro`, + kongImage, + "kong", + "docker-start", + ]); expect(runKong.status, runKong.stderr).toBe(0); - const portResult = spawnSync("docker", ["port", kongContainer, "8000/tcp"], { - encoding: "utf8", - }); - expect(portResult.status, portResult.stderr).toBe(0); - const port = Number(portResult.stdout.trim().split(":").at(-1)); - expect(port).toBeGreaterThan(0); + const port = yield* publishedPort(kongContainer, "8000/tcp"); const functionsUrl = `http://127.0.0.1:${port}/functions/v1`; const authUrl = `${functionsUrl}/test`; const diagnosticContainers = [kongContainer, runtimeContainer] as const; - const deadline = Date.now() + SERVE_OFFLINE_STARTUP_TIMEOUT_MS; - let ready = false; - let lastError: unknown; - while (Date.now() < deadline) { - try { - const response = await fetch(authUrl, { - signal: AbortSignal.timeout(SERVE_OFFLINE_ATTEMPT_TIMEOUT_MS), - }); - if (response.status === 401) { - ready = true; - break; - } - lastError = new Error(`Received ${response.status} from ${authUrl}`); - } catch (error) { - lastError = error; - } - await new Promise((resolve) => setTimeout(resolve, 250)); - } - if (!ready) { - throw new Error( - `Runtime at ${authUrl} did not become ready.\n${await containerDiagnostics(diagnosticContainers)}`, - { cause: lastError }, - ); - } + yield* awaitRuntimeReady(authUrl, diagnosticContainers); - const [customResponse, aliasResponse, nestedResponse] = await Promise.all([ - fetchColdFunction(`${functionsUrl}/custom`, diagnosticContainers, { - headers: { Origin: "http://localhost:3000" }, - }), - fetchColdFunction(`${functionsUrl}/custom-alias`, diagnosticContainers), - fetchColdFunction(`${functionsUrl}/nested-worker-path`, diagnosticContainers), - ]); + const [customResponse, aliasResponse, nestedResponse] = yield* Effect.all( + [ + fetchColdFunction(`${functionsUrl}/custom`, diagnosticContainers, { + Origin: "http://localhost:3000", + }), + fetchColdFunction(`${functionsUrl}/custom-alias`, diagnosticContainers), + fetchColdFunction(`${functionsUrl}/nested-worker-path`, diagnosticContainers), + ], + { concurrency: "unbounded" }, + ); expect(customResponse.status).toBe(200); - expect(customResponse.headers.get("x-custom-id")).toBe("abc123"); - expect(customResponse.headers.get("x-function-slug")).toBe("custom"); - expect(customResponse.headers.get("x-shared-import")).toBe("shared-import-ok"); - expect(customResponse.headers.get("x-shared")).toBe("function"); - expect(customResponse.headers.get("x-function-only")).toBe("function"); - expect(customResponse.headers.get("x-global-only")).toBe("global"); - expect(customResponse.headers.get("access-control-expose-headers")?.toLowerCase()).toBe( + expect(customResponse.headers["x-custom-id"]).toBe("abc123"); + expect(customResponse.headers["x-function-slug"]).toBe("custom"); + expect(customResponse.headers["x-shared-import"]).toBe("shared-import-ok"); + expect(customResponse.headers["x-shared"]).toBe("function"); + expect(customResponse.headers["x-function-only"]).toBe("function"); + expect(customResponse.headers["x-global-only"]).toBe("global"); + expect(customResponse.headers["access-control-expose-headers"]?.toLowerCase()).toBe( "x-custom-id", ); expect(aliasResponse.status).toBe(200); - expect(aliasResponse.headers.get("x-function-slug")).toBe("custom-alias"); - expect(aliasResponse.headers.get("x-shared-import")).toBe("shared-import-ok"); + expect(aliasResponse.headers["x-function-slug"]).toBe("custom-alias"); + expect(aliasResponse.headers["x-shared-import"]).toBe("shared-import-ok"); expect(nestedResponse.status).toBe(200); - expect(nestedResponse.headers.get("x-function-slug")).toBe("nested-worker-path"); - const earlyResponse = await fetchFunctionWithDiagnostics( - `${functionsUrl}/custom`, - diagnosticContainers, - { - method: "POST", + expect(nestedResponse.headers["x-function-slug"]).toBe("nested-worker-path"); + const earlyResponse = yield* fetchFunctionWithDiagnostics( + HttpClientRequest.post(`${functionsUrl}/custom`, { headers: { "x-reject-before-body": "true" }, - body: new Uint8Array(1024 * 1024), - signal: AbortSignal.timeout(5_000), - }, + body: HttpBody.raw(new Uint8Array(1024 * 1024)), + }), + diagnosticContainers, + Duration.seconds(5), ); expect(earlyResponse.status).toBe(400); - expect(await earlyResponse.text()).toBe("rejected"); - const runtimeLogs = containerLogs(runtimeContainer); + expect(yield* earlyResponse.text.pipe(Effect.timeout(Duration.seconds(5)))).toBe( + "rejected", + ); + const runtimeLogs = yield* containerLogs(runtimeContainer); expect(runtimeLogs).toContain("Functions config:"); expect(runtimeLogs).toContain('"custom"'); expect(runtimeLogs).not.toContain('"env"'); expect(runtimeLogs).not.toContain("must-not-appear-in-debug-logs"); - const authResponse = await fetch(authUrl, { - method: "POST", - headers: { Origin: "http://localhost:3000" }, - body: new Uint8Array(1024 * 1024), - signal: AbortSignal.timeout(5_000), - }); + const authResponse = yield* HttpClient.execute( + HttpClientRequest.post(authUrl, { + headers: { Origin: "http://localhost:3000" }, + body: HttpBody.raw(new Uint8Array(1024 * 1024)), + }), + ).pipe(Effect.timeout(Duration.seconds(5))); expect(authResponse.status).toBe(401); - expect(authResponse.headers.get("sb-error-code")).toBe("UNAUTHORIZED_NO_AUTH_HEADER"); - expect(authResponse.headers.get("access-control-expose-headers")).toBe("sb-error-code"); - expect(await authResponse.json()).toEqual({ + expect(authResponse.headers["sb-error-code"]).toBe("UNAUTHORIZED_NO_AUTH_HEADER"); + expect(authResponse.headers["access-control-expose-headers"]).toBe("sb-error-code"); + expect(yield* authResponse.json.pipe(Effect.timeout(Duration.seconds(5)))).toEqual({ code: "UNAUTHORIZED_NO_AUTH_HEADER", message: "Missing authorization header", msg: "Missing authorization header", }); - const reusedCustomResponse = await fetch(`${functionsUrl}/custom`); + const reusedCustomResponse = yield* HttpClient.execute( + HttpClientRequest.get(`${functionsUrl}/custom`), + ); expect(reusedCustomResponse.status).toBe(200); - expect(reusedCustomResponse.headers.get("x-function-slug")).toBe("custom"); - } finally { - spawnSync("docker", ["rm", "-f", kongContainer, runtimeContainer], { - stdio: "ignore", - }); - spawnSync("docker", ["network", "rm", network], { stdio: "ignore" }); - await rm(dir, { recursive: true, force: true }); - } - }, + expect(reusedCustomResponse.headers["x-function-slug"]).toBe("custom"); + }).pipe(Effect.provide(testLayer)), + SERVE_OFFLINE_TEST_TIMEOUT_MS, ); }); From 078860f422a8ad6d02a7cb4110617b54d0682eb9 Mon Sep 17 00:00:00 2001 From: 7ttp <117663341+7ttp@users.noreply.github.com> Date: Sat, 3 Oct 2026 03:29:55 +0530 Subject: [PATCH 2/2] test(cli): keep partial docker logs and network cleanup in offline e2e --- .../functions/serve-main-offline.e2e.test.ts | 45 ++++++++++++++++--- 1 file changed, 38 insertions(+), 7 deletions(-) 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 092a6db063..a98313fd90 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 @@ -9,6 +9,7 @@ import { Layer, Path, Predicate, + Ref, Result, Schedule, Stream, @@ -184,12 +185,40 @@ const authFailureCases = [ ]; const containerLogs = (container: string) => - docker(["logs", container], DOCKER_COMMAND_TIMEOUT_MS).pipe( - Effect.map(({ stdout, stderr }) => `${stdout}\n${stderr}`), - Effect.catch((error) => - Effect.succeed(`\n\n`), - ), - ); + Effect.gen(function* () { + const stdout = yield* Ref.make(""); + const stderr = yield* Ref.make(""); + const failure = yield* Effect.gen(function* () { + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const child = yield* spawner.spawn( + ChildProcess.make("docker", ["logs", container], { + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + }), + ); + yield* Effect.all( + [ + child.exitCode, + Stream.runForEach(Stream.decodeText(child.stdout), (chunk) => + Ref.update(stdout, (text) => text + chunk), + ), + Stream.runForEach(Stream.decodeText(child.stderr), (chunk) => + Ref.update(stderr, (text) => text + chunk), + ), + ], + { concurrency: "unbounded" }, + ); + }).pipe( + Effect.scoped, + Effect.timeout(DOCKER_COMMAND_TIMEOUT_MS), + Effect.match({ + onFailure: (cause) => `\n`, + onSuccess: () => "", + }), + ); + return `${yield* Ref.get(stdout)}\n${yield* Ref.get(stderr)}${failure}`; + }); const containerState = (container: string) => docker( @@ -364,8 +393,10 @@ const resolveImage = (image: string, deadline?: number) => const removeOnClose = (containers: ReadonlyArray, network?: string) => Effect.addFinalizer(() => docker(["rm", "-f", ...containers]).pipe( - Effect.andThen(network === undefined ? Effect.void : docker(["network", "rm", network])), Effect.ignore, + Effect.andThen( + network === undefined ? Effect.void : Effect.ignore(docker(["network", "rm", network])), + ), ), );