aboutsummaryrefslogtreecommitdiffstats
path: root/internal/db/db.go
diff options
context:
space:
mode:
Diffstat (limited to 'internal/db/db.go')
-rw-r--r--internal/db/db.go182
1 files changed, 170 insertions, 12 deletions
diff --git a/internal/db/db.go b/internal/db/db.go
index 44d489b..601680d 100644
--- a/internal/db/db.go
+++ b/internal/db/db.go
@@ -1,22 +1,45 @@
-// Package db opens the Postgres pool and applies embedded migrations.
+// 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/*.sql
+//go:embed migrations/control/*.sql migrations/blog/*.sql
var migrations embed.FS
-func Open(ctx context.Context, url string) (*pgxpool.Pool, error) {
- pool, err := pgxpool.New(ctx, url)
+// 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)
}
@@ -24,20 +47,155 @@ func Open(ctx context.Context, url string) (*pgxpool.Pool, error) {
pool.Close()
return nil, fmt.Errorf("ping: %w", err)
}
- return pool, nil
+ return &Cluster{controlURL: controlURL, control: pool, blogs: map[string]*pgxpool.Pool{}}, nil
}
-// Migrate applies all pending migrations using goose over database/sql.
-func Migrate(ctx context.Context, url string) error {
- sqldb, err := sql.Open("pgx", url)
+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
+}
+
+// execer is the pool or connection createDatabase runs on.
+type execer interface {
+ Exec(ctx context.Context, sql string, args ...any) (pgconn.CommandTag, error)
+}
+
+func createDatabase(ctx context.Context, db execer, 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, including the Go
+// migration that moves content out into the blog databases.
+func MigrateControl(ctx context.Context, controlURL string) error {
+ return migrate(ctx, controlURL, "migrations/control", goose.WithGoMigrations(splitMigration(controlURL)))
+}
+
+// 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, and the Go migration, apart.
+func migrate(ctx context.Context, dsn, dir string, opts ...goose.ProviderOption) error {
+ sqldb, err := sql.Open("pgx", dsn)
if err != nil {
return err
}
defer sqldb.Close()
- goose.SetBaseFS(migrations)
- goose.SetLogger(goose.NopLogger())
- if err := goose.SetDialect("postgres"); err != nil {
+ fsys, err := fs.Sub(migrations, dir)
+ if err != nil {
+ return err
+ }
+ p, err := goose.NewProvider(goose.DialectPostgres, sqldb, fsys, opts...)
+ if err != nil {
return err
}
- return goose.UpContext(ctx, sqldb, "migrations")
+ _, err = p.Up(ctx)
+ return err
}