pgx
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.
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.
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
- Szűk, sqlc-kompatibilis repository-interfész (
DBTX): Design patternek. - Unit of Work / tranzakció a contextben (
WithinTransaction,txKey{}): Design patternek.