diff --git a/modules/letta/module.json b/modules/letta/module.json index 7ad5752..077276e 100644 --- a/modules/letta/module.json +++ b/modules/letta/module.json @@ -10,7 +10,10 @@ ], "contributes": { "postgres-database": { - "name": "letta" + "name": "letta", + "extensions": [ + "vector" + ] }, "route": { "label": "letta", diff --git a/modules/postgres/README.md b/modules/postgres/README.md new file mode 100644 index 0000000..23030e6 --- /dev/null +++ b/modules/postgres/README.md @@ -0,0 +1,71 @@ +# postgres + +The mesh store (novox/hq ADR 0079): one PostgreSQL server that provides the `postgres-database` +provision to every module on any machine that requires one, and holds the `mesh-store` seat. + +Each consumer is given a database it alone owns, under the login the mesh derived for it and the +password the mesh minted (ADR 0048). The provider never invents either: it makes exactly what the mesh +handed both ends, so they agree by construction. + +## Requiring a database + +A consumer requires `postgres-database` and contributes what it wants: + +``` +"requires": ["postgres-database"], +"contributes": {"postgres-database": {"name": "letta", "extensions": ["vector"]}}, +"binds": {"postgres-database": "${dir:state}/database.json"}, +"secrets": {"postgres-database": "${dir:state}/database.secret"} +``` + +| field | meaning | +|---|---| +| `name` | what the consumer calls its database. The database and its owning login are both named by the login the mesh derives (`${bound:postgres-database:as}`), so that is the name to connect to. | +| `extensions` | optional; a list of extension names to install in the consumer's database, e.g. `["vector"]` for pgvector. | + +**Extensions are the provider's to install.** Most extensions are not trusted (pgvector is +`superuser=t, trusted=f`), so the login that owns the database cannot create them. The provider does, +as the superuser, connected to that database, with `CREATE EXTENSION IF NOT EXISTS` — on every +provisioning pass, so a consumer that adds a name to its contribution, or a database that predates +this, gets it on the next pass, and a second pass changes nothing. Only a name the server lists in +`pg_available_extensions` is installed; any other is refused, by name, in the provider's log: the +consumer's database and login are still made, but none of its extensions is installed until the +contribution is corrected. An extension is never +dropped — not when it leaves the list, not when the consumer goes. + +The minute-by-minute check (issue 120) connects as the consumer and looks for each extension it asked +for, so one removed by hand is installed again. + +## What is never done + +- **No database is ever dropped.** A consumer the mesh no longer asks for is *withdrawn*: its login is + set `NOLOGIN` and its sessions are ended, and its database stays as it was under its own name + (issue 241). A consumer that comes back is given the same database. +- **`postgres_retire_database` renames, it does not drop**: the database becomes + `_deleted_` and its owner is locked. Removing the data is a person's act, by hand. +- **A caller's statement never runs as the superuser.** `postgres_query` and the seat's `query` verb + run as `mesh_store_reader` — `pg_read_all_data` and nothing else, every session read-only by the + server's own setting, a 60 s statement timeout — with the statement sent as given (issue 193). + Without the reader's password (`own-secrets.reader`) the call is refused. + +## Tools + +| tool | answers | +|---|---| +| `postgres_list_databases` | every non-template database with its size | +| `postgres_query` `{database, sql}` | one read-only statement, as the reader; rows keyed by column, `NULL` as null | +| `postgres_retire_database` `{database, confirm}` | renames a database aside and locks its owner; `confirm` repeats the name | +| `mesh-store.databases`, `mesh-store.query` | the store seat's verbs: the same listing and read-only query | + +## Where the code lives + +One Go bundle, `cmd/postgres-provider`, launched by the node's runtime and speaking MCP over stdio +through the Go SDK (ADR 0193). Beside the tools it runs the provisioner: every five seconds it reads the +contributions file the mesh writes (`MESH_RECEIVES`), applies each consumer whose login, password or +contribution changed, and withdraws each one no longer listed; a file it cannot read withdraws nobody. +It is the TypeScript SDK's `runProvisioner` loop, carried in the module until the Go SDK has one. + +`go test ./...` runs against a fake server. `MESH_POSTGRES_LIVE=postgres://postgres:…@host:port/postgres` +also runs `live_test.go` against a real, throwaway one (see the file for a `pgvector/pgvector` +container): the extension installed and a second pass a no-op, the reader unable to write, a +withdrawn login locked out with its data kept. diff --git a/modules/postgres/client.ts b/modules/postgres/client.ts deleted file mode 100644 index cf60ea6..0000000 --- a/modules/postgres/client.ts +++ /dev/null @@ -1,339 +0,0 @@ -// 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; - /** - * The read-only login's password, which the mesh mints for this module (`own-secrets.reader`). - * Absent when the mesh has not delivered it: then a caller's statement is refused, never run as - * the admin (novox/hq issue 193). - */ - readonly readerPassword?: string; -} - -/** - * The login a caller's statement runs as (novox/hq issue 193). It may read every table and change - * nothing: `pg_read_all_data` and no other grant, and every transaction it opens is read-only by - * the server's own setting. A statement cannot climb out of a login the way it can out of a - * transaction wrapped around it as text: `COMMIT; DROP …` ended the old wrapper and ran the rest as - * the superuser, and even one read-only statement as a superuser can run a program on the server. - */ -export const READER = "mesh_store_reader"; - - -export class PostgresClient { - constructor(private readonly conn: PgConn) {} - - /** The reader is made once per process: idempotent, and repeating it re-sets a rotated password. */ - private readerReady?: Promise; - - /** - * 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; - // MESH_PROVISION_POSTGRES_PORT is the seat's twin (mesh-controller's own - // internal/envfile.Placed pattern): which port this machine actually put mesh-store at, when - // that differs from what the connection string above already says — e.g. adopted in place at - // a predecessor's port. Empty means the mesh has nothing to add and the string's own port - // stands; a placeholder the mesh never filled (a manifest ahead of the running controller) - // is treated the same way, not as a fault. - const seatPort = (env.MESH_PROVISION_POSTGRES_PORT ?? "").trim(); - const filledSeatPort = seatPort && !/^\$\{[^}]*\}$/.test(seatPort) ? seatPort : undefined; - const portSource = env.MESH_POSTGRES_PORT ?? filledSeatPort ?? url?.port ?? "5432"; - const port = Number(portSource) || 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"); - } - const readerPassword = env.MESH_POSTGRES_READER_PASSWORD ?? - readSecretFile(env.MESH_POSTGRES_READER_PASSWORD_FILE); - return new PostgresClient({ host, port, user, password, readerPassword }); - } - - 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)} VALID UNTIL 'infinity'`); - } else { - // VALID UNTIL 'infinity': a password that expired is refused like a wrong one, so the check the - // provisioner runs would report it lost, and only clearing the expiry makes applying it again work. - await this.query(`ALTER ROLE ${ident(role)} WITH LOGIN PASSWORD ${literal(password)} VALID UNTIL 'infinity'`); - } - 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)}`); - } - - /** - * Whether `role` can log in to `database` with exactly `password`: the consumer's own view of its - * credential, checked by connecting as it. Read-only. `false` only when the server says so (the - * role, the password or the database is wrong or gone); an unreachable server rejects instead, - * because being unable to ask is not evidence of loss (novox/hq issue 120). - */ - async canConnectAs(database: string, role: string, password: string): Promise { - try { - await run( - "psql", - ["-h", this.conn.host, "-p", String(this.conn.port), "-U", role, "-d", database, - "-v", "ON_ERROR_STOP=1", "--no-psqlrc", "-tAc", "SELECT 1"], - { env: { ...process.env, PGPASSWORD: password, PGCONNECT_TIMEOUT: "10" }, timeout: 20_000 }, - ); - return true; - } catch (err) { - const text = `${(err as { stderr?: string }).stderr ?? ""}`; - if (/password authentication failed|role ".*" does not exist|database ".*" does not exist|not permitted to log in|permission denied for database/i.test(text)) { - return false; - } - throw err; - } - } - - /** - * Withdraw a consumer without destroying anything (novox/hq issue 241): its login can no longer log - * in and its open connections are ended, and its database stays exactly as it was, under its own - * name. A consumer that comes back is given the same database — create sets LOGIN again — which is - * what a provider must do when "no longer asked for" turns out to be a moment's misreading. - */ - async lockRole(role: string): Promise { - const roles = await this.query("SELECT 1 FROM pg_roles WHERE rolname = " + literal(role)); - if (roles.rows.length === 0) return; - await this.query(`ALTER ROLE ${ident(role)} NOLOGIN`); - await this.query( - "SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE usename = " + - literal(role) + " AND pid <> pg_backend_pid()", - ); - } - - /** - * Take a database out of service on purpose: rename it aside to `_deleted_` and lock - * its owner. Never a drop — the data stays on the server under the new name until a person - * removes it by hand. Returns the name it now has. - */ - async retireDatabase(database: string, now: Date = new Date()): Promise { - const found = await this.query( - "SELECT pg_get_userbyid(datdba) AS owner FROM pg_database WHERE datname = " + literal(database), - ); - if (found.rows.length === 0) throw new Error(`no database named ${database}`); - const owner = String((found.rows[0] as Record).owner ?? ""); - const aside = retiredName(database, now); - const taken = await this.query("SELECT 1 FROM pg_database WHERE datname = " + literal(aside)); - if (taken.rows.length > 0) throw new Error(`${aside} already exists; retire it by hand first`); - if (owner && owner !== "postgres") await this.query(`ALTER ROLE ${ident(owner)} NOLOGIN`); - await this.query( - "SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = " + - literal(database) + " AND pid <> pg_backend_pid()", - ); - await this.query(`ALTER DATABASE ${ident(database)} RENAME TO ${ident(aside)}`); - return aside; - } - - /** 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) })); - } - - /** - * Make the read-only login, idempotently, with the password the mesh minted for it. Run as the - * admin, because only the admin can make a role. - */ - async ensureReader(): Promise { - const password = this.conn.readerPassword; - if (!password) throw readerMissing(); - const roles = await this.query("SELECT 1 FROM pg_roles WHERE rolname = " + literal(READER)); - const verb = roles.rows.length === 0 ? "CREATE" : "ALTER"; - // Every attribute stated, so an existing role someone widened is narrowed again on every start. - await this.query( - `${verb} ROLE ${ident(READER)} WITH LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOREPLICATION ` + - `NOBYPASSRLS INHERIT PASSWORD ${literal(password)} VALID UNTIL 'infinity'`, - ); - await this.query(`GRANT pg_read_all_data TO ${ident(READER)}`); - await this.query(`ALTER ROLE ${ident(READER)} SET default_transaction_read_only = on`); - await this.query(`ALTER ROLE ${ident(READER)} SET statement_timeout = '60s'`); - } - - /** - * Run a caller's statement against a named database as the read-only login, for the - * postgres_query tool and the store seat's `query` verb (novox/hq ADR 0159, issue 193). - * - * **Read-only by the login, not by text around the statement.** The statement is sent as it was - * given, as the reader, whose role can write nothing and whose transactions the server makes - * read-only. Never as the admin: without the reader's password the call is refused. - */ - async readOnlyQuery(database: string, sql: string): Promise { - this.readerReady ??= this.ensureReader().catch((err) => { - this.readerReady = undefined; // asked again next call, not failed for the process's life - throw err; - }); - await this.readerReady; - const { stdout } = await run( - "psql", - // -q: no command tags, so the output is the header and the rows and nothing else — the tags - // were what came back as rows keyed by BEGIN. - ["-h", this.conn.host, "-p", String(this.conn.port), "-U", READER, "-d", database, - "-v", "ON_ERROR_STOP=1", "--no-psqlrc", "-q", "--csv", "-c", sql], - { - env: { - ...process.env, - PGPASSWORD: this.conn.readerPassword, - // Read-only from the first statement, before the role's own setting is read. - PGOPTIONS: "-c default_transaction_read_only=on -c statement_timeout=60s", - }, - maxBuffer: 16 << 20, - }, - ); - const command = /^\s*([A-Za-z]+)/.exec(sql)?.[1]?.toUpperCase() ?? ""; - return { command, rows: parseCsvRows(stdout) }; - } -} - -function readerMissing(): Error { - return new Error( - "the read-only login's password was not delivered (own-secrets.reader, " + - "MESH_POSTGRES_READER_PASSWORD_FILE), so the statement is refused rather than run as the " + - "admin (novox/hq issue 193)", - ); -} - -/** 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; -} - -/** The name a retired database is renamed to: `_deleted_`, within postgres's 63 bytes. */ -export function retiredName(database: string, now: Date = new Date()): string { - const stamp = now.toISOString().slice(0, 10).replace(/-/g, ""); - const suffix = `_deleted_${stamp}`; - return database.slice(0, 63 - suffix.length) + suffix; -} diff --git a/modules/postgres/cmd/postgres-provider/client.go b/modules/postgres/cmd/postgres-provider/client.go new file mode 100644 index 0000000..c695423 --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/client.go @@ -0,0 +1,615 @@ +package main + +// postgres's admin client — postgres's own code, living in the module (novox/hq ADR 0039). Both this +// module's tools and its provisioner use it, and nothing outside postgres does. +// +// SQL goes over the wire protocol (pgconn), in the simple query protocol: a statement is sent as it +// was given, the way `psql -c` sent it when this module was TypeScript and could take no driver. One +// boundary, Dialer, and every method is built on it — which is also what the tests replace. + +import ( + "context" + "errors" + "fmt" + "net" + "net/url" + "os" + "regexp" + "strconv" + "strings" + "sync" + "time" + "unicode/utf8" + + "github.com/jackc/pgx/v5/pgconn" +) + +// Reader is the login a caller's statement runs as (novox/hq issue 193). It may read every table and +// change nothing: `pg_read_all_data` and no other grant, and every transaction it opens is read-only +// by the server's own setting. A statement cannot climb out of a login the way it can out of a +// transaction wrapped around it as text: `COMMIT; DROP …` ended the old wrapper and ran the rest as +// the superuser, and even one read-only statement as a superuser can run a program on the server. +const Reader = "mesh_store_reader" + +// readerOptions make the reader's session read-only from its first statement, before the role's own +// setting is read (PGOPTIONS, when this module shelled out to psql). +var readerOptions = map[string]string{ + "default_transaction_read_only": "on", + "statement_timeout": "60s", +} + +// maxResultBytes bounds what one read-only query may hand back, as psql's 16 MiB buffer did. +const maxResultBytes = 16 << 20 + +// Result is one statement's answer: the command tag's verb and the rows, each a column's text or nil. +type Result struct { + Command string + Fields []string + Rows [][]*string +} + +// Maps is the rows keyed by their columns, as the tools return them. +func (r Result) Maps() []map[string]any { + out := make([]map[string]any, 0, len(r.Rows)) + for _, row := range r.Rows { + m := map[string]any{} + for i, name := range r.Fields { + if i < len(row) && row[i] != nil { + m[name] = *row[i] + } else { + m[name] = nil + } + } + out = append(out, m) + } + return out +} + +// Session is one open connection: Run sends SQL in the simple protocol and answers every statement's +// result, in order. +type Session interface { + Run(ctx context.Context, sql string) ([]Result, error) + Close() +} + +// Login is who a session connects as, to which database, with which session settings. +type Login struct { + Database string + User string + Password string + Options map[string]string +} + +// Dialer opens a session. The real one is pgconn; a test's records what it was asked. +type Dialer func(ctx context.Context, l Login) (Session, error) + +// Conn is where the server is and who administers it. +type Conn struct { + Host string + Port int + User string + Password string + SSLMode string + // ReaderPassword is the read-only login's password, which the mesh mints for this module + // (`own-secrets.reader`). Empty when the mesh has not delivered it: then a caller's statement is + // refused, never run as the admin (novox/hq issue 193). + ReaderPassword string +} + +// Client is postgres's admin client. +type Client struct { + conn Conn + dial Dialer + + readerMu sync.Mutex + readerReady bool +} + +// NewClient is a client over a dialer; nil dials the server for real. +func NewClient(conn Conn, dial Dialer) *Client { + if dial == nil { + dial = pgDialer(conn) + } + return &Client{conn: conn, dial: dial} +} + +var unfilled = regexp.MustCompile(`^\$\{[^}]*\}$`) + +// ClientFromEnv builds the client from the module's words. MESH_POSTGRES_* first (the documented +// names), then the MESH_PROVISION_* keys the manifest sets. Fails without a host and an admin password. +func ClientFromEnv(env func(string) string) (*Client, error) { + var u *url.URL + if raw := env("MESH_PROVISION_POSTGRES"); raw != "" { + if parsed, err := url.Parse(raw); err == nil { + u = parsed + } + } + host := env("MESH_POSTGRES_HOST") + if host == "" && u != nil { + host = u.Hostname() + } + // MESH_PROVISION_POSTGRES_PORT is the seat's twin (mesh-controller's internal/envfile.Placed + // pattern): which port this machine actually put mesh-store at, when that differs from the + // connection string's. Empty, or a placeholder the mesh never filled, adds nothing. + seatPort := strings.TrimSpace(env("MESH_PROVISION_POSTGRES_PORT")) + if unfilled.MatchString(seatPort) { + seatPort = "" + } + portSource := firstOf(env("MESH_POSTGRES_PORT"), seatPort) + if portSource == "" && u != nil { + portSource = u.Port() + } + port, err := strconv.Atoi(portSource) + if err != nil || port == 0 { + port = 5432 + } + user := env("MESH_POSTGRES_USER") + if user == "" && u != nil && u.User != nil { + user = u.User.Username() + } + if user == "" { + user = "postgres" + } + password := firstOf(env("MESH_POSTGRES_PASSWORD"), readSecretFile(env("MESH_PROVISION_PASSWORD_FILE"))) + if host == "" || password == "" { + return nil, errors.New("postgres host or admin password is not set — postgres's own code cannot reach the server") + } + sslmode := "prefer" + if u != nil && u.Query().Get("sslmode") != "" { + sslmode = u.Query().Get("sslmode") + } + reader := firstOf(env("MESH_POSTGRES_READER_PASSWORD"), readSecretFile(env("MESH_POSTGRES_READER_PASSWORD_FILE"))) + return NewClient(Conn{Host: host, Port: port, User: user, Password: password, SSLMode: sslmode, ReaderPassword: reader}, nil), nil +} + +func firstOf(values ...string) string { + for _, v := range values { + if v != "" { + return v + } + } + return "" +} + +func readSecretFile(path string) string { + if path == "" { + return "" + } + b, err := os.ReadFile(path) + if err != nil { + return "" + } + return strings.TrimSpace(string(b)) +} + +// ---- the boundary ------------------------------------------------------------------------------- + +// as runs SQL as one login against one database and answers the last statement's result. +func (c *Client) as(ctx context.Context, l Login, sql string) (Result, error) { + s, err := c.dial(ctx, l) + if err != nil { + return Result{}, err + } + defer s.Close() + results, err := s.Run(ctx, sql) + if err != nil { + return Result{}, err + } + if len(results) == 0 { + return Result{}, nil + } + // The last statement that answered rows, or the last statement: what psql -c printed last. + for i := len(results) - 1; i >= 0; i-- { + if len(results[i].Fields) > 0 { + return results[i], nil + } + } + return results[len(results)-1], nil +} + +// Query runs SQL as the admin against a database ("postgres" when empty). +func (c *Client) Query(ctx context.Context, database, sql string) (Result, error) { + if database == "" { + database = "postgres" + } + return c.as(ctx, Login{Database: database, User: c.conn.User, Password: c.conn.Password}, sql) +} + +func (c *Client) exists(ctx context.Context, sql string) (bool, error) { + r, err := c.Query(ctx, "", sql) + return len(r.Rows) > 0, err +} + +// ---- what the provisioner does ----------------------------------------------------------------- + +// RoleStatement is the DDL that makes or re-sets a consumer's login. VALID UNTIL 'infinity': an +// expired password is refused like a wrong one, so the check the provisioner runs would report it +// lost, and only clearing the expiry makes applying it again work. +func RoleStatement(exists bool, role, password string) string { + verb := "CREATE" + if exists { + verb = "ALTER" + } + return fmt.Sprintf("%s ROLE %s WITH LOGIN PASSWORD %s VALID UNTIL 'infinity'", verb, Ident(role), Literal(password)) +} + +// CreateDatabaseAndRole makes a login role and a database it owns, idempotently. A database that +// exists is never recreated, and nothing here drops anything. +func (c *Client) CreateDatabaseAndRole(ctx context.Context, database, role, password string) error { + has, err := c.exists(ctx, "SELECT 1 FROM pg_roles WHERE rolname = "+Literal(role)) + if err != nil { + return err + } + if _, err := c.Query(ctx, "", RoleStatement(has, role, password)); err != nil { + return err + } + has, err = c.exists(ctx, "SELECT 1 FROM pg_database WHERE datname = "+Literal(database)) + if err != nil { + return err + } + if !has { + if _, err := c.Query(ctx, "", fmt.Sprintf("CREATE DATABASE %s OWNER %s", Ident(database), Ident(role))); err != nil { + return err + } + } + _, err = c.Query(ctx, "", fmt.Sprintf("GRANT ALL PRIVILEGES ON DATABASE %s TO %s", Ident(database), Ident(role))) + return err +} + +// Extensions reads a contribution's `extensions`: absent is none; otherwise a list of names. +func Extensions(values map[string]any) ([]string, error) { + raw, ok := values["extensions"] + if !ok || raw == nil { + return nil, nil + } + list, ok := raw.([]any) + if !ok { + return nil, fmt.Errorf("extensions must be a list of extension names, not %T", raw) + } + seen := map[string]bool{} + var out []string + for _, v := range list { + name, ok := v.(string) + name = strings.TrimSpace(name) + if !ok || name == "" { + return nil, fmt.Errorf("extensions must be a list of extension names; %v is not one", v) + } + if !seen[name] { + seen[name] = true + out = append(out, name) + } + } + return out, nil +} + +// ExtensionStatement is the DDL that installs one extension in the database it is run in. IF NOT +// EXISTS: run on every pass, and a second run changes nothing. There is no statement that removes one. +func ExtensionStatement(name string) string { + return "CREATE EXTENSION IF NOT EXISTS " + Ident(name) +} + +// Unavailable is the names asked for that the server does not offer. +func Unavailable(want []string, available map[string]bool) []string { + var missing []string + for _, name := range want { + if !available[name] { + missing = append(missing, name) + } + } + return missing +} + +// EnsureExtensions installs each named extension in the database, as the admin, connected to that +// database — most extensions (pgvector among them) are not trusted, so the consumer that owns the +// database cannot install them itself. Only names the server lists in pg_available_extensions are +// asked for; any other is refused, by name, before anything runs. An extension is never dropped. +func (c *Client) EnsureExtensions(ctx context.Context, database string, want []string) error { + if len(want) == 0 { + return nil + } + r, err := c.Query(ctx, "", "SELECT name FROM pg_available_extensions") + if err != nil { + return err + } + available := map[string]bool{} + for _, row := range r.Rows { + if len(row) > 0 && row[0] != nil { + available[*row[0]] = true + } + } + if missing := Unavailable(want, available); len(missing) > 0 { + return fmt.Errorf("database %s asks for extension(s) %s, which this server does not offer "+ + "(not in pg_available_extensions); refused, nothing installed", database, strings.Join(quoted(missing), ", ")) + } + for _, name := range want { + if _, err := c.Query(ctx, database, ExtensionStatement(name)); err != nil { + return fmt.Errorf("extension %q in %s: %w", name, database, err) + } + } + return nil +} + +func quoted(names []string) []string { + out := make([]string, len(names)) + for i, n := range names { + out[i] = strconv.Quote(n) + } + return out +} + +// lostCodes are the server's ways of saying a login, its password or its database is wrong or gone: +// invalid_password, invalid_authorization_specification (no such role, not permitted to log in), +// invalid_catalog_name (no such database), insufficient_privilege (no CONNECT). +var lostCodes = map[string]bool{"28P01": true, "28000": true, "3D000": true, "42501": true} + +// IsLost says whether an error is the server saying the credential is wrong or gone. +func IsLost(err error) bool { + var pg *pgconn.PgError + return errors.As(err, &pg) && lostCodes[pg.Code] +} + +// Holds says whether `role` can log in to `database` with exactly `password` and finds every wanted +// extension installed there: the consumer's own view, checked by connecting as it. Read-only. false +// only when the server says so; an unreachable server is an error, because being unable to ask is +// not evidence of loss (novox/hq issue 120). +func (c *Client) Holds(ctx context.Context, database, role, password string, extensions []string) (bool, error) { + ctx, cancel := context.WithTimeout(ctx, 20*time.Second) + defer cancel() + r, err := c.as(ctx, Login{Database: database, User: role, Password: password}, "SELECT extname FROM pg_extension") + if err != nil { + if IsLost(err) { + return false, nil + } + return false, err + } + installed := map[string]bool{} + for _, row := range r.Rows { + if len(row) > 0 && row[0] != nil { + installed[*row[0]] = true + } + } + return len(Unavailable(extensions, installed)) == 0, nil +} + +// LockRole withdraws a consumer without destroying anything (novox/hq issue 241): its login can no +// longer log in and its open connections are ended, and its database stays exactly as it was, under +// its own name. A consumer that comes back is given the same database — create sets LOGIN again. +func (c *Client) LockRole(ctx context.Context, role string) error { + has, err := c.exists(ctx, "SELECT 1 FROM pg_roles WHERE rolname = "+Literal(role)) + if err != nil || !has { + return err + } + if _, err := c.Query(ctx, "", fmt.Sprintf("ALTER ROLE %s NOLOGIN", Ident(role))); err != nil { + return err + } + _, err = c.Query(ctx, "", "SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE usename = "+ + Literal(role)+" AND pid <> pg_backend_pid()") + return err +} + +// ---- what the tools do -------------------------------------------------------------------------- + +// RetireDatabase takes a database out of service on purpose: renamed aside to +// `_deleted_` and its owner locked. Never a drop — the data stays on the server under the +// new name until a person removes it by hand. Answers the name it now has. +func (c *Client) RetireDatabase(ctx context.Context, database string, now time.Time) (string, error) { + found, err := c.Query(ctx, "", "SELECT pg_get_userbyid(datdba) AS owner FROM pg_database WHERE datname = "+Literal(database)) + if err != nil { + return "", err + } + if len(found.Rows) == 0 { + return "", fmt.Errorf("no database named %s", database) + } + owner := "" + if row := found.Rows[0]; len(row) > 0 && row[0] != nil { + owner = *row[0] + } + aside := RetiredName(database, now) + taken, err := c.exists(ctx, "SELECT 1 FROM pg_database WHERE datname = "+Literal(aside)) + if err != nil { + return "", err + } + if taken { + return "", fmt.Errorf("%s already exists; retire it by hand first", aside) + } + if owner != "" && owner != "postgres" { + if _, err := c.Query(ctx, "", fmt.Sprintf("ALTER ROLE %s NOLOGIN", Ident(owner))); err != nil { + return "", err + } + } + if _, err := c.Query(ctx, "", "SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = "+ + Literal(database)+" AND pid <> pg_backend_pid()"); err != nil { + return "", err + } + if _, err := c.Query(ctx, "", fmt.Sprintf("ALTER DATABASE %s RENAME TO %s", Ident(database), Ident(aside))); err != nil { + return "", err + } + return aside, nil +} + +// Database is one row of the listing. +type Database struct { + Name string `json:"name"` + SizeBytes int64 `json:"sizeBytes"` +} + +// ListDatabases is the non-template databases with their size. +func (c *Client) ListDatabases(ctx context.Context) ([]Database, error) { + r, err := c.Query(ctx, "", "SELECT datname, pg_database_size(datname) AS size FROM pg_database WHERE datistemplate = false ORDER BY datname") + if err != nil { + return nil, err + } + out := []Database{} + for _, row := range r.Rows { + if len(row) < 2 || row[0] == nil { + continue + } + d := Database{Name: *row[0]} + if row[1] != nil { + d.SizeBytes, _ = strconv.ParseInt(*row[1], 10, 64) + } + out = append(out, d) + } + return out, nil +} + +// ReaderStatements are the DDL that make the read-only login, every attribute stated so an existing +// role someone widened is narrowed again. +func ReaderStatements(exists bool, password string) []string { + verb := "CREATE" + if exists { + verb = "ALTER" + } + return []string{ + fmt.Sprintf("%s ROLE %s WITH LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOREPLICATION "+ + "NOBYPASSRLS INHERIT PASSWORD %s VALID UNTIL 'infinity'", verb, Ident(Reader), Literal(password)), + "GRANT pg_read_all_data TO " + Ident(Reader), + "ALTER ROLE " + Ident(Reader) + " SET default_transaction_read_only = on", + "ALTER ROLE " + Ident(Reader) + " SET statement_timeout = '60s'", + } +} + +var errReaderMissing = errors.New("the read-only login's password was not delivered (own-secrets.reader, " + + "MESH_POSTGRES_READER_PASSWORD_FILE), so the statement is refused rather than run as the admin (novox/hq issue 193)") + +// EnsureReader makes the read-only login, idempotently, with the password the mesh minted — as the +// admin, because only the admin can make a role. Once per process; a failure is asked again next call. +func (c *Client) EnsureReader(ctx context.Context) error { + if c.conn.ReaderPassword == "" { + return errReaderMissing + } + c.readerMu.Lock() + defer c.readerMu.Unlock() + if c.readerReady { + return nil + } + has, err := c.exists(ctx, "SELECT 1 FROM pg_roles WHERE rolname = "+Literal(Reader)) + if err != nil { + return err + } + for _, sql := range ReaderStatements(has, c.conn.ReaderPassword) { + if _, err := c.Query(ctx, "", sql); err != nil { + return err + } + } + c.readerReady = true + return nil +} + +var firstWord = regexp.MustCompile(`^\s*([A-Za-z]+)`) + +// ReadOnlyQuery runs a caller's statement against a named database as the read-only login, for the +// postgres_query tool and the store seat's `query` verb (novox/hq ADR 0159, issue 193). +// +// **Read-only by the login, not by text around the statement.** The statement is sent as it was +// given, as the reader, whose role can write nothing and whose session the server makes read-only. +// Never as the admin: without the reader's password the call is refused. +func (c *Client) ReadOnlyQuery(ctx context.Context, database, sql string) (Result, error) { + if err := c.EnsureReader(ctx); err != nil { + return Result{}, err + } + r, err := c.as(ctx, Login{Database: database, User: Reader, Password: c.conn.ReaderPassword, Options: readerOptions}, sql) + if err != nil { + return Result{}, err + } + size := 0 + for _, row := range r.Rows { + for _, v := range row { + if v != nil { + size += len(*v) + } + } + } + if size > maxResultBytes { + return Result{}, fmt.Errorf("the result is larger than %d MiB; narrow the query", maxResultBytes>>20) + } + r.Command = "" + if m := firstWord.FindStringSubmatch(sql); m != nil { + r.Command = strings.ToUpper(m[1]) + } + return r, nil +} + +// ---- quoting and names ------------------------------------------------------------------------- + +// Ident quotes a SQL identifier: double quotes, internal ones doubled. +func Ident(id string) string { return `"` + strings.ReplaceAll(id, `"`, `""`) + `"` } + +// Literal quotes a SQL string literal: single quotes, internal ones doubled (standard_conforming_strings). +func Literal(v string) string { return "'" + strings.ReplaceAll(v, "'", "''") + "'" } + +// RetiredName is `_deleted_`, within postgres's 63 bytes. +func RetiredName(database string, now time.Time) string { + suffix := "_deleted_" + now.UTC().Format("20060102") + keep := 63 - len(suffix) + if len(database) > keep { + database = database[:keep] + for !utf8.ValidString(database) { + database = database[:len(database)-1] + } + } + return database + suffix +} + +// ---- the real dialer --------------------------------------------------------------------------- + +type pgSession struct{ c *pgconn.PgConn } + +func (s pgSession) Run(ctx context.Context, sql string) ([]Result, error) { + results, err := s.c.Exec(ctx, sql).ReadAll() + if err != nil { + return nil, err + } + out := make([]Result, 0, len(results)) + for _, r := range results { + if r.Err != nil { + return nil, r.Err + } + res := Result{Command: strings.SplitN(r.CommandTag.String(), " ", 2)[0]} + for _, f := range r.FieldDescriptions { + res.Fields = append(res.Fields, f.Name) + } + for _, row := range r.Rows { + cells := make([]*string, len(row)) + for i, v := range row { + if v != nil { + s := string(v) + cells[i] = &s + } + } + res.Rows = append(res.Rows, cells) + } + out = append(out, res) + } + return out, nil +} + +func (s pgSession) Close() { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + _ = s.c.Close(ctx) +} + +func pgDialer(conn Conn) Dialer { + return func(ctx context.Context, l Login) (Session, error) { + u := url.URL{ + Scheme: "postgres", + User: url.UserPassword(l.User, l.Password), + Host: net.JoinHostPort(conn.Host, strconv.Itoa(conn.Port)), + Path: "/" + l.Database, + RawQuery: url.Values{"sslmode": {firstOf(conn.SSLMode, "prefer")}, "connect_timeout": {"10"}}.Encode(), + } + cfg, err := pgconn.ParseConfig(u.String()) + if err != nil { + return nil, err + } + for k, v := range l.Options { + cfg.RuntimeParams[k] = v + } + c, err := pgconn.ConnectConfig(ctx, cfg) + if err != nil { + return nil, err + } + return pgSession{c: c}, nil + } +} diff --git a/modules/postgres/cmd/postgres-provider/client_test.go b/modules/postgres/cmd/postgres-provider/client_test.go new file mode 100644 index 0000000..c992510 --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/client_test.go @@ -0,0 +1,325 @@ +package main + +import ( + "context" + "errors" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgconn" +) + +var ctx = context.Background() + +func TestQuoting(t *testing.T) { + if got := Ident(`we"ird`); got != `"we""ird"` { + t.Fatal(got) + } + if got := Literal(`it's`); got != `'it''s'` { + t.Fatal(got) + } + if got := RoleStatement(false, "mesh_ace_letta", "p'w"); got != `CREATE ROLE "mesh_ace_letta" WITH LOGIN PASSWORD 'p''w' VALID UNTIL 'infinity'` { + t.Fatal(got) + } + if got := RoleStatement(true, "r", "p"); !strings.HasPrefix(got, `ALTER ROLE "r" WITH LOGIN`) { + t.Fatal(got) + } +} + +func TestExtensionStatementIsIdempotentAndQuoted(t *testing.T) { + if got := ExtensionStatement("vector"); got != `CREATE EXTENSION IF NOT EXISTS "vector"` { + t.Fatal(got) + } + if got := ExtensionStatement("uuid-ossp"); got != `CREATE EXTENSION IF NOT EXISTS "uuid-ossp"` { + t.Fatal(got) + } + // A name cannot leave its quotes. + if got := ExtensionStatement(`x"; DROP DATABASE y; --`); got != `CREATE EXTENSION IF NOT EXISTS "x""; DROP DATABASE y; --"` { + t.Fatal(got) + } +} + +func TestExtensionsFromAContribution(t *testing.T) { + got, err := Extensions(map[string]any{"name": "letta"}) + if err != nil || got != nil { + t.Fatal(got, err) + } + got, err = Extensions(map[string]any{"extensions": []any{"vector", " vector ", "pg_trgm"}}) + if err != nil || strings.Join(got, ",") != "vector,pg_trgm" { + t.Fatal(got, err) + } + for _, bad := range []any{"vector", []any{"vector", 3}, []any{""}, map[string]any{}} { + if _, err := Extensions(map[string]any{"extensions": bad}); err == nil { + t.Fatalf("%v accepted", bad) + } + } +} + +func TestAnExtensionTheServerDoesNotOfferIsRefusedByName(t *testing.T) { + f, c := newFake(admin) + f.answer(`FROM pg_available_extensions`, []string{"name"}, []string{"vector"}, []string{"pg_trgm"}) + err := c.EnsureExtensions(ctx, "mesh_ace_letta", []string{"vector", "made_up", "plv8"}) + if err == nil || !strings.Contains(err.Error(), `"made_up"`) || !strings.Contains(err.Error(), `"plv8"`) { + t.Fatalf("refusal does not name the extensions: %v", err) + } + if has(f.statements(), `CREATE EXTENSION`) { + t.Fatalf("installed something although one was refused: %v", f.statements()) + } +} + +func TestExtensionsAreInstalledAsTheAdminInTheConsumersDatabase(t *testing.T) { + f, c := newFake(admin) + f.answer(`FROM pg_available_extensions`, []string{"name"}, []string{"vector"}) + if err := c.EnsureExtensions(ctx, "mesh_ace_letta", []string{"vector"}); err != nil { + t.Fatal(err) + } + last := f.calls[len(f.calls)-1] + if last.SQL != `CREATE EXTENSION IF NOT EXISTS "vector"` || last.Database != "mesh_ace_letta" || last.User != "postgres" || last.Password != "admin-secret" { + t.Fatalf("%+v", last) + } + noDrop(t, f.statements()) + + f.reset() + if err := c.EnsureExtensions(ctx, "mesh_ace_letta", nil); err != nil || len(f.calls) != 0 { + t.Fatalf("asked the server for nothing to do: %v %v", f.calls, err) + } +} + +func TestCreateIsIdempotent(t *testing.T) { + f, c := newFake(admin) + if err := c.CreateDatabaseAndRole(ctx, "mesh_ace_letta", "mesh_ace_letta", "pw"); err != nil { + t.Fatal(err) + } + sql := f.statements() + if !has(sql, `^CREATE ROLE "mesh_ace_letta"`) || !has(sql, `^CREATE DATABASE "mesh_ace_letta" OWNER "mesh_ace_letta"$`) { + t.Fatal(strings.Join(sql, "\n")) + } + + // Both exist: the password is set again, the database is not made again. + f, c = newFake(admin) + f.answer(`FROM pg_roles`, []string{"?column?"}, []string{"1"}) + f.answer(`FROM pg_database`, []string{"?column?"}, []string{"1"}) + if err := c.CreateDatabaseAndRole(ctx, "mesh_ace_letta", "mesh_ace_letta", "pw"); err != nil { + t.Fatal(err) + } + sql = f.statements() + if !has(sql, `^ALTER ROLE "mesh_ace_letta" WITH LOGIN PASSWORD 'pw' VALID UNTIL 'infinity'`) || has(sql, `CREATE DATABASE`) { + t.Fatal(strings.Join(sql, "\n")) + } + if !has(sql, `^GRANT ALL PRIVILEGES ON DATABASE "mesh_ace_letta" TO "mesh_ace_letta"$`) { + t.Fatal(strings.Join(sql, "\n")) + } + noDrop(t, sql) +} + +func TestWithdrawingLocksTheLoginAndKeepsTheDatabase(t *testing.T) { + f, c := newFake(admin) + f.answer(`FROM pg_roles`, []string{"?column?"}, []string{"1"}) + if err := c.LockRole(ctx, "mesh_anchor_mail"); err != nil { + t.Fatal(err) + } + sql := f.statements() + if !has(sql, `ALTER ROLE "mesh_anchor_mail" NOLOGIN`) || !has(sql, `pg_terminate_backend.*usename = 'mesh_anchor_mail'`) { + t.Fatal(strings.Join(sql, "\n")) + } + noDrop(t, sql) + + // No such role: nothing to lock, nothing done. + f, c = newFake(admin) + if err := c.LockRole(ctx, "gone"); err != nil || has(f.statements(), `ALTER`) { + t.Fatal(f.statements(), err) + } +} + +func TestRetiringRenamesAsideLocksTheOwnerAndDropsNothing(t *testing.T) { + f, c := newFake(admin) + f.answer(`pg_get_userbyid`, []string{"owner"}, []string{"mesh_anchor_mail"}) + aside, err := c.RetireDatabase(ctx, "mesh_anchor_mail", time.Date(2026, 10, 5, 9, 0, 0, 0, time.UTC)) + if err != nil || aside != "mesh_anchor_mail_deleted_20261005" { + t.Fatal(aside, err) + } + sql := f.statements() + if !has(sql, `ALTER DATABASE "mesh_anchor_mail" RENAME TO "mesh_anchor_mail_deleted_20261005"`) || !has(sql, `ALTER ROLE "mesh_anchor_mail" NOLOGIN`) { + t.Fatal(strings.Join(sql, "\n")) + } + noDrop(t, sql) + + f, c = newFake(admin) + f.answer(`pg_get_userbyid`, []string{"owner"}, []string{"x"}) + f.answer(`_deleted_`, []string{"?column?"}, []string{"1"}) + if _, err := c.RetireDatabase(ctx, "x", time.Now()); err == nil || has(f.statements(), `RENAME`) { + t.Fatal("renamed over a database already set aside") + } +} + +func TestARetiredNameFitsPostgres(t *testing.T) { + name := RetiredName(strings.Repeat("x", 70), time.Date(2026, 10, 5, 0, 0, 0, 0, time.UTC)) + if len(name) > 63 || !strings.HasSuffix(name, "_deleted_20261005") { + t.Fatal(name) + } +} + +func TestACallersStatementRunsAsTheReaderAsGivenReadOnly(t *testing.T) { + conn := admin + conn.ReaderPassword = "reader-secret" + f, c := newFake(conn) + f.answer(`^COMMIT; DROP`, []string{"name", "n"}, []string{"alpha", "1"}, []string{"b,eta", "2"}) + statement := "COMMIT; DROP TABLE everything" + r, err := c.ReadOnlyQuery(ctx, "inventory", statement) + if err != nil { + t.Fatal(err) + } + asked := f.calls[len(f.calls)-1] + if asked.User != Reader || asked.Password != "reader-secret" || asked.Database != "inventory" { + t.Fatalf("not as the reader: %+v", asked) + } + if asked.SQL != statement { + t.Fatalf("not sent as given: %q", asked.SQL) + } + if asked.Options["default_transaction_read_only"] != "on" || asked.Options["statement_timeout"] != "60s" { + t.Fatalf("session not read-only from its first statement: %v", asked.Options) + } + rows := r.Maps() + if r.Command != "COMMIT" || len(rows) != 2 || rows[1]["name"] != "b,eta" || rows[0]["n"] != "1" { + t.Fatalf("%+v", r) + } +} + +func TestTheReaderIsMadeAsTheAdminWithEveryAttributeOnce(t *testing.T) { + conn := admin + conn.ReaderPassword = "reader-secret" + f, c := newFake(conn) + for _, q := range []string{"SELECT 1", "SELECT 2"} { + if _, err := c.ReadOnlyQuery(ctx, "inventory", q); err != nil { + t.Fatal(err) + } + } + var ddl []string + readerCalls := 0 + for _, k := range f.calls { + switch k.User { + case "postgres": + if k.Password != "admin-secret" { + t.Fatal("admin call without the admin's password") + } + ddl = append(ddl, k.SQL) + case Reader: + readerCalls++ + } + } + creates := 0 + for _, s := range ddl { + if strings.HasPrefix(s, `CREATE ROLE "`+Reader+`"`) { + creates++ + for _, a := range []string{"LOGIN", "NOSUPERUSER", "NOCREATEDB", "NOCREATEROLE", "NOREPLICATION", "NOBYPASSRLS"} { + if !strings.Contains(s, " "+a+" ") { + t.Fatalf("%s missing from %s", a, s) + } + } + } + } + if creates != 1 || readerCalls != 2 { + t.Fatalf("made %d times, %d reader calls", creates, readerCalls) + } + if !has(ddl, `^GRANT pg_read_all_data TO "`+Reader+`"$`) || !has(ddl, `default_transaction_read_only = on`) { + t.Fatal(strings.Join(ddl, "\n")) + } +} + +func TestWithoutTheReadersPasswordNothingRunsAsTheAdmin(t *testing.T) { + f, c := newFake(admin) + _, err := c.ReadOnlyQuery(ctx, "inventory", "SELECT 1") + if err == nil || !strings.Contains(err.Error(), "refused rather than run as the admin") { + t.Fatal(err) + } + if len(f.calls) != 0 { + t.Fatalf("ran %v", f.calls) + } +} + +func TestAFailedReaderSetupIsAskedAgain(t *testing.T) { + conn := admin + conn.ReaderPassword = "r" + f, c := newFake(conn) + f.fail(`^CREATE ROLE`, errors.New("server busy")) + if _, err := c.ReadOnlyQuery(ctx, "db", "SELECT 1"); err == nil { + t.Fatal("no error") + } + f.rules = nil + if _, err := c.ReadOnlyQuery(ctx, "db", "SELECT 1"); err != nil { + t.Fatal(err) + } +} + +func TestHolds(t *testing.T) { + f, c := newFake(admin) + f.answer(`FROM pg_extension`, []string{"extname"}, []string{"plpgsql"}) + ok, err := c.Holds(ctx, "db", "db", "pw", nil) + if err != nil || !ok { + t.Fatal(ok, err) + } + if k := f.calls[0]; k.User != "db" || k.Password != "pw" || k.Database != "db" { + t.Fatalf("not checked as the consumer: %+v", k) + } + // A wanted extension missing: not held, so it is applied again. + if ok, err := c.Holds(ctx, "db", "db", "pw", []string{"vector"}); err != nil || ok { + t.Fatal(ok, err) + } + // The server saying the login is wrong or gone: not held. + for _, code := range []string{"28P01", "28000", "3D000", "42501"} { + f.dialErr = func(Login) error { return &pgconn.PgError{Code: code} } + if ok, err := c.Holds(ctx, "db", "db", "pw", nil); err != nil || ok { + t.Fatal(code, ok, err) + } + } + // Unable to ask is not evidence of loss. + f.dialErr = func(Login) error { return errors.New("connection refused") } + if _, err := c.Holds(ctx, "db", "db", "pw", nil); err == nil { + t.Fatal("an unreachable server reported as a lost login") + } +} + +func TestClientFromEnv(t *testing.T) { + dir := t.TempDir() + write := func(name, content string) string { + p := filepath.Join(dir, name) + if err := os.WriteFile(p, []byte(content), 0o600); err != nil { + t.Fatal(err) + } + return p + } + env := map[string]string{ + "MESH_PROVISION_POSTGRES": "postgres://postgres@127.0.0.1:6852/postgres?sslmode=disable", + "MESH_PROVISION_PASSWORD_FILE": write("superuser.secret", "admin\n"), + "MESH_POSTGRES_READER_PASSWORD_FILE": write("reader.secret", "from-the-file\n"), + "MESH_PROVISION_POSTGRES_PORT": "${port:mesh-store}", + } + c, err := ClientFromEnv(func(k string) string { return env[k] }) + if err != nil { + t.Fatal(err) + } + if c.conn.Host != "127.0.0.1" || c.conn.Port != 6852 || c.conn.User != "postgres" || c.conn.Password != "admin" || + c.conn.ReaderPassword != "from-the-file" || c.conn.SSLMode != "disable" { + t.Fatalf("%+v", c.conn) + } + env["MESH_PROVISION_POSTGRES_PORT"] = "7000" + c, _ = ClientFromEnv(func(k string) string { return env[k] }) + if c.conn.Port != 7000 { + t.Fatal(c.conn.Port) + } + if _, err := ClientFromEnv(func(string) string { return "" }); err == nil { + t.Fatal("no host and no password accepted") + } +} + +func TestNullsStayNull(t *testing.T) { + v := "x" + r := Result{Fields: []string{"a", "b"}, Rows: [][]*string{{&v, nil}}} + m := r.Maps()[0] + if m["a"] != "x" || m["b"] != nil { + t.Fatal(m) + } +} diff --git a/modules/postgres/cmd/postgres-provider/fake_test.go b/modules/postgres/cmd/postgres-provider/fake_test.go new file mode 100644 index 0000000..17ccc6b --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/fake_test.go @@ -0,0 +1,123 @@ +package main + +// A fake server: every session records who it connected as and what it was sent, and answers by the +// first rule whose pattern matches the statement. That the server keeps data, refuses a reader's +// write or installs an extension is proven against a real server (live_test.go), not here; this +// holds the module to asking for the right things. + +import ( + "context" + "regexp" + "strings" + "sync" + "testing" +) + +type call struct { + Login + SQL string +} + +type rule struct { + match *regexp.Regexp + fields []string + rows [][]string + err error +} + +type fakeServer struct { + mu sync.Mutex + calls []call + rules []rule + dialErr func(Login) error +} + +func (f *fakeServer) answer(pattern string, fields []string, rows ...[]string) { + f.rules = append(f.rules, rule{match: regexp.MustCompile(pattern), fields: fields, rows: rows}) +} + +func (f *fakeServer) fail(pattern string, err error) { + f.rules = append(f.rules, rule{match: regexp.MustCompile(pattern), err: err}) +} + +type fakeSession struct { + f *fakeServer + l Login +} + +func (s fakeSession) Run(_ context.Context, sql string) ([]Result, error) { + s.f.mu.Lock() + defer s.f.mu.Unlock() + s.f.calls = append(s.f.calls, call{Login: s.l, SQL: sql}) + for _, r := range s.f.rules { + if r.match.MatchString(sql) { + if r.err != nil { + return nil, r.err + } + res := Result{Command: strings.Fields(sql)[0], Fields: r.fields} + for _, row := range r.rows { + cells := make([]*string, len(row)) + for i := range row { + v := row[i] + cells[i] = &v + } + res.Rows = append(res.Rows, cells) + } + return []Result{res}, nil + } + } + return []Result{{Command: strings.Fields(sql)[0]}}, nil +} + +func (s fakeSession) Close() {} + +func (f *fakeServer) dial(_ context.Context, l Login) (Session, error) { + if f.dialErr != nil { + if err := f.dialErr(l); err != nil { + return nil, err + } + } + return fakeSession{f: f, l: l}, nil +} + +func (f *fakeServer) statements() []string { + f.mu.Lock() + defer f.mu.Unlock() + out := make([]string, len(f.calls)) + for i, c := range f.calls { + out[i] = c.SQL + } + return out +} + +func (f *fakeServer) reset() { + f.mu.Lock() + defer f.mu.Unlock() + f.calls = nil +} + +var admin = Conn{Host: "127.0.0.1", Port: 5432, User: "postgres", Password: "admin-secret"} + +func newFake(conn Conn) (*fakeServer, *Client) { + f := &fakeServer{} + return f, NewClient(conn, f.dial) +} + +func noDrop(t *testing.T, sql []string) { + t.Helper() + for _, s := range sql { + if regexp.MustCompile(`(?i)\bDROP\b`).MatchString(s) { + t.Fatalf("something was dropped:\n%s", strings.Join(sql, "\n")) + } + } +} + +func has(sql []string, pattern string) bool { + re := regexp.MustCompile(pattern) + for _, s := range sql { + if re.MatchString(s) { + return true + } + } + return false +} diff --git a/modules/postgres/cmd/postgres-provider/harness.go b/modules/postgres/cmd/postgres-provider/harness.go new file mode 100644 index 0000000..d7cceed --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/harness.go @@ -0,0 +1,340 @@ +package main + +// The reconcile loop every provider shares, as the TypeScript SDK's runProvisioner runs it +// (@novox/mesh-sdk/provisioner, 0.1.10). The Go SDK has no provisioner yet, so this module carries +// the loop itself, line for line in behaviour; when the Go SDK grows one, this file is what moves +// there (novox/hq ADR 0039: the loop is the SDK's, the adapter is the module's). +// +// Read the contributions the mesh delivered; bring each consumer's resource into being through the +// adapter, under the login and password the mesh minted; withdraw what the mesh no longer asks for. +// **A provider creates the credential the mesh minted, and seals nothing (novox/hq ADR 0048).** + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "net/url" + "os" + "strings" + "time" +) + +// Provision is one consumer's resource to bring into being — everything the mesh derived and delivered. +type Provision struct { + // As is the login the mesh derived and gave the consumer to present. + As string + // Password is the one the mesh minted, read from the file the host unsealed. + Password string + // Values are what the consumer contributed (e.g. {"name": "letta", "extensions": ["vector"]}). + Values map[string]any + // Derived is what this provider's own definition derives for the consumer (novox/hq ADR 0201). + Derived map[string]any + // At is where the consumer is; Consumer is its node. + At string + Consumer string +} + +// Adapter is the per-service half. +type Adapter interface { + Create(ctx context.Context, p Provision) error + // Remove withdraws what Create made; derived is what the mesh last derived, remembered here. + Remove(ctx context.Context, as string, derived map[string]any) error + // Holds says whether the backend still holds the consumer exactly as p says. Read-only. + Holds(ctx context.Context, p Provision) (bool, error) +} + +// Harness is the loop's settings and memory. +type Harness struct { + Resource string + Receives string + Adapter Adapter + Every time.Duration // 5s + VerifyEvery time.Duration // 60s + HoldsTimeout time.Duration // 30s + Log func(format string, args ...any) + Now func() time.Time + + verifiedAt time.Time + applied map[string]appliedEntry + lost map[string]brake + waiting map[string]int + failing map[string]failure + lastWarning string +} + +type appliedEntry struct { + hash string + derived map[string]any +} + +type brake struct { + times int + nextAt time.Time +} + +type failure struct { + text string + times int +} + +// The longest a consumer whose create keeps failing to satisfy holds waits between checks. +const maxBackoff = time.Hour + +// How many passes a secret may be unreadable before it stops being called a race (issue 225), and +// once said loudly, how often it is repeated. The same cadence quiets a create that keeps failing +// the same way. +const ( + patiently = 12 + loudlyEvery = 240 +) + +type contribution struct { + As string `json:"as"` + Secret string `json:"secret"` + Node string `json:"node"` + At string `json:"at"` + Values map[string]any `json:"values"` + Derived map[string]any `json:"derived"` +} + +func (h *Harness) init() { + if h.Every == 0 { + h.Every = 5 * time.Second + } + if h.VerifyEvery == 0 { + h.VerifyEvery = time.Minute + } + if h.HoldsTimeout == 0 { + h.HoldsTimeout = 30 * time.Second + } + if h.Now == nil { + h.Now = time.Now + } + if h.Log == nil { + h.Log = func(format string, args ...any) { fmt.Fprintf(os.Stderr, format+"\n", args...) } + } + if h.applied == nil { + h.applied = map[string]appliedEntry{} + h.lost = map[string]brake{} + h.waiting = map[string]int{} + h.failing = map[string]failure{} + } +} + +// Run reconciles until ctx ends. One consumer's failure never stops the others'. +func (h *Harness) Run(ctx context.Context) { + h.init() + for { + h.Reconcile(ctx) + select { + case <-ctx.Done(): + return + case <-time.After(h.Every): + } + } +} + +func (h *Harness) say(format string, args ...any) { + h.Log("[provisioner:"+h.Resource+"] "+format, args...) +} + +func (h *Harness) warn(why string) { + if why == h.lastWarning { + return + } + if why != "" { + h.say("%s; nothing applied or removed until it can be read", why) + } else { + h.say("contributions file readable again") + } + h.lastWarning = why +} + +// readContributions answers the consumers asked for, or nil when the file says nothing usable. +// **Nothing read is not nobody asking** (novox/hq issue 241): only a file that was read can withdraw. +func (h *Harness) readContributions() []contribution { + raw, err := os.ReadFile(h.Receives) + if err != nil { + h.warn(fmt.Sprintf("contributions file unreadable (%s): %v", h.Receives, err)) + return nil + } + var doc struct { + Requirement string `json:"requirement"` + Given json.RawMessage `json:"given"` + } + if err := json.Unmarshal(raw, &doc); err != nil { + h.warn(fmt.Sprintf("contributions file is not JSON (%s): %v", h.Receives, err)) + return nil + } + if doc.Requirement != "" && doc.Requirement != h.Resource { + h.warn(fmt.Sprintf("%s is for %s, not %s", h.Receives, doc.Requirement, h.Resource)) + return nil + } + var given []contribution + if len(doc.Given) == 0 || string(doc.Given) == "null" || json.Unmarshal(doc.Given, &given) != nil { + h.warn(fmt.Sprintf("%s has no given list", h.Receives)) + return nil + } + h.warn("") + out := []contribution{} + for _, g := range given { + // No `as` is not a credential grant: nothing to create for it. + if g.As != "" && g.Secret != "" { + out = append(out, g) + } + } + return out +} + +func hashOf(as, password string, values, derived map[string]any) string { + // Derived is in the hash: a provider that renames what it derives gave a different resource. + b, _ := json.Marshal([]any{as, password, orEmpty(values), orEmpty(derived)}) + return string(b) +} + +func orEmpty(m map[string]any) map[string]any { + if m == nil { + return map[string]any{} + } + return m +} + +// Reconcile is one pass. +func (h *Harness) Reconcile(ctx context.Context) { + h.init() + given := h.readContributions() + if given == nil { + return + } + want := map[string]bool{} + for _, g := range given { + want[g.As] = true + } + verifying := h.Now().Sub(h.verifiedAt) >= h.VerifyEvery + if verifying { + h.verifiedAt = h.Now() + } + + for _, g := range given { + raw, err := os.ReadFile(g.Secret) + if err != nil { + // A secret the host has not written yet is a race on the first pass; past a minute it is + // a person's to look at, and said so (novox/hq issue 225). + n := h.waiting[g.As] + 1 + h.waiting[g.As] = n + if n <= patiently { + h.say("%s: secret not readable yet (%s): %v", g.As, g.Secret, err) + } else if n == patiently+1 || n%loudlyEvery == 0 { + h.say("%s: CANNOT READ the secret after %d attempts (%s): %v. This is not a race any more — "+ + "nothing has been provisioned for this consumer and nothing will be until somebody looks. "+ + "Check who owns the file and who this process runs as (novox/hq issue 225)", g.As, n, g.Secret, err) + } + continue + } + delete(h.waiting, g.As) + password := strings.TrimSuffix(string(raw), "\n") + p := Provision{As: g.As, Password: password, Values: orEmpty(g.Values), Derived: orEmpty(g.Derived), At: g.At, Consumer: g.Node} + hash := hashOf(g.As, password, g.Values, g.Derived) + + reapplying := 0 + if was, ok := h.applied[g.As]; ok && was.hash == hash { + if !verifying { + continue + } + b, braked := h.lost[g.As] + if braked && h.Now().Before(b.nextAt) { + continue + } + hctx, cancel := context.WithTimeout(ctx, h.HoldsTimeout) + held, err := h.Adapter.Holds(hctx, p) + timedOut := errors.Is(hctx.Err(), context.DeadlineExceeded) + cancel() + if err != nil { + // Unable to ask is not evidence of loss. A backend that timed out will time out for + // the next consumer too, so the rest of this pass is not asked. + h.say("%s: could not check the backend, will ask again: %s", g.As, scrub(err, password)) + if timedOut { + verifying = false + } + continue + } + if held { + delete(h.lost, g.As) + continue + } + reapplying = b.times + 1 + if reapplying == 1 { + h.say("%s: the backend no longer holds it; applying again", g.As) + } else { + h.say("%s: still not held after being applied again (%d times in a row) — create does not "+ + "produce what holds checks; applying again", g.As, reapplying) + } + } + if err := h.Adapter.Create(ctx, p); err != nil { + text := scrub(err, password) + f := h.failing[g.As] + if f.text != text { + f = failure{text: text} + } + f.times++ + h.failing[g.As] = f + // Said each time it changes, and while it stays the same, as rarely as a lost secret. + if f.times == 1 || f.times%loudlyEvery == 0 { + h.say("%s: create failed, will retry: %s", g.As, text) + } + if reapplying > 0 { + h.lost[g.As] = brake{times: reapplying - 1} + } + continue + } + if f, was := h.failing[g.As]; was { + h.say("%s: created, after %d failed attempt(s)", g.As, f.times) + delete(h.failing, g.As) + } + h.applied[g.As] = appliedEntry{hash: hash, derived: p.Derived} + if reapplying == 0 { + delete(h.lost, g.As) + } else { + wait := h.VerifyEvery << (reapplying - 1) + if wait > maxBackoff || wait <= 0 { + wait = maxBackoff + } + h.lost[g.As] = brake{times: reapplying, nextAt: h.Now().Add(wait)} + if reapplying > 1 { + h.say("%s: next check in %s", g.As, wait.Round(time.Second)) + } + } + } + + // Withdraw every login this process made that the mesh no longer asks for. + for as, was := range h.applied { + if want[as] { + continue + } + h.say("%s: no longer in %s; withdrawing it from the backend", as, h.Receives) + if err := h.Adapter.Remove(ctx, as, was.derived); err != nil { + h.say("%s: remove failed, will retry: %v", as, err) + continue + } + delete(h.applied, as) + delete(h.lost, as) + } + for as := range h.failing { + if !want[as] { + delete(h.failing, as) + } + } +} + +// scrub is an error's text with the consumer's password removed, raw and URL-encoded. +func scrub(err error, password string) string { + text := err.Error() + if password == "" { + return text + } + for _, form := range []string{password, url.QueryEscape(password), url.PathEscape(password)} { + text = strings.ReplaceAll(text, form, "***") + } + return text +} diff --git a/modules/postgres/cmd/postgres-provider/harness_test.go b/modules/postgres/cmd/postgres-provider/harness_test.go new file mode 100644 index 0000000..f79aab1 --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/harness_test.go @@ -0,0 +1,232 @@ +package main + +import ( + "context" + "encoding/json" + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +type recorder struct { + created []Provision + removed []string + held bool + failing error +} + +func (r *recorder) Create(_ context.Context, p Provision) error { + if r.failing != nil { + return r.failing + } + r.created = append(r.created, p) + return nil +} + +func (r *recorder) Remove(_ context.Context, as string, _ map[string]any) error { + r.removed = append(r.removed, as) + return nil +} + +func (r *recorder) Holds(context.Context, Provision) (bool, error) { return r.held, nil } + +type world struct { + t *testing.T + dir string + receives string + now time.Time + h *Harness + a *recorder + said []string +} + +func newWorld(t *testing.T) *world { + w := &world{t: t, dir: t.TempDir(), now: time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC), a: &recorder{held: true}} + w.receives = filepath.Join(w.dir, "mesh.json") + w.h = &Harness{Resource: "postgres-database", Receives: w.receives, Adapter: w.a, + Now: func() time.Time { return w.now }, + Log: func(f string, a ...any) { w.said = append(w.said, f) }} + return w +} + +func (w *world) give(given ...map[string]any) { + for _, g := range given { + secret := filepath.Join(w.dir, g["as"].(string)+".secret") + if err := os.WriteFile(secret, []byte("pw-"+g["as"].(string)+"\n"), 0o600); err != nil { + w.t.Fatal(err) + } + g["secret"] = secret + } + if given == nil { + given = []map[string]any{} + } + raw, _ := json.Marshal(map[string]any{"requirement": "postgres-database", "given": given}) + if err := os.WriteFile(w.receives, raw, 0o600); err != nil { + w.t.Fatal(err) + } +} + +func TestAConsumerIsCreatedOnceUnderTheMeshsLoginAndPassword(t *testing.T) { + w := newWorld(t) + w.give(map[string]any{"as": "mesh_ace_letta", "node": "ace", "values": map[string]any{"name": "letta"}}) + w.h.Reconcile(ctx) + w.h.Reconcile(ctx) + if len(w.a.created) != 1 { + t.Fatalf("created %d times", len(w.a.created)) + } + p := w.a.created[0] + if p.As != "mesh_ace_letta" || p.Password != "pw-mesh_ace_letta" || p.Consumer != "ace" { + t.Fatalf("%+v", p) + } +} + +func TestAddingExtensionsAppliesTheConsumerAgain(t *testing.T) { + w := newWorld(t) + w.give(map[string]any{"as": "mesh_ace_letta", "values": map[string]any{"name": "letta"}}) + w.h.Reconcile(ctx) + w.give(map[string]any{"as": "mesh_ace_letta", "values": map[string]any{"name": "letta", "extensions": []any{"vector"}}}) + w.h.Reconcile(ctx) + if len(w.a.created) != 2 { + t.Fatalf("created %d times", len(w.a.created)) + } + if ext, _ := Extensions(w.a.created[1].Values); len(ext) != 1 || ext[0] != "vector" { + t.Fatal(ext) + } +} + +func TestAConsumerNoLongerAskedForIsWithdrawn(t *testing.T) { + w := newWorld(t) + w.give(map[string]any{"as": "a"}, map[string]any{"as": "b"}) + w.h.Reconcile(ctx) + w.give(map[string]any{"as": "a"}) + w.h.Reconcile(ctx) + if strings.Join(w.a.removed, ",") != "b" { + t.Fatal(w.a.removed) + } + // Only a file that says nobody asks withdraws everybody. + w.give() + w.h.Reconcile(ctx) + if strings.Join(w.a.removed, ",") != "b,a" { + t.Fatal(w.a.removed) + } +} + +func TestNothingReadIsNotNobodyAsking(t *testing.T) { + for name, content := range map[string]string{ + "unreadable": "", + "not JSON": "{", + "no given": `{"requirement": "postgres-database"}`, + "another": `{"requirement": "mssql-database", "given": []}`, + } { + t.Run(name, func(t *testing.T) { + w := newWorld(t) + w.give(map[string]any{"as": "a"}) + w.h.Reconcile(ctx) + if content == "" { + os.Remove(w.receives) + } else { + os.WriteFile(w.receives, []byte(content), 0o600) + } + w.h.Reconcile(ctx) + if len(w.a.removed) != 0 { + t.Fatalf("withdrew %v on a file it could not use", w.a.removed) + } + }) + } +} + +func TestALostConsumerIsAppliedAgainAndBraked(t *testing.T) { + w := newWorld(t) + w.give(map[string]any{"as": "a"}) + w.h.Reconcile(ctx) + w.a.held = false + w.now = w.now.Add(2 * time.Minute) + w.h.Reconcile(ctx) + if len(w.a.created) != 2 { + t.Fatalf("created %d times", len(w.a.created)) + } + // Still not held a minute later: applied again, then braked for two minutes. + w.now = w.now.Add(61 * time.Second) + w.h.Reconcile(ctx) + w.now = w.now.Add(61 * time.Second) + w.h.Reconcile(ctx) + if len(w.a.created) != 3 { + t.Fatalf("not braked: created %d times", len(w.a.created)) + } +} + +func TestAFailingCreateIsRetriedAndSaidOnce(t *testing.T) { + w := newWorld(t) + w.a.failing = &pgErr{"extension \"nope\" refused, password pw-a"} + w.give(map[string]any{"as": "a"}) + for i := 0; i < 5; i++ { + w.h.Reconcile(ctx) + } + n := 0 + for _, s := range w.said { + if strings.Contains(s, "create failed") { + n++ + } + } + if n != 1 { + t.Fatalf("said %d times", n) + } + w.a.failing = nil + w.h.Reconcile(ctx) + if len(w.a.created) != 1 { + t.Fatal("not retried") + } +} + +type pgErr struct{ s string } + +func (e *pgErr) Error() string { return e.s } + +func TestScrubRemovesThePassword(t *testing.T) { + if got := scrub(&pgErr{"bad pw a/b c and a%2Fb+c"}, "a/b c"); strings.Contains(got, "a/b c") || strings.Contains(got, "a%2Fb+c") { + t.Fatal(got) + } +} + +func TestProvisionerCreatesTheDatabaseThenItsExtensions(t *testing.T) { + f, c := newFake(admin) + f.answer(`FROM pg_available_extensions`, []string{"name"}, []string{"vector"}) + var events []string + a := provisioner{pg: c, announce: func(e string, _ map[string]string) { events = append(events, e) }} + p := Provision{As: "mesh_ace_letta", Password: "pw", Values: map[string]any{"name": "letta", "extensions": []any{"vector"}}} + if err := a.Create(ctx, p); err != nil { + t.Fatal(err) + } + sql := f.statements() + if !strings.HasPrefix(sql[len(sql)-1], `CREATE EXTENSION IF NOT EXISTS "vector"`) || !has(sql, `CREATE DATABASE "mesh_ace_letta"`) { + t.Fatal(strings.Join(sql, "\n")) + } + if strings.Join(events, ",") != "database.provisioned" { + t.Fatal(events) + } + + // An extension the server does not offer: refused, and no event says it was provisioned. + events = nil + p.Values = map[string]any{"extensions": []any{"nope"}} + if err := a.Create(ctx, p); err == nil || !strings.Contains(err.Error(), `"nope"`) || len(events) != 0 { + t.Fatal(err, events) + } + // A malformed list: refused before the server is asked anything. + f.reset() + p.Values = map[string]any{"extensions": "vector"} + if err := a.Create(ctx, p); err == nil || len(f.calls) != 0 { + t.Fatal(err, f.calls) + } + + f.reset() + f.answer(`FROM pg_roles`, []string{"?column?"}, []string{"1"}) + if err := a.Remove(ctx, "mesh_ace_letta", nil); err != nil { + t.Fatal(err) + } + noDrop(t, f.statements()) + if !has(f.statements(), `NOLOGIN`) { + t.Fatal(f.statements()) + } +} diff --git a/modules/postgres/cmd/postgres-provider/live_test.go b/modules/postgres/cmd/postgres-provider/live_test.go new file mode 100644 index 0000000..52c585d --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/live_test.go @@ -0,0 +1,84 @@ +package main + +// Against a real server, when one is named — skipped otherwise. A throwaway one: +// +// docker run -d --rm --name pg-test -e POSTGRES_PASSWORD=admin -p 55432:5432 pgvector/pgvector:pg17 +// MESH_POSTGRES_LIVE=postgres://postgres:admin@127.0.0.1:55432/postgres?sslmode=disable go test ./... +// +// It proves what the fakes cannot: an untrusted extension is installed by the superuser in a fresh +// consumer database and a second pass is a no-op; the consumer then holds; a withdrawn login cannot +// log in and its data is still there; the reader cannot write, even past a COMMIT. + +import ( + "net/url" + "os" + "testing" +) + +func TestLive(t *testing.T) { + raw := os.Getenv("MESH_POSTGRES_LIVE") + if raw == "" { + t.Skip("MESH_POSTGRES_LIVE names no server") + } + u, _ := url.Parse(raw) + pw, _ := u.User.Password() + env := map[string]string{"MESH_PROVISION_POSTGRES": raw, "MESH_POSTGRES_PASSWORD": pw, "MESH_POSTGRES_READER_PASSWORD": "reader-pw"} + c, err := ClientFromEnv(func(k string) string { return env[k] }) + if err != nil { + t.Fatal(err) + } + a := provisioner{pg: c, announce: func(string, map[string]string) {}} + p := Provision{As: "mesh_test_letta", Password: "consumer-pw", Values: map[string]any{"name": "letta", "extensions": []any{"vector"}}} + + // vector is not trusted: the consumer cannot install it itself. + if err := c.CreateDatabaseAndRole(ctx, p.As, p.As, p.Password); err != nil { + t.Fatal(err) + } + if _, err := c.as(ctx, Login{Database: p.As, User: p.As, Password: p.Password}, "CREATE EXTENSION IF NOT EXISTS vector"); err == nil { + t.Log("note: the consumer could install vector itself on this server") + } + for pass := 0; pass < 2; pass++ { + if err := a.Create(ctx, p); err != nil { + t.Fatalf("pass %d: %v", pass, err) + } + } + if ok, err := a.Holds(ctx, p); err != nil || !ok { + t.Fatal("not held after create:", ok, err) + } + if _, err := c.as(ctx, Login{Database: p.As, User: p.As, Password: p.Password}, + "CREATE TABLE IF NOT EXISTS kept (v vector(3)); INSERT INTO kept VALUES ('[1,2,3]')"); err != nil { + t.Fatal("the consumer cannot use the type:", err) + } + bad := p + bad.Values = map[string]any{"extensions": []any{"no_such_extension"}} + if err := a.Create(ctx, bad); err == nil { + t.Fatal("an unknown extension was accepted") + } + + r, err := c.ReadOnlyQuery(ctx, p.As, "COMMIT; DROP TABLE kept") + if err == nil { + t.Fatalf("the reader dropped a table: %+v", r) + } + r, err = c.ReadOnlyQuery(ctx, p.As, "SELECT count(*) AS n FROM kept") + if err != nil || r.Maps()[0]["n"] == nil { + t.Fatal(r, err) + } + + if err := a.Remove(ctx, p.As, nil); err != nil { + t.Fatal(err) + } + if ok, err := a.Holds(ctx, p); err != nil || ok { + t.Fatal("a withdrawn login still logs in:", ok, err) + } + r, err = c.Query(ctx, p.As, "SELECT count(*) FROM kept") + if err != nil || len(r.Rows) != 1 { + t.Fatal("withdrawing lost the data:", err) + } + // Coming back is given the same database. + if err := a.Create(ctx, p); err != nil { + t.Fatal(err) + } + if ok, _ := a.Holds(ctx, p); !ok { + t.Fatal("not held after coming back") + } +} diff --git a/modules/postgres/cmd/postgres-provider/main.go b/modules/postgres/cmd/postgres-provider/main.go new file mode 100644 index 0000000..f93b7cc --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/main.go @@ -0,0 +1,89 @@ +// postgres-provider: postgres's code, one binary the node's runtime launches and speaks MCP to over +// stdio through the Go SDK (novox/hq ADR 0188, 0193). It serves postgres's tools and the mesh-store +// seat's verbs and, beside them, runs long: the provisioner that makes postgres the provider of the +// mesh `postgres-database` interface, and an audit line for each database it provisions or withdraws. +// +// stdout is the MCP channel; everything this module says, it says on stderr. +package main + +import ( + "context" + "encoding/json" + "fmt" + "os" + "time" + + stdio "git.novox.be/novox/mesh-sdk/go" +) + +func say(format string, args ...any) { + fmt.Fprintf(os.Stderr, "[postgres] "+format+"\n", args...) +} + +func main() { + pg, err := ClientFromEnv(os.Getenv) + if err != nil { + // Without the server there is nothing to serve and nothing to provision; said, not fatal, + // so the runtime does not restart a process that cannot do better. + say("%v; serving no tools and provisioning nothing", err) + if err := stdio.Serve("", nil); err != nil { + say("%v", err) + os.Exit(1) + } + return + } + if receives := os.Getenv("MESH_RECEIVES"); receives == "" { + say("MESH_RECEIVES is not set — the provisioner cannot run without it") + } else { + h := &Harness{ + Resource: "postgres-database", + Receives: receives, + Adapter: provisioner{pg: pg, announce: announce}, + Log: func(format string, args ...any) { fmt.Fprintf(os.Stderr, format+"\n", args...) }, + } + go h.Run(context.Background()) + } + go audit() + if err := stdio.Serve("", Tools(pg)); err != nil { + say("%v", err) + os.Exit(1) + } +} + +// announce emits a lifecycle event without letting a broker hiccup fail the provisioning itself. +func announce(event string, body map[string]string) { + if err := stdio.Emit(event, body); err != nil { + say("emit %s failed: %v", event, err) + } +} + +// audit keeps a line for each database granted and withdrawn — observability the provider is best +// placed to log. The events are its own, read back from its consumer (`consumes`). +func audit() { + subscribe := func() error { + return stdio.Subscribe("postgres.database.*", func(e stdio.Envelope) error { + var body struct { + Consumer string `json:"consumer"` + Database string `json:"database"` + } + _ = json.Unmarshal(e.Body, &body) + switch e.Key { + case "postgres.database.provisioned": + say("database provisioned for %s (db %s)", body.Consumer, body.Database) + case "postgres.database.deprovisioned": + say("database withdrawn, kept (db %s)", body.Database) + } + return nil + }) + } + // Asked again until the runtime takes it: the first ask may come before Serve is running. + for wait := time.Second; ; wait = min(wait*2, time.Minute) { + if err := subscribe(); err == nil { + say("auditing database lifecycle events") + return + } else if wait >= 8*time.Second { + say("subscribing to the lifecycle events failed, will retry: %v", err) + } + time.Sleep(wait) + } +} diff --git a/modules/postgres/cmd/postgres-provider/provisioner.go b/modules/postgres/cmd/postgres-provider/provisioner.go new file mode 100644 index 0000000..d641359 --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/provisioner.go @@ -0,0 +1,64 @@ +package main + +// postgres's provisioner — the adapter that makes postgres a provider of the mesh +// `postgres-database` interface (novox/hq ADR 0039/0040/0048). A consumer connects to a database it +// alone owns, as `as` with the password the mesh minted. +// +// **The role name and password are the mesh's, not the provisioner's (ADR 0048).** postgres creates +// a role and a same-named database under exactly that login — a name the consumer cannot learn is a +// database it cannot reach. +// +// **Extensions are the provider's to install.** A contribution may name extensions +// (`"extensions": ["vector"]`); most are not trusted, so only the superuser this module holds can +// create them, in the consumer's database, on every pass, and never drops one. + +import ( + "context" +) + +// provisioner is the adapter over the client; announce emits a lifecycle event. +type provisioner struct { + pg *Client + announce func(event string, body map[string]string) +} + +func (a provisioner) Create(ctx context.Context, p Provision) error { + // Read before anything runs: a malformed list is refused without touching the server. + extensions, err := Extensions(p.Values) + if err != nil { + return err + } + // Database and owning role share the consumer's login, so the consumer owns exactly its own. + database := p.As + if err := a.pg.CreateDatabaseAndRole(ctx, database, p.As, p.Password); err != nil { + return err + } + if err := a.pg.EnsureExtensions(ctx, database, extensions); err != nil { + return err + } + a.announce("database.provisioned", map[string]string{"consumer": p.Consumer, "database": database, "user": p.As}) + return nil +} + +// Remove withdraws, never drops (novox/hq issue 241). The login is locked and the database kept under +// its own name: on 2026-10-04 a misread contributions file withdrew every consumer at once, and +// dropping made that a loss of seven databases. Taking a database out of service is a person's act — +// postgres_retire_database — and even that renames rather than drops. +func (a provisioner) Remove(ctx context.Context, as string, _ map[string]any) error { + if err := a.pg.LockRole(ctx, as); err != nil { + return err + } + a.announce("database.deprovisioned", map[string]string{"database": as, "kept": "true"}) + return nil +} + +// Holds is asked every minute: whether the consumer can still log in as the mesh gave it, and finds +// the extensions it asked for, so a login or extension lost behind the provisioner's back is made +// again (novox/hq issue 120). +func (a provisioner) Holds(ctx context.Context, p Provision) (bool, error) { + extensions, err := Extensions(p.Values) + if err != nil { + return false, err + } + return a.pg.Holds(ctx, p.As, p.As, p.Password, extensions) +} diff --git a/modules/postgres/cmd/postgres-provider/tools.go b/modules/postgres/cmd/postgres-provider/tools.go new file mode 100644 index 0000000..901759a --- /dev/null +++ b/modules/postgres/cmd/postgres-provider/tools.go @@ -0,0 +1,98 @@ +package main + +// postgres's tools, and the mesh-store seat's verbs (novox/hq ADR 0159, 0160). The seat's verbs are +// named `mesh-store.`, so the runtime serves them as the seat's wherever this module holds it +// and never lists them as postgres's own; they are scoped to what the store enables — asking what it +// holds and reading from it — so retiring a database is postgres's tool and not the store's. + +import ( + "context" + "errors" + "fmt" + "time" + + stdio "git.novox.be/novox/mesh-sdk/go" +) + +// Seat is the store seat this module claims. +const Seat = "mesh-store" + +var queryInput = map[string]any{ + "database": map[string]any{"type": "string", "description": "the database to query"}, + "sql": map[string]any{"type": "string", "description": "the SELECT (or other read-only) statement"}, +} + +func str(args map[string]any, key string) string { + if s, ok := args[key].(string); ok { + return s + } + return "" +} + +func listTool(pg *Client, name, description string) stdio.Tool { + return stdio.Tool{Name: name, Description: description, Run: func(map[string]any) (any, error) { + ctx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + dbs, err := pg.ListDatabases(ctx) + if err != nil { + return nil, err + } + return map[string]any{"databases": dbs}, nil + }} +} + +func queryTool(pg *Client, name, prefix, description string) stdio.Tool { + return stdio.Tool{Name: name, Description: description, Input: queryInput, Run: func(args map[string]any) (any, error) { + database, sql := str(args, "database"), str(args, "sql") + if database == "" { + return nil, fmt.Errorf("%s: database is required", prefix) + } + if sql == "" { + return nil, fmt.Errorf("%s: sql is required", prefix) + } + ctx, cancel := context.WithTimeout(context.Background(), 90*time.Second) + defer cancel() + r, err := pg.ReadOnlyQuery(ctx, database, sql) + if err != nil { + return nil, err + } + return map[string]any{"database": database, "command": r.Command, "rows": r.Maps()}, nil + }} +} + +// Tools are postgres's own and the seat's verbs. +func Tools(pg *Client) []stdio.Tool { + return []stdio.Tool{ + listTool(pg, "postgres_list_databases", "List the databases on the postgres server, with their on-disk size."), + queryTool(pg, "postgres_query", "postgres_query", + "Run a read-only SQL query against a named database, as a login that can read every table and change nothing."), + { + Name: "postgres_retire_database", + Description: "Take one database out of service on purpose: rename it to _deleted_ and lock its owner's login. " + + "Nothing is dropped — the data stays on the server under the new name until a person removes it by hand. " + + "Repeat the database's name in confirm.", + Input: map[string]any{ + "database": map[string]any{"type": "string", "description": "the database to retire"}, + "confirm": map[string]any{"type": "string", "description": "the same name again, to say this is meant"}, + }, + Run: func(args map[string]any) (any, error) { + database := str(args, "database") + if database == "" { + return nil, errors.New("postgres_retire_database: database is required") + } + if str(args, "confirm") != database { + return nil, errors.New("postgres_retire_database: confirm must repeat the database's name") + } + ctx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + aside, err := pg.RetireDatabase(ctx, database, time.Now()) + if err != nil { + return nil, err + } + return map[string]any{"database": database, "renamedTo": aside, "dropped": false}, nil + }, + }, + listTool(pg, Seat+".databases", "Every database the store holds, with its on-disk size."), + queryTool(pg, Seat+".query", "query", "One read-only statement against one database the store holds."), + } +} diff --git a/modules/postgres/go.mod b/modules/postgres/go.mod new file mode 100644 index 0000000..d4bf4ff --- /dev/null +++ b/modules/postgres/go.mod @@ -0,0 +1,14 @@ +module postgres + +go 1.25.0 + +require ( + git.novox.be/novox/mesh-sdk/go v0.1.7 + github.com/jackc/pgx/v5 v5.11.0 +) + +require ( + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + golang.org/x/text v0.29.0 // indirect +) diff --git a/modules/postgres/go.sum b/modules/postgres/go.sum new file mode 100644 index 0000000..79d630d --- /dev/null +++ b/modules/postgres/go.sum @@ -0,0 +1,28 @@ +git.novox.be/novox/mesh-sdk/go v0.1.7 h1:C0sTQmtTiyYH7bnqZb7PusXnqA37gKuT7Nqjn9gG47w= +git.novox.be/novox/mesh-sdk/go v0.1.7/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.11.0 h1:IzBBtyK9AHqf98cctWFifYSci2hgQR/cd56wB4p+ogg= +github.com/jackc/pgx/v5 v5.11.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug= +golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/text v0.29.0 h1:1neNs90w9YzJ9BocxfsQNHKuAT4pkghyXc4nhZ6sJvk= +golang.org/x/text v0.29.0/go.mod h1:7MhJOA9CD2qZyOKYazxdYMF85OwPdEr9jTtBpO7ydH4= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/modules/postgres/index.ts b/modules/postgres/index.ts deleted file mode 100644 index 5129f4e..0000000 --- a/modules/postgres/index.ts +++ /dev/null @@ -1,25 +0,0 @@ -// postgres's events entrypoint, loaded by the per-node tool host (the provisioner container runs -// ./provisioner separately). The database lifecycle events are EMITTED from the provisioner, where -// the lifecycle actually happens (novox/hq ADR 0041/0042): -// module.postgres.database.provisioned — a consumer's database + owning role was created -// module.postgres.database.deprovisioned — that database was removed -// Here in the tool host we react to them, keeping a lightweight audit trail of who was granted a -// database and who lost one — observability the provider itself is best placed to log. - -import { on } from "@novox/mesh-sdk/events"; - -interface DatabaseEvent { - consumer: string; - database: string; - user?: string; -} - -await on("database.provisioned", async (e) => { - console.log(`[postgres] database provisioned for ${e.body.consumer} (db ${e.body.database})`); -}); - -await on("database.deprovisioned", async (e) => { - console.log(`[postgres] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`); -}); - -console.log("[postgres] auditing database lifecycle events"); diff --git a/modules/postgres/module.json b/modules/postgres/module.json index dc1a256..5f8da0c 100644 --- a/modules/postgres/module.json +++ b/modules/postgres/module.json @@ -109,16 +109,12 @@ { "name": "code", "kind": "bundle", - "language": "typescript", - "entrypoints": [ - "index.js", - "tools/index.js", - "provisioner/index.js" - ], + "language": "go", + "system": "arch", + "from": "cmd/postgres-provider", + "binary": "postgres-provider", "loads": [ - "index.js", - "tools/index.js", - "provisioner/index.js" + "postgres-provider" ], "env": { "MESH_PROVISION_POSTGRES": "postgres://postgres@127.0.0.1:${port:5432}/postgres?sslmode=disable", diff --git a/modules/postgres/package.json b/modules/postgres/package.json deleted file mode 100644 index f6e8c55..0000000 --- a/modules/postgres/package.json +++ /dev/null @@ -1,18 +0,0 @@ -{ - "name": "@novox/module-postgres", - "version": "0.1.0", - "description": "postgres \u2014 provides the mesh postgres-database interface. Its client, provisioner, tools and events live here (novox/hq ADR 0039).", - "type": "module", - "private": true, - "scripts": { - "build": "tsc client.ts index.ts provisioner/index.ts tools/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist", - "test": "npm run build && node --test --experimental-strip-types 'test/*.test.ts'" - }, - "dependencies": { - "@novox/mesh-sdk": "^0.1.1" - }, - "devDependencies": { - "@types/node": "^22.0.0", - "typescript": "^5.6.0" - } -} diff --git a/modules/postgres/provisioner/index.ts b/modules/postgres/provisioner/index.ts deleted file mode 100644 index a2091a6..0000000 --- a/modules/postgres/provisioner/index.ts +++ /dev/null @@ -1,57 +0,0 @@ -// postgres's provisioner — the adapter that makes postgres a provider of the mesh -// `postgres-database` interface. The reconcile loop, the contributions file, and reading the mesh's -// minted password are the sdk harness's; this writes only the per-service half: how postgres creates -// and removes a consumer's database + owning role (novox/hq ADR 0039/0040/0048). -// -// The `postgres-database` interface: a consumer connects to a database it alone owns, as `as` with -// the password the mesh minted. -// -// **The role name and password are the mesh's, not the provisioner's (ADR 0048).** The mesh derives -// the login and hands it to both ends, and mints the password. postgres creates a role and a -// same-named database under exactly that login — a name the consumer cannot learn is a database it -// cannot reach. -// -// The DDL runs through PostgresClient.query(), which is the module's one pending boundary (see -// client.ts). - -import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner"; -import { emit } from "@novox/mesh-sdk/events"; -import { PostgresClient } from "../client.js"; - -const postgres = PostgresClient.fromEnv(); - -/** Emit a lifecycle event without letting a broker hiccup fail the provisioning itself. */ -async function announce(type: string, body: Record): Promise { - try { - await emit(type, body); - } catch (err) { - console.error(`[provisioner:postgres-database] emit ${type} failed: ${err}`); - } -} - -runProvisioner("postgres-database", { - async create(p: Provision): Promise { - // Database and owning role share the consumer's login, so the consumer owns exactly its own. - const database = p.as; - await postgres.createDatabaseAndRole(database, p.as, p.password); - await announce("database.provisioned", { - consumer: p.consumer ?? "", - database, - user: p.as, - }); - }, - - // **Withdrawn, never dropped** (novox/hq issue 241). The login is locked and the database kept - // under its own name: on 2026-10-04 a misread contributions file withdrew every consumer at once, - // and dropping made that a loss of seven databases. Taking a database out of service is a person's - // act — the postgres_retire_database tool — and even that renames rather than drops. - async remove(p: { as: string }): Promise { - await postgres.lockRole(p.as); - await announce("database.deprovisioned", { database: p.as, kept: "true" }); - }, - // Asked every minute by the harness: whether the backend still holds this consumer exactly as - // the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120). - async holds(p: Provision): Promise { - return postgres.canConnectAs(p.as, p.as, p.password); - }, -}); diff --git a/modules/postgres/test/reader.test.ts b/modules/postgres/test/reader.test.ts deleted file mode 100644 index daeb61e..0000000 --- a/modules/postgres/test/reader.test.ts +++ /dev/null @@ -1,112 +0,0 @@ -// What holds the store's read-only query to being read-only (novox/hq issue 193): a caller's -// statement runs as the reader login and never as the admin, is sent as given with no transaction -// wrapped around it as text, and comes back as rows keyed by their columns. The reader is made once, -// as the admin, with every attribute stated; without its password the statement is refused. -// -// psql is a fake on PATH that records each call's user, options and statement, and answers in CSV -// the way the real one does with -q. That the reader cannot write is the server's to enforce and is -// proven against a real server, not here; this holds the module to asking for it. -// Run against the compiled module (npm test builds first), the way the runtime loads it. - -import { test, before, after } from "node:test"; -import assert from "node:assert/strict"; -import { chmod, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; -import { tmpdir } from "node:os"; -import { join } from "node:path"; - -import { PostgresClient, READER } from "../dist/client.js"; - -let dir: string; -let log: string; -const originalPath = process.env.PATH; - -before(async () => { - dir = await mkdtemp(join(tmpdir(), "postgres-reader-")); - log = join(dir, "calls.jsonl"); - // Records argv, the user it connected as and the options it was given; answers a role lookup - // with no rows and anything else with a two-column result. - await writeFile(join(dir, "psql"), `#!/usr/bin/env node -const fs = require("node:fs"); -const args = process.argv.slice(2); -const at = (flag) => args[args.indexOf(flag) + 1]; -fs.appendFileSync(${JSON.stringify(log)}, JSON.stringify({ - user: at("-U"), database: at("-d"), sql: at("-c"), quiet: args.includes("-q"), - password: process.env.PGPASSWORD, options: process.env.PGOPTIONS ?? "", -}) + "\\n"); -const sql = at("-c"); -if (/FROM pg_roles/.test(sql)) process.stdout.write("?column?\\n"); -else if (/^(CREATE|ALTER|GRANT)/.test(sql)) process.stdout.write(""); -else process.stdout.write("name,n\\nalpha,1\\n\\"b,eta\\",2\\n"); -`); - await chmod(join(dir, "psql"), 0o755); - process.env.PATH = `${dir}:${originalPath}`; -}); - -after(async () => { - process.env.PATH = originalPath; - await rm(dir, { recursive: true, force: true }); -}); - -async function calls(): Promise[]> { - const text = await readFile(log, "utf8").catch(() => ""); - await writeFile(log, ""); - return text.split("\n").filter(Boolean).map((line) => JSON.parse(line)); -} - -const conn = { host: "127.0.0.1", port: 5432, user: "postgres", password: "admin-secret" }; - -test("a caller's statement runs as the reader, as given, read-only, and comes back keyed by its columns", async () => { - const client = new PostgresClient({ ...conn, readerPassword: "reader-secret" }); - const statement = "COMMIT; DROP TABLE everything"; - const result = await client.readOnlyQuery("inventory", statement); - - const made = await calls(); - const asked = made.at(-1)!; - assert.equal(asked.user, READER, "the statement never runs as the admin"); - assert.equal(asked.password, "reader-secret"); - assert.equal(asked.sql, statement, "sent as given: no transaction wrapped around it as text"); - assert.equal(asked.database, "inventory"); - assert.equal(asked.quiet, true, "no command tags, which came back as rows keyed by BEGIN"); - assert.match(String(asked.options), /default_transaction_read_only=on/); - - assert.deepEqual(result.rows, [{ name: "alpha", n: "1" }, { name: "b,eta", n: "2" }]); - assert.equal(result.command, "COMMIT"); -}); - -test("the reader is made as the admin, with every attribute stated, once per process", async () => { - const client = new PostgresClient({ ...conn, readerPassword: "reader-secret" }); - await client.readOnlyQuery("inventory", "SELECT 1"); - await client.readOnlyQuery("inventory", "SELECT 2"); - - const made = await calls(); - const asAdmin = made.filter((c) => c.user === "postgres"); - assert.ok(asAdmin.every((c) => c.password === "admin-secret")); - const ddl = asAdmin.map((c) => String(c.sql)); - const role = ddl.find((s) => s.startsWith(`CREATE ROLE "${READER}"`)); - assert.ok(role, "made when it does not exist"); - for (const attribute of ["LOGIN", "NOSUPERUSER", "NOCREATEDB", "NOCREATEROLE", "NOREPLICATION", "NOBYPASSRLS"]) { - assert.match(role!, new RegExp(`\\b${attribute}\\b`)); - } - assert.ok(ddl.includes(`GRANT pg_read_all_data TO "${READER}"`), "reads everything and is granted nothing else"); - assert.ok(ddl.some((s) => /default_transaction_read_only = on/.test(s))); - assert.equal(ddl.filter((s) => s.startsWith("CREATE ROLE")).length, 1, "made once, not per call"); - assert.equal(made.filter((c) => c.user === READER).length, 2); -}); - -test("without the reader's password the statement is refused, and nothing runs as the admin", async () => { - const client = new PostgresClient(conn); - await assert.rejects(client.readOnlyQuery("inventory", "SELECT 1"), /refused rather than run as the admin/); - assert.deepEqual(await calls(), []); -}); - -test("the reader's password is read from the file the mesh delivers", async () => { - const file = join(dir, "reader.secret"); - await writeFile(file, "from-the-file\n"); - const client = PostgresClient.fromEnv({ - MESH_POSTGRES_HOST: "127.0.0.1", MESH_POSTGRES_PASSWORD: "admin-secret", - MESH_POSTGRES_READER_PASSWORD_FILE: file, - }); - await client.readOnlyQuery("inventory", "SELECT 1"); - const asked = (await calls()).at(-1)!; - assert.equal(asked.password, "from-the-file"); -}); diff --git a/modules/postgres/test/withdraw.test.ts b/modules/postgres/test/withdraw.test.ts deleted file mode 100644 index bc97cf3..0000000 --- a/modules/postgres/test/withdraw.test.ts +++ /dev/null @@ -1,69 +0,0 @@ -// A withdrawn consumer keeps its database, and retiring one renames it — nothing here ever drops -// (novox/hq issue 241). psql is a fake on PATH that records every statement and answers the lookups -// these acts make; that the server keeps the data is proven against a real server, not here. -// Run against the compiled module (npm test builds first), the way the runtime loads it. - -import { test, before, after } from "node:test"; -import assert from "node:assert/strict"; -import { chmod, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; -import { tmpdir } from "node:os"; -import { join } from "node:path"; - -import { PostgresClient, retiredName } from "../dist/client.js"; - -let dir: string; -let log: string; -const originalPath = process.env.PATH; - -before(async () => { - dir = await mkdtemp(join(tmpdir(), "postgres-withdraw-")); - log = join(dir, "calls.jsonl"); - await writeFile(join(dir, "psql"), `#!/usr/bin/env node -const fs = require("node:fs"); -const args = process.argv.slice(2); -const sql = args[args.indexOf("-c") + 1]; -fs.appendFileSync(${JSON.stringify(log)}, JSON.stringify({ sql }) + "\\n"); -if (/FROM pg_roles/.test(sql)) process.stdout.write("?column?\\n1\\n"); -else if (/pg_get_userbyid/.test(sql)) process.stdout.write("owner\\nmesh_anchor_mail\\n"); -else if (/FROM pg_database WHERE datname = '.*_deleted_/.test(sql)) process.stdout.write("?column?\\n"); -else process.stdout.write(""); -`); - await chmod(join(dir, "psql"), 0o755); - process.env.PATH = `${dir}:${originalPath}`; -}); - -after(async () => { - process.env.PATH = originalPath; - await rm(dir, { recursive: true, force: true }); -}); - -const statements = async (): Promise => - (await readFile(log, "utf8")).trim().split("\n").map((l) => (JSON.parse(l) as { sql: string }).sql); - -const client = (): PostgresClient => - new PostgresClient({ host: "127.0.0.1", port: 5432, user: "postgres", password: "admin-secret" }); - -test("withdrawing a consumer locks its login and keeps its database", async () => { - await writeFile(log, ""); - await client().lockRole("mesh_anchor_mail"); - const sql = await statements(); - assert.ok(sql.some((s) => /ALTER ROLE "mesh_anchor_mail" NOLOGIN/.test(s)), sql.join("\n")); - assert.ok(!sql.some((s) => /\bDROP\b/i.test(s)), `a withdrawal dropped something:\n${sql.join("\n")}`); -}); - -test("retiring a database renames it aside and locks its owner, and drops nothing", async () => { - await writeFile(log, ""); - const now = new Date("2026-10-05T09:00:00Z"); - const aside = await client().retireDatabase("mesh_anchor_mail", now); - assert.equal(aside, "mesh_anchor_mail_deleted_20261005"); - const sql = await statements(); - assert.ok(sql.some((s) => /ALTER DATABASE "mesh_anchor_mail" RENAME TO "mesh_anchor_mail_deleted_20261005"/.test(s)), sql.join("\n")); - assert.ok(sql.some((s) => /ALTER ROLE "mesh_anchor_mail" NOLOGIN/.test(s))); - assert.ok(!sql.some((s) => /\bDROP\b/i.test(s)), `retiring dropped something:\n${sql.join("\n")}`); -}); - -test("a retired name fits postgres's 63 bytes", () => { - const long = "x".repeat(70); - const name = retiredName(long, new Date("2026-10-05T00:00:00Z")); - assert.ok(name.length <= 63 && name.endsWith("_deleted_20261005"), name); -}); diff --git a/modules/postgres/tools/index.ts b/modules/postgres/tools/index.ts deleted file mode 100644 index 3cdf8fe..0000000 --- a/modules/postgres/tools/index.ts +++ /dev/null @@ -1,100 +0,0 @@ -// postgres's tools — postgres's own code (novox/hq ADR 0039), importing postgres's own client. They -// return structured data; the mesh serves them through the sdk's tool harness. Both call through -// PostgresClient.query(), the module's one pending execution boundary (see client.ts): the tool -// shapes are fixed and correct, and surface the TODO honestly until that boundary is backed. - -import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; -import { PostgresClient } from "../client.js"; - -export function getPostgresTools(postgres: PostgresClient): ToolDefinition[] { - return [ - { - name: "postgres_list_databases", - description: "List the databases on the postgres server, with their on-disk size.", - input: {}, - run: async () => ({ databases: await postgres.listDatabases() }), - }, - { - name: "postgres_query", - description: "Run a read-only SQL query against a named database, as a login that can read every table and change nothing.", - input: { - database: { type: "string", description: "the database to query" }, - sql: { type: "string", description: "the SELECT (or other read-only) statement" }, - }, - run: async (args) => { - const database = String(args.database ?? ""); - const sql = String(args.sql ?? ""); - if (!database) throw new Error("postgres_query: database is required"); - if (!sql) throw new Error("postgres_query: sql is required"); - const result = await postgres.readOnlyQuery(database, sql); - return { database, command: result.command, rows: result.rows }; - }, - }, - { - name: "postgres_retire_database", - description: - "Take one database out of service on purpose: rename it to _deleted_ and lock its owner's login. Nothing is dropped — the data stays on the server under the new name until a person removes it by hand. Repeat the database's name in confirm.", - input: { - database: { type: "string", description: "the database to retire" }, - confirm: { type: "string", description: "the same name again, to say this is meant" }, - }, - run: async (args) => { - const database = String(args.database ?? ""); - if (!database) throw new Error("postgres_retire_database: database is required"); - if (String(args.confirm ?? "") !== database) { - throw new Error("postgres_retire_database: confirm must repeat the database's name"); - } - return { database, renamedTo: await postgres.retireDatabase(database), dropped: false }; - }, - }, - ]; -} - -// The store seat's verbs (novox/hq ADR 0159, 0160): the role's, not postgres's. Registered under the -// seat's name, so the runtime serves them on the seat's subjects wherever this module holds the -// seat and never lists them as postgres's own; scoped to what the store enables — asking what it -// holds and reading from it — so creating a database is postgres's tool and not the store's. -export function getStoreVerbs(postgres: PostgresClient): ToolDefinition[] { - return [ - { - name: "databases", - description: "Every database the store holds, with its on-disk size.", - input: {}, - run: async () => ({ databases: await postgres.listDatabases() }), - }, - { - name: "query", - description: "One read-only statement against one database the store holds.", - input: { - database: { type: "string", description: "the database to query" }, - sql: { type: "string", description: "the SELECT (or other read-only) statement" }, - }, - run: async (args) => { - const database = String(args.database ?? ""); - const sql = String(args.sql ?? ""); - if (!database) throw new Error("query: database is required"); - if (!sql) throw new Error("query: sql is required"); - const result = await postgres.readOnlyQuery(database, sql); - return { database, command: result.command, rows: result.rows }; - }, - }, - ]; -} - -// The tools exist only when the server can be reached from the environment; without it, postgres -// contributes none rather than failing the whole tool runtime. -registerModuleTools("postgres", (env) => { - try { - return getPostgresTools(PostgresClient.fromEnv(env)); - } catch { - return []; - } -}); - -registerModuleTools("mesh-store", (env) => { - try { - return getStoreVerbs(PostgresClient.fromEnv(env)); - } catch { - return []; - } -}); diff --git a/modules/postgres/tsconfig.json b/modules/postgres/tsconfig.json deleted file mode 100644 index 51f4046..0000000 --- a/modules/postgres/tsconfig.json +++ /dev/null @@ -1,12 +0,0 @@ -{ - "compilerOptions": { - "target": "ES2022", - "module": "NodeNext", - "moduleResolution": "NodeNext", - "strict": true, - "esModuleInterop": true, - "skipLibCheck": true, - "noEmit": true - }, - "include": ["client.ts", "index.ts", "provisioner/index.ts", "tools/index.ts"] -}