From 7a1ceed32429edff879dc9eebce4f12fd9e4c98a Mon Sep 17 00:00:00 2001 From: Alex Shapalov Date: Tue, 29 Sep 2026 16:58:05 -0700 Subject: [PATCH 1/3] Switch the PostgreSQL driver from pgx to pgGo pgbot now talks to PostgreSQL through github.com/pgrundev/pggo, a dependency-free wire-protocol client. The connection layer (pool, session pins, read-only transactions, SSH DialFunc, pooler/PgDog probes), the collectors, erd, logs, advise and the MCP tools are ported; pgx, pgpassfile, pgservicefile and puddle leave the module graph. Output is unchanged apart from pgbot's own footprint in sampled counters (see CHANGELOG), verified side by side against PostgreSQL 16-19, a streaming standby, PgBouncer and PgDog. The unit and integration suites pass on 16-19 over TLS+SCRAM, including the pooler tests. go.mod temporarily replaces pggo with ../../pggo until pggo's library release is tagged. Co-Authored-By: Claude Opus 5.5 (1M context) --- .golangci.yml | 2 +- CHANGELOG.md | 24 ++- cmd/pgbot/activity.go | 4 +- cmd/pgbot/advise.go | 34 ++--- cmd/pgbot/advise_safety_test.go | 12 +- cmd/pgbot/helpers.go | 6 +- cmd/pgbot/logs.go | 8 +- cmd/pgbot/main.go | 2 +- cmd/pgbot/main_test.go | 2 +- cmd/pgbot/mcp_tools.go | 14 +- cmd/pgbot/pgservice_test.go | 12 +- cmd/pgbot/profile_integration_test.go | 14 +- go.mod | 6 +- go.sum | 13 +- internal/advisor/engine.go | 2 +- .../activity_selfexclude_integration_test.go | 8 +- internal/collect/ash.go | 6 +- .../collect/collation_integration_test.go | 8 +- internal/collect/collector.go | 12 +- .../collect/correlate_integration_test.go | 8 +- .../collect/docverify_integration_test.go | 13 +- internal/collect/health.go | 4 +- internal/collect/integration_test.go | 2 +- .../collect/invalid_index_integration_test.go | 8 +- .../collect/pgss_schema_integration_test.go | 10 +- internal/collect/queries.go | 2 +- .../collect/readonly_role_integration_test.go | 17 +-- internal/collect/settings.go | 6 +- internal/collect/waitstudy.go | 5 +- .../collect/waitstudy_integration_test.go | 10 +- internal/conn/aurora_integration_test.go | 8 +- internal/conn/capability.go | 6 +- internal/conn/connect.go | 138 +++++++++--------- internal/conn/exclude_integration_test.go | 8 +- internal/conn/pooler.go | 24 +-- internal/conn/pooler_integration_test.go | 14 +- internal/conn/pooler_test.go | 10 +- internal/conn/sshtunnel.go | 14 +- internal/conn/sshtunnel_test.go | 2 +- internal/erd/introspect.go | 8 +- internal/pglog/sqlsource.go | 8 +- 41 files changed, 257 insertions(+), 257 deletions(-) diff --git a/.golangci.yml b/.golangci.yml index 2038293..9e95bcc 100644 --- a/.golangci.yml +++ b/.golangci.yml @@ -11,7 +11,7 @@ linters: - (io.Closer).Close - (*database/sql.Rows).Close - (*database/sql.Stmt).Close - - (*github.com/jackc/pgx/v5.Conn).Close + - (*github.com/pgrundev/pggo.Conn).Close - (*github.com/pgrundev/pgbot/internal/store.Store).Close - fmt.Fprint - fmt.Fprintf diff --git a/CHANGELOG.md b/CHANGELOG.md index 704c10d..96a1e24 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,26 @@ separately by `model.SchemaVersion` (currently 1.3.0). ## [Unreleased] ### Changed +- **PostgreSQL driver: pgx → [pgGo](https://github.com/pgrundev/pggo).** pgbot + now talks to PostgreSQL through pgGo, a dependency-free wire-protocol client + (standard library only). Output is unchanged apart from pgbot's own + footprint in the counters it samples (below) — verified by running the pgx + and pgGo builds side by side against PostgreSQL 16–19, a streaming standby, + PgBouncer (transaction mode) and PgDog. Read-only + guarantees are unchanged: the same session pins and `READ ONLY` transactions, + enforced by PostgreSQL. Effects you may notice: + - The binary is ~5 MB smaller (4 fewer modules: pgx, pgpassfile, + pgservicefile, puddle); runs need fewer round trips (`BEGIN` is pipelined + with each collector's first query). + - pgbot's own traffic is less visible in `pg_stat_database` during the + sample window, so throughput on a near-idle database is reported closer to + the truth (pgx builds counted ~15–20 tps of pgbot's own commits; with a + rate-limited 200 tps workload pgx reported ~220, pgGo ~200). On an idle + database, cache hit and rollbacks now show `—` (nothing measured) where the + pgx build graded pgbot's own reads. + - Behind a transaction pooler no protocol fallback is needed any more: pgGo + only uses unnamed statements, parsed and executed within one sync. The + prepared-statement probe still runs and still counts as a pooler signal. - **Gauge strip in the default `inspect` view.** Four vital signs sit right under the header — cache hit, lock wait (naming the culprit query when sessions are blocked), rollbacks, and idle index bytes as a share of the @@ -50,9 +70,9 @@ separately by `model.SchemaVersion` (currently 1.3.0). is passed and neither `$DATABASE_URL` nor `$PGBOT_DATABASE_URL` is set, pgbot now checks `$PGSERVICE` too, so a [connection service file](https://www.postgresql.org/docs/current/libpq-pgservice.html) - alone is enough to pick a database. pgx's `ParseConfig` already reads + alone is enough to pick a database. The driver's `ParseConfig` already reads `PGSERVICEFILE` (or the libpq default path); this just stops pgbot from - erroring out before pgx gets a chance to. + erroring out before the driver gets a chance to. - **AWS Bedrock Mantle as an `explain` / `ask` provider** (#35, contributed by @edwardsb). `PGBOT_AI_PROVIDER=bedrock` (alias `mantle`) routes `openai.*` models through the Responses API and `anthropic.*` models through the diff --git a/cmd/pgbot/activity.go b/cmd/pgbot/activity.go index 678da0b..577995b 100644 --- a/cmd/pgbot/activity.go +++ b/cmd/pgbot/activity.go @@ -8,9 +8,9 @@ import ( "text/tabwriter" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/render" + "github.com/pgrundev/pggo" "github.com/spf13/cobra" ) @@ -78,7 +78,7 @@ wait on, and the (scrubbed) SQL. Plain idle sessions are summarized, not listed if err != nil { return err } - got, err := pgx.CollectRows(rows, pgx.RowToStructByNameLax[activityRow]) + got, err := pggo.CollectStructs[activityRow](rows) if err != nil { return err } diff --git a/cmd/pgbot/advise.go b/cmd/pgbot/advise.go index 8c9ff21..4f932e1 100644 --- a/cmd/pgbot/advise.go +++ b/cmd/pgbot/advise.go @@ -7,10 +7,10 @@ import ( "os" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/advisor" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/render" + "github.com/pgrundev/pggo" "github.com/spf13/cobra" ) @@ -83,8 +83,8 @@ func adviseRun(ctx context.Context, connString string, top int, minImpr float64) // One READ ONLY transaction for the whole loop: hypopg state is per-connection, // so the hypothetical indexes and their reset must share a single connection. - err = target.ReadOnlyTx(ctx, func(tx pgx.Tx) error { - p := pgxPlanner{tx: tx, caps: target.Caps} + err = target.ReadOnlyTx(ctx, func(tx *pggo.Tx) error { + p := txPlanner{tx: tx, caps: target.Caps} defer p.ResetHypo(context.Background()) //nolint:errcheck // belt-and-braces cleanup res.Recs, res.Stats = advisor.Advise(ctx, p, inputs, advisor.Options{ MinImprovement: minImpr, StaleRelations: stale, WriteHeavy: writeHeavy, @@ -255,7 +255,7 @@ func relationStats(ctx context.Context, t *conn.Target) (stale, writeHeavy map[s return stale, writeHeavy } -// pgxPlanner implements advisor.Planner over a READ ONLY pgx transaction. Every +// txPlanner implements advisor.Planner over a READ ONLY transaction. Every // statement is plan-only or hypothetical — none executes the inspected query. // // Each fallible operation is wrapped in a SAVEPOINT: a query pgbot can't plan (or @@ -264,8 +264,8 @@ func relationStats(ctx context.Context, t *conn.Target) (stale, writeHeavy map[s // clears the error and lets the loop continue to the next query. hypopg indexes // live in backend memory, not transaction state, so they survive a savepoint // rollback — only hypopg_reset drops them. -type pgxPlanner struct { - tx pgx.Tx +type txPlanner struct { + tx *pggo.Tx caps conn.Capabilities // for hypopg's schema — its functions are called by qualified name } @@ -273,18 +273,18 @@ type pgxPlanner struct { // pg_stat_statements, hypopg lands in Supabase's "extensions" schema (issue #10) // — off a read-only role's search_path — so the fixed function names are // qualified with the namespace the probe read from pg_extension. -func (p pgxPlanner) hypo(fn string) string { return p.caps.ExtObject("hypopg", fn) } +func (p txPlanner) hypo(fn string) string { return p.caps.ExtObject("hypopg", fn) } -func (p pgxPlanner) GenericPlan(ctx context.Context, query string) ([]byte, error) { +func (p txPlanner) GenericPlan(ctx context.Context, query string) ([]byte, error) { // GENERIC_PLAN plans a normalized $N query without values; FORMAT JSON gives one - // row. This MUST use the raw simple-query protocol (PgConn.Exec): both the - // extended protocol and pgx's SimpleProtocol mode treat the $1/$2 inside the - // EXPLAIN'd query as bind parameters of the OUTER statement and demand values - // ("expected 2 arguments, got 0"). GENERIC_PLAN exists precisely to plan those - // placeholders WITHOUT values, so the SQL must reach the server byte-for-byte. + // row. This MUST use the simple-query protocol (Conn.SimpleQuery): the + // extended protocol treats the $1/$2 inside the EXPLAIN'd query as bind + // parameters of the OUTER statement and demands values. GENERIC_PLAN exists + // precisely to plan those placeholders WITHOUT values, so the SQL must reach + // the server byte-for-byte. var js []byte err := p.inSavepoint(ctx, func() error { - res, err := p.tx.Conn().PgConn().Exec(ctx, "EXPLAIN (GENERIC_PLAN, FORMAT JSON) "+query).ReadAll() + res, err := p.tx.Conn().SimpleQuery(ctx, "EXPLAIN (GENERIC_PLAN, FORMAT JSON) "+query) if err != nil { return err } @@ -299,7 +299,7 @@ func (p pgxPlanner) GenericPlan(ctx context.Context, query string) ([]byte, erro return js, err } -func (p pgxPlanner) CreateHypoIndex(ctx context.Context, ddl string) (string, int64, error) { +func (p txPlanner) CreateHypoIndex(ctx context.Context, ddl string) (string, int64, error) { var oid uint32 var name string var estBytes int64 @@ -315,14 +315,14 @@ func (p pgxPlanner) CreateHypoIndex(ctx context.Context, ddl string) (string, in return name, estBytes, err } -func (p pgxPlanner) ResetHypo(ctx context.Context) error { +func (p txPlanner) ResetHypo(ctx context.Context) error { _, err := p.tx.Exec(ctx, "SELECT "+p.hypo("hypopg_reset")+"()") return err } // inSavepoint runs fn between a SAVEPOINT and its RELEASE, rolling back to the // savepoint on error so a single failed statement doesn't poison the transaction. -func (p pgxPlanner) inSavepoint(ctx context.Context, fn func() error) error { +func (p txPlanner) inSavepoint(ctx context.Context, fn func() error) error { if _, err := p.tx.Exec(ctx, "SAVEPOINT pgbot_adv"); err != nil { return err } diff --git a/cmd/pgbot/advise_safety_test.go b/cmd/pgbot/advise_safety_test.go index 5e7c5af..de1db76 100644 --- a/cmd/pgbot/advise_safety_test.go +++ b/cmd/pgbot/advise_safety_test.go @@ -5,8 +5,8 @@ import ( "os" "testing" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/conn" + "github.com/pgrundev/pggo" ) // TestAdviseSafety_readOnlyTxBlocksInjectedWrite is the belt-and-braces proof for @@ -28,12 +28,12 @@ func TestIntegration_adviseSafety_readOnlyTxBlocksInjectedWrite(t *testing.T) { ctx := context.Background() // A raw connection (NOT pgbot-pinned) to set up and later inspect the table. - admin, err := pgx.Connect(ctx, d) + admin, err := pggo.Connect(ctx, d) if err != nil { t.Fatalf("admin connect: %v", err) } - defer admin.Close(ctx) - if _, err := admin.Exec(ctx, `DROP TABLE IF EXISTS advise_safety; CREATE TABLE advise_safety (n int)`); err != nil { + defer admin.Close() + if _, err := admin.SimpleQuery(ctx, `DROP TABLE IF EXISTS advise_safety; CREATE TABLE advise_safety (n int)`); err != nil { t.Fatalf("setup: %v", err) } @@ -47,8 +47,8 @@ func TestIntegration_adviseSafety_readOnlyTxBlocksInjectedWrite(t *testing.T) { // sanitizeQuery, exactly as a hostile pgss entry would if the sanitizer missed // it. The second statement is a write. injected := "EXPLAIN (GENERIC_PLAN, FORMAT JSON) SELECT 1; INSERT INTO advise_safety VALUES (1)" - _ = target.ReadOnlyTx(ctx, func(tx pgx.Tx) error { - _, _ = tx.Conn().PgConn().Exec(ctx, injected).ReadAll() + _ = target.ReadOnlyTx(ctx, func(tx *pggo.Tx) error { + _, _ = tx.Conn().SimpleQuery(ctx, injected) return nil }) diff --git a/cmd/pgbot/helpers.go b/cmd/pgbot/helpers.go index bfea1de..ffdcb72 100644 --- a/cmd/pgbot/helpers.go +++ b/cmd/pgbot/helpers.go @@ -25,7 +25,7 @@ func firstNonEmpty(vals ...string) string { } // pgServiceFallback lets a bare $PGSERVICE select a connection when neither an -// argument nor $DATABASE_URL/$PGBOT_DATABASE_URL is set. pgx's ParseConfig +// argument nor $DATABASE_URL/$PGBOT_DATABASE_URL is set. pggo.ParseConfig // already reads a connection service file (PGSERVICEFILE, or the libpq // default path) once it gets a "service=..." string — this just builds that // string so users who manage connections through a service file don't have @@ -69,10 +69,10 @@ func terminalWidth() int { // hostPort pulls the host/port off the pool's config for the baseline // fingerprint fallback (used only when the system identifier isn't readable). func hostPort(t *conn.Target) (string, string) { - cfg := t.Pool.Config().ConnConfig + cfg := t.Pool.Config() port := "5432" if cfg.Port != 0 { - port = strconv.Itoa(int(cfg.Port)) + port = strconv.Itoa(cfg.Port) } return cfg.Host, port } diff --git a/cmd/pgbot/logs.go b/cmd/pgbot/logs.go index 4d04912..74ed5dd 100644 --- a/cmd/pgbot/logs.go +++ b/cmd/pgbot/logs.go @@ -9,10 +9,10 @@ import ( "strings" "time" - "github.com/jackc/pgx/v5/pgconn" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/pglog" "github.com/pgrundev/pgbot/internal/render" + "github.com/pgrundev/pggo" "github.com/spf13/cobra" ) @@ -95,7 +95,7 @@ func runLogs(cmd *cobra.Command, args []string, f logsFlags) error { // register on first use, and a long --live recycles connections (5m // lifetime), minting new PIDs the whole while. connUser := "" - if cfg, cerr := pgconn.ParseConfig(connString); cerr == nil { + if cfg, cerr := pggo.ParseConfig(connString); cerr == nil { connUser = cfg.User } keep := func(e pglog.Entry) bool { @@ -159,10 +159,10 @@ func runLogs(cmd *cobra.Command, args []string, f logsFlags) error { // logsErr turns the two expected failures into their fixes: no collector, and // the one missing GRANT. func logsErr(err error, connString string) error { - var pgErr *pgconn.PgError + var pgErr *pggo.PgError if errors.As(err, &pgErr) && pgErr.Code == "42501" { // insufficient_privilege user := "your_pgbot_role" - if cfg, cerr := pgconn.ParseConfig(connString); cerr == nil && cfg.User != "" { + if cfg, cerr := pggo.ParseConfig(connString); cerr == nil && cfg.User != "" { user = cfg.User } return fmt.Errorf("reading the server log needs one grant beyond pg_monitor — run as an admin:\n\n %s\n\n(%s)", diff --git a/cmd/pgbot/main.go b/cmd/pgbot/main.go index 14d3e72..0b8d92e 100644 --- a/cmd/pgbot/main.go +++ b/cmd/pgbot/main.go @@ -21,7 +21,7 @@ var version = "dev" func main() { // A SIGINT/SIGTERM cancels the run's context instead of killing the process - // mid-flight: collectors abort at the next round trip, the pgx pool closes, + // mid-flight: collectors abort at the next round trip, the connection pool closes, // and the store finishes its write. cmd.Context() in every handler is this ctx. ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() diff --git a/cmd/pgbot/main_test.go b/cmd/pgbot/main_test.go index 1064e49..50c4d42 100644 --- a/cmd/pgbot/main_test.go +++ b/cmd/pgbot/main_test.go @@ -66,7 +66,7 @@ func TestPgServiceFallback(t *testing.T) { } // A bare $PGSERVICE resolves a connection when nothing else is set — - // pgx's ParseConfig reads PGSERVICE(FILE) itself once it sees "service=...". + // pggo.ParseConfig reads PGSERVICE(FILE) itself once it sees "service=...". t.Setenv("DATABASE_URL", "") t.Setenv("PGBOT_DATABASE_URL", "") if dsn, err := dsnFromArgs(json.RawMessage(`{}`)); err != nil || dsn != "service=mydb" { diff --git a/cmd/pgbot/mcp_tools.go b/cmd/pgbot/mcp_tools.go index 19fac09..8662943 100644 --- a/cmd/pgbot/mcp_tools.go +++ b/cmd/pgbot/mcp_tools.go @@ -7,11 +7,11 @@ import ( "strings" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/docs" "github.com/pgrundev/pgbot/internal/advisor" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/store" + "github.com/pgrundev/pggo" ) // B8 MCP tools. Every one either produces no findings (explain_plan, schema_of, @@ -57,8 +57,8 @@ func explainPlanTool(ctx context.Context, args json.RawMessage) (string, error) stmt = "EXPLAIN (GENERIC_PLAN, FORMAT JSON) " + clean } var planJSON []byte - err = target.ReadOnlyTx(ctx, func(tx pgx.Tx) error { - res, e := tx.Conn().PgConn().Exec(ctx, stmt).ReadAll() + err = target.ReadOnlyTx(ctx, func(tx *pggo.Tx) error { + res, e := tx.Conn().SimpleQuery(ctx, stmt) if e != nil { return e } @@ -129,7 +129,7 @@ func schemaOfTool(ctx context.Context, args json.RawMessage) (string, error) { out := map[string]any{"table": a.Table, "exactness": "scraped", "note": "catalog metadata only — no table data is read; the row count is the planner's estimate (reltuples)."} // The table name is passed as a bind parameter cast to regclass — never spliced. - err = target.ReadOnlyTx(ctx, func(tx pgx.Tx) error { + err = target.ReadOnlyTx(ctx, func(tx *pggo.Tx) error { var relTuples int64 var totalBytes int64 if e := tx.QueryRow(ctx, `SELECT reltuples::bigint, pg_total_relation_size($1::regclass) FROM pg_class WHERE oid = $1::regclass`, a.Table).Scan(&relTuples, &totalBytes); e != nil { @@ -226,14 +226,14 @@ func recordIndexVerdictTool(_ context.Context, args json.RawMessage) (string, er // scanRows runs a query with one bind arg and returns its rows as []map, column // names from the result description (best-effort; a query error yields nil). -func scanRows(ctx context.Context, tx pgx.Tx, sql string, arg any) []map[string]any { +func scanRows(ctx context.Context, tx *pggo.Tx, sql string, arg any) []map[string]any { rows, err := tx.Query(ctx, sql, arg) if err != nil { return nil } defer rows.Close() var out []map[string]any - fields := rows.FieldDescriptions() + fields := rows.Columns() for rows.Next() { vals, err := rows.Values() if err != nil { @@ -241,7 +241,7 @@ func scanRows(ctx context.Context, tx pgx.Tx, sql string, arg any) []map[string] } m := make(map[string]any, len(fields)) for i, fd := range fields { - m[string(fd.Name)] = vals[i] + m[fd.Name] = vals[i] } out = append(out, m) } diff --git a/cmd/pgbot/pgservice_test.go b/cmd/pgbot/pgservice_test.go index e20f5d2..b298141 100644 --- a/cmd/pgbot/pgservice_test.go +++ b/cmd/pgbot/pgservice_test.go @@ -5,10 +5,10 @@ import ( "path/filepath" "testing" - "github.com/jackc/pgx/v5" + "github.com/pgrundev/pggo" ) -// The fallback hands pgx a bare "service=" and relies on pgx to read the +// The fallback hands the driver a bare "service=" and relies on it to read the // connection service file. Pin that end to end: a service file named through // PGSERVICEFILE must supply host, port, user, and database, and a service name // missing from the file must be an error rather than a silent localhost. @@ -28,19 +28,19 @@ func TestPgServiceFallback_resolvesThroughServiceFile(t *testing.T) { } dsn := firstNonEmpty("", os.Getenv("DATABASE_URL"), os.Getenv("PGBOT_DATABASE_URL"), pgServiceFallback()) - cfg, err := pgx.ParseConfig(dsn) + cfg, err := pggo.ParseConfig(dsn) if err != nil { - t.Fatalf("pgx.ParseConfig(%q): %v", dsn, err) + t.Fatalf("pggo.ParseConfig(%q): %v", dsn, err) } if cfg.Host != "db.internal" || cfg.Port != 6432 || cfg.User != "pgbot_ro" || cfg.Database != "appdb" { t.Fatalf("service file not applied: host=%q port=%d user=%q db=%q", cfg.Host, cfg.Port, cfg.User, cfg.Database) } - if cfg.TLSConfig == nil { + if cfg.SSLMode != "require" { t.Fatal("sslmode=require from the service file was not applied") } t.Setenv("PGSERVICE", "does-not-exist") - if _, err := pgx.ParseConfig(pgServiceFallback()); err == nil { + if _, err := pggo.ParseConfig(pgServiceFallback()); err == nil { t.Fatal("an unknown service name should fail to parse, not fall through to defaults") } } diff --git a/cmd/pgbot/profile_integration_test.go b/cmd/pgbot/profile_integration_test.go index f6ef5f5..a5ed748 100644 --- a/cmd/pgbot/profile_integration_test.go +++ b/cmd/pgbot/profile_integration_test.go @@ -8,10 +8,10 @@ import ( "testing" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/collect" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // TestIntegration_schemaProfile_acceptance is the D3 acceptance test (DoD 9 & 10): @@ -27,11 +27,11 @@ func TestIntegration_schemaProfile_acceptance(t *testing.T) { ctx := context.Background() const dbName = "pgbot_d3_accept" - admin, err := pgx.Connect(ctx, su) + admin, err := pggo.Connect(ctx, su) if err != nil { t.Fatalf("admin connect: %v", err) } - defer admin.Close(ctx) + defer admin.Close() _, _ = admin.Exec(ctx, `DROP DATABASE IF EXISTS `+dbName+` WITH (FORCE)`) if _, err := admin.Exec(ctx, `CREATE DATABASE `+dbName); err != nil { t.Fatalf("create db: %v", err) @@ -41,11 +41,11 @@ func TestIntegration_schemaProfile_acceptance(t *testing.T) { dsn := swapDatabase(t, su, dbName) // A sound schema: bigint keys, the FK is indexed, no invalid/redundant indexes. - fixture, err := pgx.Connect(ctx, dsn) + fixture, err := pggo.Connect(ctx, dsn) if err != nil { t.Fatalf("fixture connect: %v", err) } - defer fixture.Close(ctx) + defer fixture.Close() mustExec(t, ctx, fixture, `CREATE TABLE users (id bigserial PRIMARY KEY, email text)`, `CREATE TABLE orders (id bigserial PRIMARY KEY, user_id bigint REFERENCES users(id), total numeric)`, @@ -86,7 +86,7 @@ func TestIntegration_schemaProfile_acceptance(t *testing.T) { func swapDatabase(t *testing.T, dsn, db string) string { t.Helper() - cfg, err := pgx.ParseConfig(dsn) + cfg, err := pggo.ParseConfig(dsn) if err != nil { t.Fatalf("parse dsn: %v", err) } @@ -100,7 +100,7 @@ func swapDatabase(t *testing.T, dsn, db string) string { return u.String() } -func mustExec(t *testing.T, ctx context.Context, c *pgx.Conn, stmts ...string) { +func mustExec(t *testing.T, ctx context.Context, c *pggo.Conn, stmts ...string) { t.Helper() tctx, cancel := context.WithTimeout(ctx, 15*time.Second) defer cancel() diff --git a/go.mod b/go.mod index 2ce2072..25fe237 100644 --- a/go.mod +++ b/go.mod @@ -8,9 +8,9 @@ require ( github.com/BurntSushi/toml v1.6.0 github.com/charmbracelet/lipgloss v1.1.0 github.com/invopop/jsonschema v0.14.0 - github.com/jackc/pgx/v5 v5.11.0 github.com/kevinburke/ssh_config v1.6.0 github.com/owenrumney/go-sarif/v2 v2.3.3 + github.com/pgrundev/pggo v0.1.0 github.com/spf13/cobra v1.10.2 golang.org/x/crypto v0.56.0 golang.org/x/sync v0.23.0 @@ -29,9 +29,6 @@ require ( github.com/dustin/go-humanize v1.0.1 // indirect github.com/google/uuid v1.6.0 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect - github.com/jackc/pgpassfile v1.0.0 // indirect - github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect - github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/lucasb-eyer/go-colorful v1.2.0 // indirect github.com/mattn/go-isatty v0.0.24 // indirect github.com/mattn/go-runewidth v0.0.16 // indirect @@ -44,7 +41,6 @@ require ( github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect go.yaml.in/yaml/v4 v4.0.0-rc.2 // indirect golang.org/x/sys v0.48.0 // indirect - golang.org/x/text v0.41.0 // indirect modernc.org/libc v1.75.6 // indirect modernc.org/mathutil v1.7.1 // indirect modernc.org/memory v1.12.1 // indirect diff --git a/go.sum b/go.sum index 34c8d3b..9a9ebe5 100644 --- a/go.sum +++ b/go.sum @@ -36,14 +36,6 @@ github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2 github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= github.com/invopop/jsonschema v0.14.0 h1:MHQqLhvpNUZfw+hM3AZDYK7jxO8FZoQeQM77g8iyZjg= github.com/invopop/jsonschema v0.14.0/go.mod h1:ygm6C2EaVNMBDPpaPlnOA2pFAxBnxGjFlMZABxm9n2I= -github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= -github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= -github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= -github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= -github.com/jackc/pgx/v5 v5.11.0 h1:IzBBtyK9AHqf98cctWFifYSci2hgQR/cd56wB4p+ogg= -github.com/jackc/pgx/v5 v5.11.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4= -github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= -github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= github.com/kevinburke/ssh_config v1.6.0 h1:J1FBfmuVosPHf5GRdltRLhPJtJpTlMdKTBjRgTaQBFY= github.com/kevinburke/ssh_config v1.6.0/go.mod h1:q2RIzfka+BXARoNexmF9gkxEX7DmvbW9P4hIVx2Kg4M= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= @@ -64,6 +56,8 @@ github.com/owenrumney/go-sarif/v2 v2.3.3 h1:ubWDJcF5i3L/EIOER+ZyQ03IfplbSU1BLOE2 github.com/owenrumney/go-sarif/v2 v2.3.3/go.mod h1:MSqMMx9WqlBSY7pXoOZWgEsVB4FDNfhcaXDA1j6Sr+w= github.com/pb33f/ordered-map/v2 v2.3.1 h1:5319HDO0aw4DA4gzi+zv4FXU9UlSs3xGZ40wcP1nBjY= github.com/pb33f/ordered-map/v2 v2.3.1/go.mod h1:qxFQgd0PkVUtOMCkTapqotNgzRhMPL7VvaHKbd1HnmQ= +github.com/pgrundev/pggo v0.1.0 h1:SORSc6h+sq6k2qSrk7YFghShgb60fsuIIyL9/2u4PTU= +github.com/pgrundev/pggo v0.1.0/go.mod h1:PqAWvEdWAMoePDrIR3D8jnBm+zy56/kEK4TDzNRDEJo= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE= @@ -77,7 +71,6 @@ github.com/spf13/cobra v1.10.2/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiT github.com/spf13/pflag v1.0.9 h1:9exaQaMOCwffKiiiYk6/BndUBv+iRViNW+4lEMi0PvY= github.com/spf13/pflag v1.0.9/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= -github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= @@ -108,8 +101,6 @@ golang.org/x/term v0.46.0/go.mod h1:+K02xbkittuwc0Am4abfA3Fc+XRGXkvBXNO88NCXPoc= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk= golang.org/x/text v0.3.5/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= -golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= -golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE= golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk= diff --git a/internal/advisor/engine.go b/internal/advisor/engine.go index a5edec7..815e8f8 100644 --- a/internal/advisor/engine.go +++ b/internal/advisor/engine.go @@ -6,7 +6,7 @@ import ( ) // Planner is the minimal database surface the advisor drives. Implemented over a -// pgx READ ONLY transaction in the command; mocked in tests. Every method is a +// READ ONLY transaction in the command; mocked in tests. Every method is a // plan-only or hypothetical operation — none executes the inspected query or // writes anything. type Planner interface { diff --git a/internal/collect/activity_selfexclude_integration_test.go b/internal/collect/activity_selfexclude_integration_test.go index f3540df..c17059e 100644 --- a/internal/collect/activity_selfexclude_integration_test.go +++ b/internal/collect/activity_selfexclude_integration_test.go @@ -5,10 +5,10 @@ import ( "testing" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/collect" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // End-to-end wiring guard for the observer-exclusion class, found walking the @@ -46,16 +46,16 @@ func TestIntegration_selfExclusion_wiring(t *testing.T) { // Hold a real idle-in-transaction session; pgbot must count it. Counting is a // lower bound (background activity can only add), so this stays robust while // still failing if the activity collector stopped counting real sessions. - cfg, err := pgx.ParseConfig(d) + cfg, err := pggo.ParseConfig(d) if err != nil { t.Fatalf("parse dsn: %v", err) } cfg.RuntimeParams["application_name"] = "pgbot_selftest_app" - held, err := pgx.ConnectConfig(ctx, cfg) + held, err := pggo.ConnectConfig(ctx, cfg) if err != nil { t.Fatalf("hold connect: %v", err) } - defer held.Close(ctx) + defer held.Close() if _, err := held.Exec(ctx, "BEGIN READ ONLY"); err != nil { t.Fatalf("begin: %v", err) } diff --git a/internal/collect/ash.go b/internal/collect/ash.go index 00ad210..feb5a37 100644 --- a/internal/collect/ash.go +++ b/internal/collect/ash.go @@ -6,9 +6,9 @@ import ( "sort" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // WaitSample is one observation of one active backend at one instant. Nil @@ -71,7 +71,7 @@ type ashResult struct { // interval: at the default 10 Hz that would be 100 ms, and a poll over a // normal-latency link (a laptop or CI reaching RDS/Neon/Supabase at 30–100 ms // RTT) cannot complete in that — every poll timed out, every timeout tore down -// its pool connection (pgx closes a connection whose context expires mid-query) +// its pool connection (the driver closes a connection whose context expires mid-query) // and the profile came back "all N polls errored" while the report said nothing. // With a fixed budget the sampler runs at min(hz, what the round trip allows) // instead of 0 (PR#1). @@ -108,7 +108,7 @@ func sampleWaitsOpt(ctx context.Context, t *conn.Target, caps conn.Capabilities, res.failures++ // drop this poll, keep going return } - got, err := pgx.CollectRows(rows, pgx.RowToStructByNameLax[WaitSample]) + got, err := pggo.CollectStructs[WaitSample](rows) if err != nil { res.failures++ return diff --git a/internal/collect/collation_integration_test.go b/internal/collect/collation_integration_test.go index 4779a9f..13eb221 100644 --- a/internal/collect/collation_integration_test.go +++ b/internal/collect/collation_integration_test.go @@ -6,11 +6,11 @@ import ( "testing" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/collect" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/findings" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // A collation version mismatch can't be produced by upgrading glibc inside a @@ -27,11 +27,11 @@ func TestIntegration_collationVersionMismatch(t *testing.T) { ro := dsn(t) ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) defer cancel() - admin, err := pgx.Connect(ctx, su) + admin, err := pggo.Connect(ctx, su) if err != nil { t.Fatalf("admin connect: %v", err) } - t.Cleanup(func() { admin.Close(context.Background()) }) + t.Cleanup(func() { admin.Close() }) var vnum int if err := admin.QueryRow(ctx, `SELECT current_setting('server_version_num')::int`).Scan(&vnum); err != nil { @@ -49,7 +49,7 @@ func TestIntegration_collationVersionMismatch(t *testing.T) { if recorded == nil { t.Skip("this database's collation records no version (C/POSIX) — nothing can drift") } - refresh := `ALTER DATABASE ` + pgx.Identifier{db}.Sanitize() + ` REFRESH COLLATION VERSION` + refresh := `ALTER DATABASE ` + pggo.QuoteIdentifier(db) + ` REFRESH COLLATION VERSION` if _, err := admin.Exec(ctx, `UPDATE pg_database SET datcollversion = '0.0-pgbot-test' WHERE datname = current_database()`); err != nil { t.Fatalf("forge a stale datcollversion: %v", err) } diff --git a/internal/collect/collector.go b/internal/collect/collector.go index 843243f..b620ed9 100644 --- a/internal/collect/collector.go +++ b/internal/collect/collector.go @@ -8,9 +8,9 @@ import ( "context" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // Kind decides how the runner samples a collector. @@ -156,12 +156,12 @@ func newContext(caps conn.Capabilities, tB time.Time, dt time.Duration) *model.C func queryOne[T any](ctx context.Context, t *conn.Target, sql string, args ...any) (T, error) { var out T - err := t.ReadOnlyTx(ctx, func(tx pgx.Tx) error { + err := t.ReadOnlyTx(ctx, func(tx *pggo.Tx) error { rows, err := tx.Query(ctx, sql, args...) if err != nil { return err } - out, err = pgx.CollectExactlyOneRow(rows, pgx.RowToStructByNameLax[T]) + out, err = pggo.CollectOneStruct[T](rows) return err }) return out, err @@ -175,7 +175,7 @@ func queryMany[T any](ctx context.Context, t *conn.Target, sql string, args ...a // (SET LOCAL … — they end with the transaction, so the session's pins stay). func queryManyLocal[T any](ctx context.Context, t *conn.Target, setLocal []string, sql string, args ...any) ([]T, error) { var out []T - err := t.ReadOnlyTx(ctx, func(tx pgx.Tx) error { + err := t.ReadOnlyTx(ctx, func(tx *pggo.Tx) error { for _, s := range setLocal { if _, err := tx.Exec(ctx, s); err != nil { return err @@ -185,7 +185,7 @@ func queryManyLocal[T any](ctx context.Context, t *conn.Target, setLocal []strin if err != nil { return err } - out, err = pgx.CollectRows(rows, pgx.RowToStructByNameLax[T]) + out, err = pggo.CollectStructs[T](rows) return err }) return out, err @@ -193,7 +193,7 @@ func queryManyLocal[T any](ctx context.Context, t *conn.Target, setLocal []strin func scalar[T any](ctx context.Context, t *conn.Target, sql string, args ...any) (T, error) { var out T - err := t.ReadOnlyTx(ctx, func(tx pgx.Tx) error { + err := t.ReadOnlyTx(ctx, func(tx *pggo.Tx) error { return tx.QueryRow(ctx, sql, args...).Scan(&out) }) return out, err diff --git a/internal/collect/correlate_integration_test.go b/internal/collect/correlate_integration_test.go index 56d8d58..08f6703 100644 --- a/internal/collect/correlate_integration_test.go +++ b/internal/collect/correlate_integration_test.go @@ -6,11 +6,11 @@ import ( "testing" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/collect" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/correlate" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // TestIntegration_indexCorrelation exercises the new index attributes (method, @@ -26,15 +26,15 @@ func TestIntegration_indexCorrelation(t *testing.T) { t.Skip("set PGBOT_TEST_SUPERUSER_DSN (a superuser DSN) to run the index-correlation test") } ctx := context.Background() - admin, err := pgx.Connect(ctx, d) + admin, err := pggo.Connect(ctx, d) if err != nil { t.Fatalf("admin connect: %v", err) } - defer admin.Close(ctx) + defer admin.Close() // A table big enough that its indexes cross pgbot's 16 KB collector floor, with // one of each shape the classifier distinguishes. - if _, err := admin.Exec(ctx, ` + if _, err := admin.SimpleQuery(ctx, ` DROP TABLE IF EXISTS public.corr_job, public.corr_ri; CREATE TABLE public.corr_job (id bigint primary key, "externalIdNormalized" text, tags jsonb, status text, customer_id bigint); INSERT INTO public.corr_job diff --git a/internal/collect/docverify_integration_test.go b/internal/collect/docverify_integration_test.go index cb3332b..45dd56a 100644 --- a/internal/collect/docverify_integration_test.go +++ b/internal/collect/docverify_integration_test.go @@ -10,9 +10,8 @@ import ( "strings" "testing" - "github.com/jackc/pgx/v5" - "github.com/jackc/pgx/v5/pgconn" "github.com/pgrundev/pgbot/internal/conn" + "github.com/pgrundev/pggo" ) // TestDocVerifyQueries_run executes every "How to verify it yourself" SQL query @@ -38,11 +37,11 @@ func TestIntegration_docVerifyQueries(t *testing.T) { } ctx := context.Background() - admin, err := pgx.Connect(ctx, d) + admin, err := pggo.Connect(ctx, d) if err != nil { t.Fatalf("admin connect: %v", err) } - defer admin.Close(ctx) + defer admin.Close() var vnum int if err := admin.QueryRow(ctx, "SELECT current_setting('server_version_num')::int").Scan(&vnum); err != nil { @@ -54,7 +53,7 @@ func TestIntegration_docVerifyQueries(t *testing.T) { // Example objects the pages reference by name (public.issues.last_seen_at, // public.orders). Everything else is catalog views with literal filters. - if _, err := admin.Exec(ctx, ` + if _, err := admin.SimpleQuery(ctx, ` CREATE EXTENSION IF NOT EXISTS pgstattuple; CREATE TABLE IF NOT EXISTS public.orders (id bigserial primary key, customer_id int, status int, amount numeric, note text); CREATE TABLE IF NOT EXISTS public.issues (id bigserial primary key, last_seen_at timestamptz, project_id int)`); err != nil { @@ -76,7 +75,7 @@ func TestIntegration_docVerifyQueries(t *testing.T) { } for _, stmt := range splitStatements(b.sql) { ran++ - err := target.ReadOnlyTx(ctx, func(tx pgx.Tx) error { + err := target.ReadOnlyTx(ctx, func(tx *pggo.Tx) error { _, e := tx.Exec(ctx, stmt) return e }) @@ -160,7 +159,7 @@ func splitStatements(block string) []string { // tolerated reports whether an error is the known pg_stat_bgwriter column move // (columns relocated to pg_stat_checkpointer in PG17) rather than a real defect. func tolerated(stmt string, err error) bool { - var pg *pgconn.PgError + var pg *pggo.PgError if !errors.As(err, &pg) { return false } diff --git a/internal/collect/health.go b/internal/collect/health.go index fe54d38..d9aac7b 100644 --- a/internal/collect/health.go +++ b/internal/collect/health.go @@ -5,10 +5,10 @@ import ( _ "embed" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/model" "github.com/pgrundev/pgbot/internal/rate" + "github.com/pgrundev/pggo" ) //go:embed sql/health.sql @@ -51,7 +51,7 @@ func (healthCollector) Sample(ctx context.Context, t *conn.Target, _ conn.Capabi if err != nil { return healthSample{}, err } - return pgx.CollectExactlyOneRow(rows, pgx.RowToStructByNameLax[healthSample]) + return pggo.CollectOneStruct[healthSample](rows) } func (healthCollector) Assemble(c *model.Context, _ conn.Capabilities, s sampled, dt time.Duration, _ Options) { diff --git a/internal/collect/integration_test.go b/internal/collect/integration_test.go index 301398f..d24ed7e 100644 --- a/internal/collect/integration_test.go +++ b/internal/collect/integration_test.go @@ -26,7 +26,7 @@ func TestIntegration_cancelMidRun(t *testing.T) { defer target.Close() // Warm the pool with one full run first, so the goroutine baseline reflects - // steady state (pgx pool goroutines up, that run's sampler already drained) — + // steady state (pool goroutines up, that run's sampler already drained) — // otherwise we'd be measuring pool warm-up, not a sampler leak. if _, err := collect.Run(context.Background(), target, collect.Options{Interval: 200 * time.Millisecond, ASHHz: 10, ASHWindow: 200 * time.Millisecond}); err != nil { diff --git a/internal/collect/invalid_index_integration_test.go b/internal/collect/invalid_index_integration_test.go index 26ed87a..e8ab886 100644 --- a/internal/collect/invalid_index_integration_test.go +++ b/internal/collect/invalid_index_integration_test.go @@ -7,11 +7,11 @@ import ( "testing" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/collect" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/findings" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // Issue #11: a CREATE INDEX CONCURRENTLY that fails during the build (here: a @@ -29,13 +29,13 @@ func TestIntegration_invalidIndexDebrisIsNotCritical(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) defer cancel() - admin, err := pgx.Connect(ctx, su) + admin, err := pggo.Connect(ctx, su) if err != nil { t.Fatalf("admin connect: %v", err) } - t.Cleanup(func() { admin.Close(context.Background()) }) + t.Cleanup(func() { admin.Close() }) - if _, err := admin.Exec(ctx, ` + if _, err := admin.SimpleQuery(ctx, ` DROP TABLE IF EXISTS public.pgbot_it_dup; CREATE TABLE public.pgbot_it_dup (id serial primary key, v int); INSERT INTO public.pgbot_it_dup (v) SELECT g % 100 FROM generate_series(1, 20000) g`); err != nil { diff --git a/internal/collect/pgss_schema_integration_test.go b/internal/collect/pgss_schema_integration_test.go index 37358ae..2c556db 100644 --- a/internal/collect/pgss_schema_integration_test.go +++ b/internal/collect/pgss_schema_integration_test.go @@ -6,10 +6,10 @@ import ( "testing" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/collect" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // Issue #10: pg_stat_statements installed outside public (Supabase's @@ -29,13 +29,13 @@ func TestIntegration_pgssInNonDefaultSchema(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) defer cancel() - admin, err := pgx.Connect(ctx, su) + admin, err := pggo.Connect(ctx, su) if err != nil { t.Fatalf("admin connect: %v", err) } // Registered before the restore cleanup below: cleanups run last-in-first-out, // so the extension is moved back on a still-open connection. - t.Cleanup(func() { admin.Close(context.Background()) }) + t.Cleanup(func() { admin.Close() }) var installed bool if err := admin.QueryRow(ctx, `SELECT count(*) > 0 FROM pg_extension WHERE extname = 'pg_stat_statements'`).Scan(&installed); err != nil { @@ -53,12 +53,12 @@ func TestIntegration_pgssInNonDefaultSchema(t *testing.T) { // the read-only role can see the objects — the point is that they are found // by *qualified* name, not via search_path. const relocated = "pgbot_ext_test" - if _, err := admin.Exec(ctx, `CREATE SCHEMA IF NOT EXISTS `+relocated+`; GRANT USAGE ON SCHEMA `+relocated+` TO PUBLIC; ALTER EXTENSION pg_stat_statements SET SCHEMA `+relocated); err != nil { + if _, err := admin.SimpleQuery(ctx, `CREATE SCHEMA IF NOT EXISTS `+relocated+`; GRANT USAGE ON SCHEMA `+relocated+` TO PUBLIC; ALTER EXTENSION pg_stat_statements SET SCHEMA `+relocated); err != nil { t.Fatalf("relocate extension: %v", err) } t.Cleanup(func() { c := context.Background() - if _, err := admin.Exec(c, `ALTER EXTENSION pg_stat_statements SET SCHEMA `+pgx.Identifier{origSchema}.Sanitize()+`; DROP SCHEMA IF EXISTS `+relocated); err != nil { + if _, err := admin.SimpleQuery(c, `ALTER EXTENSION pg_stat_statements SET SCHEMA `+pggo.QuoteIdentifier(origSchema)+`; DROP SCHEMA IF EXISTS `+relocated); err != nil { t.Errorf("restore extension schema: %v", err) } }) diff --git a/internal/collect/queries.go b/internal/collect/queries.go index 154ee72..7e2b0a6 100644 --- a/internal/collect/queries.go +++ b/internal/collect/queries.go @@ -65,7 +65,7 @@ func (queriesCollector) Sample(ctx context.Context, t *conn.Target, caps conn.Ca // (caps.Pgss — discovered from pg_extension at connect). Supabase installs the // extension in "extensions", off a read-only role's search_path; the bare // name there raised 42P01 while the capability list still said "present" - // (issue #10). The names are the fixed allowlisted objects, quoted by pgx. + // (issue #10). The names are the fixed allowlisted objects, quoted by pggo.QuoteIdentifier. rows, err := queryManyLocal[queryRow](ctx, t, []string{"SET LOCAL work_mem = '64MB'"}, fmt.Sprintf(sqlQueries, caps.StatStatementsTotalCol(), caps.Pgss("pg_stat_statements"))) if err != nil { diff --git a/internal/collect/readonly_role_integration_test.go b/internal/collect/readonly_role_integration_test.go index 5b36dc5..c69a6bb 100644 --- a/internal/collect/readonly_role_integration_test.go +++ b/internal/collect/readonly_role_integration_test.go @@ -10,11 +10,10 @@ import ( "testing" "time" - "github.com/jackc/pgx/v5" - "github.com/jackc/pgx/v5/pgconn" "github.com/pgrundev/pgbot/internal/collect" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // The read-only guarantee's real boundary is the role: a pg_monitor role with no @@ -30,19 +29,19 @@ func TestIntegration_readOnlyRole_runsFullPipelineAndCannotWrite(t *testing.T) { } ctx := context.Background() - cfg, err := pgx.ParseConfig(su) + cfg, err := pggo.ParseConfig(su) if err != nil { t.Fatalf("parse superuser dsn: %v", err) } const roUser, roPass = "pgbot_ro_test", "ro_test_pw" - admin, err := pgx.Connect(ctx, su) + admin, err := pggo.Connect(ctx, su) if err != nil { t.Fatalf("admin connect: %v", err) } - defer admin.Close(ctx) + defer admin.Close() - db := pgx.Identifier{cfg.Database}.Sanitize() + db := pggo.QuoteIdentifier(cfg.Database) // Idempotent (re)provision: least privilege — pg_monitor + CONNECT, nothing more. _, _ = admin.Exec(ctx, `DROP OWNED BY `+roUser) _, _ = admin.Exec(ctx, `DROP ROLE IF EXISTS `+roUser) @@ -90,16 +89,16 @@ func TestIntegration_readOnlyRole_runsFullPipelineAndCannotWrite(t *testing.T) { // 2. The role itself must be unable to write — a raw connection with NO pgbot // read-only pinning still cannot INSERT, because the grant simply isn't there. - raw, err := pgx.Connect(ctx, roDSN) + raw, err := pggo.Connect(ctx, roDSN) if err != nil { t.Fatalf("raw connect as %s: %v", roUser, err) } - defer raw.Close(ctx) + defer raw.Close() _, werr := raw.Exec(ctx, `INSERT INTO ro_probe VALUES (1)`) if werr == nil { t.Fatal("SAFETY: a pg_monitor role was able to INSERT — it has write access it must not have") } - var pgErr *pgconn.PgError + var pgErr *pggo.PgError if !strings.Contains(werr.Error(), "permission denied") && !(errors.As(werr, &pgErr) && pgErr.Code == "42501") { t.Errorf("write should be denied for insufficient privilege (42501), got: %v", werr) diff --git a/internal/collect/settings.go b/internal/collect/settings.go index 52e5617..79cdf46 100644 --- a/internal/collect/settings.go +++ b/internal/collect/settings.go @@ -5,9 +5,9 @@ import ( _ "embed" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) //go:embed sql/settings.sql @@ -35,7 +35,7 @@ func (settingsCollector) Sample(ctx context.Context, t *conn.Target, _ conn.Capa // database. UnpinLocal reverts them for this one transaction; the query also // drops session/client-sourced rows from the override set. var rows []settingRow - err := t.ReadOnlyTx(ctx, func(tx pgx.Tx) error { + err := t.ReadOnlyTx(ctx, func(tx *pggo.Tx) error { if err := conn.UnpinLocal(ctx, tx); err != nil { return err } @@ -43,7 +43,7 @@ func (settingsCollector) Sample(ctx context.Context, t *conn.Target, _ conn.Capa if err != nil { return err } - rows, err = pgx.CollectRows(r, pgx.RowToStructByNameLax[settingRow]) + rows, err = pggo.CollectStructs[settingRow](r) return err }) return rows, err diff --git a/internal/collect/waitstudy.go b/internal/collect/waitstudy.go index 2cd3967..992046d 100644 --- a/internal/collect/waitstudy.go +++ b/internal/collect/waitstudy.go @@ -6,10 +6,9 @@ import ( "sort" "time" - "github.com/jackc/pgx/v5" - "github.com/pgrundev/pgbot/internal/conn" "github.com/pgrundev/pgbot/internal/model" + "github.com/pgrundev/pggo" ) // LockEdge is one blocked→holder observation from one slow-plane snapshot: @@ -312,7 +311,7 @@ func RunWaitStudy(ctx context.Context, t *conn.Target, caps conn.Capabilities, o snapFails++ return } - got, err := pgx.CollectRows(rows, pgx.RowToStructByNameLax[lockEdgeRow]) + got, err := pggo.CollectStructs[lockEdgeRow](rows) if err != nil { snapFails++ return diff --git a/internal/collect/waitstudy_integration_test.go b/internal/collect/waitstudy_integration_test.go index 61a2fe3..a8a5104 100644 --- a/internal/collect/waitstudy_integration_test.go +++ b/internal/collect/waitstudy_integration_test.go @@ -7,8 +7,8 @@ import ( "testing" "time" - "github.com/jackc/pgx/v5" "github.com/pgrundev/pgbot/internal/conn" + "github.com/pgrundev/pggo" ) // A real two-session lock conflict must produce a SUSTAINED blocker with the @@ -24,11 +24,11 @@ func TestIntegration_waitStudy_namesTheBlocker(t *testing.T) { ctx := context.Background() // Holder: an open transaction owning the advisory lock. - holder, err := pgx.Connect(ctx, dsn) + holder, err := pggo.Connect(ctx, dsn) if err != nil { t.Fatalf("connect: %v", err) } - defer holder.Close(context.Background()) + defer holder.Close() tx, err := holder.Begin(ctx) if err != nil { t.Fatal(err) @@ -42,11 +42,11 @@ func TestIntegration_waitStudy_namesTheBlocker(t *testing.T) { victimDone := make(chan struct{}) go func() { defer close(victimDone) - v, err := pgx.Connect(ctx, dsn) + v, err := pggo.Connect(ctx, dsn) if err != nil { return } - defer v.Close(context.Background()) + defer v.Close() vctx, cancel := context.WithTimeout(ctx, 20*time.Second) defer cancel() _, _ = v.Exec(vctx, `SELECT pg_advisory_xact_lock(987654321012345)`) diff --git a/internal/conn/aurora_integration_test.go b/internal/conn/aurora_integration_test.go index 91cf591..745ef3f 100644 --- a/internal/conn/aurora_integration_test.go +++ b/internal/conn/aurora_integration_test.go @@ -6,7 +6,7 @@ import ( "testing" "time" - "github.com/jackc/pgx/v5/pgxpool" + "github.com/pgrundev/pggo" ) // TestIntegration_auroraInstances is the opt-in end-to-end check for @@ -36,15 +36,15 @@ func TestIntegration_auroraInstances(t *testing.T) { if err != nil { t.Fatalf("discover: %v", err) } - cfg, err := pgxpool.ParseConfig(dsn) + cfg, err := pggo.ParseConfig(dsn) if err != nil { t.Fatal(err) } - endpoint, err := CanonicalRDSHost(cfg.ConnConfig.Host) + endpoint, err := CanonicalRDSHost(cfg.Host) if err != nil { t.Fatalf("endpoint: %v", err) } - t.Logf("entry endpoint %s → %s; %d instance(s)", cfg.ConnConfig.Host, endpoint, len(instances)) + t.Logf("entry endpoint %s → %s; %d instance(s)", cfg.Host, endpoint, len(instances)) for _, inst := range instances { host, err := AuroraInstanceHost(endpoint, inst.ID) diff --git a/internal/conn/capability.go b/internal/conn/capability.go index 6a8ce2f..52d6088 100644 --- a/internal/conn/capability.go +++ b/internal/conn/capability.go @@ -3,7 +3,7 @@ package conn import ( "time" - "github.com/jackc/pgx/v5" + "github.com/pgrundev/pggo" ) // Capabilities is built once at connect from server_version_num + a probe of @@ -39,7 +39,7 @@ func (c Capabilities) ExtensionSchema(ext string) string { return c.ExtensionSch // allowlisted object (view / function / table) belonging to an extension — // e.g. "extensions"."pg_stat_statements" — so a read works whatever the // session's search_path is. The schema comes from the catalog, never user input, -// but it is quoted anyway (pgx.Identifier) so an unusual namespace name can't +// but it is quoted anyway (pggo.QuoteIdentifier) so an unusual namespace name can't // break the statement. When the schema is unknown it degrades to the bare name, // i.e. exactly the pre-discovery behaviour. func (c Capabilities) ExtObject(ext, object string) string { @@ -47,7 +47,7 @@ func (c Capabilities) ExtObject(ext, object string) string { if schema == "" { return object } - return pgx.Identifier{schema, object}.Sanitize() + return pggo.QuoteIdentifier(schema, object) } // Pgss is ExtObject for pg_stat_statements' objects: the view, the diff --git a/internal/conn/connect.go b/internal/conn/connect.go index 8813c98..506273f 100644 --- a/internal/conn/connect.go +++ b/internal/conn/connect.go @@ -6,26 +6,40 @@ package conn import ( "context" "fmt" + "net/url" "os" + "strings" "time" - "github.com/jackc/pgx/v5" - "github.com/jackc/pgx/v5/pgxpool" + "github.com/pgrundev/pggo" ) -// clientOnlyParams are libpq connection parameters that pgx/pgconn does NOT -// recognize and so forwards to the server as startup GUCs — where the server -// rejects them ("unrecognized configuration parameter"). Managed providers -// (pgrun, Neon) ship channel_binding in their default connection strings. pgx -// can't implement SCRAM channel binding anyway, so we drop it; TLS from sslmode -// still applies — channel binding was hardening on top of that. +// clientOnlyParams are libpq connection parameters the driver accepts but cannot +// honor. Managed providers (pgrun, Neon) ship channel_binding in their default +// connection strings; pgGo does not implement SCRAM channel binding, so it is +// dropped (never sent as a server GUC) and we say so. TLS from sslmode still +// applies — channel binding was hardening on top of that. var clientOnlyParams = []string{"channel_binding"} +// hasConnParam reports whether a URL or keyword=value connection string sets key. +func hasConnParam(connString, key string) bool { + if strings.HasPrefix(connString, "postgres://") || strings.HasPrefix(connString, "postgresql://") { + u, err := url.Parse(connString) + return err == nil && u.Query().Has(key) + } + for _, f := range strings.Fields(connString) { + if strings.HasPrefix(f, key+"=") { + return true + } + } + return false +} + // Target is a configured, capability-probed connection to one database. The // pool is small (max 4) so a burst of concurrent collectors can't itself become // a connection storm on the database it was invoked to inspect. type Target struct { - Pool *pgxpool.Pool + Pool *pggo.Pool Caps Capabilities Pooler PoolerInfo self *selfPIDs // backend PIDs of our own pool connections; see ExcludeSelf @@ -56,26 +70,19 @@ func ConnectDBAt(ctx context.Context, connString, database, host string) (*Targe } func connect(ctx context.Context, connString, database, host string) (*Target, error) { - cfg, err := pgxpool.ParseConfig(connString) + cfg, err := pggo.ParseConfig(connString) if err != nil { return nil, fmt.Errorf("parse connection string: %w", err) } if database != "" { - cfg.ConnConfig.Database = database + cfg.Database = database } if host != "" { - cfg.ConnConfig.Host = host - cfg.ConnConfig.Fallbacks = nil // a multi-host DSN names the cluster, not this member - if tc := cfg.ConnConfig.TLSConfig; tc != nil { - tc = tc.Clone() - tc.ServerName = host - cfg.ConnConfig.TLSConfig = tc - } + // The TLS server name follows Config.Host, so verify-full checks this + // member's own certificate. + cfg.Host = host } - cfg.MaxConns = maxConns - cfg.MinConns = 0 - cfg.MaxConnLifetime = 5 * time.Minute - cfg.ConnConfig.RuntimeParams["application_name"] = "pgbot" + cfg.RuntimeParams["application_name"] = "pgbot" // Route the TCP leg through the SSH jump host when one is configured. This has // to happen before probe(): the probe connection dials too, and it must take @@ -83,32 +90,26 @@ func connect(ctx context.Context, connString, database, host string) (*Target, e // to a local forward is what keeps sslmode= and .pgpass matching on the real // hostname — see sshtunnel.go. if dial := sshDialFunc(); dial != nil { - cfg.ConnConfig.DialFunc = dial - // Let the SSH server resolve database hostnames. - cfg.ConnConfig.LookupFunc = func(_ context.Context, host string) ([]string, error) { - return []string{host}, nil - } + // pgGo never resolves the host itself: DialFunc receives host:port, so + // the SSH server resolves database hostnames. + cfg.DialFunc = dial } - // Drop client-only params pgx forwarded into RuntimeParams (it would send them - // as server GUCs, which the server rejects). See clientOnlyParams. + // The driver never sends client-only params to the server. See clientOnlyParams. for _, p := range clientOnlyParams { - if _, ok := cfg.ConnConfig.RuntimeParams[p]; ok { - delete(cfg.ConnConfig.RuntimeParams, p) + if hasConnParam(connString, p) { fmt.Fprintf(os.Stderr, "pgbot: ignoring connection param %q — the driver can't honor it; TLS from sslmode still applies\n", p) } } // Probe capabilities + pooler on a throwaway connection first, so AfterConnect - // applies only the GUCs this server understands and the pool uses the right - // wire protocol. - caps, pooler, probePID, err := probe(ctx, cfg.ConnConfig.Copy()) + // applies only the GUCs this server understands. (Behind a transaction pooler + // no protocol switch is needed: pgGo only uses the unnamed statement, parsed + // and executed within one Sync, which poolers route as a unit.) + caps, pooler, probePID, err := probe(ctx, cfg.Copy()) if err != nil { return nil, err } - if pooler.SimpleProtocol { - cfg.ConnConfig.DefaultQueryExecMode = pgx.QueryExecModeSimpleProtocol - } // Track our own backend PIDs so collectors can exclude every pgbot connection // from pg_stat_activity, not just the one running a given query (ExcludeSelf). @@ -116,21 +117,21 @@ func connect(ctx context.Context, connString, database, host string) (*Target, e // filters historical log lines by PID, and the probe wrote some. self := newSelfPIDs() self.add(probePID) - cfg.AfterConnect = func(ctx context.Context, c *pgx.Conn) error { - if err := applySessionSetup(ctx, c, caps); err != nil { - return err - } - self.add(c.PgConn().PID()) - return nil - } - cfg.BeforeClose = func(c *pgx.Conn) { - self.remove(c.PgConn().PID()) - } - - pool, err := pgxpool.NewWithConfig(ctx, cfg) - if err != nil { - return nil, fmt.Errorf("open pool: %w", err) - } + pool := pggo.NewPool(pggo.PoolConfig{ + Config: cfg, + MaxConns: maxConns, + MaxConnLifetime: 5 * time.Minute, + AfterConnect: func(ctx context.Context, c *pggo.Conn) error { + if err := applySessionSetup(ctx, c, caps); err != nil { + return err + } + self.add(c.PID()) + return nil + }, + BeforeClose: func(c *pggo.Conn) { + self.remove(c.PID()) + }, + }) return &Target{Pool: pool, Caps: caps, Pooler: pooler, self: self}, nil } @@ -152,7 +153,7 @@ func (t *Target) Warm(ctx context.Context) { if t.Pool == nil { return } - held := make([]*pgxpool.Conn, 0, maxConns) + held := make([]*pggo.PoolConn, 0, maxConns) for i := 0; i < maxConns; i++ { c, err := t.Pool.Acquire(ctx) if err != nil { @@ -185,7 +186,7 @@ var sessionPins = []struct{ name, value string }{ // rather than pgbot. Only the settings collector needs it. stats_fetch_consistency // (PG15+, also pinned) is left alone: it isn't a tuning parameter, and a SET LOCAL // of an unknown GUC would abort the transaction on PG < 15. -func UnpinLocal(ctx context.Context, tx pgx.Tx) error { +func UnpinLocal(ctx context.Context, tx *pggo.Tx) error { for _, p := range sessionPins { if _, err := tx.Exec(ctx, "SET LOCAL "+p.name+" = DEFAULT"); err != nil { return fmt.Errorf("unpin %s: %w", p.name, err) @@ -197,7 +198,7 @@ func UnpinLocal(ctx context.Context, tx pgx.Tx) error { // applySessionSetup pins every physical connection. statement_timeout and // lock_timeout are mandatory: pgbot must never become the incident it was // invoked to diagnose. -func applySessionSetup(ctx context.Context, c *pgx.Conn, caps Capabilities) error { +func applySessionSetup(ctx context.Context, c *pggo.Conn, caps Capabilities) error { stmts := []string{"SET application_name = 'pgbot'"} for _, p := range sessionPins { stmts = append(stmts, "SET "+p.name+" = "+p.value) @@ -220,21 +221,16 @@ func applySessionSetup(ctx context.Context, c *pgx.Conn, caps Capabilities) erro // probe reads server_version_num, installed extensions, role membership, and // the system identifier in one round trip (with a best-effort fallback for the // identifier, which needs elevated read access on some managed providers). -func probe(ctx context.Context, cc *pgx.ConnConfig) (Capabilities, PoolerInfo, uint32, error) { - c, err := pgx.ConnectConfig(ctx, cc) +func probe(ctx context.Context, cc *pggo.Config) (Capabilities, PoolerInfo, uint32, error) { + c, err := pggo.ConnectConfig(ctx, cc) if err != nil { return Capabilities{}, PoolerInfo{}, 0, fmt.Errorf("connect: %w", err) } - defer c.Close(ctx) - probePID := c.PgConn().PID() + defer c.Close() + probePID := c.PID() - // Detect the pooler first — if prepared statements are broken, every later - // query on this probe connection must use the simple protocol too. + // Detect the pooler first (named prepared statements, session persistence). pooler := detectPooler(ctx, c, cc) - mode := []any{} - if pooler.SimpleProtocol { - mode = []any{pgx.QueryExecModeSimpleProtocol} - } var caps Capabilities var mk providerMarkers @@ -258,7 +254,7 @@ func probe(ctx context.Context, cc *pgx.ConnConfig) (Capabilities, PoolerInfo, u -- counter pgbot reports. to_regprocedure does neither. to_regprocedure('aurora_version()') IS NOT NULL OR to_regprocedure('aurora_replica_status()') IS NOT NULL` - err = c.QueryRow(ctx, q, mode...).Scan(&caps.VersionNum, &caps.VersionText, &caps.Database, + err = c.QueryRow(ctx, q).Scan(&caps.VersionNum, &caps.VersionText, &caps.Database, &caps.StartedAt, &caps.HasStatStatements, &caps.HasHypopg, &caps.HasPgMonitor, &mk.HasRDS, &mk.HasCloudSQL, &mk.HasAzure, &caps.InRecovery, &mk.IsAurora) if err != nil { @@ -271,7 +267,7 @@ func probe(ctx context.Context, cc *pgx.ConnConfig) (Capabilities, PoolerInfo, u // system_identifier makes the baseline fingerprint survive a restore/rename; // it needs pg_monitor/superuser on some providers, so it's best-effort. var sysID int64 - if err := c.QueryRow(ctx, `SELECT system_identifier FROM pg_control_system()`, mode...).Scan(&sysID); err == nil { + if err := c.QueryRow(ctx, `SELECT system_identifier FROM pg_control_system()`).Scan(&sysID); err == nil { caps.SystemIdentifier = fmt.Sprintf("%d", sysID) } @@ -280,12 +276,12 @@ func probe(ctx context.Context, cc *pgx.ConnConfig) (Capabilities, PoolerInfo, u // by qualified name — Supabase and friends install them in "extensions", off // the read-only role's search_path (issue #10). Best-effort: on failure the // map stays empty and callers fall back to bare names. - if rows, err := c.Query(ctx, `SELECT e.extname, n.nspname FROM pg_extension e JOIN pg_namespace n ON n.oid = e.extnamespace ORDER BY e.extname`, mode...); err == nil { + if rows, err := c.Query(ctx, `SELECT e.extname, n.nspname FROM pg_extension e JOIN pg_namespace n ON n.oid = e.extnamespace ORDER BY e.extname`); err == nil { type extRow struct { Name string Schema string } - if exts, err := pgx.CollectRows(rows, pgx.RowToStructByPos[extRow]); err == nil { + if exts, err := pggo.CollectStructsByPos[extRow](rows); err == nil { caps.ExtensionSchemas = make(map[string]string, len(exts)) for _, e := range exts { caps.Extensions = append(caps.Extensions, e.Name) @@ -299,8 +295,8 @@ func probe(ctx context.Context, cc *pgx.ConnConfig) (Capabilities, PoolerInfo, u // ReadOnlyTx runs fn inside its own short READ ONLY transaction and always rolls // back. Each collector sample gets a fresh transaction — that, plus // stats_fetch_consistency='none', is what keeps double-sampled rates non-zero. -func (t *Target) ReadOnlyTx(ctx context.Context, fn func(pgx.Tx) error) error { - tx, err := t.Pool.BeginTx(ctx, pgx.TxOptions{AccessMode: pgx.ReadOnly}) +func (t *Target) ReadOnlyTx(ctx context.Context, fn func(*pggo.Tx) error) error { + tx, err := t.Pool.BeginTx(ctx, pggo.TxOptions{ReadOnly: true}) if err != nil { return err } diff --git a/internal/conn/exclude_integration_test.go b/internal/conn/exclude_integration_test.go index 2e6e277..9be7fc8 100644 --- a/internal/conn/exclude_integration_test.go +++ b/internal/conn/exclude_integration_test.go @@ -5,7 +5,7 @@ import ( "os" "testing" - "github.com/jackc/pgx/v5" + "github.com/pgrundev/pggo" ) // Self-exclusion must key on backend PID, not the 'pgbot' application_name: a @@ -32,16 +32,16 @@ func TestIntegration_excludeSelf_byPIDNotLabel(t *testing.T) { } // An impostor: an external connection labelled application_name='pgbot'. - cfg, err := pgx.ParseConfig(d) + cfg, err := pggo.ParseConfig(d) if err != nil { t.Fatalf("parse dsn: %v", err) } cfg.RuntimeParams["application_name"] = "pgbot" - imp, err := pgx.ConnectConfig(ctx, cfg) + imp, err := pggo.ConnectConfig(ctx, cfg) if err != nil { t.Fatalf("impostor connect: %v", err) } - defer imp.Close(ctx) + defer imp.Close() var impPID uint32 if err := imp.QueryRow(ctx, "SELECT pg_backend_pid()").Scan(&impPID); err != nil { t.Fatalf("impostor pid: %v", err) diff --git a/internal/conn/pooler.go b/internal/conn/pooler.go index 9274914..02767f1 100644 --- a/internal/conn/pooler.go +++ b/internal/conn/pooler.go @@ -7,7 +7,7 @@ import ( "math/big" "strings" - "github.com/jackc/pgx/v5" + "github.com/pgrundev/pggo" ) // PoolerInfo records whether the connection routes through a transaction pooler @@ -18,7 +18,7 @@ import ( // zeroed rates that LOOK like a healthy run. We must detect this, not guess it. type PoolerInfo struct { Detected bool - SimpleProtocol bool // prepared statements failed; use the simple protocol + SimpleProtocol bool // named prepared statements failed (a transaction-pooler signal) Hint string // provider-specific fix, for the error message } @@ -33,23 +33,24 @@ type PoolerInfo struct { // pooler the two statements may land on different backends, so it doesn't. // // Host/port heuristics are only used to make the error message friendlier. -func detectPooler(ctx context.Context, c *pgx.Conn, cc *pgx.ConnConfig) PoolerInfo { +func detectPooler(ctx context.Context, c *pggo.Conn, cc *pggo.Config) PoolerInfo { info := PoolerInfo{Hint: poolerHint(cc)} // (1) prepared-statement probe psName := fmt.Sprintf("pgbot_ps_%d", nonce()) - if _, err := c.Prepare(ctx, psName, "SELECT 1"); err != nil { + if err := c.Prepare(ctx, psName, "SELECT 1"); err != nil { info.SimpleProtocol = true } else { _ = c.Deallocate(ctx, psName) } - // (2) session-persistence probe, forced through the simple protocol so a - // prepared-statement failure doesn't masquerade as non-persistence. + // (2) session-persistence probe. pgGo sends it as an unnamed statement in one + // Sync, which works even where named prepared statements fail, so a + // prepared-statement failure can't masquerade as non-persistence. want := fmt.Sprintf("pgbot_probe_%d", nonce()) if _, err := c.Exec(ctx, "SET application_name = '"+want+"'"); err == nil { var got string - err := c.QueryRow(ctx, "SELECT current_setting('application_name')", pgx.QueryExecModeSimpleProtocol).Scan(&got) + err := c.QueryRow(ctx, "SELECT current_setting('application_name')").Scan(&got) if err == nil && got != want { info.Detected = true } @@ -87,7 +88,7 @@ func detectPooler(ctx context.Context, c *pgx.Conn, cc *pgx.ConnConfig) PoolerIn // LOCAL inside an explicit transaction is forwarded verbatim), and shard 0 // exists in every deployment. pgdog.sharding_key must never be used here — on // some configs it errors and poisons the session for every later statement. -func detectPgDog(ctx context.Context, c *pgx.Conn) bool { +func detectPgDog(ctx context.Context, c *pggo.Conn) bool { want := fmt.Sprintf("pgbot_%d", nonce()) if _, err := c.Exec(ctx, "SET pgbot.probe = '"+want+"'"); err != nil { return false @@ -98,8 +99,7 @@ func detectPgDog(ctx context.Context, c *pgx.Conn) bool { } var control, shard *string err := c.QueryRow(ctx, - "SELECT current_setting('pgbot.probe', true), current_setting('pgdog.shard', true)", - pgx.QueryExecModeSimpleProtocol).Scan(&control, &shard) + "SELECT current_setting('pgbot.probe', true), current_setting('pgdog.shard', true)").Scan(&control, &shard) // Clean up with SET … TO DEFAULT, never RESET: SET of a dotted name is // accepted even by a backend that never saw the parameter, while RESET // errors there — and pgbot must not book server errors it would then report. @@ -121,11 +121,11 @@ func pgdogVerdict(control, shard *string, want string) bool { return shard == nil || *shard != "0" } -func isKnownPoolerEndpoint(cc *pgx.ConnConfig) bool { +func isKnownPoolerEndpoint(cc *pggo.Config) bool { return strings.Contains(strings.ToLower(cc.Host), "-pooler") || cc.Port == 6543 } -func poolerHint(cc *pgx.ConnConfig) string { +func poolerHint(cc *pggo.Config) string { host := strings.ToLower(cc.Host) switch { // "-pooler" alone is not Neon (issue #22): PgDog and self-hosted poolers diff --git a/internal/conn/pooler_integration_test.go b/internal/conn/pooler_integration_test.go index 54090b5..9dafd98 100644 --- a/internal/conn/pooler_integration_test.go +++ b/internal/conn/pooler_integration_test.go @@ -5,7 +5,7 @@ import ( "os" "testing" - "github.com/jackc/pgx/v5" + "github.com/pgrundev/pggo" ) // Direct PostgreSQL must never be identified as PgDog (issue #22): the probe's @@ -18,15 +18,15 @@ func TestIntegration_detectPgDog_notOnRealPostgres(t *testing.T) { t.Skip("set PGBOT_TEST_DSN to run integration tests") } ctx := context.Background() - cfg, err := pgx.ParseConfig(d) + cfg, err := pggo.ParseConfig(d) if err != nil { t.Fatalf("parse dsn: %v", err) } - c, err := pgx.ConnectConfig(ctx, cfg) + c, err := pggo.ConnectConfig(ctx, cfg) if err != nil { t.Fatalf("connect: %v", err) } - defer c.Close(ctx) + defer c.Close() if detectPgDog(ctx, c) { t.Error("real PostgreSQL misidentified as PgDog") @@ -56,15 +56,15 @@ func TestIntegration_detectPgDog_throughPgDog(t *testing.T) { t.Skip("set PGBOT_PGDOG_TEST_DSN (a DSN through a PgDog pooler) to run") } ctx := context.Background() - cfg, err := pgx.ParseConfig(d) + cfg, err := pggo.ParseConfig(d) if err != nil { t.Fatalf("parse dsn: %v", err) } - c, err := pgx.ConnectConfig(ctx, cfg) + c, err := pggo.ConnectConfig(ctx, cfg) if err != nil { t.Fatalf("connect: %v", err) } - defer c.Close(ctx) + defer c.Close() if !detectPgDog(ctx, c) { t.Error("PgDog not identified by the behavioral probe") diff --git a/internal/conn/pooler_test.go b/internal/conn/pooler_test.go index b1fbeac..4a10f6b 100644 --- a/internal/conn/pooler_test.go +++ b/internal/conn/pooler_test.go @@ -4,7 +4,7 @@ import ( "strings" "testing" - "github.com/jackc/pgx/v5" + "github.com/pgrundev/pggo" ) func TestIsKnownPoolerEndpoint(t *testing.T) { @@ -21,9 +21,9 @@ func TestIsKnownPoolerEndpoint(t *testing.T) { {"127.0.0.1", 6432, false}, // generic PgBouncer on a nonstandard port — undetectable by signature } for _, c := range cases { - cc := &pgx.ConnConfig{} + cc := &pggo.Config{} cc.Host = c.host - cc.Port = c.port + cc.Port = int(c.port) if got := isKnownPoolerEndpoint(cc); got != c.want { t.Errorf("isKnownPoolerEndpoint(%s:%d) = %v, want %v", c.host, c.port, got, c.want) } @@ -31,8 +31,8 @@ func TestIsKnownPoolerEndpoint(t *testing.T) { } func TestPoolerHint(t *testing.T) { - mk := func(host string, port uint16) *pgx.ConnConfig { - cc := &pgx.ConnConfig{} + mk := func(host string, port int) *pggo.Config { + cc := &pggo.Config{} cc.Host, cc.Port = host, port return cc } diff --git a/internal/conn/sshtunnel.go b/internal/conn/sshtunnel.go index 98cfbcb..c3b2ba8 100644 --- a/internal/conn/sshtunnel.go +++ b/internal/conn/sshtunnel.go @@ -4,9 +4,9 @@ package conn // DSN's host resolves to, which leaves out every database that only answers from // inside a bastion, a VPN-routed jump host, or a private VPC subnet. // -// The tunnel is installed as pgx's DialFunc rather than as a local port forward. -// That distinction matters: pgconn documents DialFunc as running BEFORE TLS is -// established, so the DSN keeps naming the REAL host all the way through. +// The tunnel is installed as the driver's DialFunc (pggo.Config.DialFunc) rather +// than as a local port forward. That distinction matters: DialFunc runs BEFORE +// TLS is negotiated, so the DSN keeps naming the REAL host all the way through. // sslmode=verify-full still validates against that hostname, and .pgpass still // matches on it. A `ssh -L` forward would force the DSN to say 127.0.0.1 and // silently break both, besides leaving a port open to every local user. @@ -77,8 +77,8 @@ func CloseSSHTunnel() { } // sshDialFunc returns a dialer that opens the database connection as a channel on -// the SSH connection, or nil when no tunnel is configured (pgx then keeps its own -// default dialer, timeouts included). +// the SSH connection, or nil when no tunnel is configured (the driver then keeps +// its own default dialer, timeouts included). func sshDialFunc() func(context.Context, string, string) (net.Conn, error) { if !SSHTunnelActive() { return nil @@ -559,8 +559,8 @@ func expandTilde(p string) string { return p } -// warnOnce prints a per-key diagnostic a single time. pgx dials more than once -// (probe, then each pool connection, then any fallback host), and repeating the +// warnOnce prints a per-key diagnostic a single time. The driver dials more than +// once (probe, then each pool connection, then any cancel request), and repeating the // same "skipping key" line four times reads like four different problems. var warned sync.Map diff --git a/internal/conn/sshtunnel_test.go b/internal/conn/sshtunnel_test.go index 8bee158..0aca054 100644 --- a/internal/conn/sshtunnel_test.go +++ b/internal/conn/sshtunnel_test.go @@ -46,7 +46,7 @@ func TestSplitTunnelSpec(t *testing.T) { } } -// A tunnel that was never configured must leave pgx's own dialer in place — +// A tunnel that was never configured must leave the driver's own dialer in place — // otherwise every direct connection would start paying for this feature. func TestSSHDialFunc_nilWhenUnconfigured(t *testing.T) { SetSSHTunnel("") diff --git a/internal/erd/introspect.go b/internal/erd/introspect.go index a2a5bef..970fcdd 100644 --- a/internal/erd/introspect.go +++ b/internal/erd/introspect.go @@ -4,13 +4,13 @@ import ( "context" "fmt" - "github.com/jackc/pgx/v5" + "github.com/pgrundev/pggo" ) -// Querier is the one pgx capability introspection needs (satisfied by -// *pgxpool.Pool and *pgx.Conn). +// Querier is the one driver capability introspection needs (satisfied by +// *pggo.Pool, *pggo.Conn and *pggo.Tx). type Querier interface { - Query(ctx context.Context, sql string, args ...any) (pgx.Rows, error) + Query(ctx context.Context, sql string, args ...any) (*pggo.Rows, error) } // Structure only — names, types, key membership. No table data is ever read; diff --git a/internal/pglog/sqlsource.go b/internal/pglog/sqlsource.go index cad21f5..4ba6036 100644 --- a/internal/pglog/sqlsource.go +++ b/internal/pglog/sqlsource.go @@ -6,13 +6,13 @@ import ( "fmt" "path" - "github.com/jackc/pgx/v5" + "github.com/pgrundev/pggo" ) -// RowQuerier is the one pgx capability the SQL source needs (satisfied by -// *pgxpool.Pool and *pgx.Conn). +// RowQuerier is the one driver capability the SQL source needs (satisfied by +// *pggo.Pool and *pggo.Conn). type RowQuerier interface { - QueryRow(ctx context.Context, sql string, args ...any) pgx.Row + QueryRow(ctx context.Context, sql string, args ...any) *pggo.Row } // ErrNoCollector means the server writes no logfile pgbot can address — From c453287711d81e02e7909d9c27362c46663e948b Mon Sep 17 00:00:00 2001 From: Alex Shapalov Date: Tue, 29 Sep 2026 20:02:11 -0700 Subject: [PATCH 2/3] Depend on the released pggo v0.1.0; add PostgreSQL 19 to the CI matrix Drops the local replace directive. The integration matrix gains 19beta1, where the pgGo migration was verified. Co-Authored-By: Claude Opus 5.5 (1M context) --- .github/workflows/ci.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 6020d0f..07eb5e6 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -185,7 +185,7 @@ jobs: strategy: fail-fast: false matrix: - pg: [14, 15, 16, 17, 18] + pg: [14, 15, 16, 17, 18, 19beta1] steps: - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 - uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0 From 86cb1fad5b14973fcd691f26d5c0226d6f4d119c Mon Sep 17 00:00:00 2001 From: Alex Shapalov Date: Tue, 29 Sep 2026 20:14:02 -0700 Subject: [PATCH 3/3] ci: generate the rate-guard write load with pgbench, not a DO loop Since PostgreSQL 15 a backend flushes its counters to pg_stat_database only when it goes idle, so the commits inside one long DO block stay invisible until the block ends. On 15+ the 'non-zero TPS under load' guard therefore only ever saw pgbot's own traffic: the pgx build passed by counting itself (and would have passed with the stats-caching bug it exists to catch); the pgGo build, which generates less self-traffic in the sample window, correctly saw 0 TPS. pgbench runs the same INSERTs as client transactions, which are counted. Co-Authored-By: Claude Opus 5.5 (1M context) --- .github/workflows/ci.yml | 15 ++++++++++----- 1 file changed, 10 insertions(+), 5 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 07eb5e6..189082f 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -214,9 +214,13 @@ jobs: PGBOT_TEST_SUPERUSER_DSN: postgres://postgres:pw@127.0.0.1:5432/postgres PGBOT_TEST_LOAD: "1" run: | - # concurrent write load so the stats-caching / rate guards are meaningful - docker exec pg psql -U postgres -q -c \ - "DO \$\$ BEGIN FOR i IN 1..200000 LOOP INSERT INTO people(email) VALUES('l@example.com'); IF i%50=0 THEN COMMIT; END IF; END LOOP; END \$\$;" & + # Concurrent write load so the stats-caching / rate guards are meaningful. + # Client transactions (pgbench), not a DO loop: since PG15 a backend + # flushes its counters to pg_stat_database only when it goes idle, so + # commits inside one long DO block stay invisible until it ends, and the + # guard would only have seen pgbot's own traffic. + echo "INSERT INTO people(email) VALUES ('l@example.com');" \ + | docker exec -i pg pgbench -U postgres -n -c 2 -T 180 -f /dev/stdin postgres >/dev/null & go test ./internal/collect/ -run Integration -v # cmd/pgbot integration: the advisor's read-only-blocks-a-write proof and # the MCP explain_plan/schema_of tools, against the same server. @@ -262,8 +266,9 @@ jobs: PGBOT_POOLER_DSN: postgres://postgres:pw@127.0.0.1:6543/postgres PGBOT_TEST_LOAD: "1" run: | - docker exec pg psql -U postgres -q -c \ - "DO \$\$ BEGIN FOR i IN 1..200000 LOOP INSERT INTO t(note) VALUES('x'); IF i%50=0 THEN COMMIT; END IF; END LOOP; END \$\$;" & + # client transactions, visible in pg_stat_database (see the integration job) + echo "INSERT INTO t(note) VALUES ('x');" \ + | docker exec -i pg pgbench -U postgres -n -c 2 -T 120 -f /dev/stdin postgres >/dev/null & go test ./internal/collect/ -run Integration_poolerRatesStayCorrect -v npm: