diff options
Diffstat (limited to 'internal/db')
| -rw-r--r-- | internal/db/db.go | 182 | ||||
| -rw-r--r-- | internal/db/db_test.go | 25 | ||||
| -rw-r--r-- | internal/db/migrations/blog/00001_init.sql | 91 | ||||
| -rw-r--r-- | internal/db/migrations/control/00001_init.sql (renamed from internal/db/migrations/00001_init.sql) | 0 | ||||
| -rw-r--r-- | internal/db/migrations/control/00002_sections.sql (renamed from internal/db/migrations/00002_sections.sql) | 0 | ||||
| -rw-r--r-- | internal/db/migrations/control/00003_layout.sql (renamed from internal/db/migrations/00003_layout.sql) | 0 | ||||
| -rw-r--r-- | internal/db/migrations/control/00004_section_placement.sql (renamed from internal/db/migrations/00004_section_placement.sql) | 0 | ||||
| -rw-r--r-- | internal/db/migrations/control/00005_per_blog_databases.sql | 12 | ||||
| -rw-r--r-- | internal/db/migrations/control/00007_drop_content.sql | 9 | ||||
| -rw-r--r-- | internal/db/split.go | 152 |
10 files changed, 459 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 } diff --git a/internal/db/db_test.go b/internal/db/db_test.go new file mode 100644 index 0000000..b4fa3da --- /dev/null +++ b/internal/db/db_test.go @@ -0,0 +1,25 @@ +package db + +import ( + "strings" + "testing" +) + +func TestDBName(t *testing.T) { + for sub, want := range map[string]string{"alice": "blog_alice", "my-blog": "blog_my_blog", "a-b-c": "blog_a_b_c", "www": "blog_www"} { + if got := DBName(sub); got != want || !dbNameRe.MatchString(got) { + t.Errorf("DBName(%q) = %q, want %q", sub, got, want) + } + } + // the longest subdomain the admin form accepts still fits a Postgres identifier + if n := len(DBName(strings.Repeat("a-", 29))); n > 63 { // 58 chars, the admin form's limit + t.Errorf("database name too long: %d", n) + } +} + +func TestWithDatabase(t *testing.T) { + got := withDatabase("postgres://u:p@db:5432/blogspace?sslmode=disable", "blog_alice") + if want := "postgres://u:p@db:5432/blog_alice?sslmode=disable"; got != want { + t.Errorf("got %q, want %q", got, want) + } +} diff --git a/internal/db/migrations/blog/00001_init.sql b/internal/db/migrations/blog/00001_init.sql new file mode 100644 index 0000000..4649935 --- /dev/null +++ b/internal/db/migrations/blog/00001_init.sql @@ -0,0 +1,91 @@ +-- +goose Up +-- One database per blog: nothing here carries a blog id, the database is the scope. +CREATE TABLE settings ( + id boolean PRIMARY KEY DEFAULT true CHECK (id), -- exactly one row + title text NOT NULL, + tagline text NOT NULL DEFAULT '', + theme jsonb NOT NULL DEFAULT '{}'::jsonb, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now() +); + +CREATE TABLE pages ( + id bigserial PRIMARY KEY, + slug text NOT NULL UNIQUE, + title text NOT NULL, + intro_md text NOT NULL DEFAULT '', + intro_html text NOT NULL DEFAULT '', + nav_order integer NOT NULL DEFAULT 0, + is_home boolean NOT NULL DEFAULT false, + created_at timestamptz NOT NULL DEFAULT now() +); +CREATE UNIQUE INDEX pages_one_home ON pages ((true)) WHERE is_home; + +CREATE TABLE posts ( + id bigserial PRIMARY KEY, + page_id bigint NOT NULL REFERENCES pages(id) ON DELETE CASCADE, + slug text NOT NULL, + title text NOT NULL, + body_md text NOT NULL DEFAULT '', + body_html text NOT NULL DEFAULT '', + published boolean NOT NULL DEFAULT true, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now(), + UNIQUE (page_id, slug) +); +CREATE INDEX posts_page_created ON posts (page_id, created_at DESC); + +CREATE TABLE images ( + id uuid PRIMARY KEY, + filename text NOT NULL, + content_type text NOT NULL, + size integer NOT NULL, + data bytea NOT NULL, + created_at timestamptz NOT NULL DEFAULT now() +); +CREATE INDEX images_created ON images (created_at DESC); + +-- Announcements: blog-wide notices shown on every page and post. +CREATE TABLE sections ( + id bigserial PRIMARY KEY, + title text NOT NULL DEFAULT '', + body_md text NOT NULL DEFAULT '', + body_html text NOT NULL DEFAULT '', + placement text NOT NULL DEFAULT 'main-top' + CHECK (placement IN ('left-top', 'left-bottom', 'main-top', 'main-bottom', 'right-top', 'right-bottom')), + style text NOT NULL DEFAULT 'note' CHECK (style IN ('plain', 'note', 'warning')), + enabled boolean NOT NULL DEFAULT true, + sort_order integer NOT NULL DEFAULT 0, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now() +); +CREATE INDEX sections_order ON sections (sort_order, id); + +-- Layout modules: what each area of a blog (header, columns, footer) shows. +CREATE TABLE modules ( + id bigserial PRIMARY KEY, + area text NOT NULL CHECK (area IN ('header', 'left', 'right', 'above', 'below', 'footer')), + kind text NOT NULL CHECK (kind IN ('title', 'logo', 'menu', 'archive', 'recent', 'html', 'rss', 'text', 'sitemap')), + title text NOT NULL DEFAULT '', + body text NOT NULL DEFAULT '', + count integer NOT NULL DEFAULT 5, + sort_order integer NOT NULL DEFAULT 0, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now() +); +CREATE INDEX modules_area ON modules (area, sort_order, id); + +-- The menu: blog pages and custom links in one ordered list. +CREATE TABLE menu_items ( + id bigserial PRIMARY KEY, + page_id bigint REFERENCES pages(id) ON DELETE CASCADE, + label text NOT NULL DEFAULT '', + url text NOT NULL DEFAULT '', + sort_order integer NOT NULL DEFAULT 0, + CHECK (page_id IS NOT NULL OR url <> '') +); +CREATE UNIQUE INDEX menu_items_page ON menu_items (page_id) WHERE page_id IS NOT NULL; +CREATE INDEX menu_items_order ON menu_items (sort_order, id); + +-- +goose Down +DROP TABLE menu_items, modules, sections, images, posts, pages, settings; diff --git a/internal/db/migrations/00001_init.sql b/internal/db/migrations/control/00001_init.sql index e27dad4..e27dad4 100644 --- a/internal/db/migrations/00001_init.sql +++ b/internal/db/migrations/control/00001_init.sql diff --git a/internal/db/migrations/00002_sections.sql b/internal/db/migrations/control/00002_sections.sql index 388cb17..388cb17 100644 --- a/internal/db/migrations/00002_sections.sql +++ b/internal/db/migrations/control/00002_sections.sql diff --git a/internal/db/migrations/00003_layout.sql b/internal/db/migrations/control/00003_layout.sql index e8fff23..e8fff23 100644 --- a/internal/db/migrations/00003_layout.sql +++ b/internal/db/migrations/control/00003_layout.sql diff --git a/internal/db/migrations/00004_section_placement.sql b/internal/db/migrations/control/00004_section_placement.sql index c36787e..c36787e 100644 --- a/internal/db/migrations/00004_section_placement.sql +++ b/internal/db/migrations/control/00004_section_placement.sql diff --git a/internal/db/migrations/control/00005_per_blog_databases.sql b/internal/db/migrations/control/00005_per_blog_databases.sql new file mode 100644 index 0000000..92d3fa2 --- /dev/null +++ b/internal/db/migrations/control/00005_per_blog_databases.sql @@ -0,0 +1,12 @@ +-- +goose Up +-- Each blog gets its own database; the registry remembers which one. +-- db_name stays nullable until the Go migration 00006 has moved the content. +ALTER TABLE blogs ADD COLUMN db_name text UNIQUE; +-- "blog_" + subdomain must fit in a 63-char Postgres identifier. +ALTER TABLE blogs DROP CONSTRAINT blogs_subdomain_check; +ALTER TABLE blogs ADD CONSTRAINT blogs_subdomain_check CHECK (subdomain ~ '^[a-z0-9](-?[a-z0-9]){0,57}$'); + +-- +goose Down +ALTER TABLE blogs DROP CONSTRAINT blogs_subdomain_check; +ALTER TABLE blogs ADD CONSTRAINT blogs_subdomain_check CHECK (subdomain ~ '^[a-z0-9](-?[a-z0-9]){0,62}$'); +ALTER TABLE blogs DROP COLUMN db_name; diff --git a/internal/db/migrations/control/00007_drop_content.sql b/internal/db/migrations/control/00007_drop_content.sql new file mode 100644 index 0000000..9daa1ac --- /dev/null +++ b/internal/db/migrations/control/00007_drop_content.sql @@ -0,0 +1,9 @@ +-- +goose Up +-- The content now lives in the per-blog databases (see 00006 in split.go); +-- the control database keeps only users and the blog registry. +DROP TABLE menu_items, modules, sections, images, posts, pages; +ALTER TABLE blogs DROP COLUMN title, DROP COLUMN tagline, DROP COLUMN theme, DROP COLUMN updated_at; +ALTER TABLE blogs ALTER COLUMN db_name SET NOT NULL; + +-- +goose Down +-- Not reversible: the content is gone from this database. Restore from a backup instead. diff --git a/internal/db/split.go b/internal/db/split.go new file mode 100644 index 0000000..668d255 --- /dev/null +++ b/internal/db/split.go @@ -0,0 +1,152 @@ +package db + +import ( + "context" + "database/sql" + "errors" + "fmt" + "log" + + "github.com/jackc/pgx/v5" + "github.com/pressly/goose/v3" +) + +// splitMigration is control migration 00006: it moves every blog's content out +// of the control database into a database of its own. It runs in the control +// transaction, so either every blog is moved and marked with its db_name or +// nothing changes; the blog databases it created are then dropped by hand +// (the error says which) and the next start retries. 00007 drops the old +// tables afterwards. +func splitMigration(controlURL string) *goose.Migration { + up := func(ctx context.Context, tx *sql.Tx) error { + return splitBlogs(ctx, tx, controlURL) + } + down := func(ctx context.Context, tx *sql.Tx) error { + return errors.New("the per-blog split cannot be undone; restore the control database from a backup") + } + return goose.NewGoMigration(6, &goose.GoFunc{RunTx: up}, &goose.GoFunc{RunTx: down}) +} + +// splitTable is one content table to copy: the query selects the blog's rows +// in the column order of the new table. +type splitTable struct { + name string + cols []string + query string // $1 = blog id + seq bool // bigserial id to bump after the copy +} + +var splitTables = []splitTable{ + {"pages", []string{"id", "slug", "title", "intro_md", "intro_html", "nav_order", "is_home", "created_at"}, + `SELECT id, slug, title, intro_md, intro_html, nav_order, is_home, created_at FROM pages WHERE blog_id=$1`, true}, + {"posts", []string{"id", "page_id", "slug", "title", "body_md", "body_html", "published", "created_at", "updated_at"}, + `SELECT p.id, p.page_id, p.slug, p.title, p.body_md, p.body_html, p.published, p.created_at, p.updated_at + FROM posts p JOIN pages g ON g.id=p.page_id WHERE g.blog_id=$1`, true}, + {"images", []string{"id", "filename", "content_type", "size", "data", "created_at"}, + `SELECT id, filename, content_type, size, data, created_at FROM images WHERE blog_id=$1`, false}, + {"sections", []string{"id", "title", "body_md", "body_html", "placement", "style", "enabled", "sort_order", "created_at", "updated_at"}, + `SELECT id, title, body_md, body_html, placement, style, enabled, sort_order, created_at, updated_at FROM sections WHERE blog_id=$1`, true}, + {"modules", []string{"id", "area", "kind", "title", "body", "count", "sort_order", "created_at", "updated_at"}, + `SELECT id, area, kind, title, body, count, sort_order, created_at, updated_at FROM modules WHERE blog_id=$1`, true}, + {"menu_items", []string{"id", "page_id", "label", "url", "sort_order"}, + `SELECT id, page_id, label, url, sort_order FROM menu_items WHERE blog_id=$1`, true}, +} + +func splitBlogs(ctx context.Context, tx *sql.Tx, controlURL string) error { + type blog struct { + id int64 + sub, title, tagline string + theme []byte + createdAt, updatedAt any + } + rows, err := tx.QueryContext(ctx, `SELECT id, subdomain, title, tagline, theme, created_at, updated_at FROM blogs WHERE db_name IS NULL ORDER BY id`) + if err != nil { + return err + } + var blogs []blog + for rows.Next() { + var b blog + if err := rows.Scan(&b.id, &b.sub, &b.title, &b.tagline, &b.theme, &b.createdAt, &b.updatedAt); err != nil { + rows.Close() + return err + } + blogs = append(blogs, b) + } + rows.Close() + if err := rows.Err(); err != nil { + return err + } + if len(blogs) == 0 { + return nil + } + + // DDL needs its own connection: CREATE DATABASE refuses to run in a transaction. + admin, err := pgx.Connect(ctx, controlURL) + if err != nil { + return err + } + defer admin.Close(ctx) + var created []string + for _, b := range blogs { + name := DBName(b.sub) + log.Printf("moving blog %q into database %s", b.sub, name) + if err := createDatabase(ctx, admin, name); err != nil { + return fmt.Errorf("%w — a leftover of an earlier failed split must be dropped by hand (created so far: %v)", err, created) + } + created = append(created, name) + if err := MigrateBlog(ctx, withDatabase(controlURL, name)); err != nil { + return fmt.Errorf("migrate %s: %w (drop the databases %v before retrying)", name, err, created) + } + if err := copyBlog(ctx, tx, withDatabase(controlURL, name), b.id, b.title, b.tagline, b.theme, b.createdAt, b.updatedAt); err != nil { + return fmt.Errorf("copy blog %q: %w (drop the databases %v before retrying)", b.sub, err, created) + } + if _, err := tx.ExecContext(ctx, `UPDATE blogs SET db_name=$1 WHERE id=$2`, name, b.id); err != nil { + return err + } + } + return nil +} + +// copyBlog streams one blog's rows from the control transaction into its new database. +func copyBlog(ctx context.Context, tx *sql.Tx, blogURL string, blogID int64, title, tagline string, theme []byte, createdAt, updatedAt any) error { + conn, err := pgx.Connect(ctx, blogURL) + if err != nil { + return err + } + defer conn.Close(ctx) + if _, err := conn.Exec(ctx, `INSERT INTO settings (title, tagline, theme, created_at, updated_at) VALUES ($1,$2,$3,$4,$5)`, + title, tagline, theme, createdAt, updatedAt); err != nil { + return fmt.Errorf("settings: %w", err) + } + for _, t := range splitTables { + rows, err := tx.QueryContext(ctx, t.query, blogID) + if err != nil { + return fmt.Errorf("%s: %w", t.name, err) + } + n := len(t.cols) + src := pgx.CopyFromFunc(func() ([]any, error) { + if !rows.Next() { + return nil, rows.Err() + } + vals := make([]any, n) + ptrs := make([]any, n) + for i := range vals { + ptrs[i] = &vals[i] + } + // database/sql hands back int64/string/bool/[]byte/time.Time, + // which pgx encodes for the matching column types. + return vals, rows.Scan(ptrs...) + }) + _, err = conn.CopyFrom(ctx, pgx.Identifier{t.name}, t.cols, src) + rows.Close() + if err != nil { + return fmt.Errorf("%s: %w", t.name, err) + } + if t.seq { + if _, err := conn.Exec(ctx, `SELECT setval(pg_get_serial_sequence($1,'id'), coalesce(max(id),0)+1, false) FROM `+pgx.Identifier{t.name}.Sanitize(), t.name); err != nil { + return fmt.Errorf("%s sequence: %w", t.name, err) + } + } + } + return nil +} |
