dbx

pgx

A pgx-natív adatbázis-stack: pool, a DBTX query executor és a contextet figyelő DB.

A dbx/pg a dbx alapértelmezett, pgx-natív implementációja: pgx-re épül, a Go legszélesebb körben használt PostgreSQL-driverére. Három dolgot ad: pool építését a dbx.Config alapján, egy query executort (DB), amely figyeli, fut-e éppen tranzakció, és egy dbx.Transactor implementációt, ami pgx.Tx-et hordoz a contextben. Önálló package: semmihez sem kell hozzá a gpsystem kit.

import (
    "github.com/gp-system/dbx"
    "github.com/gp-system/dbx/pg"
)

Nincs wrapper-API: a csomag pgxpool.Pool-t, pgx.Rows-t és pgconn.CommandTag-et ad vissza. Ami a pgx dokumentációjában áll, az itt is igaz. A dbx/pg csak a pool-építést, az instrumentációt és a tranzakció-plumbingot központosítja.

Pool

pool, err := pg.NewPool(ctx, cfg.DB)  // DSN a konfigból, connect, ping
pool := pg.MustNewPool(ctx, cfg.DB)   // ugyanez, hibára panic, main()-be

A NewPool a cfg.DB.DSN()-ből parse-olja a pool-konfigot, beállítja a MaxConns-t (DB_MAX_CONNS, default 10), csatlakozik, és pinggel ellenőrzi a kapcsolatot, így a rossz jelszó vagy elérhetetlen adatbázis már startupkor bukik, nem az első kérésnél. A MustNewPool ugyanez panic-kel; a generált main.go-k ezt hívják.

A kitre épülő projektben a generált belépési pont a lezárást a server-életciklusra köti:

err := server.Run(ctx, cfg.Server, register,
    server.WithCloser("pgxpool", func(context.Context) error {
        pool.Close()
        return nil
    }))

A poolt az otelpgx instrumentálja: minden query egy OTel spant bocsát ki a statementtel és az idejével, az aktív request/task span alá ágyazva. No-op tracer providerrel (dev mód, telemetria ki) ez gyakorlatilag ingyenes; bekapcsolt telemetriával a query-k megjelennek a trace waterfallban és a Sentry hiba-eseményeinél is.

DBTX: a repositoryk query-felülete

type DBTX interface {
    Exec(ctx context.Context, sql string, args ...any) (pgconn.CommandTag, error)
    Query(ctx context.Context, sql string, args ...any) (pgx.Rows, error)
    QueryRow(ctx context.Context, sql string, args ...any) pgx.Row
}

Ettől a három metódustól függ minden repository. A visszatérési típusok szándékosan nyers pgx-típusok (pgconn.CommandTag, pgx.Rows, pgx.Row), nem kit-wrapperek: a scanelés, a CollectRows, a RowsAffected() pontosan úgy működik, ahogy a pgx dokumentációjában olvasod.

A DBTX ráadásul metódusról metódusra azonos a sqlc által generált DBTX interfésszel. A kit generátorai kézzel írt repositorykat vázolnak fel, de ha egy projekt sqlc-vel generálna query-kódot, annak structjai változtatás nélkül elfogadják a pg.NewDB(pool)-t, és ugyanazokban a context-hordozta tranzakciókban vesznek részt.

DB: az executor, ami a contextet figyeli

db := pg.NewDB(pool) // DBTX-et implementál

A pg.DB a DBTX egyetlen kit-beli implementációja, és egyetlen trükkje van: minden hívás előbb tranzakciót keres a contextben. Ha a WithinTransaction nyitott egyet, a hívás azon fut; ha nincs, a poolon. A repository sosem tudja, és nem is kell tudnia, melyik eset áll fenn.

Ez az a pont, ahol a kit tranzakció-modellje összeér: a service WithinTransaction-t hív, a repositoryk pedig változatlan szignatúrával, automatikusan csatlakoznak.

Egy shop repository

A shop mintaalkalmazás modul-szintű termék-repositoryja (internal/modules/shop/repository/, minden surface számára közös). A repository egyetlen pg.DBTX mezőt fog. A main.go a pg.NewDB(pool)-t injektálja a Dependencies-en át:

package repository

type ProductRepository struct {
    db pg.DBTX
}

func NewProductRepository(db pg.DBTX) *ProductRepository {
    return &ProductRepository{db: db}
}

A terméklistázás cursor-lapozással (keyset seek, nem OFFSET):

const productSort = "created_at DESC, id DESC" // egyezik az ORDER BY-jal

func (r *ProductRepository) ListProducts(
    ctx context.Context, p paginate.CursorParams,
) (paginate.CursorPage[Product], error) {
    args := []any{p.FetchLimit()}
    query := `SELECT id, name, price_cents, stock, created_at FROM products`
    if p.HasCursor() {
        var createdAt time.Time
        var id int64
        if err := p.DecodeKey(&createdAt, &id); err != nil {
            return paginate.CursorPage[Product]{}, err
        }
        query += ` WHERE (created_at, id) < ($2, $3)`
        args = append(args, createdAt, id)
    }
    query += ` ORDER BY ` + productSort + ` LIMIT $1`

    rows, err := r.db.Query(ctx, query, args...)
    if err != nil {
        return paginate.CursorPage[Product]{}, errs.Wrap(err, "products: list")
    }
    products, err := pgx.CollectRows(rows, pgx.RowToStructByName[Product])
    if err != nil {
        return paginate.CursorPage[Product]{}, errs.Wrap(err, "products: scan")
    }
    return paginate.NewCursorPage(products, p, func(last Product) []any {
        return []any{last.CreatedAt, last.ID}
    })
}

És a készletcsökkentés, amit a PlaceOrder tranzakcióból hív a service:

func (r *ProductRepository) DecrementStock(ctx context.Context, productID int64, qty int) error {
    tag, err := r.db.Exec(ctx,
        `UPDATE products SET stock = stock - $2 WHERE id = $1 AND stock >= $2`,
        productID, qty)
    if err != nil {
        return errs.Wrap(err, "products: decrement stock")
    }
    if tag.RowsAffected() == 0 {
        return errs.New("insufficient stock",
            errs.Code("insufficient_stock"), errs.Public("A termék elfogyott."))
    }
    return nil
}

Figyeld meg, mi nincs a kódban: tranzakció-paraméter. A ListProducts a publikus listázó endpointból poolon fut; a DecrementStock a PlaceOrder tranzakcióján belül hívva ugyanabba a pgx.Tx-be csatlakozik, mindkettő ugyanazzal a szignatúrával. A pgx.CollectRows + RowToStructByName a pgx beépített scanelése; a WHERE ... AND stock >= $2 + RowsAffected() páros pedig az atomikus készlet-ellenőrzés, SELECT-majd-UPDATE verseny nélkül.

Transactor

tx := pg.NewTransactor(pool) // dbx.Transactor

err := tx.WithinTransaction(ctx, func(ctx context.Context) error {
    if err := s.orders.Insert(ctx, o); err != nil {
        return err
    }
    return s.products.DecrementStock(ctx, o.ProductID, o.Qty) // ugyanaz a tx
})

A WithinTransaction tranzakciót nyit a poolon, a pgx.Tx-et a contextbe injektálja, és az fn lefutása után commitol. Bármilyen hiba, vagy panic, ami rollback után újradobódik, visszagörgeti az egészet. Ha a context már hordoz tranzakciót, az fn ahhoz csatlakozik új tranzakció helyett: a nesting lapos, savepoint nincs, a commitot/rollbacket mindig a legkülső hívás birtokolja. A minta teljes tárgyalása, a shop PlaceOrder példájával, a Tranzakciók oldalon.

Escape hatch: a context-plumbing

Saját executorhoz (pl. batch-íróhoz vagy egy kézzel instrumentált query-úthoz) a tranzakció-hordozás exportált:

ctx = pg.ContextWithTx(ctx, tx)   // tranzakció a contextbe
tx, ok := pg.TxFromContext(ctx)   // tranzakció a contextből

Alkalmazáskód ezeket normál esetben sosem hívja: a pg.DB és a Transactor együtt lefedi. Akkor kellenek, ha a DBTX-nél gazdagabb pgx-felületre van szükséged (pl. CopyFrom, SendBatch) egy tranzakción belül.

A tranzakció-szemantikát (commit/rollback/nesting/panic) a dbx saját integrációs tesztjei rögzítik valós PostgreSQL ellen: go test -tags=integration ./dbx/pg/.

Önálló példa

A dbx/pg-hez semmi nem kell a kitből: egy pool, egy DBTX-alapú repository és egy transactor önmagában is elég.

main.go
package main

import (
    "context"
    "log"

    "github.com/gp-system/dbx"
    "github.com/gp-system/dbx/pg"
)

func main() {
    ctx := context.Background()
    pool := pg.MustNewPool(ctx, dbx.Config{
        Host: "localhost", Port: 5432,
        User: "shop", Password: "secret", Database: "shop",
    })
    defer pool.Close()

    db := pg.NewDB(pool)
    tx := pg.NewTransactor(pool)

    err := tx.WithinTransaction(ctx, func(ctx context.Context) error {
        _, err := db.Exec(ctx, `UPDATE products SET stock = stock - 1 WHERE id = $1`, "prod_42")
        return err
    })
    if err != nil {
        log.Fatal(err)
    }
}

Használt patternek

Copyright © 2026