From 5c23db8885c7070f2aee34ea73ac948782aaf695 Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Mon, 31 Aug 2026 22:44:10 +1000 Subject: [PATCH 1/2] Add read/write splitting to the pg module MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit pg.Cluster wraps a primary pool and an optional read-replica pool. Routing is explicit — Read(), Write(), Primary() — with no SQL inspection: statement text misroutes CTE writes and SELECT ... FOR UPDATE in both directions. An unhealthy replica falls back to the primary behind a ping-based circuit that retries on a doubling backoff window, and migrations always run through Write(). pg.ClusterDSNsFromEnv resolves both endpoints from the environment, naming the read-only host explicitly rather than deriving it from the primary's. --- README.md | 107 +++++++++++++- pg/cluster.go | 239 ++++++++++++++++++++++++++++++ pg/cluster_test.go | 319 +++++++++++++++++++++++++++++++++++++++++ pg/dsn.go | 106 +++++++++++--- pg/dsn_test.go | 162 +++++++++++++++++++++ pg/integration_test.go | 53 +++++++ pg/pg.go | 5 +- 7 files changed, 971 insertions(+), 20 deletions(-) create mode 100644 pg/cluster.go create mode 100644 pg/cluster_test.go diff --git a/README.md b/README.md index 29fa7a4..a9fe51f 100644 --- a/README.md +++ b/README.md @@ -11,7 +11,7 @@ imports it inherits them. | Import | What it does | | --- | --- | -| `git.unkin.net/unkin/golib/pg` | Postgres: DSN from the environment, pgxpool construction, and the estate's migration runner. | +| `git.unkin.net/unkin/golib/pg` | Postgres: DSN from the environment, pgxpool construction, read/write splitting across a primary and a replica, and the estate's migration runner. | | `git.unkin.net/unkin/golib/pg/pgtest` | A throwaway Postgres container for a consumer's own `_test.go` files. Test-only. | ### pg @@ -65,6 +65,111 @@ rather than on the first query. `pg.NewMigrated` does that and then migrates. `pg.LockKey(name)` exposes the derivation, so a service migrating off a hardcoded key can assert the two agree before switching over. +#### Read/write splitting + +`pg.Cluster` wraps a primary pool and an optional read-replica pool: + +```go +primaryDSN, replicaDSN, err := pg.ClusterDSNsFromEnv("ENCAPI_") +if err != nil { + return err +} +db, err := pg.NewCluster(ctx, pg.ClusterConfig{ + PrimaryDSN: primaryDSN, + ReplicaDSN: replicaDSN, + Logger: log, +}) +if err != nil { + return err +} +defer db.Close() + +if err := db.Migrate(ctx, migrations.FS, pg.MigrateOptions{LockName: "encapi-migrations"}); err != nil { + return err +} + +rows, err := db.Read().Query(ctx, "SELECT id, name FROM nodes") // replica when healthy +tag, err := db.Write().Exec(ctx, "UPDATE nodes SET ...") // always the primary +row := db.Primary().QueryRow(ctx, "SELECT ... WHERE id = $1", id) // read-your-writes +``` + +| Call | Goes to | +| --- | --- | +| `Write()` | The primary, always. | +| `Read()` | The replica when one is configured and healthy; the primary otherwise. | +| `Primary()` | The primary. Same pool as `Write()`, named for reads that must not be stale. | + +**Routing is explicit; nothing inspects SQL.** Deciding from the statement text +gets it wrong in both directions: a CTE with an `INSERT` in it reads as a +`SELECT`, and `SELECT ... FOR UPDATE` takes row locks a replica cannot grant. +The caller knows which it wants, so the caller picks. + +**Replica lag is real.** A replica serves a slightly stale snapshot, so a read +that has to observe a write this process just made goes to `Primary()`. The +usual shape is a handler that writes and then re-reads what it wrote, or a +redirect straight into a `GET` of the row just created — both need the primary. +Everything else (list endpoints, dashboards, reports, background aggregation) +can take the replica. + +**A missing replica is not an error.** With `ReplicaDSN` empty, `Read()` returns +the primary and the `Cluster` is an ordinary single-pool handle, so a service +can use `Cluster` unconditionally and let the deployment decide whether reads +split. A `ReplicaDSN` equal to `PrimaryDSN` does the same rather than opening a +second pool to the same place. + +**An unhealthy replica falls back.** The `Cluster` keeps a small circuit over +the replica: + +- A replica that fails its startup ping does not fail startup; reads begin on + the primary and move over once it answers a probe. +- `db.ReportReplicaError(err)` — pass the error from a query run on the pool + `Read()` handed out — trips the circuit, and reads move to the primary. +- While tripped, `Read()` issues no probes until the backoff window expires; the + next `Read()` after that pays for one probe that either promotes the replica + or doubles the window. The window runs from `ReplicaRetryMin` to + `ReplicaRetryMax` (5s to 2m by default) and resets on recovery. A healthy + replica is never probed at all, so the split costs nothing on the read path. +- Reporting is advisory. A caller that never reports still routes correctly; it + just does not react to a replica that dies mid-flight. + +**Migrations always target the primary.** `Cluster.Migrate` runs through +`Write()`. A replica is physically read-only, and a schema change has to +originate on the primary to reach the replica at all. + +`pg.ClusterDSNsFromEnv(prefix)` returns both connection strings. The primary +follows `DSNFromEnv` exactly. The replica resolves, highest first: + +1. `DATABASE_RO_URL` — used verbatim. +2. `DATABASE_RO_URL` — likewise. +3. `DB_RO_HOST`, or bare `DB_RO_HOST` — the primary's port, user, + password, database and sslmode with that host substituted. +4. Nothing set — an empty replica DSN, so the `Cluster` runs single-pool. + +Nothing is derived. The read-only host is never rewritten out of the primary's, +because a wrong guess silently sends reads somewhere unintended. Setting only +`DB_RO_HOST` while the primary comes from a whole `DATABASE_URL` is an error, +not a guess: the fields to substitute into are not known. + +CloudNativePG publishes exactly the two endpoints this expects — `-rw` +routes to the primary and `-ro` to the replicas — so a Deployment names +both: + +```yaml +env: + - name: ENCAPI_DBHOST + value: encapi-db-rw + - name: ENCAPI_DB_RO_HOST + value: encapi-db-ro + - name: ENCAPI_DBNAME + value: encapi + - name: ENCAPI_DBUSER + valueFrom: { secretKeyRef: { name: encapi-db-app, key: username } } + - name: ENCAPI_DBPASS + valueFrom: { secretKeyRef: { name: encapi-db-app, key: password } } +``` + +Dropping `ENCAPI_DB_RO_HOST` turns the split off without a code change. + ### pgtest `pgtest` starts `postgres:17-alpine` via testcontainers. Import it only from diff --git a/pg/cluster.go b/pg/cluster.go new file mode 100644 index 0000000..f8163a5 --- /dev/null +++ b/pg/cluster.go @@ -0,0 +1,239 @@ +package pg + +import ( + "context" + "fmt" + "io/fs" + "log/slog" + "sync" + "time" + + "github.com/jackc/pgx/v5/pgxpool" +) + +// Replica health defaults. The retry window is deliberately coarse: a replica +// that just fell over is not coming back within a request, and probing it more +// often only moves the failure onto the read path. +const ( + defaultReplicaRetryMin = 5 * time.Second + defaultReplicaRetryMax = 2 * time.Minute + replicaProbeTimeout = 5 * time.Second +) + +// ClusterConfig configures a Cluster. +type ClusterConfig struct { + // PrimaryDSN is the read-write endpoint. Required. + PrimaryDSN string + + // ReplicaDSN is the read-only endpoint. Empty puts the Cluster in + // single-pool mode, where Read returns the primary. A ReplicaDSN equal to + // PrimaryDSN does the same rather than opening a second pool to the same + // place. + ReplicaDSN string + + // Logger receives replica health transitions. Nil discards them. + Logger *slog.Logger + + // ReplicaRetryMin and ReplicaRetryMax bound the backoff window between + // probes of an unhealthy replica; the window doubles from Min up to Max. + // Zero means the defaults, 5s and 2m. + ReplicaRetryMin time.Duration + ReplicaRetryMax time.Duration +} + +// Cluster routes queries across a primary and an optional read replica. +// +// Routing is explicit: the caller picks Read or Write per query. Nothing +// inspects SQL to decide, because statement inspection gets it wrong in both +// directions — a CTE with an INSERT in it reads as a SELECT, and +// SELECT ... FOR UPDATE takes row locks a replica cannot grant. +// +// A Cluster with no replica is a working single-pool Cluster, so a service can +// use it unconditionally and a deployment decides whether reads are split. +type Cluster struct { + primary *pgxpool.Pool + // replica is nil in single-pool mode. + replica *pgxpool.Pool + health *replicaCircuit +} + +// NewCluster opens the primary and, when configured, the replica. +// +// A primary that cannot be reached is fatal: the service has nowhere to write. +// A replica that cannot be reached is not — the Cluster starts with reads on +// the primary and promotes the replica once it answers a probe, which is the +// same degradation a replica failing later gets. +// +// The caller owns the Cluster and must Close it. +func NewCluster(ctx context.Context, cfg ClusterConfig) (*Cluster, error) { + primary, err := New(ctx, cfg.PrimaryDSN, cfg.Logger) + if err != nil { + return nil, err + } + c, err := newCluster(ctx, primary, cfg) + if err != nil { + primary.Close() + return nil, err + } + return c, nil +} + +// newCluster attaches the replica half to an already-open primary pool. +func newCluster(ctx context.Context, primary *pgxpool.Pool, cfg ClusterConfig) (*Cluster, error) { + log := logger(cfg.Logger) + c := &Cluster{primary: primary} + if cfg.ReplicaDSN == "" || cfg.ReplicaDSN == cfg.PrimaryDSN { + log.Debug("postgres cluster in single-pool mode, reads go to the primary") + return c, nil + } + + replica, err := pgxpool.New(ctx, cfg.ReplicaDSN) + if err != nil { + return nil, fmt.Errorf("connect postgres replica: %w", err) + } + c.replica = replica + c.health = newReplicaCircuit(replica, log, cfg.ReplicaRetryMin, cfg.ReplicaRetryMax) + if err := replica.Ping(ctx); err != nil { + c.health.trip(err) + } else { + log.Debug("postgres replica pool ready", "host", replica.Config().ConnConfig.Host) + } + return c, nil +} + +// Write returns the primary pool. Every statement that changes data, takes row +// locks, or must be read back immediately goes here. +func (c *Cluster) Write() *pgxpool.Pool { return c.primary } + +// Read returns the pool to run a read-only query on: the replica when one is +// configured and healthy, the primary otherwise. +// +// Reads served by a replica see a slightly stale snapshot. Use Primary for a +// read that must observe a write this process just made. +func (c *Cluster) Read() *pgxpool.Pool { + if c.replica == nil || !c.health.use() { + return c.primary + } + return c.replica +} + +// Primary returns the primary pool, whatever the replica's state. It is the +// same pool as Write, named for the read-your-writes case: a read that must not +// be served stale asks for the primary explicitly. +func (c *Cluster) Primary() *pgxpool.Pool { return c.primary } + +// ReportReplicaError tells the Cluster a query on the replica failed, so reads +// move to the primary until the replica answers a probe again. Pass the error +// from a query run on the pool Read returned; a nil error, or a Cluster with no +// replica, is a no-op. +// +// It is advisory. A caller that does not report anything still gets correct +// routing, just no reaction to a replica that dies mid-flight. +func (c *Cluster) ReportReplicaError(err error) { + if err == nil || c.replica == nil { + return + } + c.health.trip(err) +} + +// Migrate brings the schema up to date on the primary. Migrations never run +// against a replica: a replica is physically read-only, and a schema change has +// to originate on the primary to reach it at all. +func (c *Cluster) Migrate(ctx context.Context, fsys fs.FS, opts MigrateOptions) error { + return Migrate(ctx, c.Write(), fsys, opts) +} + +// Close closes both pools. +func (c *Cluster) Close() { + if c.replica != nil { + c.replica.Close() + } + c.primary.Close() +} + +// pinger is the slice of *pgxpool.Pool the health circuit drives, so the +// circuit is testable without a database. +type pinger interface { + Ping(ctx context.Context) error +} + +// replicaCircuit tracks whether the replica is worth sending reads to. +// +// While healthy it costs nothing: Read consults a boolean and no probe is +// issued. A reported error or a failed startup ping trips it, and reads go to +// the primary until the backoff window expires, at which point the next Read +// pays for one probe that either promotes the replica or extends the window. +type replicaCircuit struct { + probe pinger + log *slog.Logger + now func() time.Time + min time.Duration + max time.Duration + + mu sync.Mutex + tripped bool + retryAt time.Time + backoff time.Duration +} + +func newReplicaCircuit(probe pinger, log *slog.Logger, retryMin, retryMax time.Duration) *replicaCircuit { + if retryMin <= 0 { + retryMin = defaultReplicaRetryMin + } + if retryMax < retryMin { + retryMax = defaultReplicaRetryMax + } + if retryMax < retryMin { + retryMax = retryMin + } + return &replicaCircuit{probe: probe, log: log, now: time.Now, min: retryMin, max: retryMax} +} + +// use reports whether the replica may serve the next read, probing it when the +// backoff window has expired. The probe runs under the mutex so a burst of +// concurrent reads issues one probe between them rather than one each. +func (c *replicaCircuit) use() bool { + c.mu.Lock() + defer c.mu.Unlock() + if !c.tripped { + return true + } + if c.now().Before(c.retryAt) { + return false + } + + ctx, cancel := context.WithTimeout(context.Background(), replicaProbeTimeout) + defer cancel() + if err := c.probe.Ping(ctx); err != nil { + c.log.Debug("postgres replica still unhealthy", "err", err) + c.backOffLocked() + return false + } + c.log.Info("postgres replica healthy again, reads return to it") + c.tripped = false + c.backoff = 0 + return true +} + +// trip marks the replica unhealthy. An already-tripped circuit keeps the window +// it has, so a flood of reported errors cannot push the retry out indefinitely. +func (c *replicaCircuit) trip(err error) { + c.mu.Lock() + defer c.mu.Unlock() + if c.tripped { + return + } + c.log.Warn("postgres replica unhealthy, reads move to the primary", "err", err) + c.backOffLocked() +} + +func (c *replicaCircuit) backOffLocked() { + c.tripped = true + switch c.backoff { + case 0: + c.backoff = c.min + default: + c.backoff = min(c.backoff*2, c.max) + } + c.retryAt = c.now().Add(c.backoff) +} diff --git a/pg/cluster_test.go b/pg/cluster_test.go new file mode 100644 index 0000000..0ee7e1f --- /dev/null +++ b/pg/cluster_test.go @@ -0,0 +1,319 @@ +package pg + +import ( + "context" + "errors" + "strings" + "testing" + "testing/fstest" + "time" + + "github.com/jackc/pgx/v5/pgxpool" +) + +// A second unreachable DSN, distinct from deadDSN, so a test can tell the two +// pools of a Cluster apart by pointer. +const deadReplicaDSN = "postgres://u:p@127.0.0.1:2/d?sslmode=disable&connect_timeout=2" + +// fakePinger answers probes from a scripted sequence of results, so recovery +// can be driven without a database. +type fakePinger struct { + results []error + calls int +} + +func (p *fakePinger) Ping(context.Context) error { + p.calls++ + if len(p.results) == 0 { + return nil + } + i := min(p.calls-1, len(p.results)-1) + return p.results[i] +} + +// fakeClock is a manually advanced clock for the backoff window. +type fakeClock struct{ t time.Time } + +func (c *fakeClock) now() time.Time { return c.t } +func (c *fakeClock) advance(d time.Duration) { c.t = c.t.Add(d) } + +// lazyPool opens a pool without connecting: pgxpool is lazy, so this needs no +// server and lets a test build a Cluster with two distinguishable pools. +func lazyPool(t *testing.T, dsn string) *pgxpool.Pool { + t.Helper() + pool, err := pgxpool.New(context.Background(), dsn) + if err != nil { + t.Fatalf("pgxpool.New(%q): %v", dsn, err) + } + t.Cleanup(pool.Close) + return pool +} + +// testCluster builds a Cluster with a replica and a circuit under a fake clock. +func testCluster(t *testing.T, probe pinger) (*Cluster, *fakeClock) { + t.Helper() + c := &Cluster{ + primary: lazyPool(t, deadDSN), + replica: lazyPool(t, deadReplicaDSN), + } + clock := &fakeClock{t: time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC)} + c.health = newReplicaCircuit(probe, testLogger(), time.Second, 4*time.Second) + c.health.now = clock.now + return c, clock +} + +// The routing table: what Read and Write return in each replica state. +func TestCluster_Routing(t *testing.T) { + single := &Cluster{primary: lazyPool(t, deadDSN)} + if single.Read() != single.primary { + t.Error("with no replica, Read must return the primary") + } + if single.Write() != single.primary || single.Primary() != single.primary { + t.Error("Write and Primary must return the primary") + } + + c, clock := testCluster(t, &fakePinger{}) + if c.Read() != c.replica { + t.Fatal("a healthy replica must serve reads") + } + if c.Write() != c.primary { + t.Fatal("writes must always go to the primary") + } + if c.Primary() != c.primary { + t.Fatal("Primary must return the primary even with a healthy replica") + } + + c.ReportReplicaError(errors.New("replica exploded")) + if c.Read() != c.primary { + t.Fatal("an unhealthy replica must not serve reads") + } + + // Inside the backoff window the decision stands without a probe. + clock.advance(500 * time.Millisecond) + if c.Read() != c.primary { + t.Fatal("reads must stay on the primary inside the backoff window") + } + + // Once the window expires the next read probes, and a healthy probe + // promotes the replica again. + clock.advance(time.Second) + if c.Read() != c.replica { + t.Fatal("a recovered replica must serve reads again") + } +} + +// A nil error and a Cluster without a replica must both be no-ops rather than a +// nil-pointer dereference on the circuit. +func TestCluster_ReportReplicaErrorIgnoresNoOps(t *testing.T) { + single := &Cluster{primary: lazyPool(t, deadDSN)} + single.ReportReplicaError(errors.New("no replica to blame")) + + c, _ := testCluster(t, &fakePinger{}) + c.ReportReplicaError(nil) + if c.Read() != c.replica { + t.Fatal("a nil error must not trip the circuit") + } +} + +func TestReplicaCircuit_BackoffDoublesToTheCeiling(t *testing.T) { + probe := &fakePinger{results: []error{errors.New("down")}} + clock := &fakeClock{t: time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC)} + c := newReplicaCircuit(probe, testLogger(), time.Second, 4*time.Second) + c.now = clock.now + + c.trip(errors.New("down")) + // 1s, then 2s, then 4s, then held at the 4s ceiling. + for _, want := range []time.Duration{time.Second, 2 * time.Second, 4 * time.Second, 4 * time.Second} { + if c.backoff != want { + t.Fatalf("backoff = %s, want %s", c.backoff, want) + } + // Just short of the window: no probe, still tripped. + before := probe.calls + clock.advance(want - time.Millisecond) + if c.use() { + t.Fatal("the circuit reopened before its window expired") + } + if probe.calls != before { + t.Fatal("the circuit probed inside its window") + } + clock.advance(time.Millisecond) + if c.use() { + t.Fatal("a failing probe must leave the circuit closed") + } + if probe.calls != before+1 { + t.Fatalf("probe calls = %d, want %d", probe.calls, before+1) + } + } +} + +// Reported errors while already tripped must not push the retry out: a busy +// service reporting on every read would otherwise never probe again. +func TestReplicaCircuit_RepeatedTripsKeepTheWindow(t *testing.T) { + clock := &fakeClock{t: time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC)} + c := newReplicaCircuit(&fakePinger{}, testLogger(), time.Second, 4*time.Second) + c.now = clock.now + + c.trip(errors.New("first")) + retryAt := c.retryAt + clock.advance(900 * time.Millisecond) + c.trip(errors.New("second")) + if !c.retryAt.Equal(retryAt) { + t.Fatalf("retryAt moved from %s to %s on a repeated trip", retryAt, c.retryAt) + } +} + +// After a recovery the backoff restarts at the minimum rather than resuming +// where the previous outage left off. +func TestReplicaCircuit_RecoveryResetsTheBackoff(t *testing.T) { + probe := &fakePinger{results: []error{errors.New("down"), nil}} + clock := &fakeClock{t: time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC)} + c := newReplicaCircuit(probe, testLogger(), time.Second, 4*time.Second) + c.now = clock.now + + c.trip(errors.New("down")) + clock.advance(time.Second) + if c.use() { + t.Fatal("the first probe fails, so the circuit stays closed") + } + clock.advance(2 * time.Second) + if !c.use() { + t.Fatal("the second probe succeeds, so the circuit must reopen") + } + if c.backoff != 0 { + t.Fatalf("backoff = %s after recovery, want 0", c.backoff) + } + + c.trip(errors.New("down again")) + if c.backoff != time.Second { + t.Fatalf("backoff = %s on the next outage, want the 1s minimum", c.backoff) + } +} + +func TestNewReplicaCircuit_Defaults(t *testing.T) { + c := newReplicaCircuit(&fakePinger{}, testLogger(), 0, 0) + if c.min != defaultReplicaRetryMin || c.max != defaultReplicaRetryMax { + t.Fatalf("defaults = %s/%s, want %s/%s", c.min, c.max, defaultReplicaRetryMin, defaultReplicaRetryMax) + } + // A minimum above the default ceiling must not produce max < min. + c = newReplicaCircuit(&fakePinger{}, testLogger(), time.Hour, 0) + if c.max < c.min { + t.Fatalf("max %s is below min %s", c.max, c.min) + } +} + +func TestNewCluster_SinglePoolModes(t *testing.T) { + tests := []struct { + name string + cfg ClusterConfig + }{ + {"no replica configured", ClusterConfig{PrimaryDSN: deadDSN}}, + {"replica DSN equal to the primary", ClusterConfig{PrimaryDSN: deadDSN, ReplicaDSN: deadDSN}}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + c, err := newCluster(shortCtx(t), lazyPool(t, deadDSN), tc.cfg) + if err != nil { + t.Fatalf("newCluster: %v", err) + } + if c.replica != nil { + t.Fatal("expected single-pool mode") + } + if c.Read() != c.Write() { + t.Fatal("in single-pool mode Read and Write must be the same pool") + } + }) + } +} + +// A replica that will not answer must degrade to primary reads, not fail +// startup: the service can still serve everything, just without the split. +func TestNewCluster_UnreachableReplicaStartsDegraded(t *testing.T) { + c, err := newCluster(shortCtx(t), lazyPool(t, deadDSN), ClusterConfig{ + PrimaryDSN: deadDSN, + ReplicaDSN: deadReplicaDSN, + Logger: testLogger(), + }) + if err != nil { + t.Fatalf("newCluster: %v", err) + } + t.Cleanup(func() { c.replica.Close() }) + if c.replica == nil { + t.Fatal("the replica pool must still be opened") + } + if !c.health.tripped { + t.Fatal("a replica that failed its startup ping must start tripped") + } + if c.Read() != c.primary { + t.Fatal("reads must start on the primary while the replica is down") + } +} + +func TestNewCluster_RejectsAnUnparseableReplicaDSN(t *testing.T) { + _, err := newCluster(shortCtx(t), lazyPool(t, deadDSN), ClusterConfig{ + PrimaryDSN: deadDSN, + ReplicaDSN: "://not a dsn", + }) + if err == nil { + t.Fatal("expected an error for an unparseable replica DSN") + } + if !strings.Contains(err.Error(), "connect postgres replica") { + t.Fatalf("error %q does not identify the failing step", err) + } +} + +func TestNewCluster_PropagatesPrimaryFailure(t *testing.T) { + c, err := NewCluster(shortCtx(t), ClusterConfig{PrimaryDSN: deadDSN, ReplicaDSN: deadReplicaDSN}) + if err == nil { + c.Close() + t.Fatal("expected NewCluster to fail against an unreachable primary") + } + if !strings.Contains(err.Error(), "ping postgres") { + t.Fatalf("error %q does not identify the failing step", err) + } +} + +// Migrations must reach the primary even when the replica is the pool Read +// would hand out. A closed replica pool fails Acquire with a distinctive error, +// so routing the run to it would be visible here. +func TestCluster_MigrateTargetsThePrimary(t *testing.T) { + c, _ := testCluster(t, &fakePinger{}) + if c.Read() != c.replica { + t.Fatal("the replica must be healthy for this test to mean anything") + } + c.replica.Close() + + err := c.Migrate(shortCtx(t), testFS(), MigrateOptions{LockName: testLockName}) + if err == nil { + t.Fatal("expected the migration to fail against an unreachable primary") + } + if strings.Contains(err.Error(), "closed pool") { + t.Fatalf("the migration ran against the replica: %v", err) + } + if !strings.Contains(err.Error(), "acquire migration connection") { + t.Fatalf("error %q does not identify the failing step", err) + } +} + +func TestCluster_MigrateRequiresALockName(t *testing.T) { + c := &Cluster{primary: lazyPool(t, deadDSN)} + if err := c.Migrate(shortCtx(t), fstest.MapFS{}, MigrateOptions{}); err == nil { + t.Fatal("expected Migrate to reject an empty LockName") + } +} + +func TestCluster_CloseClosesBothPools(t *testing.T) { + c := &Cluster{ + primary: lazyPool(t, deadDSN), + replica: lazyPool(t, deadReplicaDSN), + } + c.health = newReplicaCircuit(c.replica, testLogger(), 0, 0) + c.Close() + + // Acquiring from a closed pool fails immediately; a live one would dial. + for name, pool := range map[string]*pgxpool.Pool{"primary": c.primary, "replica": c.replica} { + if _, err := pool.Acquire(context.Background()); err == nil || !strings.Contains(err.Error(), "closed pool") { + t.Errorf("%s pool was not closed: %v", name, err) + } + } + // Close is registered again by lazyPool's cleanup; pgxpool tolerates it. +} diff --git a/pg/dsn.go b/pg/dsn.go index a97321c..28d5fd8 100644 --- a/pg/dsn.go +++ b/pg/dsn.go @@ -32,6 +32,10 @@ var ( fieldPass = dsnField{"DBPASS", "PGPASSWORD"} fieldName = dsnField{"DBNAME", "PGDATABASE"} fieldSSL = dsnField{"DBSSL", "PGSSLMODE"} + + // The replica host has no libpq counterpart; the second name is the bare + // variable a prefixed lookup falls back to. + fieldROHost = dsnField{"DB_RO_HOST", "DB_RO_HOST"} ) // DSNFromEnv builds a libpq/pgx connection string from the environment. @@ -58,35 +62,103 @@ var ( // User and database name have no default: DSNFromEnv reports an error naming // the variables it looked at rather than connecting somewhere unintended. func DSNFromEnv(prefix string) (string, error) { - if v := os.Getenv(prefix + "DATABASE_URL"); v != "" { + if v := urlFromEnv(prefix, "DATABASE_URL"); v != "" { return v, nil } - if v := os.Getenv("DATABASE_URL"); v != "" { - return v, nil + p, err := connPartsFromEnv(prefix) + if err != nil { + return "", err + } + return p.dsn(), nil +} + +// ClusterDSNsFromEnv resolves the primary and replica connection strings for a +// Cluster. The primary follows DSNFromEnv exactly; the replica resolves, +// highest first: +// +// 1. DATABASE_RO_URL — used verbatim. +// 2. DATABASE_RO_URL — likewise. +// 3. DB_RO_HOST, or bare DB_RO_HOST — the primary's port, user, +// password, database and sslmode with that host substituted, which is the +// CNPG shape where the -rw and -ro services differ only +// in hostname. +// 4. Nothing set: an empty replica DSN, so the Cluster runs single-pool. +// +// Nothing is derived. A deployment that wants split reads names the read-only +// endpoint; the replica host is never rewritten out of the primary's, because a +// wrong guess silently sends reads somewhere unintended. +// +// A host-only replica variable set alongside a primary given as a whole URL is +// an error rather than a guess: the fields to substitute into are not known. +func ClusterDSNsFromEnv(prefix string) (primary, replica string, err error) { + replica = urlFromEnv(prefix, "DATABASE_RO_URL") + roHost := lookup(prefix, fieldROHost, "") + + if u := urlFromEnv(prefix, "DATABASE_URL"); u != "" { + if replica == "" && roHost != "" { + return "", "", fmt.Errorf("%s%s is set but the primary comes from a connection URL: set %sDATABASE_RO_URL instead", + prefix, fieldROHost.suffix, prefix) + } + return u, replica, nil } - host := lookup(prefix, fieldHost, defaultHost) + p, err := connPartsFromEnv(prefix) + if err != nil { + return "", "", err + } + primary = p.dsn() + if replica != "" || roHost == "" { + return primary, replica, nil + } + p.host = roHost + return primary, p.dsn(), nil +} + +// urlFromEnv reads a whole-connection-string variable: prefixed, then bare. +func urlFromEnv(prefix, name string) string { + if v := os.Getenv(prefix + name); v != "" { + return v + } + return os.Getenv(name) +} + +// connParts is one resolved set of connection fields. +type connParts struct { + host string + port int + user string + pass string + name string + ssl string +} + +func (p connParts) dsn() string { return DSN(p.host, p.port, p.user, p.pass, p.name, p.ssl) } + +// connPartsFromEnv resolves levels 3 to 5 of DSNFromEnv's precedence. +func connPartsFromEnv(prefix string) (connParts, error) { portStr := lookup(prefix, fieldPort, strconv.Itoa(defaultPort)) - user := lookup(prefix, fieldUser, "") - pass := lookup(prefix, fieldPass, "") - name := lookup(prefix, fieldName, "") - ssl := lookup(prefix, fieldSSL, defaultSSLMode) - port, err := strconv.Atoi(portStr) if err != nil { - return "", fmt.Errorf("invalid %s%s: %w", prefix, fieldPort.suffix, err) + return connParts{}, fmt.Errorf("invalid %s%s: %w", prefix, fieldPort.suffix, err) } if port < 1 || port > 65535 { - return "", fmt.Errorf("invalid %s%s: port %d out of range", prefix, fieldPort.suffix, port) + return connParts{}, fmt.Errorf("invalid %s%s: port %d out of range", prefix, fieldPort.suffix, port) } - if user == "" { - return "", fmt.Errorf("no database user: set %s%s or %s", prefix, fieldUser.suffix, fieldUser.libpq) + p := connParts{ + host: lookup(prefix, fieldHost, defaultHost), + port: port, + user: lookup(prefix, fieldUser, ""), + pass: lookup(prefix, fieldPass, ""), + name: lookup(prefix, fieldName, ""), + ssl: lookup(prefix, fieldSSL, defaultSSLMode), } - if name == "" { - return "", fmt.Errorf("no database name: set %s%s or %s", prefix, fieldName.suffix, fieldName.libpq) + if p.user == "" { + return connParts{}, fmt.Errorf("no database user: set %s%s or %s", prefix, fieldUser.suffix, fieldUser.libpq) } - - return DSN(host, port, user, pass, name, ssl), nil + if p.name == "" { + return connParts{}, fmt.Errorf("no database name: set %s%s or %s", prefix, fieldName.suffix, fieldName.libpq) + } + return p, nil } // DSN renders the estate's standard connection string. User and password are diff --git a/pg/dsn_test.go b/pg/dsn_test.go index 4aed4f5..28c68ee 100644 --- a/pg/dsn_test.go +++ b/pg/dsn_test.go @@ -14,6 +14,7 @@ var dsnVars = []string{ "PGHOST", "PGPORT", "PGUSER", "PGPASSWORD", "PGDATABASE", "PGSSLMODE", "APP_DATABASE_URL", "APP_DBHOST", "APP_DBPORT", "APP_DBUSER", "APP_DBPASS", "APP_DBNAME", "APP_DBSSL", + "DATABASE_RO_URL", "DB_RO_HOST", "APP_DATABASE_RO_URL", "APP_DB_RO_HOST", } // setEnv clears every variable DSNFromEnv consults, then sets the given ones. @@ -227,6 +228,167 @@ func TestDSNFromEnv_ErrorNamesPrefixedVar(t *testing.T) { } } +func TestClusterDSNsFromEnv(t *testing.T) { + // The primary every field-based case resolves to, spelled once. + const wantPrimary = "postgres://app:pw@db-rw:5432/appdb?sslmode=require" + fieldEnv := map[string]string{ + "APP_DBHOST": "db-rw", "APP_DBUSER": "app", "APP_DBPASS": "pw", + "APP_DBNAME": "appdb", "APP_DBSSL": "require", + } + withFields := func(extra map[string]string) map[string]string { + env := map[string]string{} + for k, v := range fieldEnv { + env[k] = v + } + for k, v := range extra { + env[k] = v + } + return env + } + + tests := []struct { + name string + prefix string + env map[string]string + wantPrimary string + wantReplica string + }{ + { + name: "no replica variables leaves the replica empty", + prefix: "APP_", + env: withFields(nil), + wantPrimary: wantPrimary, + wantReplica: "", + }, + { + name: "prefixed DB_RO_HOST substitutes only the host", + prefix: "APP_", + env: withFields(map[string]string{"APP_DB_RO_HOST": "db-ro"}), + wantPrimary: wantPrimary, + wantReplica: "postgres://app:pw@db-ro:5432/appdb?sslmode=require", + }, + { + name: "bare DB_RO_HOST fills in for a prefixed lookup", + prefix: "APP_", + env: withFields(map[string]string{"DB_RO_HOST": "db-ro"}), + wantPrimary: wantPrimary, + wantReplica: "postgres://app:pw@db-ro:5432/appdb?sslmode=require", + }, + { + name: "prefixed DB_RO_HOST wins over the bare one", + prefix: "APP_", + env: withFields(map[string]string{ + "APP_DB_RO_HOST": "db-ro", "DB_RO_HOST": "ignored", + }), + wantPrimary: wantPrimary, + wantReplica: "postgres://app:pw@db-ro:5432/appdb?sslmode=require", + }, + { + name: "DATABASE_RO_URL passes through verbatim and wins over the host", + prefix: "APP_", + env: withFields(map[string]string{ + "APP_DATABASE_RO_URL": "postgres://ro@ro-host/db?application_name=reader", + "APP_DB_RO_HOST": "ignored", + }), + wantPrimary: wantPrimary, + wantReplica: "postgres://ro@ro-host/db?application_name=reader", + }, + { + name: "prefixed DATABASE_RO_URL wins over the bare one", + prefix: "APP_", + env: withFields(map[string]string{ + "APP_DATABASE_RO_URL": "postgres://app-ro@app-ro-host/db", + "DATABASE_RO_URL": "postgres://bare-ro@bare-ro-host/db", + }), + wantPrimary: wantPrimary, + wantReplica: "postgres://app-ro@app-ro-host/db", + }, + { + name: "both endpoints as whole URLs", + prefix: "", + env: map[string]string{ + "DATABASE_URL": "postgres://u:p@rw/db", + "DATABASE_RO_URL": "postgres://u:p@ro/db", + }, + wantPrimary: "postgres://u:p@rw/db", + wantReplica: "postgres://u:p@ro/db", + }, + { + name: "an unprefixed deployment reads the bare names", + prefix: "", + env: map[string]string{ + "DBHOST": "cnpg-rw", "DBUSER": "u", "DBNAME": "d", + "DB_RO_HOST": "cnpg-ro", + }, + wantPrimary: "postgres://u:@cnpg-rw:5432/d?sslmode=disable", + wantReplica: "postgres://u:@cnpg-ro:5432/d?sslmode=disable", + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + setEnv(t, tc.env) + primary, replica, err := ClusterDSNsFromEnv(tc.prefix) + if err != nil { + t.Fatalf("ClusterDSNsFromEnv: %v", err) + } + if primary != tc.wantPrimary { + t.Errorf("primary = %q, want %q", primary, tc.wantPrimary) + } + if replica != tc.wantReplica { + t.Errorf("replica = %q, want %q", replica, tc.wantReplica) + } + }) + } +} + +func TestClusterDSNsFromEnv_Errors(t *testing.T) { + tests := []struct { + name string + prefix string + env map[string]string + wantSub string + }{ + { + name: "a broken primary is reported", + prefix: "", + env: map[string]string{"DBUSER": "u"}, + wantSub: "set DBNAME or PGDATABASE", + }, + { + // The fields to substitute the read-only host into are unknown, so + // the alternative to an error is guessing at the connection. + name: "a host-only replica alongside a URL primary is ambiguous", + prefix: "APP_", + env: map[string]string{ + "APP_DATABASE_URL": "postgres://u:p@rw/db", + "APP_DB_RO_HOST": "db-ro", + }, + wantSub: "set APP_DATABASE_RO_URL instead", + }, + { + name: "a broken port is reported when the replica is resolved", + prefix: "", + env: map[string]string{ + "DBUSER": "u", "DBNAME": "d", "DBPORT": "70000", "DB_RO_HOST": "ro", + }, + wantSub: "out of range", + }, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + setEnv(t, tc.env) + primary, replica, err := ClusterDSNsFromEnv(tc.prefix) + if err == nil { + t.Fatalf("expected an error, got %q / %q", primary, replica) + } + if !strings.Contains(err.Error(), tc.wantSub) { + t.Fatalf("error %q does not mention %q", err, tc.wantSub) + } + }) + } +} + func TestDSN_EmptyPasswordMatchesLegacyFormat(t *testing.T) { // The Sprintf builders rendered an unset password as an empty string // between the colon and the "@"; keep that shape so DSNs do not churn. diff --git a/pg/integration_test.go b/pg/integration_test.go index 5895f37..017586e 100644 --- a/pg/integration_test.go +++ b/pg/integration_test.go @@ -2,6 +2,7 @@ package pg_test import ( "context" + "errors" "net/url" "strings" "testing" @@ -167,6 +168,58 @@ func TestDSNFromEnv_ConnectsToRealPostgres(t *testing.T) { pool.Close() } +// A Cluster against a real server. One container plays both roles — what is +// under test is the routing, not Postgres' own replication — so the read pool +// is the same server reached through a second, distinct DSN. +func TestCluster_AgainstRealPostgres(t *testing.T) { + ctx := testCtx(t) + dsn := pgtest.MustStartPostgres(ctx, t) + + c, err := pg.NewCluster(ctx, pg.ClusterConfig{ + PrimaryDSN: dsn, + ReplicaDSN: dsn + "&application_name=reader", + // Short enough that the recovery probe lands inside the test. + ReplicaRetryMin: 10 * time.Millisecond, + ReplicaRetryMax: 10 * time.Millisecond, + }) + if err != nil { + t.Fatalf("NewCluster: %v", err) + } + t.Cleanup(c.Close) + + if c.Read() == c.Write() { + t.Fatal("a configured, healthy replica must be a distinct pool") + } + if err := c.Migrate(ctx, migrations, pg.MigrateOptions{LockName: "golib-pg-cluster-integration"}); err != nil { + t.Fatalf("Migrate: %v", err) + } + if _, err := c.Write().Exec(ctx, "INSERT INTO widgets (name) VALUES ($1)", "sprocket"); err != nil { + t.Fatalf("insert on the primary: %v", err) + } + + var n int + if err := c.Read().QueryRow(ctx, "SELECT count(*) FROM widgets").Scan(&n); err != nil { + t.Fatalf("read: %v", err) + } + if err := c.Primary().QueryRow(ctx, "SELECT count(*) FROM widgets").Scan(&n); err != nil { + t.Fatalf("read-your-writes: %v", err) + } + if n != 1 { + t.Fatalf("widgets has %d rows, want 1", n) + } + + // A reported replica failure moves reads to the primary; a replica that + // still answers earns them back on the next probe. + c.ReportReplicaError(errors.New("simulated replica failure")) + if c.Read() != c.Write() { + t.Fatal("a reported replica failure must move reads to the primary") + } + time.Sleep(20 * time.Millisecond) + if c.Read() == c.Write() { + t.Fatal("a replica that answers its probe must get reads back") + } +} + // dsnParts is the container DSN split back into the fields DSNFromEnv reads. type dsnParts struct{ host, port, user, pass, name string } diff --git a/pg/pg.go b/pg/pg.go index 7aa1a9d..74c9cdb 100644 --- a/pg/pg.go +++ b/pg/pg.go @@ -1,6 +1,7 @@ // Package pg holds the estate's shared Postgres plumbing: environment-driven -// DSN construction, pgxpool construction, and the migration runner every -// service uses to bring its own schema up to date at startup. +// DSN construction, pgxpool construction, read/write splitting across a primary +// and a read replica, and the migration runner every service uses to bring its +// own schema up to date at startup. package pg import ( -- 2.47.3 From 02233919d7cbedd88ef3f75e6f9fdcd3899f1c4d Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Mon, 31 Aug 2026 23:00:24 +1000 Subject: [PATCH 2/2] Address review findings on the pg read/write split Name the replica host variable that was actually resolved. lookup falls back from the prefixed name to the bare one, but the ambiguous-replica error printed the prefixed form unconditionally, so a deployment setting APP_DATABASE_URL alongside a bare DB_RO_HOST was told to go and look at APP_DB_RO_HOST, which is not in its environment. lookupNamed reports the variable the value came from and the error names that one. Build the replica pool before the primary in NewCluster. pgxpool connects lazily, so this dials nothing, but it parses the replica DSN: a mistyped one now fails before the primary is opened, which removes the cleanup branch that closed a primary pool nobody could observe and makes the failure testable through the public constructor. A primary that cannot be reached closes the lazy replica pool on the way out. --- pg/cluster.go | 46 ++++++++++++++++++++++++++++++++-------------- pg/cluster_test.go | 32 ++++++++++++++++++++++---------- pg/dsn.go | 20 ++++++++++++++------ pg/dsn_test.go | 21 ++++++++++++++++++++- 4 files changed, 88 insertions(+), 31 deletions(-) diff --git a/pg/cluster.go b/pg/cluster.go index f8163a5..bc7dc79 100644 --- a/pg/cluster.go +++ b/pg/cluster.go @@ -59,6 +59,10 @@ type Cluster struct { // NewCluster opens the primary and, when configured, the replica. // +// The replica pool is built first. pgxpool connects lazily, so that dials +// nothing and costs nothing, but it does parse the replica DSN: a mistyped one +// fails before the primary is opened, leaving no pool behind to clean up. +// // A primary that cannot be reached is fatal: the service has nowhere to write. // A replica that cannot be reached is not — the Cluster starts with reads on // the primary and promotes the replica once it answers a probe, which is the @@ -66,31 +70,45 @@ type Cluster struct { // // The caller owns the Cluster and must Close it. func NewCluster(ctx context.Context, cfg ClusterConfig) (*Cluster, error) { + replica, err := openReplica(ctx, cfg) + if err != nil { + return nil, err + } primary, err := New(ctx, cfg.PrimaryDSN, cfg.Logger) if err != nil { + if replica != nil { + replica.Close() + } return nil, err } - c, err := newCluster(ctx, primary, cfg) - if err != nil { - primary.Close() - return nil, err - } - return c, nil + return newCluster(ctx, primary, replica, cfg), nil } -// newCluster attaches the replica half to an already-open primary pool. -func newCluster(ctx context.Context, primary *pgxpool.Pool, cfg ClusterConfig) (*Cluster, error) { - log := logger(cfg.Logger) - c := &Cluster{primary: primary} +// openReplica builds the replica pool, or nil when the Cluster is to run +// single-pool: no replica configured, or one pointed at the primary, which is +// answered with the pool already open rather than a second one to the same +// place. +func openReplica(ctx context.Context, cfg ClusterConfig) (*pgxpool.Pool, error) { if cfg.ReplicaDSN == "" || cfg.ReplicaDSN == cfg.PrimaryDSN { - log.Debug("postgres cluster in single-pool mode, reads go to the primary") - return c, nil + return nil, nil } - replica, err := pgxpool.New(ctx, cfg.ReplicaDSN) if err != nil { return nil, fmt.Errorf("connect postgres replica: %w", err) } + return replica, nil +} + +// newCluster joins an already-open primary to an already-open replica, nil for +// single-pool mode, and probes the replica once to set its starting health. +func newCluster(ctx context.Context, primary, replica *pgxpool.Pool, cfg ClusterConfig) *Cluster { + log := logger(cfg.Logger) + c := &Cluster{primary: primary} + if replica == nil { + log.Debug("postgres cluster in single-pool mode, reads go to the primary") + return c + } + c.replica = replica c.health = newReplicaCircuit(replica, log, cfg.ReplicaRetryMin, cfg.ReplicaRetryMax) if err := replica.Ping(ctx); err != nil { @@ -98,7 +116,7 @@ func newCluster(ctx context.Context, primary *pgxpool.Pool, cfg ClusterConfig) ( } else { log.Debug("postgres replica pool ready", "host", replica.Config().ConnConfig.Host) } - return c, nil + return c } // Write returns the primary pool. Every statement that changes data, takes row diff --git a/pg/cluster_test.go b/pg/cluster_test.go index 0ee7e1f..50e6bd3 100644 --- a/pg/cluster_test.go +++ b/pg/cluster_test.go @@ -211,10 +211,14 @@ func TestNewCluster_SinglePoolModes(t *testing.T) { } for _, tc := range tests { t.Run(tc.name, func(t *testing.T) { - c, err := newCluster(shortCtx(t), lazyPool(t, deadDSN), tc.cfg) + replica, err := openReplica(shortCtx(t), tc.cfg) if err != nil { - t.Fatalf("newCluster: %v", err) + t.Fatalf("openReplica: %v", err) } + if replica != nil { + t.Fatal("no replica pool may be opened in single-pool mode") + } + c := newCluster(shortCtx(t), lazyPool(t, deadDSN), replica, tc.cfg) if c.replica != nil { t.Fatal("expected single-pool mode") } @@ -228,14 +232,12 @@ func TestNewCluster_SinglePoolModes(t *testing.T) { // A replica that will not answer must degrade to primary reads, not fail // startup: the service can still serve everything, just without the split. func TestNewCluster_UnreachableReplicaStartsDegraded(t *testing.T) { - c, err := newCluster(shortCtx(t), lazyPool(t, deadDSN), ClusterConfig{ - PrimaryDSN: deadDSN, - ReplicaDSN: deadReplicaDSN, - Logger: testLogger(), - }) + cfg := ClusterConfig{PrimaryDSN: deadDSN, ReplicaDSN: deadReplicaDSN, Logger: testLogger()} + replica, err := openReplica(shortCtx(t), cfg) if err != nil { - t.Fatalf("newCluster: %v", err) + t.Fatalf("openReplica: %v", err) } + c := newCluster(shortCtx(t), lazyPool(t, deadDSN), replica, cfg) t.Cleanup(func() { c.replica.Close() }) if c.replica == nil { t.Fatal("the replica pool must still be opened") @@ -248,19 +250,29 @@ func TestNewCluster_UnreachableReplicaStartsDegraded(t *testing.T) { } } -func TestNewCluster_RejectsAnUnparseableReplicaDSN(t *testing.T) { - _, err := newCluster(shortCtx(t), lazyPool(t, deadDSN), ClusterConfig{ +// An unparseable replica DSN must be reported without the primary ever being +// opened, so the failure path has no pool to leak. The primary here is +// unreachable too: had it been opened first, its ping would have failed and +// masked the replica error. +func TestNewCluster_RejectsAnUnparseableReplicaDSNBeforeOpeningThePrimary(t *testing.T) { + c, err := NewCluster(shortCtx(t), ClusterConfig{ PrimaryDSN: deadDSN, ReplicaDSN: "://not a dsn", }) if err == nil { + c.Close() t.Fatal("expected an error for an unparseable replica DSN") } if !strings.Contains(err.Error(), "connect postgres replica") { t.Fatalf("error %q does not identify the failing step", err) } + if strings.Contains(err.Error(), "ping postgres") { + t.Fatalf("the primary was opened before the replica DSN was parsed: %v", err) + } } +// The mirror case: the replica DSN parses, so its pool exists when the primary +// turns out to be unreachable, and NewCluster must close it on the way out. func TestNewCluster_PropagatesPrimaryFailure(t *testing.T) { c, err := NewCluster(shortCtx(t), ClusterConfig{PrimaryDSN: deadDSN, ReplicaDSN: deadReplicaDSN}) if err == nil { diff --git a/pg/dsn.go b/pg/dsn.go index 28d5fd8..bbb30ac 100644 --- a/pg/dsn.go +++ b/pg/dsn.go @@ -92,12 +92,12 @@ func DSNFromEnv(prefix string) (string, error) { // an error rather than a guess: the fields to substitute into are not known. func ClusterDSNsFromEnv(prefix string) (primary, replica string, err error) { replica = urlFromEnv(prefix, "DATABASE_RO_URL") - roHost := lookup(prefix, fieldROHost, "") + roHost, roHostVar := lookupNamed(prefix, fieldROHost, "") if u := urlFromEnv(prefix, "DATABASE_URL"); u != "" { if replica == "" && roHost != "" { - return "", "", fmt.Errorf("%s%s is set but the primary comes from a connection URL: set %sDATABASE_RO_URL instead", - prefix, fieldROHost.suffix, prefix) + return "", "", fmt.Errorf("%s is set but the primary comes from a connection URL: set %sDATABASE_RO_URL instead", + roHostVar, prefix) } return u, replica, nil } @@ -187,11 +187,19 @@ func hostPort(host string, port int) string { // lookup resolves one field: prefixed variable, then libpq variable, then def. func lookup(prefix string, f dsnField, def string) string { + v, _ := lookupNamed(prefix, f, def) + return v +} + +// lookupNamed resolves one field and reports the variable the value came from, +// empty when the default was used. An error message names the variable the +// operator actually set rather than the prefixed form they may never have used. +func lookupNamed(prefix string, f dsnField, def string) (value, name string) { if v := os.Getenv(prefix + f.suffix); v != "" { - return v + return v, prefix + f.suffix } if v := os.Getenv(f.libpq); v != "" { - return v + return v, f.libpq } - return def + return def, "" } diff --git a/pg/dsn_test.go b/pg/dsn_test.go index 28c68ee..47e4491 100644 --- a/pg/dsn_test.go +++ b/pg/dsn_test.go @@ -348,6 +348,8 @@ func TestClusterDSNsFromEnv_Errors(t *testing.T) { prefix string env map[string]string wantSub string + // notSub, when set, must not appear in the error. + notSub string }{ { name: "a broken primary is reported", @@ -364,7 +366,21 @@ func TestClusterDSNsFromEnv_Errors(t *testing.T) { "APP_DATABASE_URL": "postgres://u:p@rw/db", "APP_DB_RO_HOST": "db-ro", }, - wantSub: "set APP_DATABASE_RO_URL instead", + wantSub: "APP_DB_RO_HOST is set but the primary comes from a connection URL: set APP_DATABASE_RO_URL instead", + }, + { + // The same ambiguity, but the host came from the bare variable the + // prefixed lookup falls back to. The error must name the variable + // the operator actually set: naming APP_DB_RO_HOST sends them + // looking for something that is not in their environment. + name: "the bare replica host variable is named, not the prefixed one", + prefix: "APP_", + env: map[string]string{ + "APP_DATABASE_URL": "postgres://u:p@rw/db", + "DB_RO_HOST": "db-ro", + }, + wantSub: "DB_RO_HOST is set but", + notSub: "APP_DB_RO_HOST", }, { name: "a broken port is reported when the replica is resolved", @@ -385,6 +401,9 @@ func TestClusterDSNsFromEnv_Errors(t *testing.T) { if !strings.Contains(err.Error(), tc.wantSub) { t.Fatalf("error %q does not mention %q", err, tc.wantSub) } + if tc.notSub != "" && strings.Contains(err.Error(), tc.notSub) { + t.Fatalf("error %q names %q, which is not the variable that was set", err, tc.notSub) + } }) } } -- 2.47.3