Skip to content
Open
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
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -190,6 +190,10 @@ jobs:
-c "CREATE TABLE people(id serial primary key, email text);" \
-c "INSERT INTO people(email) SELECT 'user'||g||'@example.com' FROM generate_series(1,2000) g;" \
-c "CREATE ROLE pgbot_ro LOGIN PASSWORD 'ro'; GRANT pg_monitor TO pgbot_ro; GRANT CONNECT ON DATABASE postgres TO pgbot_ro;"
- name: install HypoPG for the generic-plan advisor fixture
if: matrix.pg >= 16
run: |
docker exec --user root pg sh -eu -c 'apt-get update; apt-get install -y --no-install-recommends postgresql-${{ matrix.pg }}-hypopg'
- name: run integration suite (full pipeline + PII + rates under load)
env:
PGBOT_TEST_DSN: postgres://pgbot_ro:ro@127.0.0.1:5432/postgres
Expand Down
5 changes: 5 additions & 0 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,11 @@ docker run -d -p 5432:5432 -e POSTGRES_PASSWORD=pw postgres:18 -c shared_preload
PGBOT_TEST_SUPERUSER_DSN=postgres://postgres:pw@127.0.0.1:5432/postgres go test ./internal/collect/ -run Integration
```

The advisor schema integration test additionally requires PostgreSQL 16+ and the
server-side HypoPG package (for example, `postgresql-18-hypopg`). CI installs it
before the integration suite; the test creates the extension in its own disposable
database and does not skip a missing extension on supported PostgreSQL versions.

CI also runs `gofmt -l`, `go vet`, `go test -race`, `golangci-lint`, and
`govulncheck`; keep them green.

Expand Down
7 changes: 4 additions & 3 deletions cmd/pgbot/advise.go
Original file line number Diff line number Diff line change
Expand Up @@ -276,15 +276,16 @@ type pgxPlanner struct {
func (p pgxPlanner) hypo(fn string) string { return p.caps.ExtObject("hypopg", fn) }

func (p pgxPlanner) 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
// GENERIC_PLAN plans a normalized $N query without values; VERBOSE preserves
// each resolved relation's schema in FORMAT JSON, which 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.
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().PgConn().Exec(ctx, "EXPLAIN (GENERIC_PLAN, VERBOSE, FORMAT JSON) "+query).ReadAll()
if err != nil {
return err
}
Expand Down
99 changes: 99 additions & 0 deletions cmd/pgbot/advise_safety_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,13 @@ package main

import (
"context"
"fmt"
"os"
"testing"
"time"

"github.com/jackc/pgx/v5"
"github.com/pgrundev/pgbot/internal/advisor"
"github.com/pgrundev/pgbot/internal/conn"
)

Expand Down Expand Up @@ -62,3 +65,99 @@ func TestIntegration_adviseSafety_readOnlyTxBlocksInjectedWrite(t *testing.T) {
}
_, _ = admin.Exec(ctx, `DROP TABLE IF EXISTS advise_safety`)
}

func TestIntegration_advisorTargetsResolvedNonPublicSchema(t *testing.T) {
superuserDSN := os.Getenv("PGBOT_TEST_SUPERUSER_DSN")
if superuserDSN == "" {
t.Skip("set PGBOT_TEST_SUPERUSER_DSN to run the advisor schema integration test")
}
ctx, cancel := context.WithTimeout(context.Background(), 45*time.Second)
defer cancel()

admin, err := pgx.Connect(ctx, superuserDSN)
if err != nil {
t.Fatalf("admin connect: %v", err)
}
defer admin.Close(context.Background())
var serverVersion int
if err := admin.QueryRow(ctx, `SELECT current_setting('server_version_num')::int`).Scan(&serverVersion); err != nil {
t.Fatal(err)
}
if serverVersion < 160000 {
t.Skip("generic-plan advisor requires PostgreSQL 16+")
}
dbName := fmt.Sprintf("pgbot_advisor_schema_%d", time.Now().UnixNano())
quotedDB := pgx.Identifier{dbName}.Sanitize()
dropDatabase := func() {
cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cleanupCancel()
if _, err := admin.Exec(cleanupCtx, `DROP DATABASE `+quotedDB+` WITH (FORCE)`); err != nil {
t.Errorf("remove owned fixture database: %v", err)
}
}
if _, err := admin.Exec(ctx, `CREATE DATABASE `+quotedDB); err != nil {
t.Fatalf("create fixture database: %v", err)
}
defer dropDatabase()

dsn := swapDatabase(t, superuserDSN, dbName)
fixture, err := pgx.Connect(ctx, dsn)
if err != nil {
t.Fatalf("fixture connect: %v", err)
}
defer fixture.Close(ctx)

for _, stmt := range []string{
`CREATE EXTENSION IF NOT EXISTS hypopg`,
`CREATE SCHEMA pgbot_advise_schema_it`,
`CREATE TABLE pgbot_advise_schema_it.pgbot_same_name_it (customer_id integer NOT NULL, payload text)`,
`CREATE TABLE public.pgbot_same_name_it (customer_id integer NOT NULL, payload text)`,
`INSERT INTO pgbot_advise_schema_it.pgbot_same_name_it SELECT n, repeat('x', 80) FROM generate_series(1, 50000) AS n`,
`INSERT INTO public.pgbot_same_name_it VALUES (1, 'public control')`,
`ANALYZE pgbot_advise_schema_it.pgbot_same_name_it`,
`ANALYZE public.pgbot_same_name_it`,
} {
if _, err := fixture.Exec(ctx, stmt); err != nil {
t.Fatalf("fixture statement failed: %v", err)
}
}

target, err := conn.Connect(ctx, dsn)
if err != nil {
t.Fatalf("pgbot connect: %v", err)
}
defer target.Close()
if !target.Caps.HasHypopg {
t.Fatal("fixture installed hypopg but pgbot did not detect it")
}

var recs []advisor.Recommendation
err = target.ReadOnlyTx(ctx, func(tx pgx.Tx) error {
planner := pgxPlanner{tx: tx, caps: target.Caps}
defer func() {
cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cleanupCancel()
_ = planner.ResetHypo(cleanupCtx)
}()
recs, _ = advisor.Advise(ctx, planner, []advisor.QueryInput{{
QueryID: 1,
Text: "SELECT * FROM pgbot_advise_schema_it.pgbot_same_name_it " +
"WHERE customer_id = $1",
Scrubbed: "SELECT * FROM pgbot_advise_schema_it.pgbot_same_name_it " +
"WHERE customer_id = $1",
Calls: 100,
SharePct: 80,
}}, advisor.Options{MinImprovement: 0.5})
return nil
})
if err != nil {
t.Fatalf("advisor transaction: %v", err)
}
if len(recs) != 1 {
t.Fatalf("expected one non-public recommendation, got %+v", recs)
}
const want = "CREATE INDEX ON pgbot_advise_schema_it.pgbot_same_name_it (customer_id)"
if recs[0].Schema != "pgbot_advise_schema_it" || recs[0].IndexDDL != want {
t.Fatalf("advisor targeted the wrong same-named relation: %+v", recs[0])
}
}
4 changes: 3 additions & 1 deletion internal/advisor/advisor.go
Original file line number Diff line number Diff line change
Expand Up @@ -133,9 +133,11 @@ func candidatesFromPlan(root planNode) []Candidate {
if n.NodeType != "Seq Scan" || n.RelationName == "" || n.Filter == "" {
return
}
// PostgreSQL emits Schema only for a VERBOSE plan. Without it there is
// no safe way to distinguish same-named relations resolved by search_path.
schema := n.Schema
if schema == "" {
schema = "public"
return
}
if systemSchemas[schema] {
return
Expand Down
29 changes: 29 additions & 0 deletions internal/advisor/advisor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,35 @@ func TestCandidatesFromPlan_seqScanOnly(t *testing.T) {
}
}

func TestCandidatesFromPlan_requiresResolvedSchema(t *testing.T) {
t.Run("schema-less plan fails closed", func(t *testing.T) {
root, err := parsePlan([]byte(`[{
"Plan":{"Node Type":"Seq Scan","Relation Name":"orders",
"Filter":"(customer_id = $1)","Total Cost":4040.0}
}]`))
if err != nil {
t.Fatal(err)
}
if got := candidatesFromPlan(root); len(got) != 0 {
t.Fatalf("schema-less plan must not invent a target schema, got %+v", got)
}
})

t.Run("verbose qualified filter retains non-public schema", func(t *testing.T) {
root, err := parsePlan([]byte(`[{
"Plan":{"Node Type":"Seq Scan","Relation Name":"orders","Schema":"tenant_data",
"Filter":"(orders.customer_id = $1)","Total Cost":4040.0}
}]`))
if err != nil {
t.Fatal(err)
}
got := candidatesFromPlan(root)
if len(got) != 1 || got[0].DDL() != "CREATE INDEX ON tenant_data.orders (customer_id)" {
t.Fatalf("wrong candidate from verbose plan: %+v", got)
}
})
}

func TestUsesIndex(t *testing.T) {
js := []byte(`[{"Plan":{"Node Type":"Index Scan","Index Name":"<13337>btree_orders_customer_id","Total Cost":52.0}}]`)
root, _ := parsePlan(js)
Expand Down
2 changes: 1 addition & 1 deletion internal/advisor/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ import (
// plan-only or hypothetical operation — none executes the inspected query or
// writes anything.
type Planner interface {
// GenericPlan returns EXPLAIN (GENERIC_PLAN, FORMAT JSON) <query> as raw JSON.
// GenericPlan returns EXPLAIN (GENERIC_PLAN, VERBOSE, FORMAT JSON) <query> as raw JSON.
// GENERIC_PLAN (PG16+) plans a normalized $N query without values or execution.
GenericPlan(ctx context.Context, query string) ([]byte, error)
// CreateHypoIndex creates a hypothetical index from a CREATE INDEX statement
Expand Down
Loading