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
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
integration test proves that, with replication paused via
`pg_wal_replay_pause()`, a write followed by a read never observes a stale
value under strict mode. Rationale in `docs/adr/0008-lsn-fencing.md`.
- Routing policy engine (Phase 8): `internal/router` makes replica selection
pluggable behind a `Policy` interface, choosing among the replicas the fence
has already deemed eligible. Three policies ship — `round-robin`,
`least-in-flight` (the default), and `scored`, which ranks replicas by
estimated completion time using an EWMA of each query shape's latency keyed by
pg_query fingerprint, steering expensive shapes away from busy replicas.
`routing.policy` selects one. A deterministic discrete-event simulator
compares the policies on a synthetic mixed workload with no database, and an
integration test asserts reads spread across both replicas. Rationale in
`docs/adr/0009-routing-policy-engine.md`.

### Dependencies

Expand Down
31 changes: 26 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,8 +42,9 @@ The goal is to do one thing — correct, observable read/write routing — well.
Early development, built in phases (see the roadmap). Not production-ready yet.
pgpilot now **routes**: it authenticates each client with SCRAM-SHA-256, pools
connections, classifies each query, and sends writes to the primary and reads to
a replica — enforcing read-your-writes with a per-session LSN fence. A routing
*policy* (scoring replicas by lag and load) and observability come next.
a replica — enforcing read-your-writes with a per-session LSN fence — and
balances reads across eligible replicas with a selectable routing policy
(round-robin, least-in-flight, or latency-scored). Observability comes next.

## Roadmap

Expand All @@ -57,8 +58,8 @@ a replica — enforcing read-your-writes with a per-session LSN fence. A routing
| 5 | Query classification (read vs. write via pg_query) | done |
| 6 | Replica registry, health polling, circuit breakers | done |
| 7 | LSN fencing | done |
| 8 | Routing policy engine | next |
| 9 | Observability (Prometheus, structured logs, pprof) | |
| 8 | Routing policy engine | done |
| 9 | Observability (Prometheus, structured logs, pprof) | next |
| 10 | Fault-injection harness | |
| 11 | Benchmarks vs. direct connection and pgbouncer | |
| 12 | Docs and the v0.1.0 release | |
Expand Down Expand Up @@ -106,7 +107,8 @@ pgpilot reads a JSON config file (see [`pgpilot.example.json`](pgpilot.example.j
"users": [{"name": "pgpilot", "password": "pgpilot"}],
"pool": {"mode": "session", "max_size": 10, "acquire_timeout": "5s", "idle_timeout": "5m"},
"health": {"interval": "1s", "failure_threshold": 3, "base_backoff": "1s", "max_backoff": "30s"},
"fencing": {"mode": "strict", "bounded_ms": 100}
"fencing": {"mode": "strict", "bounded_ms": 100},
"routing": {"policy": "least-in-flight"}
}
```

Expand Down Expand Up @@ -142,6 +144,25 @@ primary). `fencing.mode` selects the trade-off:
stale value under strict mode. Design in
[`docs/adr/0008-lsn-fencing.md`](docs/adr/0008-lsn-fencing.md).

### Routing policies

Fencing decides which replicas *may* serve a read; `routing.policy` decides which
one *does* when more than one qualifies:

- **`round-robin`** — even rotation across eligible replicas.
- **`least-in-flight`** (default) — the eligible replica with the fewest reads
outstanding, so a slow or overloaded replica sheds load until it drains.
- **`scored`** — ranks replicas by estimated completion time,
`(inFlight + 1) * ewmaLatency(addr, shape) + lagPenalty * lag`, learning each
query shape's cost per replica (keyed by pg_query fingerprint) so it steers
expensive shapes away from busy replicas. It costs one fingerprint parse per
read, which is why it is opt-in.

`go test -run WorkloadComparison -v ./internal/router` runs a deterministic
simulation comparing the policies on a synthetic mixed workload with no database.
Design in
[`docs/adr/0009-routing-policy-engine.md`](docs/adr/0009-routing-policy-engine.md).

### Health and replication lag

A background poller (`internal/registry`) tracks each backend's role and
Expand Down
7 changes: 7 additions & 0 deletions cmd/pgpilot/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import (
"github.com/sachhg/pgpilot/internal/config"
"github.com/sachhg/pgpilot/internal/proxy"
"github.com/sachhg/pgpilot/internal/registry"
"github.com/sachhg/pgpilot/internal/router"
)

// version is the build version, overridden via -ldflags "-X main.version=...".
Expand Down Expand Up @@ -91,11 +92,17 @@ func run(args []string) error {
reg.Start(ctx, backendsFrom(cfg))
go reloadOnHUP(ctx, logger, *configPath, reg)

policy, err := router.New(cfg.Routing.Policy)
if err != nil {
return err
}

srv := proxy.New(proxy.Config{
ListenAddr: cfg.Listen,
Users: cfg,
Manager: mgr,
Registry: reg,
Policy: policy,
Logger: logger,
})
addr, err := srv.Listen()
Expand Down
79 changes: 79 additions & 0 deletions docs/adr/0009-routing-policy-engine.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
# 9. Routing policy engine

Date: 2026-07-23

## Status

Accepted.

## Context

Phase 7 made pgpilot route reads to replicas that have replayed past the
session's fence, but among the replicas that qualify it simply took the first
one. On a cluster with more than one eligible replica that concentrates every
read on the same backend, wasting the rest and ignoring that replicas differ in
load, speed, and how expensive a given query is on each. This phase decides
*which* eligible replica serves a read.

Eligibility (healthy, and fresh enough under the fencing rules) already lives in
the proxy and is a correctness question. Selection among the eligible is a
performance question with real alternatives, so it belongs behind an interface.

## Decision

**A `Policy` interface, separate from eligibility.** `internal/router` defines
`Policy.Choose(candidates, fingerprint) addr` over the replicas the proxy has
already deemed eligible. Because eligibility is decided upstream, a policy never
needs a database or the registry: it is pure Go, exercised in unit tests and a
deterministic simulator with no cluster. Every `Choose` is paired with one
`Release(addr, fingerprint, latency)` so a load-aware policy tracks in-flight
work and observed latency; the single policy instance is shared across all
sessions so it sees the whole proxy's load.

**Three implementations, increasing in sophistication.**

- **round-robin** — even rotation. Fair over a stable candidate set, keeps no
state, ignores load and query cost.
- **least-in-flight** — send each read to the eligible replica with the fewest
outstanding reads. Adapts to a slow or briefly overloaded replica for free,
because its reads pile up and it stops being chosen until it drains. This is
the default: strong, zero-config, and it costs no per-read parse.
- **scored** — rank candidates by estimated completion time,
`(inFlight + 1) * ewmaLatency(addr, fingerprint) + lagPenalty * lagSeconds`,
and take the minimum. Keying the EWMA by pg_query fingerprint is the point: it
lets the policy steer an *expensive* query shape away from a replica already
busy with that shape while still sending cheap queries there, and it learns
per-replica speed differences as slower backends accrue higher averages.

**least-in-flight is the default, scored is opt-in.** The scored policy needs a
pg_query fingerprint per read — a cgo parse — and only pays off when replicas or
query costs are heterogeneous. Least-in-flight already handles slow and
overloaded replicas without that cost, so it is the safe default; `routing.policy`
selects another, validated at load time against the router's own names.

## Consequences

- A deterministic discrete-event simulator (seeded arrivals, a heavy-tailed
shape mix, one 3x-slow replica, each replica a single FIFO server) compares
the policies with no database. It shows the failure it is meant to: blind
round-robin forces a third of all reads onto the slow replica and its queue
explodes (p95 into the seconds), while least-in-flight and scored back off it —
and scored, alone knowing heavy shapes are dear there, sends it the fewest
reads and wins the tail. Micro-benchmarks confirm per-decision cost stays
sub-microsecond.
- The scored policy trades slightly higher median latency for a better tail: it
concentrates load on the fast replicas to spare the slow one, so the fast
replicas queue a little more at p50 while p95/p99 improve. That is the right
trade for a router, but it is a trade.
- The EWMA folds in queueing time, not pure service time, because that is what
the proxy can measure — the latency from dispatch to ReadyForQuery. Under load
the estimate therefore runs a little hot, which reinforces the correct steering
but is not a clean per-shape cost model. A server-reported cost would be
cleaner and is future work.
- The scored weights are constants tuned for comparable replicas. A cluster with
a large staleness budget or very heterogeneous hardware may want them tunable
from config; they are exported so that remains an additive change.
- Rejected: random choice (no load awareness), lag-only weighting (blind to load
and query cost — the very thing that overloads a replica), and power-of-two-
choices (a good cheap approximation of least-in-flight, but with only a handful
of replicas exact least-in-flight is affordable and strictly better).
12 changes: 12 additions & 0 deletions internal/classify/classify.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,18 @@ var volatileFuncs = map[string]struct{}{
"clock_timestamp": {}, "timeofday": {}, "pg_sleep": {},
}

// Fingerprint returns pg_query's stable fingerprint of a query — a hash of its
// parse tree with literals normalized away, so queries that differ only in their
// constants share a fingerprint. Load-aware routing uses it to key per-shape
// latency estimates. An unparseable query has no fingerprint and returns "".
func Fingerprint(sql string) string {
fp, err := pg.Fingerprint(sql)
if err != nil {
return ""
}
return fp
}

// Classify returns the routing class of a simple-query string, which may contain
// several statements. It is a Write if any statement is a Write.
func Classify(sql string) Class {
Expand Down
43 changes: 43 additions & 0 deletions internal/classify/fingerprint_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
package classify_test

import (
"testing"

"github.com/sachhg/pgpilot/internal/classify"
)

func TestFingerprint_StableAcrossLiterals(t *testing.T) {
// Queries differing only in constants share a fingerprint.
a := classify.Fingerprint("SELECT * FROM users WHERE id = 1")
b := classify.Fingerprint("SELECT * FROM users WHERE id = 42")
if a == "" {
t.Fatal("fingerprint was empty for a valid query")
}
if a != b {
t.Errorf("fingerprints differ across literals: %q vs %q", a, b)
}
}

func TestFingerprint_DistinguishesShapes(t *testing.T) {
a := classify.Fingerprint("SELECT * FROM users WHERE id = 1")
b := classify.Fingerprint("SELECT * FROM orders WHERE id = 1")
if a == b {
t.Errorf("different table shapes shared a fingerprint: %q", a)
}
}

func TestFingerprint_UnparseableIsEmpty(t *testing.T) {
if fp := classify.Fingerprint("SELECT FROM WHERE ("); fp != "" {
t.Errorf("unparseable query returned fingerprint %q, want empty", fp)
}
}

func TestFingerprint_Deterministic(t *testing.T) {
sql := "SELECT a, b FROM t JOIN u ON t.id = u.id WHERE t.x > 3 ORDER BY a"
first := classify.Fingerprint(sql)
for i := 0; i < 5; i++ {
if got := classify.Fingerprint(sql); got != first {
t.Fatalf("fingerprint not deterministic: %q vs %q", got, first)
}
}
}
22 changes: 22 additions & 0 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ import (
"fmt"
"os"
"time"

"github.com/sachhg/pgpilot/internal/router"
)

// Config is the top-level pgpilot configuration.
Expand All @@ -21,6 +23,17 @@ type Config struct {
Pool Pool `json:"pool"`
Health Health `json:"health"`
Fencing Fencing `json:"fencing"`
Routing Routing `json:"routing"`
}

// Routing configures how reads are balanced across the replicas eligible to
// serve them (those healthy and fresh enough under the fencing rules).
type Routing struct {
// Policy is "round-robin" (even rotation), "least-in-flight" (fewest
// outstanding reads wins), or "scored" (rank by estimated completion time,
// learning per-query-shape latency). "scored" costs one pg_query fingerprint
// per read, so "least-in-flight" is the default.
Policy string `json:"policy"`
}

// Fencing-mode values.
Expand Down Expand Up @@ -157,6 +170,9 @@ func (c *Config) applyDefaults() {
if c.Fencing.BoundedMs == 0 {
c.Fencing.BoundedMs = 100
}
if c.Routing.Policy == "" {
c.Routing.Policy = router.PolicyLeastInFlight
}
}

func (c *Config) validate() error {
Expand Down Expand Up @@ -191,6 +207,12 @@ func (c *Config) validate() error {
return fmt.Errorf("config: fencing.mode must be %q, %q, or %q, got %q",
FenceStrict, FenceBounded, FenceRelaxed, c.Fencing.Mode)
}
switch c.Routing.Policy {
case router.PolicyRoundRobin, router.PolicyLeastInFlight, router.PolicyScored:
default:
return fmt.Errorf("config: routing.policy must be %q, %q, or %q, got %q",
router.PolicyRoundRobin, router.PolicyLeastInFlight, router.PolicyScored, c.Routing.Policy)
}
return nil
}

Expand Down
29 changes: 29 additions & 0 deletions internal/config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,35 @@ func TestLoad_Fencing(t *testing.T) {
}
}

func TestLoad_Routing(t *testing.T) {
path := writeConfig(t, `{
"listen": "x", "primary": "y",
"users": [{"name": "u", "password": "p"}],
"routing": {"policy": "scored"}
}`)
c, err := config.Load(path)
if err != nil {
t.Fatalf("Load: %v", err)
}
if c.Routing.Policy != "scored" {
t.Errorf("routing policy = %q, want scored", c.Routing.Policy)
}

def := writeConfig(t, `{"listen":"x","primary":"y","users":[{"name":"u","password":"p"}]}`)
dc, err := config.Load(def)
if err != nil {
t.Fatalf("Load default: %v", err)
}
if dc.Routing.Policy != "least-in-flight" {
t.Errorf("default routing policy = %q, want least-in-flight", dc.Routing.Policy)
}

bad := writeConfig(t, `{"listen":"x","primary":"y","users":[{"name":"u","password":"p"}],"routing":{"policy":"nope"}}`)
if _, err := config.Load(bad); err == nil {
t.Error("expected an error for an invalid routing policy")
}
}

func TestLoad_HealthDefaults(t *testing.T) {
path := writeConfig(t, `{"listen":"x","primary":"y","users":[{"name":"u","password":"p"}]}`)
c, err := config.Load(path)
Expand Down
31 changes: 28 additions & 3 deletions internal/proxy/proxy.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import (
"github.com/sachhg/pgpilot/internal/config"
"github.com/sachhg/pgpilot/internal/protocol"
"github.com/sachhg/pgpilot/internal/registry"
"github.com/sachhg/pgpilot/internal/router"
)

// Config configures a proxy Server.
Expand All @@ -32,15 +33,19 @@ type Config struct {
// Registry reports backend health and lag for read routing. Nil disables
// routing (every session is served by the primary).
Registry *registry.Registry
// Policy selects which eligible replica serves each read. Nil resolves the
// policy named in Users.Routing.Policy, defaulting to least-in-flight.
Policy router.Policy
// Logger receives structured logs. Nil selects slog.Default.
Logger *slog.Logger
}

// Server is pgpilot's client-facing proxy. Construct it with New, then call
// Listen and Serve.
type Server struct {
cfg Config
log *slog.Logger
cfg Config
log *slog.Logger
policy router.Policy

sessions atomic.Uint64
wg sync.WaitGroup
Expand All @@ -55,7 +60,26 @@ func New(cfg Config) *Server {
if logger == nil {
logger = slog.Default()
}
return &Server{cfg: cfg, log: logger}
return &Server{cfg: cfg, log: logger, policy: resolvePolicy(cfg)}
}

// resolvePolicy returns cfg.Policy, or constructs the policy named in the config
// (defaulting to least-in-flight) when the caller did not inject one. The single
// returned instance is shared across every session, since load-aware policies
// track in-flight work across the whole proxy.
func resolvePolicy(cfg Config) router.Policy {
if cfg.Policy != nil {
return cfg.Policy
}
name := router.PolicyLeastInFlight
if cfg.Users != nil && cfg.Users.Routing.Policy != "" {
name = cfg.Users.Routing.Policy
}
p, err := router.New(name)
if err != nil {
return router.NewLeastInFlight()
}
return p
}

// Listen binds the configured listen address and returns the resolved address.
Expand Down Expand Up @@ -128,6 +152,7 @@ func (s *Server) handle(ctx context.Context, client net.Conn) {
cfg: s.cfg.Users,
manager: s.cfg.Manager,
registry: s.cfg.Registry,
policy: s.policy,
log: log,
tracker: &protocol.TxTracker{},
}
Expand Down
Loading
Loading