dbconn

package
v0.3.3 Latest Latest
Warning

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

Go to latest
Published: Sep 16, 2026 License: Apache-2.0 Imports: 18 Imported by: 1

Documentation

Overview

Package dbconn is the engine's database connectivity layer: pgx pool construction with safe session defaults (lock_timeout, statement_timeout), RDS/Aurora TLS, bounded retries for transient errors, and a helper to terminate backends blocking a session's lock acquisition.

Index

Constants

View Source
const (
	DefaultLockTimeout      = 3 * time.Second
	DefaultStatementTimeout = 30 * time.Second
	// DefaultConnectTimeout bounds each dial attempt.
	DefaultConnectTimeout = 10 * time.Second
)

Defaults for the session timeouts every pooled connection runs under. Every statement the engine issues is bounded: lock_timeout keeps the engine from sitting at the head of the lock queue (the lock-queue pile-up), and statement_timeout bounds runaway work.

Variables

View Source
var ErrNoSessionAffinity = errors.New("the connection does not keep one PostgreSQL session")

ErrNoSessionAffinity reports a connection that does not keep one server session, so nothing set on it — least of all an execution bound — can be relied on to still be there for the next statement.

View Source
var ErrSessionSettingsDiscarded = errors.New("the connection does not keep pg-sprite's session settings")

ErrSessionSettingsDiscarded reports a connection that did not keep the session settings the pool asked for. The engine's execution bounds are session settings, so a connection that drops them removes the bounds without removing the work they were bounding.

Functions

func IsRDSHost

func IsRDSHost(host string) bool

IsRDSHost reports whether host is an Amazon RDS/Aurora endpoint.

func LocalSearchPath added in v0.3.2

func LocalSearchPath(schemas ...string) string

LocalSearchPath returns the SET LOCAL statement that scopes the current transaction's search_path to schemas, in order. SET cannot take bind parameters, so each schema is quoted as an identifier, and the path then receives the same rewrite as every pooled session's: a pg_catalog entry that another schema precedes is dropped. A transaction-local path replaces the session path for the transaction, so it must uphold the same guarantee (INV: CO-9); every transaction-local search_path pg-sprite sets is built here, and TestLocalSearchPathIsTheOnlySearchPathWriter keeps it that way.

func NewPool

func NewPool(ctx context.Context, cfg Config) (*pgxpool.Pool, error)

NewPool builds a pgx pool from cfg, applies the session defaults, and verifies connectivity with a ping before returning. Each physical connection has any pg_catalog entry that another schema precedes removed from its search_path (see unshadowCatalog), so a schema listed ahead of an explicit pg_catalog no longer shadows the catalog while every other entry stays as configured.

It proves session affinity before returning, which needs a second connection to the same server for the length of the proof. The server — or the pooler in front of it — must have one connection to spare beyond this pool's own, or NewPool fails rather than opening a pool whose bounds were never proven to hold.

func ProveSessionAffinity added in v0.3.3

func ProveSessionAffinity(ctx context.Context, conn *pgxpool.Conn, pinner *pgxpool.Pool, bound time.Duration) error

ProveSessionAffinity proves on conn the property every session-scoped advisory lock rests on: that this connection keeps one server session, so a lock taken on it is a lock other connections cannot take. It returns an error wrapping ErrNoSessionAffinity when it proves the property does not hold, so a caller fails closed instead of running without the exclusion it believes it has. pinner supplies a second connection to the same server and is not otherwise disturbed.

The proof takes a lock on conn and then makes a second connection hold a transaction open across the check. A transaction-mode pooler pins a backend for a transaction's duration, so the second connection takes the backend conn's single-statement acquire just released, and conn's next statement lands somewhere else. That turns a rebind that would otherwise depend on load into one the proof can observe, and it gives three independent readings of the same failure: the second connection takes a lock conn holds, conn no longer appears in pg_locks as the session holding its own lock, and conn cannot release what it took.

Against a direct connection, and against a pooler that hands out a session per client connection, all three readings are the healthy one. There is no false positive to trade off, because each reading is a fact about the connection in hand rather than a guess about what sits behind it.

bound caps the whole proof; a non-positive bound leaves it on the caller's context alone.

It is one-sided in the other direction: an idle transaction-mode pooler with spare backends can answer every reading the healthy way, so a clean proof is evidence and not certainty. It is a guard against the configuration an operator lands on by following a hosted platform's default connection string, not a substitute for pointing the engine at a direct endpoint.

The probe key is freshly random on every call, so concurrent proofs never contend and a probe lock stranded on an unreachable backend can never block anything later.

INV: LK-2 — every strong lock acquisition is bounded. The bounds are session settings, so they only hold where the session does. An advisory lock is the instrument here rather than the subject: it is the one piece of session state whose loss a client can observe directly, which makes it the way to prove the session is stable enough to carry the timeouts.

func Retry

func Retry(ctx context.Context, attempts int, backoff time.Duration, fn func(context.Context) error) error

Retry runs fn up to attempts times with linear backoff, retrying only errors Retryable classifies as transient. Non-transient errors return immediately; context cancellation always wins.

func Retryable

func Retryable(err error) bool

Retryable reports whether err is transient: a bounded lock wait that timed out, a deadlock or serialization failure, or a connection-level error.

func ServerMajor

func ServerMajor(ctx context.Context, pool *pgxpool.Pool) (int, error)

ServerMajor reads the connected server's numeric major version.

func ServerVersion

func ServerVersion(ctx context.Context, pool *pgxpool.Pool) (string, error)

ServerVersion reads the connected server's server_version setting. Plan reports carry it because classification is version-sensitive: a stored report names the server whose rules produced it.

func TerminateBlockers

func TerminateBlockers(ctx context.Context, q Querier, pid int) ([]int, error)

TerminateBlockers terminates every backend currently blocking pid's lock acquisition (per pg_blocking_pids) and returns the pids it terminated.

This is the bounded-cutover escape hatch: it targets only the backends standing in front of a specific waiting session (e.g. the cutover swap), never a broad sweep. Callers decide whether evicting those backends is acceptable; this function only does the targeted termination.

Types

type AdvisoryLockHolder added in v0.3.3

type AdvisoryLockHolder interface {
	QueryRow(ctx context.Context, sql string, args ...any) pgx.Row
}

AdvisoryLockHolder is the query surface a lock confirmation needs; *pgxpool.Conn and *pgx.Conn both satisfy it. The confirmation must run on the session under test, so this is deliberately not the pool.

type Config

type Config struct {
	// URL is a libpq connection string or URL (postgres://...).
	URL string
	// LockTimeout is applied as the session lock_timeout on every connection.
	// Zero means DefaultLockTimeout.
	LockTimeout time.Duration
	// StatementTimeout is applied as the session statement_timeout on every
	// connection. Zero means DefaultStatementTimeout.
	StatementTimeout time.Duration
	// ConnectTimeout bounds each dial attempt. The session affinity proof
	// runs against the same server and opens a connection of its own, so its
	// budget is this plus the proof's own floor. Zero means
	// DefaultConnectTimeout.
	ConnectTimeout time.Duration
	// CACertPath, when set, enables verify-full TLS using the given CA bundle
	// (e.g. the RDS/Aurora global bundle). Unset, RDS/Aurora endpoints are
	// auto-verified with the embedded bundle (see rds.go).
	CACertPath string

	// Pool sizing and lifecycle. Zero values keep pgxpool's defaults.
	//
	// NOTE: the advisory-lock connection (LK-1) must NOT come from this pool:
	// session-scoped locks die with their session, and lifetime/idle
	// recycling would silently release the lock. The lock helper owns a
	// dedicated single-connection pool exempt from recycling.
	MaxConns              int32
	MinConns              int32
	MaxConnLifetime       time.Duration
	MaxConnLifetimeJitter time.Duration
	MaxConnIdleTime       time.Duration
	HealthCheckPeriod     time.Duration

	// QueryExecMode overrides pgx's default protocol usage — e.g.
	// pgx.QueryExecModeExec when a transaction-pooling proxy that cannot
	// handle prepared statements sits in front of the pool. Zero keeps pgx's
	// default (statement caching).
	QueryExecMode pgx.QueryExecMode
	// Logger, when set, enables statement-level tracing (pgx tracelog) at
	// debug level through the given slog logger.
	Logger *slog.Logger
	// BeforeConnect, when set, can mutate each new connection's config just
	// before dialing — the hook for short-lived credentials such as RDS IAM
	// authentication tokens.
	BeforeConnect func(context.Context, *pgx.ConnConfig) error
}

Config describes a connection target. Zero values keep sensible defaults (ours for the session timeouts, pgxpool's for pool sizing and lifecycle) — decisions, not options.

type IndexBuildProgress

type IndexBuildProgress struct {
	Phase            string
	BlocksDone       uint64
	BlocksTotal      uint64
	TuplesDone       uint64
	TuplesTotal      uint64
	LockersTotal     uint64
	LockersDone      uint64
	CurrentLockerPID uint32
}

IndexBuildProgress is one server observation of a concurrent index build.

func ConcurrentIndexProgress

func ConcurrentIndexProgress(ctx context.Context, session RowQuerier, backendPID uint32) (IndexBuildProgress, bool, error)

ConcurrentIndexProgress reads the active build owned by backendPID. The boolean is false when PostgreSQL has not published the row yet or the build has already left the progress view.

type Querier

type Querier interface {
	Query(ctx context.Context, sql string, args ...any) (pgx.Rows, error)
}

Querier is the query surface TerminateBlockers needs; *pgxpool.Pool and *pgx.Conn both satisfy it.

type RowQuerier

type RowQuerier interface {
	QueryRow(context.Context, string, ...any) pgx.Row
}

RowQuerier is the session capability needed for a progress observation.

type TableLock added in v0.3.2

type TableLock struct {
	// contains filtered or unexported fields
}

TableLock proves this process holds LK-1's per-table advisory lock on a dedicated session. Its zero value is forgeable; consumers must reject it when Table is empty.

func (TableLock) Key added in v0.3.2

func (l TableLock) Key() int64

Key returns the advisory lock key.

func (TableLock) Schema added in v0.3.2

func (l TableLock) Schema() string

Schema returns the locked table's schema.

func (TableLock) Table added in v0.3.2

func (l TableLock) Table() string

Table returns the locked table name.

Jump to

Keyboard shortcuts

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