diff options
Diffstat (limited to 'internal/db/split.go')
| -rw-r--r-- | internal/db/split.go | 152 |
1 files changed, 0 insertions, 152 deletions
diff --git a/internal/db/split.go b/internal/db/split.go deleted file mode 100644 index 668d255..0000000 --- a/internal/db/split.go +++ /dev/null @@ -1,152 +0,0 @@ -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 -} |
