// postgres's admin client — postgres's own code, living in the module (novox/hq ADR 0039). Both this // module's tools and its provisioner import it, and nothing outside postgres does. // // SQL is executed through `psql`, not a wire-protocol driver: the module may take NO npm dependency // beyond @novox/mesh-sdk, and hand-rolling startup + SCRAM auth + the query protocol is more surface // than this should carry — so it shells out to the client the postgres tools ship, the same way // minio drives itself through `mc` and mailu through doveadm. One boundary, `query()`, and every // method is built on it. import { randomBytes } from "node:crypto"; import { readFileSync } from "node:fs"; import { execFile } from "node:child_process"; import { promisify } from "node:util"; const run = promisify(execFile); export interface QueryResult { /** The command tag postgres returns, e.g. "SELECT", "CREATE DATABASE". */ readonly command: string; readonly rows: Record[]; } export interface PgConn { readonly host: string; readonly port: number; readonly user: string; readonly password: string; } export class PostgresClient { constructor(private readonly conn: PgConn) {} /** * Build from the module's resolved environment. Reads MESH_POSTGRES_* first (the documented * names), falling back to the MESH_PROVISION_* keys the manifest already sets on the provisioner * container. Throws if it cannot find a host and an admin password. */ static fromEnv(env: NodeJS.ProcessEnv = process.env): PostgresClient { const url = env.MESH_PROVISION_POSTGRES ? safeUrl(env.MESH_PROVISION_POSTGRES) : undefined; const host = env.MESH_POSTGRES_HOST ?? url?.hostname; const port = Number(env.MESH_POSTGRES_PORT ?? url?.port ?? "5432") || 5432; const user = env.MESH_POSTGRES_USER ?? url?.username ?? "postgres"; const password = env.MESH_POSTGRES_PASSWORD ?? readSecretFile(env.MESH_PROVISION_PASSWORD_FILE); if (!host || !password) { throw new Error("postgres host or admin password is not set — postgres's own code cannot reach the server"); } return new PostgresClient({ host, port, user, password }); } get host(): string { return this.conn.host; } get port(): number { return this.conn.port; } /** Execute SQL against a database as the admin and return its rows, through `psql` (see header). */ async query(sql: string, database = "postgres"): Promise { // Executed through `psql`, the way minio drives itself through `mc` and mailu through doveadm: // node has no postgres wire client without an npm dependency, and the module owns its own code // (ADR 0039), so it shells out to the client the postgres tools ship. CSV so the rows come back // structured; ON_ERROR_STOP so a failed statement is an error here, not a success with a warning. const { stdout } = await run( "psql", ["-h", this.conn.host, "-p", String(this.conn.port), "-U", this.conn.user, "-d", database, "-v", "ON_ERROR_STOP=1", "--no-psqlrc", "--csv", "-c", sql], { env: { ...process.env, PGPASSWORD: this.conn.password }, maxBuffer: 16 << 20 }, ); const rows = parseCsvRows(stdout); return { command: sql.trimStart().split(/\s+/)[0]?.toUpperCase() ?? "", rows }; } /** * Create a login role and a database it owns, idempotently. The DDL is the full, correct shape, run through query(). Extensions can be requested per * database and are created as the admin (a plain owner cannot install most of them). */ async createDatabaseAndRole(database: string, role: string, password: string): Promise { const roles = await this.query("SELECT 1 FROM pg_roles WHERE rolname = " + literal(role)); if (roles.rows.length === 0) { await this.query(`CREATE ROLE ${ident(role)} WITH LOGIN PASSWORD ${literal(password)}`); } else { await this.query(`ALTER ROLE ${ident(role)} WITH LOGIN PASSWORD ${literal(password)}`); } const dbs = await this.query("SELECT 1 FROM pg_database WHERE datname = " + literal(database)); if (dbs.rows.length === 0) { await this.query(`CREATE DATABASE ${ident(database)} OWNER ${ident(role)}`); } await this.query(`GRANT ALL PRIVILEGES ON DATABASE ${ident(database)} TO ${ident(role)}`); } /** Drop a database and its owning role, idempotently, after evicting live connections. */ async dropDatabaseAndRole(database: string, role: string): Promise { await this.query( "SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = " + literal(database) + " AND pid <> pg_backend_pid()", ); await this.query(`DROP DATABASE IF EXISTS ${ident(database)}`); await this.query(`DROP ROLE IF EXISTS ${ident(role)}`); } /** List the non-template databases, with size, for the postgres_list_databases tool. */ async listDatabases(): Promise<{ name: string; sizeBytes: number }[]> { const res = await this.query( "SELECT datname, pg_database_size(datname) AS size FROM pg_database WHERE datistemplate = false ORDER BY datname", ); return res.rows.map((r) => ({ name: String(r.datname), sizeBytes: Number(r.size) })); } /** Run a read-only SQL statement against a named database, for the postgres_query tool. */ async readOnlyQuery(database: string, sql: string): Promise { // The read-only guarantee is a wrapping transaction the server honours. return this.query(`BEGIN TRANSACTION READ ONLY; ${sql}; ROLLBACK;`, database); } } /** Generate a URL-safe password. */ export function generatePassword(): string { return randomBytes(24).toString("base64url"); } /** Quote a SQL identifier (double quotes, doubled internal quotes). */ export function ident(id: string): string { return '"' + id.replace(/"/g, '""') + '"'; } /** Quote a SQL string literal (single quotes, doubled internal quotes). */ export function literal(val: string): string { return "'" + val.replace(/'/g, "''") + "'"; } function readSecretFile(path: string | undefined): string | undefined { if (!path) return undefined; try { return readFileSync(path, "utf8").trim(); } catch { return undefined; } } function safeUrl(raw: string): URL | undefined { try { return new URL(raw); } catch { return undefined; } } /** Parse psql --csv output into row objects. RFC-4180: fields may be quoted, an embedded quote is * doubled, and a quoted field may span newlines. Empty output (a DDL statement) yields no rows. */ function parseCsvRows(csv: string): Record[] { const records = parseCsv(csv); if (records.length === 0) return []; const [header, ...rows] = records; return rows.map((cells) => { const row: Record = {}; header.forEach((name, i) => (row[name] = cells[i] ?? null)); return row; }); } function parseCsv(text: string): string[][] { const records: string[][] = []; let field = ""; let record: string[] = []; let inQuotes = false; let started = false; const endRecord = (): void => { if (started || field.length > 0 || record.length > 0) { record.push(field); records.push(record); } field = ""; record = []; started = false; }; for (let i = 0; i < text.length; i++) { const c = text[i]; if (inQuotes) { if (c === '"') { if (text[i + 1] === '"') { field += '"'; i++; } else inQuotes = false; } else field += c; } else if (c === '"') { inQuotes = true; started = true; } else if (c === ",") { record.push(field); field = ""; started = true; } else if (c === "\n" || c === "\r") { if (c === "\r" && text[i + 1] === "\n") i++; endRecord(); } else { field += c; started = true; } } endRecord(); return records; }