// Package db opens the Postgres pools and applies the embedded migrations. // // There are two kinds of database: the control database (users and the blog // registry, named by DATABASE_URL) and one database per blog on the same // server, created by the app. A Cluster hands out pools for both. package db import ( "context" "database/sql" "embed" "errors" "fmt" "io/fs" "net/url" "regexp" "strings" "sync" "time" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgconn" "github.com/jackc/pgx/v5/pgxpool" _ "github.com/jackc/pgx/v5/stdlib" "github.com/pressly/goose/v3" ) //go:embed migrations/control/*.sql migrations/blog/*.sql var migrations embed.FS // Cluster is the connection to one Postgres server: the control pool plus a // lazily opened, cached pool per blog database. type Cluster struct { controlURL string control *pgxpool.Pool mu sync.Mutex blogs map[string]*pgxpool.Pool } // Open connects to the control database. func Open(ctx context.Context, controlURL string) (*Cluster, error) { pool, err := pgxpool.New(ctx, controlURL) if err != nil { return nil, fmt.Errorf("connect: %w", err) } if err := pool.Ping(ctx); err != nil { pool.Close() return nil, fmt.Errorf("ping: %w", err) } return &Cluster{controlURL: controlURL, control: pool, blogs: map[string]*pgxpool.Pool{}}, nil } func (c *Cluster) Control() *pgxpool.Pool { return c.control } // Blog returns the pool for a blog database, opening it on first use. Blog // pools are small and drop idle connections quickly, so a hundred quiet blogs // cost nothing; only the busy ones hold connections. func (c *Cluster) Blog(ctx context.Context, dbName string) (*pgxpool.Pool, error) { c.mu.Lock() defer c.mu.Unlock() if p, ok := c.blogs[dbName]; ok { return p, nil } cfg, err := pgxpool.ParseConfig(c.BlogURL(dbName)) if err != nil { return nil, err } cfg.MaxConns = 4 cfg.MinConns = 0 cfg.MaxConnIdleTime = 2 * time.Minute p, err := pgxpool.NewWithConfig(ctx, cfg) if err != nil { return nil, err } c.blogs[dbName] = p return p, nil } // forget closes and drops the cached pool for a blog database, if any. func (c *Cluster) forget(dbName string) { c.mu.Lock() p, ok := c.blogs[dbName] delete(c.blogs, dbName) c.mu.Unlock() if ok { p.Close() } } // Close closes every pool. func (c *Cluster) Close() { c.mu.Lock() defer c.mu.Unlock() for name, p := range c.blogs { p.Close() delete(c.blogs, name) } c.control.Close() } // BlogURL is the control URL pointed at another database on the same server. func (c *Cluster) BlogURL(dbName string) string { return withDatabase(c.controlURL, dbName) } func withDatabase(dsn, dbName string) string { u, err := url.Parse(dsn) if err != nil { return dsn } u.Path = "/" + dbName return u.String() } var dbNameRe = regexp.MustCompile(`^[a-z0-9_]{1,63}$`) // DBName is the database a blog lives in: "blog_" + subdomain, dashes as // underscores so the name needs no quoting in psql or pg_dump. Subdomains only // allow [a-z0-9-], so the mapping is one-to-one. func DBName(subdomain string) string { return "blog_" + strings.ReplaceAll(subdomain, "-", "_") } // ErrDatabaseExists is returned by CreateBlogDB when the name is taken — a // leftover from a deleted blog or a failed attempt that must be dropped by hand. var ErrDatabaseExists = errors.New("database already exists") // CreateBlogDB creates an empty blog database and applies the blog migrations. // CREATE DATABASE cannot run inside a transaction, so this always uses its own // connection; callers doing registry work in a transaction must order it so a // failure here rolls the transaction back. func (c *Cluster) CreateBlogDB(ctx context.Context, dbName string) error { if !dbNameRe.MatchString(dbName) { return fmt.Errorf("invalid database name %q", dbName) } if err := createDatabase(ctx, c.control, dbName); err != nil { return err } if err := MigrateBlog(ctx, c.BlogURL(dbName)); err != nil { c.control.Exec(ctx, `DROP DATABASE IF EXISTS `+pgx.Identifier{dbName}.Sanitize()+` WITH (FORCE)`) return fmt.Errorf("migrate %s: %w", dbName, err) } return nil } func createDatabase(ctx context.Context, db *pgxpool.Pool, dbName string) error { _, err := db.Exec(ctx, `CREATE DATABASE `+pgx.Identifier{dbName}.Sanitize()) var pgErr *pgconn.PgError if errors.As(err, &pgErr) && pgErr.Code == "42P04" { // duplicate_database return fmt.Errorf("%w: %s", ErrDatabaseExists, dbName) } if err != nil { return fmt.Errorf("create database %s: %w", dbName, err) } return nil } // DropBlogDB closes the blog's pool and drops its database, kicking any other // session still connected to it. func (c *Cluster) DropBlogDB(ctx context.Context, dbName string) error { if !dbNameRe.MatchString(dbName) { return fmt.Errorf("invalid database name %q", dbName) } c.forget(dbName) _, err := c.control.Exec(ctx, `DROP DATABASE IF EXISTS `+pgx.Identifier{dbName}.Sanitize()+` WITH (FORCE)`) return err } // ---- migrations ------------------------------------------------------------ // MigrateControl applies the control database migrations. func MigrateControl(ctx context.Context, controlURL string) error { return migrate(ctx, controlURL, "migrations/control") } // MigrateBlog applies the blog schema migrations to one blog database. func MigrateBlog(ctx context.Context, blogURL string) error { return migrate(ctx, blogURL, "migrations/blog") } // migrate runs goose over database/sql. A Provider (rather than the package // globals) keeps the two migration sets apart. func migrate(ctx context.Context, dsn, dir string) error { sqldb, err := sql.Open("pgx", dsn) if err != nil { return err } defer sqldb.Close() fsys, err := fs.Sub(migrations, dir) if err != nil { return err } p, err := goose.NewProvider(goose.DialectPostgres, sqldb, fsys) if err != nil { return err } _, err = p.Up(ctx) return err }