Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 11 additions & 6 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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:
Expand Down
2 changes: 1 addition & 1 deletion .golangci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
24 changes: 22 additions & 2 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions cmd/pgbot/activity.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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
}
Expand Down
34 changes: 17 additions & 17 deletions cmd/pgbot/advise.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand All @@ -264,27 +264,27 @@ 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
}

// hypo returns the schema-qualified name of a hypopg function. Like
// 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
}
Expand All @@ -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
Expand All @@ -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
}
Expand Down
12 changes: 6 additions & 6 deletions cmd/pgbot/advise_safety_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
}

Expand All @@ -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
})

Expand Down
6 changes: 3 additions & 3 deletions cmd/pgbot/helpers.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
}
8 changes: 4 additions & 4 deletions cmd/pgbot/logs.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)",
Expand Down
2 changes: 1 addition & 1 deletion cmd/pgbot/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
2 changes: 1 addition & 1 deletion cmd/pgbot/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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" {
Expand Down
14 changes: 7 additions & 7 deletions cmd/pgbot/mcp_tools.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -226,22 +226,22 @@ 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 {
return out
}
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)
}
Expand Down
Loading
Loading