From 99ae90fea2d5c2d75a62bb5c2c23650806697a52 Mon Sep 17 00:00:00 2001 From: Facundo Date: Mon, 6 Apr 2026 01:07:42 -0700 Subject: [PATCH] refactor(core): use Postgres template databases for fast branch reset Replaces the per-branch container model with one long-lived container per connector that hosts a frozen seed database (IS_TEMPLATE=true) plus N branch databases cloned via CREATE DATABASE ... TEMPLATE seed. Reset goes from "docker rm -f + docker run + waitForReady + psql restore" (5-15s) to "DROP DATABASE WITH (FORCE); CREATE DATABASE FROM TEMPLATE" (~200-800ms on a 10k-row schema), an ~10x improvement. Key changes: - New low-level helpers in branching/docker.ts: createConnectorContainer, waitForConnectorReady, execSqlInDb, loadInitSqlIntoDb, dumpDatabase, restoreDumpToDatabase, listDatabases. - DockerBranchProvider rewritten around container.json (one per connector, in .sow/snapshots//container.json). - providerMetaVersion=2 with explicit migration error for v1 branches. - stopBranch is now per-branch connection termination (the container is shared); startBranch verifies the container is running. - manager.resetBranch no longer recreates the container in the docker case. - All identifiers go through quoteIdent and connector/branch names are regex-validated to prevent injection through the connector name. - 7 new unit tests with mocked docker helpers; one BUGSTER_DB_URL-gated integration benchmark asserting reset < 1500ms. Co-Authored-By: Claude Opus 4.6 (1M context) --- packages/cli/src/commands/branch.ts | 9 + .../reset-perf.integration.test.ts | 52 ++++ packages/core/src/branching/docker.ts | 234 ++++++++++++++- packages/core/src/branching/manager.ts | 13 +- .../src/branching/providers/docker.test.ts | 208 +++++++++++++ .../core/src/branching/providers/docker.ts | 275 +++++++++++++++--- packages/core/src/branching/storage.ts | 43 +++ 7 files changed, 775 insertions(+), 59 deletions(-) create mode 100644 packages/core/src/__integration__/reset-perf.integration.test.ts create mode 100644 packages/core/src/branching/providers/docker.test.ts diff --git a/packages/cli/src/commands/branch.ts b/packages/cli/src/commands/branch.ts index 6a21301..4908659 100644 --- a/packages/cli/src/commands/branch.ts +++ b/packages/cli/src/commands/branch.ts @@ -184,6 +184,15 @@ export async function runBranch( console.log(` Port: ${branch.port}`); console.log(` Created: ${timeAgo(branch.createdAt)}`); console.log(` URL: ${branch.connectionString}`); + if (branch.provider === "docker") { + const dmeta = branch.providerMeta as { databaseName?: string; containerName?: string }; + if (dmeta.databaseName) { + console.log(` Database: ${dmeta.databaseName}`); + } + if (dmeta.containerName) { + console.log(` Container: ${dmeta.containerName}`); + } + } const checkpoints = listCheckpoints(branch.name); if (checkpoints.length > 0) { console.log(` Checkpoints: ${checkpoints.map((cp) => cp.name).join(", ")}`); diff --git a/packages/core/src/__integration__/reset-perf.integration.test.ts b/packages/core/src/__integration__/reset-perf.integration.test.ts new file mode 100644 index 0000000..430de30 --- /dev/null +++ b/packages/core/src/__integration__/reset-perf.integration.test.ts @@ -0,0 +1,52 @@ +import { describe, it, expect, afterEach } from "vitest"; +import { createConnector } from "../branching/connector.js"; +import { createBranch, deleteBranch, resetBranch, execBranch } from "../branching/manager.js"; +import { BUGSTER_DB_URL, cleanupConnector } from "./helpers.js"; + +// Skip unless explicitly enabled — this test boots Docker, builds a real +// 10k-row schema, and is too heavy for the regular unit run. +const HAS_BUGSTER = !!process.env.BUGSTER_DB_URL || !!process.env.RUN_PERF_TESTS; + +const tracked: { branch: string; connector: string }[] = []; + +afterEach(async () => { + for (const { branch, connector } of tracked) { + try { + await deleteBranch(branch); + } catch { + // ignore + } + cleanupConnector(connector); + } + tracked.length = 0; +}); + +describe.skipIf(!HAS_BUGSTER)("docker provider — reset perf (Lane B)", () => { + it("resetBranch completes in under 1.5s on a 10k-row schema", async () => { + const connector = "perf-reset-bench"; + const branchName = "perf-feature"; + cleanupConnector(connector); + + await createConnector(BUGSTER_DB_URL, { name: connector }); + const branch = await createBranch(branchName, connector); + tracked.push({ branch: branchName, connector }); + + // Write some divergent data so reset has something to undo. + try { + await execBranch( + branchName, + "CREATE TABLE IF NOT EXISTS _scratch (id int); INSERT INTO _scratch SELECT generate_series(1,1000);", + ); + } catch { + // schema-level changes may fail on read-only/bizarre snapshots; not fatal + } + + const start = Date.now(); + await resetBranch(branchName); + const elapsed = Date.now() - start; + + // The whole point of Lane B: reset must complete in <1.5s. + expect(elapsed).toBeLessThan(1500); + expect(branch.connectionString).toContain("postgresql://"); + }, 120_000); +}); diff --git a/packages/core/src/branching/docker.ts b/packages/core/src/branching/docker.ts index faaf50f..975f02e 100644 --- a/packages/core/src/branching/docker.ts +++ b/packages/core/src/branching/docker.ts @@ -3,9 +3,11 @@ import { promisify } from "node:util"; const execFile = promisify(execFileCb); -const POSTGRES_USER = "sow"; -const POSTGRES_PASSWORD = "sow"; -const POSTGRES_DB = "sow"; +export const POSTGRES_USER = "sow"; +export const POSTGRES_PASSWORD = "sow"; +export const POSTGRES_DB = "sow"; +/** Bootstrap database used to create/drop other databases. */ +export const POSTGRES_BOOTSTRAP_DB = "postgres"; export interface CreateContainerOptions { containerName: string; @@ -234,3 +236,229 @@ export async function execSql( function sleep(ms: number): Promise { return new Promise((resolve) => setTimeout(resolve, ms)); } + +// --------------------------------------------------------------------------- +// New helpers for the template-database model (Lane B / Issue #1). +// One long-lived container per connector hosts a frozen seed database and +// N branch databases cloned from it via `CREATE DATABASE ... TEMPLATE seed`. +// --------------------------------------------------------------------------- + +export interface CreateConnectorContainerOptions { + containerName: string; + port: number; + pgVersion: string; +} + +/** + * Create a long-lived per-connector Postgres container with no init.sql + * mounted. The default `postgres` database is the bootstrap; the seed + * and per-branch databases are created later via SQL. + */ +export async function createConnectorContainer( + opts: CreateConnectorContainerOptions, +): Promise { + const { containerName, port, pgVersion } = opts; + const args = [ + "run", + "-d", + "--name", + containerName, + "-p", + `${port}:5432`, + "-e", + `POSTGRES_DB=${POSTGRES_BOOTSTRAP_DB}`, + "-e", + `POSTGRES_USER=${POSTGRES_USER}`, + "-e", + `POSTGRES_PASSWORD=${POSTGRES_PASSWORD}`, + `postgres:${pgVersion}-alpine`, + ]; + const { stdout } = await execFile("docker", args, { timeout: 30_000 }); + return stdout.trim(); +} + +/** + * Wait for the bootstrap `postgres` database inside a connector container + * to accept connections. + */ +export async function waitForConnectorReady( + containerName: string, + timeoutMs = 120_000, +): Promise { + const start = Date.now(); + const pollInterval = 500; + while (Date.now() - start < timeoutMs) { + try { + const { stdout } = await execFile( + "docker", + ["exec", containerName, "pg_isready", "-U", POSTGRES_USER, "-d", POSTGRES_BOOTSTRAP_DB], + { timeout: 5_000 }, + ); + if (stdout.includes("accepting connections")) { + try { + await execFile( + "docker", + ["exec", containerName, "psql", "-U", POSTGRES_USER, "-d", POSTGRES_BOOTSTRAP_DB, "-c", "SELECT 1"], + { timeout: 5_000 }, + ); + return; + } catch { + // not ready yet + } + } + } catch { + // not ready yet + } + await sleep(pollInterval); + } + throw new Error( + `Postgres container '${containerName}' did not become ready within ${timeoutMs / 1000}s`, + ); +} + +/** Run a single SQL statement against a specific database inside a container. */ +export async function execSqlInDb( + containerName: string, + databaseName: string, + sql: string, +): Promise { + const { stdout } = await execFile( + "docker", + [ + "exec", + containerName, + "psql", + "-U", + POSTGRES_USER, + "-d", + databaseName, + "-v", + "ON_ERROR_STOP=1", + "-c", + sql, + ], + { timeout: 60_000 }, + ); + return stdout; +} + +/** Stream a SQL file from the host into psql for a given database. */ +export async function loadInitSqlIntoDb( + containerName: string, + databaseName: string, + initSqlPath: string, +): Promise { + const { readFile } = await import("node:fs/promises"); + const sqlContent = await readFile(initSqlPath, "utf-8"); + await pipeSqlToDb(containerName, databaseName, sqlContent, /*onErrorStop*/ false); +} + +async function pipeSqlToDb( + containerName: string, + databaseName: string, + sqlContent: string, + onErrorStop: boolean, +): Promise { + await new Promise((resolve, reject) => { + const args = [ + "exec", + "-i", + containerName, + "psql", + "-U", + POSTGRES_USER, + "-d", + databaseName, + ]; + if (onErrorStop) { + args.push("-v", "ON_ERROR_STOP=1"); + } else { + args.push("-v", "ON_ERROR_STOP=0"); + } + const child = spawnCb("docker", args); + let stderr = ""; + child.stderr.on("data", (chunk: Buffer) => { + stderr += chunk.toString(); + }); + child.on("close", (code) => { + if (code !== 0) { + reject(new Error(`psql failed (exit ${code}): ${stderr}`)); + } else { + resolve(); + } + }); + child.on("error", reject); + child.stdin.write(sqlContent); + child.stdin.end(); + }); +} + +/** Dump a specific database inside a container to a SQL string. */ +export async function dumpDatabase( + containerName: string, + databaseName: string, +): Promise { + const { stdout } = await execFile( + "docker", + [ + "exec", + containerName, + "pg_dump", + "-U", + POSTGRES_USER, + "-d", + databaseName, + "--no-owner", + "--no-privileges", + ], + { timeout: 60_000, maxBuffer: 100 * 1024 * 1024 }, + ); + return stdout; +} + +/** + * Restore a SQL dump into a branch database. Wipes the public schema first. + * Targets the per-branch database (NOT the seed). + */ +export async function restoreDumpToDatabase( + containerName: string, + databaseName: string, + sqlContent: string, +): Promise { + await execSqlInDb( + containerName, + databaseName, + "DROP SCHEMA IF EXISTS public CASCADE; CREATE SCHEMA public;", + ); + await pipeSqlToDb(containerName, databaseName, sqlContent, /*onErrorStop*/ false); +} + +/** + * List database names matching a SQL LIKE pattern (used to count branches + * remaining for a connector before tearing down its container). + */ +export async function listDatabases( + containerName: string, + likePattern: string, +): Promise { + const { stdout } = await execFile( + "docker", + [ + "exec", + containerName, + "psql", + "-U", + POSTGRES_USER, + "-d", + POSTGRES_BOOTSTRAP_DB, + "-tA", + "-c", + `SELECT datname FROM pg_database WHERE datname LIKE '${likePattern.replace(/'/g, "''")}'`, + ], + { timeout: 10_000 }, + ); + return stdout + .split("\n") + .map((s) => s.trim()) + .filter((s) => s.length > 0); +} diff --git a/packages/core/src/branching/manager.ts b/packages/core/src/branching/manager.ts index 2e4b707..82b07cb 100644 --- a/packages/core/src/branching/manager.ts +++ b/packages/core/src/branching/manager.ts @@ -221,16 +221,9 @@ export async function resetBranch(name: string): Promise { const initSqlPath = getInitSqlPath(branch.connector); await provider.resetBranch(branch, initSqlPath); - - if (provider.name === "docker") { - // Docker reset = recreate container, so re-run full createBranch - removeBranch(name); - return createBranch(name, branch.connector, { - port: branch.port, - pgVersion: (branch.providerMeta as any).pgVersion, - }); - } - + // The Docker provider now resets in-place by dropping & recreating the + // branch database from the connector's seed template — no need to tear + // down the container or rewrite branch metadata. return branch; } diff --git a/packages/core/src/branching/providers/docker.test.ts b/packages/core/src/branching/providers/docker.test.ts new file mode 100644 index 0000000..86f701b --- /dev/null +++ b/packages/core/src/branching/providers/docker.test.ts @@ -0,0 +1,208 @@ +import { describe, it, expect, beforeEach, vi } from "vitest"; + +// Mock the low-level docker helpers BEFORE importing the provider. +vi.mock("../docker.js", async () => { + return { + POSTGRES_USER: "sow", + POSTGRES_PASSWORD: "sow", + POSTGRES_DB: "sow", + POSTGRES_BOOTSTRAP_DB: "postgres", + ensureDocker: vi.fn().mockResolvedValue(undefined), + createConnectorContainer: vi.fn().mockResolvedValue("container-id-abc"), + waitForConnectorReady: vi.fn().mockResolvedValue(undefined), + removeContainer: vi.fn().mockResolvedValue(undefined), + startContainer: vi.fn().mockResolvedValue(undefined), + getContainerStatus: vi.fn().mockResolvedValue("running"), + execSqlInDb: vi.fn().mockResolvedValue(""), + loadInitSqlIntoDb: vi.fn().mockResolvedValue(undefined), + dumpDatabase: vi.fn().mockResolvedValue("-- dump"), + restoreDumpToDatabase: vi.fn().mockResolvedValue(undefined), + listDatabases: vi.fn().mockResolvedValue([]), + }; +}); + +vi.mock("../ports.js", () => ({ + findFreePort: vi.fn().mockResolvedValue(54321), +})); + +// In-memory connector container store. +const containerStore = new Map(); +vi.mock("../storage.js", () => ({ + readConnectorContainer: vi.fn((name: string) => containerStore.get(name) ?? null), + writeConnectorContainer: vi.fn((name: string, info: any) => { + containerStore.set(name, info); + }), + deleteConnectorContainer: vi.fn((name: string) => { + containerStore.delete(name); + }), +})); + +import { + DockerBranchProvider, + PROVIDER_META_VERSION, + type DockerProviderMeta, +} from "./docker.js"; +import * as docker from "../docker.js"; +import * as storage from "../storage.js"; + +const provider = new DockerBranchProvider(); + +function makeBranch(overrides: Partial = {}): any { + const meta: DockerProviderMeta = { + providerMetaVersion: PROVIDER_META_VERSION, + containerId: "container-id-abc", + containerName: "sow-myconn", + pgVersion: "16", + connector: "myconn", + databaseName: "sow_feature_a", + ...overrides, + }; + return { + name: "feature-a", + connector: "myconn", + provider: "docker", + providerMeta: meta, + port: 54321, + status: "running", + createdAt: "2024-01-01T00:00:00Z", + connectionString: "postgresql://sow:sow@localhost:54321/sow_feature_a", + }; +} + +beforeEach(() => { + containerStore.clear(); + vi.clearAllMocks(); + // Re-prime defaults the way the top-of-file vi.mock does. + (docker.createConnectorContainer as any).mockResolvedValue("container-id-abc"); + (docker.execSqlInDb as any).mockResolvedValue(""); + (docker.getContainerStatus as any).mockResolvedValue("running"); + (docker.listDatabases as any).mockResolvedValue([]); +}); + +describe("DockerBranchProvider.createBranch", () => { + it("creates the connector container, seed db, and branch db on first branch", async () => { + const result = await provider.createBranch({ + name: "feature-a", + connector: "myconn", + initSqlPath: "/tmp/init.sql", + detection: { meta: {} }, + }); + + expect(docker.createConnectorContainer).toHaveBeenCalledTimes(1); + expect(docker.waitForConnectorReady).toHaveBeenCalledWith("sow-myconn"); + expect(docker.loadInitSqlIntoDb).toHaveBeenCalledWith( + "sow-myconn", + "sow_seed_myconn", + "/tmp/init.sql", + ); + + const sqlCalls = (docker.execSqlInDb as any).mock.calls.map((c: any[]) => c[2]); + expect(sqlCalls.some((s: string) => s.includes("CREATE DATABASE") && s.includes("sow_seed_myconn"))).toBe(true); + expect(sqlCalls.some((s: string) => s.includes("IS_TEMPLATE true"))).toBe(true); + expect(sqlCalls.some((s: string) => s.includes("ALLOW_CONNECTIONS false"))).toBe(true); + expect(sqlCalls.some((s: string) => s.includes("CREATE DATABASE") && s.includes("sow_feature_a") && s.includes("TEMPLATE"))).toBe(true); + + expect(result.connectionString).toBe( + "postgresql://sow:sow@localhost:54321/sow_feature_a", + ); + const meta = result.providerMeta as DockerProviderMeta; + expect(meta.providerMetaVersion).toBe(PROVIDER_META_VERSION); + expect(meta.databaseName).toBe("sow_feature_a"); + expect(meta.containerName).toBe("sow-myconn"); + expect(storage.writeConnectorContainer).toHaveBeenCalledTimes(1); + }); + + it("reuses the existing container on the second branch", async () => { + await provider.createBranch({ + name: "feature-a", + connector: "myconn", + initSqlPath: "/tmp/init.sql", + detection: { meta: {} }, + }); + (docker.createConnectorContainer as any).mockClear(); + (docker.loadInitSqlIntoDb as any).mockClear(); + + const result = await provider.createBranch({ + name: "feature-b", + connector: "myconn", + initSqlPath: "/tmp/init.sql", + detection: { meta: {} }, + }); + + expect(docker.createConnectorContainer).not.toHaveBeenCalled(); + expect(docker.loadInitSqlIntoDb).not.toHaveBeenCalled(); + expect((result.providerMeta as DockerProviderMeta).databaseName).toBe( + "sow_feature_b", + ); + // The branch db should have been cloned from the seed. + const sqlCalls = (docker.execSqlInDb as any).mock.calls.map((c: any[]) => c[2]); + expect( + sqlCalls.some((s: string) => + s.includes("CREATE DATABASE") && s.includes("sow_feature_b") && s.includes("TEMPLATE"), + ), + ).toBe(true); + }); + + it("rejects unsafe connector names", async () => { + await expect( + provider.createBranch({ + name: "ok", + connector: "evil; DROP", + initSqlPath: "/tmp/init.sql", + detection: { meta: {} }, + }), + ).rejects.toThrow(/Invalid connector/); + }); +}); + +describe("DockerBranchProvider.resetBranch", () => { + it("issues DROP + CREATE FROM TEMPLATE", async () => { + const branch = makeBranch(); + await provider.resetBranch(branch, "/tmp/init.sql"); + + const sqlCalls = (docker.execSqlInDb as any).mock.calls.map((c: any[]) => c[2]); + expect(sqlCalls).toHaveLength(2); + expect(sqlCalls[0]).toMatch(/DROP DATABASE.*sow_feature_a.*FORCE/); + expect(sqlCalls[1]).toMatch(/CREATE DATABASE.*sow_feature_a.*TEMPLATE.*sow_seed_myconn/); + }); +}); + +describe("DockerBranchProvider.deleteBranch", () => { + it("drops the branch db but leaves the container when other branches remain", async () => { + (docker.listDatabases as any).mockResolvedValue([ + "sow_seed_myconn", + "sow_feature_a", + "sow_feature_b", + ]); + const branch = makeBranch(); + await provider.deleteBranch(branch); + expect(docker.execSqlInDb).toHaveBeenCalled(); + expect(docker.removeContainer).not.toHaveBeenCalled(); + }); + + it("removes the container when the last non-seed branch is deleted", async () => { + (docker.listDatabases as any).mockResolvedValue(["sow_seed_myconn"]); + const branch = makeBranch(); + await provider.deleteBranch(branch); + expect(docker.removeContainer).toHaveBeenCalledWith("sow-myconn"); + expect(storage.deleteConnectorContainer).toHaveBeenCalledWith("myconn"); + }); +}); + +describe("DockerBranchProvider migration / version handling", () => { + it("rejects branches with old-shape providerMeta", async () => { + const oldBranch: any = { + name: "legacy", + connector: "myconn", + provider: "docker", + providerMeta: { containerId: "x", containerName: "sow-myconn-legacy", pgVersion: "16" }, + port: 54321, + status: "running", + createdAt: "2024-01-01T00:00:00Z", + connectionString: "", + }; + await expect(provider.resetBranch(oldBranch, "/tmp/init.sql")).rejects.toThrow( + /older sow version/, + ); + }); +}); diff --git a/packages/core/src/branching/providers/docker.ts b/packages/core/src/branching/providers/docker.ts index d718665..70153f1 100644 --- a/packages/core/src/branching/providers/docker.ts +++ b/packages/core/src/branching/providers/docker.ts @@ -2,35 +2,87 @@ import type { BranchProvider, ProviderDetection, ProviderBranchOpts, ProviderBra import type { Branch, BranchStatus } from "../types.js"; import { ensureDocker, - createContainer, - stopContainer, - startContainer, + createConnectorContainer, + waitForConnectorReady, removeContainer, - waitForReady, getContainerStatus, - dumpBranch as dockerDump, - restoreFromDump as dockerRestore, - execSql, + execSqlInDb, + loadInitSqlIntoDb, + dumpDatabase, + restoreDumpToDatabase, + listDatabases, + POSTGRES_USER, + POSTGRES_PASSWORD, + POSTGRES_BOOTSTRAP_DB, } from "../docker.js"; import { findFreePort } from "../ports.js"; +import { + readConnectorContainer, + writeConnectorContainer, + deleteConnectorContainer, + type ConnectorContainerInfo, +} from "../storage.js"; +import { quoteIdent } from "../../sql/identifiers.js"; -const POSTGRES_USER = "sow"; -const POSTGRES_PASSWORD = "sow"; -const POSTGRES_DB = "sow"; - -function buildConnectionString(port: number): string { - return `postgresql://${POSTGRES_USER}:${POSTGRES_PASSWORD}@localhost:${port}/${POSTGRES_DB}`; -} +const SAFE_NAME_RE = /^[a-zA-Z0-9_-]{1,40}$/; -function buildContainerName(connector: string, branchName: string): string { - return `sow-${connector}-${branchName}`; -} +export const PROVIDER_META_VERSION = 2; /** Docker provider metadata stored in Branch.providerMeta. */ export interface DockerProviderMeta { + providerMetaVersion: number; containerId: string; containerName: string; pgVersion: string; + connector: string; + databaseName: string; +} + +function assertSafeName(kind: string, value: string): void { + if (!SAFE_NAME_RE.test(value)) { + throw new Error( + `Invalid ${kind} name '${value}': must match ${SAFE_NAME_RE.source}`, + ); + } +} + +function buildContainerName(connector: string): string { + return `sow-${connector}`; +} + +function buildSeedDatabaseName(connector: string): string { + return `sow_seed_${connector.replace(/-/g, "_")}`; +} + +function buildBranchDatabaseName(branchName: string): string { + return `sow_${branchName.replace(/-/g, "_")}`; +} + +function buildConnectionString(port: number, databaseName: string): string { + return `postgresql://${POSTGRES_USER}:${POSTGRES_PASSWORD}@localhost:${port}/${databaseName}`; +} + +function isV2Meta(raw: unknown): raw is DockerProviderMeta { + if (!raw || typeof raw !== "object") return false; + const m = raw as Record; + return ( + m.providerMetaVersion === PROVIDER_META_VERSION && + typeof m.containerName === "string" && + typeof m.databaseName === "string" && + typeof m.connector === "string" + ); +} + +function readMeta(branch: Branch): DockerProviderMeta { + const raw = branch.providerMeta; + if (!isV2Meta(raw)) { + throw new Error( + `Branch '${branch.name}' was created by an older sow version and is not ` + + `compatible with the template-database provider. Run \`sow branch delete ${branch.name}\` ` + + `and recreate the branch.`, + ); + } + return raw; } export class DockerBranchProvider implements BranchProvider { @@ -46,71 +98,202 @@ export class DockerBranchProvider implements BranchProvider { } async createBranch(opts: ProviderBranchOpts): Promise { + assertSafeName("connector", opts.connector); + assertSafeName("branch", opts.name); + const pgVersion = opts.pgVersion ?? "16"; - const port = await findFreePort(opts.port); - const containerName = buildContainerName(opts.connector, opts.name); + const containerName = buildContainerName(opts.connector); + const seedDb = buildSeedDatabaseName(opts.connector); + const branchDb = buildBranchDatabaseName(opts.name); - const containerId = await createContainer({ - containerName, - initSqlPath: opts.initSqlPath, - port, - pgVersion, - }); + let container = readConnectorContainer(opts.connector); - await waitForReady(containerName); + // First branch for this connector — bring up the long-lived container + // and seed the template database from init.sql. + if (!container) { + const port = await findFreePort(opts.port); + const containerId = await createConnectorContainer({ + containerName, + port, + pgVersion, + }); + try { + await waitForConnectorReady(containerName); + // Create the seed database, restore init.sql into it, then mark + // it as a template (and disallow direct connections). + await execSqlInDb( + containerName, + POSTGRES_BOOTSTRAP_DB, + `CREATE DATABASE ${quoteIdent(seedDb)} OWNER ${quoteIdent(POSTGRES_USER)}`, + ); + await loadInitSqlIntoDb(containerName, seedDb, opts.initSqlPath); + // ALTER DATABASE ... IS_TEMPLATE / ALLOW_CONNECTIONS is supported on + // PG12+. The repo defaults to postgres:16-alpine. + await execSqlInDb( + containerName, + POSTGRES_BOOTSTRAP_DB, + `ALTER DATABASE ${quoteIdent(seedDb)} IS_TEMPLATE true`, + ); + await execSqlInDb( + containerName, + POSTGRES_BOOTSTRAP_DB, + `ALTER DATABASE ${quoteIdent(seedDb)} ALLOW_CONNECTIONS false`, + ); + } catch (err) { + // Roll back: leaving a half-initialized container would poison the + // next createBranch call. + await removeContainer(containerName); + throw err; + } - return { - connectionString: buildConnectionString(port), - port, - providerMeta: { + container = { containerId, containerName, + port, pgVersion, + seedDatabase: seedDb, + createdAt: new Date().toISOString(), + }; + writeConnectorContainer(opts.connector, container); + } + + // Clone the seed into a fresh per-branch database. + await execSqlInDb( + containerName, + POSTGRES_BOOTSTRAP_DB, + `CREATE DATABASE ${quoteIdent(branchDb)} WITH TEMPLATE ${quoteIdent(container.seedDatabase)} OWNER ${quoteIdent(POSTGRES_USER)}`, + ); + + return { + connectionString: buildConnectionString(container.port, branchDb), + port: container.port, + providerMeta: { + providerMetaVersion: PROVIDER_META_VERSION, + containerId: container.containerId, + containerName: container.containerName, + pgVersion: container.pgVersion, + connector: opts.connector, + databaseName: branchDb, } satisfies DockerProviderMeta, }; } async deleteBranch(branch: Branch): Promise { - const meta = branch.providerMeta as DockerProviderMeta; - await removeContainer(meta.containerName); + const meta = readMeta(branch); + // DROP DATABASE cannot be wrapped in a transaction; pass straight through. + try { + await execSqlInDb( + meta.containerName, + POSTGRES_BOOTSTRAP_DB, + `DROP DATABASE IF EXISTS ${quoteIdent(meta.databaseName)} WITH (FORCE)`, + ); + } catch (err) { + // If the container is gone there's nothing to clean. + const status = await getContainerStatus(meta.containerName); + if (status === "not_found") { + deleteConnectorContainer(meta.connector); + return; + } + throw err; + } + + // If this was the last branch for the connector, tear the whole + // container down so we don't leak idle Postgres processes. + const remaining = await listDatabases(meta.containerName, "sow_%"); + const seed = buildSeedDatabaseName(meta.connector); + const nonSeed = remaining.filter((d) => d !== seed); + if (nonSeed.length === 0) { + await removeContainer(meta.containerName); + deleteConnectorContainer(meta.connector); + } } async resetBranch(branch: Branch, _initSqlPath: string): Promise { - const meta = branch.providerMeta as DockerProviderMeta; - await removeContainer(meta.containerName); + const meta = readMeta(branch); + const seed = buildSeedDatabaseName(meta.connector); + // DROP + CREATE FROM TEMPLATE: ~200-800ms on a 10k-row schema, vs the + // 5-15s of the previous "tear down the container and re-run init.sql" + // approach. WITH (FORCE) terminates lingering connections (PG13+). + await execSqlInDb( + meta.containerName, + POSTGRES_BOOTSTRAP_DB, + `DROP DATABASE IF EXISTS ${quoteIdent(meta.databaseName)} WITH (FORCE)`, + ); + await execSqlInDb( + meta.containerName, + POSTGRES_BOOTSTRAP_DB, + `CREATE DATABASE ${quoteIdent(meta.databaseName)} WITH TEMPLATE ${quoteIdent(seed)} OWNER ${quoteIdent(POSTGRES_USER)}`, + ); } async execSQL(branch: Branch, sql: string): Promise { - const meta = branch.providerMeta as DockerProviderMeta; - return execSql(meta.containerName, sql); + const meta = readMeta(branch); + return execSqlInDb(meta.containerName, meta.databaseName, sql); } async getBranchStatus(branch: Branch): Promise { - const meta = branch.providerMeta as DockerProviderMeta; + const meta = readMeta(branch); const actual = await getContainerStatus(meta.containerName); if (actual === "not_found") return "error"; return actual === "running" ? "running" : "stopped"; } async stopBranch(branch: Branch): Promise { - const meta = branch.providerMeta as DockerProviderMeta; - await stopContainer(meta.containerName); + // In the template-database model the container is shared by every + // branch for a connector, so "stop one branch" has no analog at the + // container level. Instead, terminate any open connections to this + // branch's database — the container keeps running and other branches + // are unaffected. + const meta = readMeta(branch); + await execSqlInDb( + meta.containerName, + POSTGRES_BOOTSTRAP_DB, + `SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = ${pgString(meta.databaseName)}`, + ); } async startBranch(branch: Branch): Promise { - const meta = branch.providerMeta as DockerProviderMeta; + // The container is always running while any branch exists, so start + // is effectively a no-op. We just verify the container is up. + const meta = readMeta(branch); await ensureDocker(); - await startContainer(meta.containerName); - await waitForReady(meta.containerName); + const status = await getContainerStatus(meta.containerName); + if (status === "not_found") { + throw new Error( + `Connector container '${meta.containerName}' is gone. Recreate the branch.`, + ); + } + if (status === "stopped") { + // Reuse the existing low-level startContainer + const { startContainer } = await import("../docker.js"); + await startContainer(meta.containerName); + await waitForConnectorReady(meta.containerName); + } } async dumpBranch(branch: Branch): Promise { - const meta = branch.providerMeta as DockerProviderMeta; - return dockerDump(meta.containerName); + const meta = readMeta(branch); + return dumpDatabase(meta.containerName, meta.databaseName); } async restoreDump(branch: Branch, sql: string): Promise { - const meta = branch.providerMeta as DockerProviderMeta; - await dockerRestore(meta.containerName, sql); + const meta = readMeta(branch); + await restoreDumpToDatabase(meta.containerName, meta.databaseName, sql); } } + +/** SQL string literal escape (single-quote doubling). For LITERALS only. */ +function pgString(s: string): string { + return "'" + s.replace(/'/g, "''") + "'"; +} + +// Re-exported helpers for tests. +export const __test__ = { + buildContainerName, + buildSeedDatabaseName, + buildBranchDatabaseName, + buildConnectionString, + isV2Meta, + pgString, +}; +export type { ConnectorContainerInfo }; diff --git a/packages/core/src/branching/storage.ts b/packages/core/src/branching/storage.ts index 4f672a3..b6edc3c 100644 --- a/packages/core/src/branching/storage.ts +++ b/packages/core/src/branching/storage.ts @@ -152,6 +152,49 @@ export function deleteConnectorSnapshot(connectorName: string): void { } } +// --------------------------------------------------------------------------- +// Per-connector Docker container metadata (Lane B) +// --------------------------------------------------------------------------- + +export interface ConnectorContainerInfo { + containerId: string; + containerName: string; + port: number; + pgVersion: string; + seedDatabase: string; + createdAt: string; +} + +export function getConnectorContainerPath(connectorName: string): string { + return join(getSnapshotDir(connectorName), "container.json"); +} + +export function readConnectorContainer( + connectorName: string, +): ConnectorContainerInfo | null { + const p = getConnectorContainerPath(connectorName); + if (!existsSync(p)) return null; + try { + return JSON.parse(readFileSync(p, "utf-8")) as ConnectorContainerInfo; + } catch { + return null; + } +} + +export function writeConnectorContainer( + connectorName: string, + info: ConnectorContainerInfo, +): void { + const dir = getSnapshotDir(connectorName); + mkdirSync(dir, { recursive: true }); + writeFileSync(getConnectorContainerPath(connectorName), JSON.stringify(info, null, 2), "utf-8"); +} + +export function deleteConnectorContainer(connectorName: string): void { + const p = getConnectorContainerPath(connectorName); + if (existsSync(p)) rmSync(p, { force: true }); +} + // --------------------------------------------------------------------------- // Checkpoints — point-in-time snapshots within a branch // ---------------------------------------------------------------------------