diff options
Diffstat (limited to 'internal/db/db.go')
| -rw-r--r-- | internal/db/db.go | 182 |
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 } |
