pgxext

package module
v0.2.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 22, 2026 License: MIT Imports: 11 Imported by: 0

README

 ███████████    █████████  █████ █████    ██████████ █████ █████ ███████████
░░███░░░░░███  ███░░░░░███░░███ ░░███    ░░███░░░░░█░░███ ░░███ ░█░░░███░░░█
 ░███    ░███ ███     ░░░  ░░███ ███      ░███  █ ░  ░░███ ███  ░   ░███  ░
 ░██████████ ░███           ░░█████       ░██████     ░░█████       ░███
 ░███░░░░░░  ░███    █████   ███░███      ░███░░█      ███░███      ░███
 ░███        ░░███  ░░███   ███ ░░███     ░███ ░   █  ███ ░░███     ░███
 █████        ░░█████████  █████ █████    ██████████ █████ █████    █████
░░░░░          ░░░░░░░░░  ░░░░░ ░░░░░    ░░░░░░░░░░ ░░░░░ ░░░░░    ░░░░░
    

A thin pgx/v5 connection-pool wrapper with migrations and a generic query builder.


Overview

pgxext wraps pgxpool to provide:

  • Config — fluent builder for connection settings, TLS, pool parameters, and runtime params.
  • DataSource — a *pgxpool.Pool wrapper exposing Query, Exec, QueryRow, transactions, and batch operations.
  • Migration — transactional SQL migration runner backed by a migrations table.
  • Repository — generic, type-safe query builder (SELECT / INSERT / UPDATE / DELETE) with JOIN and COUNT support.
  • Notification — trigger builder and LISTEN consumer for JSON PostgreSQL notifications.
  • Functional — builders for PostgreSQL views, materialized views, and functions.
  • TimescaleDB — safe hypertable, frame/gapfill, continuous-aggregate, policy, and metadata primitives.

Quick start

Config & DataSource

ds, err := pgxext.NewDataSource(ctx,
    pgxext.NewConfig().
        WithHost("localhost").
        WithPort(5432).
        WithDatabase("mydb").
        WithUser("alice").
        WithPassword("secret").
        WithMaxConns(10),
)

Or from a URL:

cfg, err := pgxext.NewConfig().WithURL("postgres://alice:secret@localhost/mydb?pool_max_conns=10")
ds, err := pgxext.NewDataSource(ctx, cfg)

Migrations

var schema = migration.MigrationSet{
    {
        Name:      "001_create_users",
        UpQuery:   `CREATE TABLE users (id SERIAL PRIMARY KEY, name TEXT NOT NULL)`,
        DownQuery: `DROP TABLE users`,
    },
}

m := migration.NewMigrator(ctx, ds)
m.Up(schema)   // apply
m.Down(schema) // revert

Repository

Define a model with db struct tags (pgx RowToAddrOfStructByName convention):

type User struct {
    ID    int    `db:"id"`
    Name  string `db:"name"`
    Email string `db:"email"`
}

Create a repository once and reuse it:

repo := repository.NewRepository[User](ds, "users")

Select

users, err := repo.Select().
    Where("name", repository.Like, "%alice%").
    OrderBy("name", repository.ASC).
    Limit(20).
    Execute(ctx)

Count

n, err := repo.Select().
    Where("email", repository.Like, "%@example.com").
    Count(ctx)

Insert

n, err := repo.Insert().
    Set("name", "Alice").
    Set("email", "alice@example.com").
    Execute(ctx)

Update

n, err := repo.Update().
    Set("email", "new@example.com").
    Where("id", repository.Equals, 42).
    Execute(ctx)

Delete

n, err := repo.Delete().
    Where("id", repository.Equals, 42).
    Execute(ctx)

Join

type UserOrder struct {
    UserID  int    `db:"id"`
    Name    string `db:"name"`
    OrderID int    `db:"order_id"`
}

repo := repository.NewRepository[UserOrder](ds, "users")
rows, err := repo.Select().
    Join("orders", "users.id", repository.Equals, "orders.user_id").
    Where("orders.status", repository.Equals, "paid").
    Execute(ctx)

Notifications

Create a trigger that sends JSON payloads on selected row operations:

notificationSQL := notification.NewNotification("users", "user.changed").
    On(notification.Insert, notification.Update, notification.Delete).
    WithPayloadProperties(
        notification.FromState,
        notification.ToState,
        notification.TableName,
        notification.RowID,
        notification.CreatedAt,
    )

err := notificationSQL.Apply(ctx, ds)

The generated payload uses these top-level properties:

{
  "fromState": {},
  "toState": {},
  "tableName": "users",
  "rowId": 42,
  "createdAt": "2026-05-25T12:00:00Z"
}

Use EmptyPayload() or omit WithPayloadProperties to emit {}.

Consume notifications with a dedicated LISTEN connection:

consumer := notification.NewConsumer(ds, "user.changed")
err := consumer.Listen(ctx, func(ctx context.Context, event notification.Event) error {
    fmt.Println(event.Payload.TableName, event.Payload.RowID)
    return nil
})

Database objects

Build and apply database-level objects directly or embed the generated SQL in migrations:

view := functional.NewView("public.active_users").
    OrReplace().
    As(`SELECT id, email FROM users WHERE active`)

err := view.Apply(ctx, ds)
matView := functional.NewMaterializedView("public.user_totals").
    IfNotExists().
    As(`SELECT user_id, count(*) AS total FROM orders GROUP BY user_id`).
    WithNoData()

err := matView.Apply(ctx, ds)
err = matView.Refresh(ctx, ds, true, true)
fn := functional.NewFunction("public.normalize_email").
    OrReplace().
    WithArguments("value text").
    Returns("text").
    Language("sql").
    WithVolatility(functional.Immutable).
    Strict().
    Body(`SELECT lower(trim(value))`)

err := fn.Apply(ctx, ds)

TimescaleDB

The additive timescaledb package targets TimescaleDB 2.20+ on PostgreSQL 15–17, and TimescaleDB 2.23+ on PostgreSQL 18. It uses pgx/v5, supports TIMESTAMPTZ time dimensions, and takes *pgxext.DataSource directly for runtime operations. Run every package against the PostgreSQL 16/17/18 and compatible TimescaleDB matrix with just test-integration (or the shorter just test alias). Each matrix row writes a package coverage profile under .coverage/.

source := timescaledb.Relation{Schema: "metrics", Name: "observations"}

hypertable := timescaledb.NewHypertable(source, "observed_at").
    ChunkInterval(timescaledb.FixedInterval(24 * time.Hour)).
    IfNotExists()

if err := hypertable.Apply(ctx, ds); err != nil {
    return err
}

frames, err := timescaledb.NewFrameQuery[Frame](source, "observed_at").
    Bucket(timescaledb.FixedInterval(15 * time.Second)).
    Dimension("series_id").
    Measure(timescaledb.Avg("value").As("average")).
    Between(from, to).
    Where(timescaledb.Equal("series_id", seriesID)).
    Execute(ctx, ds)

Extension installation is explicit: use timescaledb.CreateExtensionSQL() in an application migration or call timescaledb.CreateExtension from an authorized bootstrap. The package never installs or drops the extension during initialization.

Continuous aggregates, real-time aggregation, gapfill, retention, and columnstore depend on the installed Timescale edition and version. See the complete package guide for compatibility, safe DDL, gapfill modes, policy planning, manual refresh, lifecycle warnings, transactions, metadata, integration tests, and non-goals.

Utilities

// Scan one row into *T (nil on no rows)
user, err := pgxext.CollectOneRow[User](rows)

// Scan all rows into []*T
users, err := pgxext.CollectRows[User](rows)

// Inspect a Postgres error code / constraint
if pgErr, ok := pgxext.IsPostgresError(err); ok {
    fmt.Println(pgErr.Code, pgErr.ConstraintName)
}

Documentation

Overview

Package pgxext provides pgx/v5 helpers.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func CollectOneRow

func CollectOneRow[T any](rows pgx.Rows) (*T, error)

CollectOneRow scans one row into *T.

func CollectRows

func CollectRows[T any](rows pgx.Rows) ([]*T, error)

CollectRows scans rows into []*T.

func IsPostgresError

func IsPostgresError(err error) (*pgconn.PgError, bool)

IsPostgresError unwraps a PostgreSQL server error.

Types

type Config

type Config struct {
	*pgxpool.Config
}

Config wraps pgxpool.Config.

func NewConfig

func NewConfig() *Config

NewConfig creates a default Config.

func (*Config) WithAfterConnect

func (config *Config) WithAfterConnect(fn func(context.Context, *pgx.Conn) error) *Config

WithAfterConnect sets AfterConnect.

func (*Config) WithAfterRelease

func (config *Config) WithAfterRelease(fn func(*pgx.Conn) bool) *Config

WithAfterRelease sets AfterRelease.

func (*Config) WithApplicationName

func (config *Config) WithApplicationName(name string) *Config

WithApplicationName sets application_name.

func (*Config) WithBeforeAcquire deprecated

func (config *Config) WithBeforeAcquire(fn func(context.Context, *pgx.Conn) bool) *Config

WithBeforeAcquire sets BeforeAcquire.

Deprecated: prefer WithPrepareConn.

func (*Config) WithBeforeClose

func (config *Config) WithBeforeClose(fn func(*pgx.Conn)) *Config

WithBeforeClose sets BeforeClose.

func (*Config) WithBeforeConnect

func (config *Config) WithBeforeConnect(fn func(context.Context, *pgx.ConnConfig) error) *Config

WithBeforeConnect sets BeforeConnect.

func (*Config) WithConnectTimeout

func (config *Config) WithConnectTimeout(timeout time.Duration) *Config

WithConnectTimeout sets the connect timeout.

func (*Config) WithDatabase

func (config *Config) WithDatabase(database string) *Config

WithDatabase sets the database.

func (*Config) WithDialFunc

func (config *Config) WithDialFunc(fn pgconn.DialFunc) *Config

WithDialFunc sets DialFunc.

func (*Config) WithHealthCheckPeriod

func (config *Config) WithHealthCheckPeriod(d time.Duration) *Config

WithHealthCheckPeriod sets HealthCheckPeriod.

func (*Config) WithHost

func (config *Config) WithHost(host string) *Config

WithHost sets the host.

func (*Config) WithLookupFunc

func (config *Config) WithLookupFunc(fn pgconn.LookupFunc) *Config

WithLookupFunc sets LookupFunc.

func (*Config) WithMaxConnIdleTime

func (config *Config) WithMaxConnIdleTime(d time.Duration) *Config

WithMaxConnIdleTime sets MaxConnIdleTime.

func (*Config) WithMaxConnLifetime

func (config *Config) WithMaxConnLifetime(d time.Duration) *Config

WithMaxConnLifetime sets MaxConnLifetime.

func (*Config) WithMaxConns

func (config *Config) WithMaxConns(n int32) *Config

WithMaxConns sets MaxConns.

func (*Config) WithMinConns

func (config *Config) WithMinConns(n int32) *Config

WithMinConns sets MinConns.

func (*Config) WithOnNotice

func (config *Config) WithOnNotice(fn pgconn.NoticeHandler) *Config

WithOnNotice sets OnNotice.

func (*Config) WithOnNotification

func (config *Config) WithOnNotification(fn pgconn.NotificationHandler) *Config

WithOnNotification sets OnNotification.

func (*Config) WithOnPgError

func (config *Config) WithOnPgError(fn pgconn.PgErrorHandler) *Config

WithOnPgError sets OnPgError.

func (*Config) WithPassword

func (config *Config) WithPassword(password string) *Config

WithPassword sets the password.

func (*Config) WithPort

func (config *Config) WithPort(port uint16) *Config

WithPort sets the port.

func (*Config) WithPrepareConn

func (config *Config) WithPrepareConn(fn func(context.Context, *pgx.Conn) (bool, error)) *Config

WithPrepareConn sets PrepareConn.

func (*Config) WithSSLClientCert

func (config *Config) WithSSLClientCert(certFile, keyFile, password string) (*Config, error)

WithSSLClientCert loads a client certificate.

func (*Config) WithSSLMode

func (config *Config) WithSSLMode(mode string) *Config

WithSSLMode sets the PostgreSQL sslmode.

func (*Config) WithSSLRootCert

func (config *Config) WithSSLRootCert(rootCertFile string) (*Config, error)

WithSSLRootCert loads a root certificate.

func (*Config) WithSearchPath

func (config *Config) WithSearchPath(path string) *Config

WithSearchPath sets search_path.

func (*Config) WithShouldPing

func (config *Config) WithShouldPing(fn func(context.Context, pgxpool.ShouldPingParams) bool) *Config

WithShouldPing sets ShouldPing.

func (*Config) WithTimezone

func (config *Config) WithTimezone(tz string) *Config

WithTimezone sets TimeZone.

func (*Config) WithURL

func (config *Config) WithURL(rawURL string) (*Config, error)

WithURL parses a PostgreSQL URL.

func (*Config) WithUser

func (config *Config) WithUser(user string) *Config

WithUser sets the user.

func (*Config) WithValidateConnect

func (config *Config) WithValidateConnect(fn pgconn.ValidateConnectFunc) *Config

WithValidateConnect sets ValidateConnect.

type DataSource

type DataSource struct {
	*pgxpool.Pool
}

DataSource wraps a pgxpool.Pool.

func NewDataSource

func NewDataSource(ctx context.Context, config *Config) (ds *DataSource, err error)

NewDataSource opens a connection pool.

func (*DataSource) Exec

func (ds *DataSource) Exec(ctx context.Context, sql string, args ...interface{}) (pgconn.CommandTag, error)

Exec runs a SQL statement.

func (*DataSource) NewBatch

func (ds *DataSource) NewBatch() *pgx.Batch

NewBatch returns an empty batch.

func (*DataSource) NewCustomTransaction

func (ds *DataSource) NewCustomTransaction(ctx context.Context, options pgx.TxOptions) (pgx.Tx, error)

NewCustomTransaction starts a transaction with options.

func (*DataSource) NewTransaction

func (ds *DataSource) NewTransaction(ctx context.Context) (pgx.Tx, error)

NewTransaction starts a transaction.

func (*DataSource) Query

func (ds *DataSource) Query(ctx context.Context, sql string, args ...any) (pgx.Rows, error)

Query runs a SQL query.

func (*DataSource) QueryRow

func (ds *DataSource) QueryRow(ctx context.Context, sql string, args ...any) (pgx.Row, error)

QueryRow runs a SQL query and returns one row.

func (*DataSource) SendBatch

func (ds *DataSource) SendBatch(ctx context.Context, batch *pgx.Batch) (pgx.BatchResults, error)

SendBatch sends a batch.

Directories

Path Synopsis
Package database provides PostgreSQL database-level object builders.
Package database provides PostgreSQL database-level object builders.
Package migration provides SQL migrations.
Package migration provides SQL migrations.
Package notification provides PostgreSQL LISTEN/NOTIFY helpers.
Package notification provides PostgreSQL LISTEN/NOTIFY helpers.
Package repository provides a generic PostgreSQL query builder.
Package repository provides a generic PostgreSQL query builder.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL