Port postgres to Go and install the extensions a consumer asks for
letta crash-loops on 'type "vector" does not exist': pgvector is not a trusted extension, so only the provider's superuser can create it, and the provisioner never did. A contribution may now name extensions; the provider creates each (IF NOT EXISTS, available ones only) in the consumer's database on every pass. Go per the standing rule for a TypeScript module that changes. letta asks for vector.
This commit is contained in:
@@ -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
|
||||
// `<name>_deleted_<date>` 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 `<name>_deleted_<yyyymmdd>`, 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
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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())
|
||||
}
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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.<verb>`, 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 <name>_deleted_<date> 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."),
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user