Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 2 additions & 4 deletions packages/stack/src/runtime/Container.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1523,10 +1523,8 @@ describe("container process adapter", { timeout: 120_000 }, () => {
if (unpinned[0] === "version") return Effect.succeed(handle(0, "27.0.0|27.0.0"));
if (unpinned[0] === "inspect")
return Effect.succeed(handle(0, '{"Ports":{},"Status":"running","ExitCode":0}'));
if (unpinned[0] === "run")
return Effect.succeed(
handle(0, unpinned.includes("-i") ? "supabase-helper-ready\n" : "abcdef012345"),
);
if (unpinned[0] === "run") return Effect.succeed(handle(0, "abcdef012345"));
if (unpinned[0] === "start") return Effect.succeed(handle(0, "supabase-helper-ready\n"));
// Two call sites share "ps": the launched container's own `--rm` removal poll (filtered
// by name, always confirmed gone) and the label-based orphan discovery below.
if (unpinned[0] === "ps" && unpinned.some((arg) => arg.includes("com.supabase.stack="))) {
Expand Down
227 changes: 182 additions & 45 deletions packages/stack/src/storage/DockerDatabaseStorage.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,11 @@ import {
import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process";
import * as TestClock from "effect/testing/TestClock";
import { postgresVersion, resolveArtifact } from "../Artifacts.ts";
import { makeContainerRuntime, type EngineTarget } from "../runtime/Container.ts";
import {
makeContainerRuntime,
removeStackContainers,
type EngineTarget,
} from "../runtime/Container.ts";
import { makeDatabaseSnapshots } from "../services/DatabaseSnapshot.ts";
import { makeDockerDatabaseStorage } from "./DockerDatabaseStorage.ts";
import { makeDockerHelperRegistry } from "./DockerHelperRegistry.ts";
Expand Down Expand Up @@ -60,7 +64,7 @@ const dockerStoragePublishLoopFixture = fileURLToPath(

const helperMirror = "registry.test/supabase/postgres:17";
const fakeHelperEngine = (
onRun: () => Effect.Effect<void> = () => Effect.void,
onCreate: () => Effect.Effect<void> = () => Effect.void,
{ volumePresent = false }: { readonly volumePresent?: boolean } = {},
) => {
const commands: string[][] = [];
Expand Down Expand Up @@ -97,16 +101,17 @@ const fakeHelperEngine = (
return handle(0);
}
if (args[0] === "inspect") return handle(1, "", "no such container");
if (args[0] === "run") {
if (args[0] === "create") {
const shellIndex = args.indexOf("/bin/sh");
const image = args[shellIndex - 1];
if (image === undefined || !local.has(image))
return handle(1, "", `Unable to find image '${image ?? ""}' locally`);
// Fires once the create command the claim must precede is observed, so the claim's
// presence at this exact point is a deterministic fact, not a timing guess.
yield* onRun();
return handle(0, args.includes("-i") ? "supabase-helper-ready\n" : "abcdef0123456789");
yield* onCreate();
return handle(0, "abcdef0123456789");
}
if (args[0] === "start") return handle(0, "supabase-helper-ready\n");
if (args[0] === "exec" || args[0] === "rm") return handle(0, "done");
return handle(1, "", `unexpected engine command: ${args[0] ?? ""}`);
}),
Expand Down Expand Up @@ -183,12 +188,15 @@ const awaitDestroyed = Effect.fn("DockerStorageTest.awaitDestroyed")((id: string
),
);

const ownerDeathScenario = (shared: boolean) =>
const ownerDeathScenario = (mode: "instance" | "shared" | "stall-before-start") =>
Effect.scoped(
Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const crypto = yield* Crypto.Crypto;
const spawner = yield* ChildProcessSpawner.ChildProcessSpawner;
const shared = mode === "shared";
const stalled = mode === "stall-before-start";
const root = yield* fs.makeTempDirectoryScoped({ prefix: "docker-helper-owner-death-" });
const stackId = `storage-owner-death-${yield* crypto.randomUUIDv4}`;
const filter = `label=com.supabase.stack=${stackId}`;
Expand All @@ -202,7 +210,7 @@ const ownerDeathScenario = (shared: boolean) =>
fileURLToPath(new URL("../../tests/helper-owner-fixture.ts", import.meta.url)),
stackId,
root,
...(shared ? ["shared"] : []),
...(mode === "instance" ? [] : [mode]),
],
{ stdin: "ignore" },
),
Expand Down Expand Up @@ -238,7 +246,7 @@ const ownerDeathScenario = (shared: boolean) =>
if (
Exit.isFailure(ready) ||
Option.isNone(ready.value) ||
ready.value.value !== "HELPER_READY"
ready.value.value !== (stalled ? "HELPER_CREATED" : "HELPER_READY")
) {
yield* child.kill({ killSignal: "SIGKILL" }).pipe(Effect.ignore);
const [stderr, exited] = yield* Effect.all(
Expand All @@ -251,6 +259,7 @@ const ownerDeathScenario = (shared: boolean) =>
}
const id = yield* engine([
"ps",
...(stalled ? ["--all", "--filter", "status=created"] : []),
"--filter",
filter,
"--no-trunc",
Expand All @@ -260,7 +269,7 @@ const ownerDeathScenario = (shared: boolean) =>
Effect.flatMap((value) =>
value.length > 0
? Effect.succeed(value)
: new DockerTestError({ message: "Owner helper was not running" }),
: new DockerTestError({ message: "Owner helper was not found" }),
),
);
if (shared)
Expand All @@ -272,18 +281,33 @@ const ownerDeathScenario = (shared: boolean) =>
id,
]),
).toBe("volume");
const since = Math.floor((yield* Clock.currentTimeMillis) / 1000) - 1;
const destroyed = yield* awaitDestroyed(id, since).pipe(
Effect.forkChild({ startImmediately: true }),
);
expect(yield* child.isRunning).toBe(true);
yield* child.kill({ killSignal: "SIGKILL" });
yield* Fiber.join(destroyed).pipe(
Effect.timeoutOrElse({
duration: "1 minute",
orElse: () => new DockerTestError({ message: `Orphaned database helper: ${id}` }),
}),
);
if (stalled) {
expect(yield* child.isRunning).toBe(true);
yield* child.kill({ killSignal: "SIGKILL" });
yield* child.exitCode.pipe(Effect.exit);
/* Nothing attached, so the helper outlives its owner until the stack sweep removes it. */
expect(yield* engine(["ps", "-aq", "--no-trunc", "--filter", filter])).toBe(id);
expect(
yield* removeStackContainers({
target: dockerTarget,
stackId,
stackRoot: path.resolve(path.join(root, "state", "stack", "data")),
}),
).toEqual([]);
} else {
const since = Math.floor((yield* Clock.currentTimeMillis) / 1000) - 1;
const destroyed = yield* awaitDestroyed(id, since).pipe(
Effect.forkChild({ startImmediately: true }),
);
expect(yield* child.isRunning).toBe(true);
yield* child.kill({ killSignal: "SIGKILL" });
yield* Fiber.join(destroyed).pipe(
Effect.timeoutOrElse({
duration: "1 minute",
orElse: () => new DockerTestError({ message: `Orphaned database helper: ${id}` }),
}),
);
}
expect(yield* engine(["ps", "-aq", "--filter", filter])).toBe("");
}),
).pipe(Effect.provide(NodeServices.layer));
Expand All @@ -296,12 +320,17 @@ describe.runIf(testEngine === "docker")(
() => {
it.live.skipIf(process.platform === "win32")(
"removes a database helper when its owner is killed",
() => ownerDeathScenario(false),
() => ownerDeathScenario("instance"),
);

it.live.skipIf(process.platform === "win32")(
"removes a shared volume helper when its owner is killed",
() => ownerDeathScenario(true),
() => ownerDeathScenario("shared"),
);

it.live.skipIf(process.platform === "win32")(
"sweeps a helper whose owner was killed between create and start",
() => ownerDeathScenario("stall-before-start"),
);

it.live("waits for helper creation to settle before cleaning up an interrupted operation", () =>
Expand All @@ -316,31 +345,33 @@ describe.runIf(testEngine === "docker")(
const started = yield* Deferred.make<void>();
const release = yield* Deferred.make<void>();
const present = yield* Ref.make(false);
const removals = yield* Ref.make(0);
const removals = yield* Ref.make<ReadonlyArray<boolean>>([]);
const spawner = ChildProcessSpawner.make((command) =>
Effect.gen(function* () {
if (!ChildProcess.isStandardCommand(command))
return yield* Effect.die("Unexpected command");
const creating = command.args[0] === "run";
const creating = command.args[0] === "create";
if (creating) yield* Deferred.succeed(started, undefined);
/* Each removal records whether the create had settled into a container by then. */
if (command.args[0] === "rm") {
yield* Ref.update(removals, (count) => count + 1);
yield* Ref.set(present, false);
const existed = yield* Ref.getAndSet(present, false);
yield* Ref.update(removals, (all) => [...all, existed]);
}
return ChildProcessSpawner.makeHandle({
pid: ChildProcessSpawner.ProcessId(0),
exitCode: Effect.succeed(ChildProcessSpawner.ExitCode(0)),
exitCode: creating
? Deferred.await(release).pipe(
Effect.andThen(Ref.set(present, true)),
Effect.as(ChildProcessSpawner.ExitCode(0)),
)
: Effect.succeed(ChildProcessSpawner.ExitCode(0)),
isRunning: Effect.succeed(false),
kill: () => Effect.void,
stdin: Sink.drain,
stdout: creating
? Stream.fromEffect(
Deferred.await(release).pipe(
Effect.andThen(Ref.set(present, true)),
Effect.as(new TextEncoder().encode("supabase-helper-ready\n")),
),
)
: Stream.empty,
stdout:
command.args[0] === "start"
? Stream.succeed(new TextEncoder().encode("supabase-helper-ready\n"))
: Stream.empty,
stderr: Stream.empty,
all: Stream.empty,
getInputFd: () => Sink.drain,
Expand Down Expand Up @@ -372,12 +403,9 @@ describe.runIf(testEngine === "docker")(
const operation = yield* storage.removeData("17").pipe(Effect.forkScoped);
yield* Deferred.await(started);
yield* Effect.sync(() => operation.interruptUnsafe());
yield* Effect.yieldNow;
expect(yield* Ref.get(removals)).toBe(0);
yield* Deferred.succeed(release, undefined);
yield* Fiber.await(operation);
expect(yield* Ref.get(removals)).toBe(1);
expect(yield* Ref.get(present)).toBe(false);
expect(yield* Ref.get(removals)).toEqual([true]);
}),
).pipe(Effect.provide(NodeServices.layer)),
);
Expand All @@ -390,25 +418,29 @@ describe.runIf(testEngine === "docker")(
const crypto = yield* Crypto.Crypto;
const root = yield* fs.makeTempDirectoryScoped({ prefix: "docker-helper-timeout-" });
const waiting = yield* Deferred.make<void>();
const commands: string[][] = [];
const warning =
"WARNING: The requested image's platform does not match the host platform";
const spawner = ChildProcessSpawner.make((command) =>
Effect.gen(function* () {
if (!ChildProcess.isStandardCommand(command))
return yield* Effect.die("Unexpected command");
const creating = command.args[0] === "run";
commands.push([...command.args]);
const creating = command.args[0] === "create";
const attaching = command.args[0] === "start";
return ChildProcessSpawner.makeHandle({
pid: ChildProcessSpawner.ProcessId(0),
exitCode: Effect.succeed(ChildProcessSpawner.ExitCode(0)),
isRunning: Effect.succeed(false),
kill: () => Effect.void,
stdin: Sink.drain,
stdout: creating
stdout: attaching
? Stream.fromEffect(Deferred.succeed(waiting, undefined)).pipe(
Stream.drain,
Stream.concat(Stream.never),
)
: Stream.empty,
// Docker prints the platform warning while creating the container, before start.
stderr: creating
? Stream.succeed(new TextEncoder().encode(`${warning}\n`))
: Stream.empty,
Expand Down Expand Up @@ -446,6 +478,111 @@ describe.runIf(testEngine === "docker")(
expect(failure.message).toBe(
`Database helper did not become ready within 30 seconds: ${warning}`,
);
const create = commands.find((args) => args[0] === "create");
const name = create?.[create.indexOf("--name") + 1];
expect(create?.slice(0, 3)).toEqual(["create", "--rm", "-i"]);
expect(name).toMatch(/^supabase-db-helper-/u);
const order = commands.map((args) => args[0]);
expect(order.indexOf("create")).toBeLessThan(order.indexOf("start"));
expect(commands.find((args) => args[0] === "start")).toEqual([
"start",
"--attach",
"--interactive",
name,
]);
expect(commands.filter((args) => args[0] === "rm")).toEqual([["rm", "-f", name]]);
}),
).pipe(Effect.provide([NodeServices.layer, TestClock.layer()])),
);

it.effect("bounds helper creation and still removes the helper by name afterwards", () =>
Effect.scoped(
Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const crypto = yield* Crypto.Crypto;
const root = yield* fs.makeTempDirectoryScoped({
prefix: "docker-helper-create-timeout-",
});
const creating = yield* Deferred.make<void>();
const commands: string[][] = [];
const killed = yield* Deferred.make<void>();
const present = yield* Ref.make(false);
const warning =
"WARNING: The requested image's platform does not match the host platform";
const spawner = ChildProcessSpawner.make((command) =>
Effect.gen(function* () {
if (!ChildProcess.isStandardCommand(command))
return yield* Effect.die("Unexpected command");
commands.push([...command.args]);
const create = command.args[0] === "create";
if (create) yield* Deferred.succeed(creating, undefined);
const missing = command.args[0] === "rm" && !(yield* Ref.getAndSet(present, false));
return ChildProcessSpawner.makeHandle({
pid: ChildProcessSpawner.ProcessId(0),
// The daemon never answers the create.
exitCode: create
? Effect.never
: Effect.succeed(ChildProcessSpawner.ExitCode(missing ? 1 : 0)),
isRunning: Effect.succeed(create),
kill: () => (create ? Deferred.succeed(killed, undefined) : Effect.void),
stdin: Sink.drain,
stdout: Stream.empty,
stderr: create
? Stream.succeed(new TextEncoder().encode(`${warning}\n`)).pipe(
Stream.concat(Stream.fromEffect(Deferred.await(killed)).pipe(Stream.drain)),
)
: missing
? Stream.succeed(new TextEncoder().encode("No such container\n"))
: Stream.empty,
all: Stream.empty,
getInputFd: () => Sink.drain,
getOutputFd: () => Stream.empty,
unref: Effect.succeed(Effect.void),
});
}),
);
const owner = yield* Scope.fork(yield* Scope.Scope);
const storage = yield* makeDockerDatabaseStorage({
runtime: "docker",
target: dockerTarget,
stackId: "helper-create-timeout",
instanceId: "database",
instanceRoot: root,
root,
cacheRoot: path.join(root, "cache"),
fs,
path,
crypto,
spawner,
container: {
prepare: () => Effect.void,
prepareImage: (image) => Effect.succeed(image),
launch: () => Effect.die("unused"),
launchCommand: () => Effect.die("unused"),
},
}).pipe(Scope.provide(owner));
yield* storage.prepare("17");
const operation = yield* storage.removeData("17").pipe(Effect.flip, Effect.forkScoped);
yield* Deferred.await(creating);
yield* TestClock.adjust("30 seconds");
const failure = yield* Fiber.join(operation);
expect(failure.message).toBe(
`Database helper was not created within 30 seconds: ${warning}`,
);
expect(yield* Deferred.isDone(killed)).toBe(true);
const create = commands.find((args) => args[0] === "create");
const name = create?.[create.indexOf("--name") + 1];
expect(name).toMatch(/^supabase-db-helper-/u);
expect(commands.some((args) => args[0] === "start")).toBe(false);
/* The daemon finishes the create after the failure cleanup found nothing to remove. */
yield* Ref.set(present, true);
yield* Scope.close(owner, Exit.void);
expect(commands.filter((args) => args[0] === "rm")).toEqual([
["rm", "-f", name],
["rm", "-f", name],
]);
expect(yield* Ref.get(present)).toBe(false);
}),
).pipe(Effect.provide([NodeServices.layer, TestClock.layer()])),
);
Expand Down Expand Up @@ -486,9 +623,9 @@ describe.runIf(testEngine === "docker")(
yield* storage.prepare("17");
yield* storage.removeData("17");

const run = engine.commands.find((args) => args[0] === "run");
const shellIndex = run?.indexOf("/bin/sh") ?? -1;
expect(run?.[shellIndex - 1]).toBe(helperMirror);
const create = engine.commands.find((args) => args[0] === "create");
const shellIndex = create?.indexOf("/bin/sh") ?? -1;
expect(create?.[shellIndex - 1]).toBe(helperMirror);
expect(
engine.commands.filter((args) => args[0] === "pull").map((args) => args.at(-1)),
).toEqual([expect.stringContaining("ghcr.io/supabase/cli/postgres:"), helperMirror]);
Expand Down
Loading
Loading