aboutsummaryrefslogtreecommitdiffstats
path: root/internal/db
diff options
context:
space:
mode:
Diffstat (limited to 'internal/db')
-rw-r--r--internal/db/db.go182
-rw-r--r--internal/db/db_test.go25
-rw-r--r--internal/db/migrations/blog/00001_init.sql91
-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.sql12
-rw-r--r--internal/db/migrations/control/00007_drop_content.sql9
-rw-r--r--internal/db/split.go152
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
+}