Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -67,8 +67,11 @@ it again to apply those changes.
`storage`, `functions`, `studio`, `mail`, `analytics`, and `pooler`. Database cannot be excluded.
Storage includes its Imgproxy companion, Studio includes Pgmeta, and Analytics includes Vector.
Studio requires REST; excluding REST while keeping Studio fails before stopping the composition.
Vector runs a stack-owned default configuration that enables its health API and forwards no service
logs; log collection into Analytics is not implemented yet.
Vector runs a stack-owned default configuration that enables its health API. With Docker or Podman,
Vector mounts the engine socket read-only (or reaches a TCP `DOCKER_HOST` through
`host.docker.internal`) and ships the stack's Auth, REST, Realtime, Storage, Functions, and database
container logs to Analytics. The native runtime keeps service output in memory, so native Vector
forwards no service logs.

After an explicit stop, start compares the project configuration with the saved composition through
the stack package's composition plan, ignoring values the composition and stack credentials supply.
Expand Down
83 changes: 80 additions & 3 deletions packages/stack/src/runtime/Container.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,10 @@ interface ContainerSpec {
readonly image: string;
readonly stackId: string;
readonly instanceId: string;
/** Labels the container with its service kind so stack log collectors can route it. */
readonly service?: string;
/** Exposes the engine API at the Docker-compatible default socket or `DOCKER_HOST`. */
readonly engineApi?: boolean;
readonly env: Readonly<Record<string, string>>;
readonly args?: ReadonlyArray<string>;
readonly entrypoint?: string;
Expand Down Expand Up @@ -124,6 +128,49 @@ const PublishedPorts = Schema.Record(
),
);

interface EngineApiAccess {
readonly env: Readonly<Record<string, string>>;
readonly mounts: NonNullable<ContainerSpec["mounts"]>;
readonly securityOpt: ReadonlyArray<string>;
}

const DEFAULT_ENGINE_SOCKET = "/var/run/docker.sock";

// Docker Desktop and Colima expose a per-user socket on the host but serve the rootful daemon
// socket at the default path inside their VM, which is the only one bind mounts can reach.
const vmDefaultSocket = (socket: string) =>
/\/\.docker\/(run|desktop)\/docker\.sock$/u.test(socket) ||
(socket.includes("/.colima/") && socket.endsWith("/docker.sock"));

const engineApiAccess = (endpoint: string): Effect.Effect<EngineApiAccess, ContainerError> => {
if (endpoint.startsWith("unix://")) {
const socket = endpoint.slice("unix://".length);
const vmDefault = vmDefaultSocket(socket);
return Effect.succeed({
env: {},
mounts: [
{
source: vmDefault ? DEFAULT_ENGINE_SOCKET : socket,
target: DEFAULT_ENGINE_SOCKET,
readOnly: true,
},
],
securityOpt: vmDefault ? [] : ["label=disable"],
});
}
const tcp = /^tcp:\/\/[^/]*?(?::(\d+))?(?:\/.*)?$/u.exec(endpoint);
// A named pipe cannot be mounted, so containers reach it only through a TCP-exposed daemon.
const port =
tcp === null ? (endpoint.startsWith("npipe://") ? "2375" : undefined) : (tcp[1] ?? "2375");
if (port !== undefined)
return Effect.succeed({
env: { DOCKER_HOST: `http://host.docker.internal:${port}` },
mounts: [],
securityOpt: [],
});
return Effect.fail(errorFor("engine-api", `Unsupported engine endpoint ${endpoint}`));
};

const mountField = (key: string, value: string) => {
const field = `${key}=${value}`;
return /[,"\n\r]/u.test(field) ? `"${field.replaceAll('"', '""')}"` : field;
Expand Down Expand Up @@ -258,13 +305,41 @@ export const makeContainerRuntime = (options: {
prepareImage(image).pipe(Effect.asVoid),
);

const engineApi = Effect.gen(function* () {
if (options.engine === "podman") {
// Podman runs containers without its API service, but reports the socket path regardless.
const [socket = "", exists] = (yield* run([
"info",
"--format",
"{{.Host.RemoteSocket.Path}} {{.Host.RemoteSocket.Exists}}",
])).split(" ");
if (socket.length === 0 || exists !== "true")
return yield* errorFor(
"engine-api",
"Podman's API socket is not active; enable it with `systemctl --user enable --now podman.socket`",
);
return socket.includes("://") ? socket : `unix://${socket}`;
}
// Reports DOCKER_HOST when set, otherwise the current context's endpoint.
const endpoint = yield* run([
"context",
"inspect",
"--format",
"{{.Endpoints.docker.Host}}",
]).pipe(Effect.orElseSucceed(() => ""));
return endpoint.length > 0 ? endpoint : `unix://${DEFAULT_ENGINE_SOCKET}`;
}).pipe(Effect.flatMap(engineApiAccess));

const launch = Effect.fn("Container.launch")(function* (
spec: ContainerSpec,
interactive = false,
) {
const owner = yield* Scope.Scope;
const image = (yield* Ref.get(mirrored)).get(spec.image) ?? spec.image;
for (const [key, value] of Object.entries(spec.env)) {
const access = spec.engineApi === true ? yield* engineApi : undefined;
const env = { ...spec.env, ...access?.env };
const mounts = [...(spec.mounts ?? []), ...(access?.mounts ?? [])];
for (const [key, value] of Object.entries(env)) {
if (!/^[A-Za-z_][A-Za-z0-9_]*$/u.test(key) || /[\0\r\n]/u.test(value)) {
return yield* errorFor(
"environment",
Expand All @@ -283,7 +358,7 @@ export const makeContainerRuntime = (options: {
yield* fs
.writeFileString(
envPath,
Object.entries(spec.env)
Object.entries(env)
.map(([key, value]) => `${key}=${value}`)
.join("\n"),
{ mode: 0o600 },
Expand All @@ -309,9 +384,11 @@ export const makeContainerRuntime = (options: {
`com.supabase.instance=${spec.instanceId}`,
"--label",
`com.supabase.stack-root=${stackRoot}`,
...(spec.service === undefined ? [] : ["--label", `com.supabase.service=${spec.service}`]),
...(access?.securityOpt ?? []).flatMap((option) => ["--security-opt", option]),
"--env-file",
envPath,
...(spec.mounts ?? []).flatMap((mount) => [
...mounts.flatMap((mount) => [
"--mount",
[
`type=${mount.type ?? "bind"}`,
Expand Down
1 change: 1 addition & 0 deletions packages/stack/src/services/Database.ts
Original file line number Diff line number Diff line change
Expand Up @@ -720,6 +720,7 @@ export const makeDatabase = (
image: image.image,
stackId: String(options.stackId),
instanceId: options.instanceId,
service: "database",
env: {
PGDATA: "/var/lib/postgresql/data",
PGSODIUM_KEY_FILE: "/etc/postgresql-custom/pgsodium_root.key",
Expand Down
3 changes: 3 additions & 0 deletions packages/stack/src/services/ProcessRecipe.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,7 @@ export interface ProcessRecipeSpec<C extends RecipeCreation<ServiceKind, unknown
readonly enabledPort?: (creation: C, name: string) => boolean;
readonly containerPort?: (creation: C, name: string, port: number) => number;
readonly containerEntrypoint?: (creation: C) => string | undefined;
readonly engineApi?: boolean;
readonly prepare?: (creation: C) => Effect.Effect<void, ServiceError>;
readonly removeData?: (creation: C) => Effect.Effect<void, ServiceError>;
}
Expand Down Expand Up @@ -780,6 +781,8 @@ export const makeProcessRecipe = <C extends RecipeCreation<ServiceKind, unknown>
image: resolved.image,
stackId: options.stackId,
instanceId: options.instanceId,
service: spec.service,
engineApi: spec.engineApi === true,
env: yield* spec.env(context.config, containerDesired, true),
entrypoint: spec.containerEntrypoint?.(context.config),
args: yield* spec.args(context.config, containerDesired, { container: true }),
Expand Down
125 changes: 124 additions & 1 deletion packages/stack/src/services/Vector.integration.test.ts
Original file line number Diff line number Diff line change
@@ -1,10 +1,14 @@
import { PgClient } from "@effect/sql-pg";
import { NodeHttpClient, NodeServices } from "@effect/platform-node";
import { describe, expect, it } from "@effect/vitest";
import { tmpdir } from "node:os";
import { Effect, FileSystem, Layer } from "effect";
import { Context, Crypto, Effect, FileSystem, Layer, Redacted, Schedule } from "effect";
import { HttpClient, HttpClientRequest } from "effect/unstable/http";
import { makeService } from "../Service.ts";
import { makeServiceRecipe } from "./Catalog.ts";
import type { ServiceEndpoint } from "./Recipe.ts";
import { makeDockerHttpRelay, makeDockerTcpRelay } from "../../tests/docker-relay.ts";
import { makeDockerDatabaseRoot } from "../../tests/docker-fixture.ts";

const options = (root: string, runtime: "docker" | "native") => ({
stackId: "catalog-test",
Expand All @@ -14,6 +18,24 @@ const options = (root: string, runtime: "docker" | "native") => ({
runtime,
});

const password = Redacted.make("postgres");

const query = <A extends object>(endpoint: ServiceEndpoint, database: string, statement: string) =>
Effect.scoped(
Effect.gen(function* () {
const services = yield* Layer.build(
PgClient.layer({
host: endpoint.host ?? "127.0.0.1",
port: endpoint.port,
database,
username: "supabase_admin",
password,
}),
);
return yield* Context.get(services, PgClient.PgClient).unsafe<A>(statement);
}),
);

describe("vector recipe", () => {
for (const runtime of ["docker", "native"] as const) {
it.live(
Expand Down Expand Up @@ -119,4 +141,105 @@ describe("vector recipe", () => {
}),
).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))),
);

it.live(
"ships the stack's service logs to Analytics with its default pipeline",
() =>
Effect.scoped(
Effect.gen(function* () {
const id = yield* Crypto.Crypto.pipe(Effect.flatMap((crypto) => crypto.randomUUIDv4));
const marker = `vector-pipeline-${id.slice(0, 8)}`;
const options = {
stackId: marker,
instanceId: "instance",
root: yield* makeDockerDatabaseRoot("vector-pipeline-", marker),
runtime: "docker" as const,
};
const catalogOptions = { ...options, cacheRoot: `${options.root}/cache` };
const databaseRecipe = yield* makeServiceRecipe(
{
service: "database",
config: {
version: "17",
databasePassword: password,
jwtSecret: Redacted.make("vector-pipeline-secret-with-at-least-32-chars"),
jwtExpiry: 3600,
},
},
catalogOptions,
);
const database = yield* makeService(databaseRecipe.definition, {
id: "database",
config: databaseRecipe.creation,
});
yield* database.start;
yield* database.ready;
const sql = yield* databaseRecipe.endpoint("sql");
const databaseRelay = yield* makeDockerTcpRelay(databaseRecipe.endpoint("sql"));

const analyticsRecipe = yield* makeServiceRecipe(
{
service: "analytics",
config: {
databaseUrl: `postgresql://supabase_admin:postgres@${databaseRelay.host}:${databaseRelay.port}/_supabase`,
backend: "postgres",
apiKey: "vector-pipeline-key",
},
},
catalogOptions,
);
const analytics = yield* makeService(analyticsRecipe.definition, {
id: "analytics",
config: analyticsRecipe.creation,
});
yield* analytics.start;
yield* analytics.ready;
const analyticsRelay = yield* makeDockerHttpRelay(analyticsRecipe.endpoint("http"));

const vectorRecipe = yield* makeServiceRecipe(
{
service: "vector",
config: {
analyticsUrl: `http://${analyticsRelay.host}:${analyticsRelay.port}`,
apiKey: "vector-pipeline-key",
},
},
catalogOptions,
);
const vector = yield* makeService(vectorRecipe.definition, {
id: "vector",
config: vectorRecipe.creation,
});
yield* vector.start;
yield* vector.ready;

yield* query(sql, "postgres", `DO $$ BEGIN RAISE LOG '${marker}'; END $$`);

const [source] = yield* query<{ readonly token: string }>(
sql,
"_supabase",
"SELECT replace(token::text, '-', '_') AS token FROM _analytics.sources WHERE name = 'postgres.logs'",
);
const events = `_analytics."log_events_${source?.token}"`;
// Vector batches and Logflare inserts asynchronously, with no completion signal to await.
const shipped = yield* query<{ readonly message: string; readonly severity: string }>(
sql,
"_supabase",
`SELECT body->>'event_message' AS message, body->'metadata'->'parsed'->>'error_severity' AS severity FROM ${events} WHERE body->>'event_message' LIKE '%LOG: ${marker}'`,
).pipe(
Effect.filterOrFail((rows) => rows.length > 0),
Effect.retry(Schedule.spaced("500 millis")),
Effect.timeout("60 seconds"),
);
expect(shipped).toEqual([
{ message: expect.stringContaining(`LOG: ${marker}`), severity: "LOG" },
]);

yield* vector.stop;
yield* analytics.stop;
yield* database.stop;
}),
).pipe(Effect.provide(Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp))),
{ timeout: 180_000 },
);
});
Loading