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
4 changes: 4 additions & 0 deletions .markdownlint.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
MD013: false
MD033:
allowed_elements: [p, img]
MD041: false
8 changes: 8 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,7 @@ All flags are bound to environment variables with the `EP_` prefix. For example,
| `--liveness-interval` | `EP_LIVENESS_INTERVAL` | `5s` | Heartbeat probe interval for liveness detection |
| `--liveness-threshold` | `EP_LIVENESS_THRESHOLD` | `15s` | Maximum time since last heartbeat before liveness fails |
| `--drain-timeout` | `EP_DRAIN_TIMEOUT` | `30s` | Grace period for connections to removed endpoints before force close (`0` closes immediately) |
| `--upstream-selection` | `EP_UPSTREAM_SELECTION` | `random` | How new connections pick an upstream: `random` (uniform among pickable upstreams) or `latency` (prefer the lowest-latency tier, see below) |

### Validation rules

Expand All @@ -153,6 +154,7 @@ All flags are bound to environment variables with the `EP_` prefix. For example,
- `--liveness-interval` and `--liveness-threshold` must be at least 1 second
- `--liveness-threshold` must be greater than `--liveness-interval`
- `--drain-timeout` must not be negative
- `--upstream-selection` must be `random` or `latency`

## Examples

Expand All @@ -174,6 +176,10 @@ The TCP load balancer is implemented in `internal/proxy`. It accepts a channel o

When an endpoint disappears from the upstream list (node cordon, API server rolling restart), the proxy stops routing new connections to it but lets existing connections finish within `--drain-timeout`. Remaining connections are force-closed when the timeout expires. Re-adding the endpoint before the timeout aborts the drain. A timeout of `0` closes connections to removed endpoints immediately. A new connection whose endpoint starts draining while it is being dialed gets one retry on another pickable endpoint, if there is one.

### Latency-aware selection

With `--upstream-selection=latency`, every successful health check feeds its connect time into a smoothed average (EWMA) per upstream; failed checks do not. For an endpoint given as a hostname, the connect time includes name resolution. Each sample is capped at four times the current average (but never below 2ms), so one slow check, such as a retransmitted SYN, moves the average by a bounded step instead of jumping to the retransmission time, while a lasting increase still shows up within a few checks. A sample below a quarter of the average replaces it, so a first check that took more than four times the usual connect time is undone by the next normal one. Upstreams with a sample are grouped into tiers by latency alone, healthy or not: tier 0 holds every upstream whose average is at most `max(best + 2ms, 2 × best)`, where `best` is the lowest average of all; tier 1 applies the same rule to the rest, and so on. An upstream leaves its tier only when its average exceeds the limit by a further 25%, so an upstream near the boundary does not flap. Membership follows the group rather than the tier number, so a tier appearing or disappearing above does not split a group. New connections are picked at random among the pickable upstreams of the lowest tier that has one, together with pickable upstreams that have no sample yet. Unhealthy and draining upstreams keep their tier, so losing some upstreams of a tier does not let a lower tier in, and losing all of them hands over to the next tier. An upstream removed from the endpoint list leaves the tiers once its drain completes, and the tiers are rebuilt without it. A new upstream takes picks before its first successful check; if it is unreachable, it keeps an equal share of that tier's picks, not of all upstreams as in `random` mode, until the consecutive-failure threshold excludes it. Every upstream within 2ms of the best one is in tier 0, so a single site with sub-millisecond differences behaves like `random`. Only health checks feed the average, so tiers follow latency changes at the pace of `--health-interval`. Existing connections stay where they are.

### Endpoint discovery

The merged discovery provider runs all configured sub-providers concurrently and combines their results:
Expand Down Expand Up @@ -205,6 +211,8 @@ The health server also exposes Prometheus metrics at `/metrics`. It binds to `--
| `extractedprism_connection_errors_total` | Counter | `upstream` | Failed connection attempts (`none` when no upstream was available); failures caused by proxy shutdown are not counted. The series of an upstream is deleted once it is drained and removed, and failures reported after that are not counted until the address is added again |
| `extractedprism_discovery_updates_total` | Counter | `provider` | Endpoint list updates received (`static`, `kubernetes`) |
| `extractedprism_discovery_errors_total` | Counter | `provider` | Discovery errors: provider failures, failed Watch calls, watch error events and failed re-lists. A watch stream ending is not counted, including on a dropped connection, and neither is 410 Gone expiry; a dead upstream shows in `health_check_status` instead |
| `extractedprism_upstream_rtt_seconds` | Gauge | `upstream` | Smoothed health check connect time, including name resolution for hostname endpoints; only with `--upstream-selection=latency`, after the first successful check |
| `extractedprism_upstream_latency_tier` | Gauge | `upstream` | Latency tier of the upstream, 0 being the closest group; new connections go to the lowest tier with a healthy, non-draining upstream, together with upstreams that have no sample yet. Only with `--upstream-selection=latency`, after the first successful check |
| `extractedprism_health_check_status` | Gauge | `upstream` | Upstream health state after the consecutive-failure threshold, fed by health checks and client dials (1 healthy, 0 unhealthy). Health checks stop while an upstream drains, so the value holds until it is removed or re-added |

### Graceful shutdown
Expand Down
3 changes: 3 additions & 0 deletions cmd/extractedprism/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ func registerFlags() {
flags.Duration("liveness-interval", base.LivenessInterval, "heartbeat probe interval for liveness detection")
flags.Duration("liveness-threshold", base.LivenessThreshold, "maximum time since last heartbeat before liveness fails")
flags.Duration("drain-timeout", base.DrainTimeout, "grace period for connections to removed endpoints before force close (0 = immediate)")
flags.String("upstream-selection", base.UpstreamSelection, "how new connections pick an upstream (random, latency)")
}

func bindEnvVars() {
Expand All @@ -78,6 +79,7 @@ func bindEnvVars() {
mustBindPFlag("liveness_interval", flags.Lookup("liveness-interval"))
mustBindPFlag("liveness_threshold", flags.Lookup("liveness-threshold"))
mustBindPFlag("drain_timeout", flags.Lookup("drain-timeout"))
mustBindPFlag("upstream_selection", flags.Lookup("upstream-selection"))

viper.AutomaticEnv()
}
Expand Down Expand Up @@ -143,6 +145,7 @@ func buildConfig() *config.Config {
cfg.LivenessInterval = viper.GetDuration("liveness_interval")
cfg.LivenessThreshold = viper.GetDuration("liveness_threshold")
cfg.DrainTimeout = viper.GetDuration("drain_timeout")
cfg.UpstreamSelection = viper.GetString("upstream_selection")

return cfg
}
Expand Down
17 changes: 17 additions & 0 deletions cmd/extractedprism/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -99,3 +99,20 @@ func TestBuildConfig_DrainTimeoutFromViper(t *testing.T) {
assert.Equal(t, 7*time.Second, cfg.DrainTimeout,
"DrainTimeout must be read from viper")
}

func TestBuildConfig_UpstreamSelectionFromViper(t *testing.T) {
setValidViperDefaults()
t.Cleanup(viper.Reset)

viper.Set("upstream_selection", "latency")

cfg := buildConfig()
assert.Equal(t, "latency", cfg.UpstreamSelection,
"UpstreamSelection must be read from viper")
}

func TestRegisterFlags_UpstreamSelectionDefaultsToRandom(t *testing.T) {
flag := rootCmd.PersistentFlags().Lookup("upstream-selection")
require.NotNil(t, flag, "--upstream-selection must be registered")
assert.Equal(t, "random", flag.DefValue)
}
22 changes: 21 additions & 1 deletion internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,12 @@ const (
defaultLivenessThreshold = 15 * time.Second
defaultDrainTimeout = 30 * time.Second

// UpstreamSelectionRandom picks uniformly among pickable upstreams.
UpstreamSelectionRandom = "random"
// UpstreamSelectionLatency prefers the upstreams with the lowest
// smoothed health check connect time.
UpstreamSelectionLatency = "latency"

minPort = 1
maxPort = 65535
)
Expand All @@ -39,6 +45,7 @@ var (
ErrInvalidLivenessTiming = errors.New("liveness threshold must be greater than liveness interval")
ErrInvalidHealthBindAddress = errors.New("invalid health bind address")
ErrInvalidDrainTimeout = errors.New("drain timeout must not be negative")
ErrInvalidUpstreamSelection = errors.New("invalid upstream selection")
)

const minDuration = 1 * time.Second
Expand Down Expand Up @@ -68,6 +75,8 @@ type Config struct {
// DrainTimeout bounds how long connections to a removed endpoint may
// finish before being force-closed. Zero closes them immediately.
DrainTimeout time.Duration
// UpstreamSelection is UpstreamSelectionRandom or UpstreamSelectionLatency.
UpstreamSelection string
}

// NewBaseConfig returns a Config populated with sensible defaults for optional
Expand All @@ -85,6 +94,7 @@ func NewBaseConfig() *Config {
LivenessInterval: defaultLivenessInterval,
LivenessThreshold: defaultLivenessThreshold,
DrainTimeout: defaultDrainTimeout,
UpstreamSelection: UpstreamSelectionRandom,
}
}

Expand Down Expand Up @@ -159,7 +169,17 @@ func (cfg *Config) Validate() error {
return errors.Wrapf(ErrInvalidDrainTimeout, "drain timeout %s", cfg.DrainTimeout)
}

return nil
return validateUpstreamSelection(cfg.UpstreamSelection)
}

func validateUpstreamSelection(selection string) error {
switch selection {
case UpstreamSelectionRandom, UpstreamSelectionLatency:
return nil
default:
return errors.Wrapf(ErrInvalidUpstreamSelection, "%q: must be %s or %s",
selection, UpstreamSelectionRandom, UpstreamSelectionLatency)
}
}

// validateAddress checks that addr is a valid IP address or a syntactically
Expand Down
36 changes: 36 additions & 0 deletions internal/config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -790,3 +790,39 @@ func TestNewBaseConfig_DrainTimeoutDefault(t *testing.T) {
cfg := config.NewBaseConfig()
assert.Equal(t, 30*time.Second, cfg.DrainTimeout)
}

func TestNewBaseConfig_UpstreamSelectionDefaultsToRandom(t *testing.T) {
cfg := config.NewBaseConfig()
assert.Equal(t, config.UpstreamSelectionRandom, cfg.UpstreamSelection)
}

func TestValidate_UpstreamSelection(t *testing.T) {
tests := []struct {
name string
selection string
wantErr error
}{
{name: "random is valid", selection: "random", wantErr: nil},
{name: "latency is valid", selection: "latency", wantErr: nil},
{name: "empty returns ErrInvalidUpstreamSelection", selection: "", wantErr: config.ErrInvalidUpstreamSelection},
{name: "unknown mode returns ErrInvalidUpstreamSelection", selection: "fastest", wantErr: config.ErrInvalidUpstreamSelection},
{name: "mode names are case sensitive", selection: "Latency", wantErr: config.ErrInvalidUpstreamSelection},
{name: "surrounding whitespace is not trimmed", selection: " random", wantErr: config.ErrInvalidUpstreamSelection},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
cfg := validTestConfig()
cfg.UpstreamSelection = tt.selection

err := cfg.Validate()

if tt.wantErr == nil {
require.NoError(t, err)
} else {
require.Error(t, err)
assert.True(t, errors.Is(err, tt.wantErr))
}
})
}
}
40 changes: 33 additions & 7 deletions internal/metrics/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package metrics

import (
"net/http"
"time"

"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"
Expand All @@ -28,6 +29,8 @@ type Metrics struct {
upstreamsActive prometheus.Gauge
upstreamsTotal prometheus.Gauge
backendHealth *prometheus.GaugeVec
backendRTT *prometheus.GaugeVec
backendTier *prometheus.GaugeVec
}

// New creates a Metrics with all series registered on a private registry.
Expand Down Expand Up @@ -72,11 +75,12 @@ func New() *Metrics {
Name: "upstreams_total",
Help: "Total known upstreams, including unhealthy and draining ones.",
}),
backendHealth: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: namespace,
Name: "health_check_status",
Help: "Upstream health state after the consecutive-failure threshold, fed by health checks and client dials (1 healthy, 0 unhealthy). Health checks stop while an upstream drains, so the value holds until it is removed or re-added.",
}, []string{labelUpstream}),
backendHealth: newUpstreamGaugeVec("health_check_status",
"Upstream health state after the consecutive-failure threshold, fed by health checks and client dials (1 healthy, 0 unhealthy). Health checks stop while an upstream drains, so the value holds until it is removed or re-added."),
backendRTT: newUpstreamGaugeVec("upstream_rtt_seconds",
"Smoothed connect time of successful health checks, including name resolution for hostname endpoints. Only exported with --upstream-selection=latency, after the first successful check."),
backendTier: newUpstreamGaugeVec("upstream_latency_tier",
"Latency tier of the upstream, 0 being the closest group. New connections go to the lowest tier with a healthy, non-draining upstream, together with upstreams that have no sample yet. Only exported with --upstream-selection=latency, after the first successful check."),
}

m.registry.MustRegister(
Expand All @@ -88,11 +92,21 @@ func New() *Metrics {
m.upstreamsActive,
m.upstreamsTotal,
m.backendHealth,
m.backendRTT,
m.backendTier,
)

return m
}

func newUpstreamGaugeVec(name, help string) *prometheus.GaugeVec {
return prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: namespace,
Name: name,
Help: help,
}, []string{labelUpstream})
}

// Handler returns an HTTP handler exposing the registry in Prometheus format.
func (m *Metrics) Handler() http.Handler {
return promhttp.HandlerFor(m.registry, promhttp.HandlerOpts{})
Expand Down Expand Up @@ -143,9 +157,21 @@ func (m *Metrics) SetBackendHealth(upstream string, healthy bool) {
m.backendHealth.WithLabelValues(upstream).Set(value)
}

// RemoveBackend deletes the health and connection error series of a removed
// upstream so dead series do not accumulate as endpoints change.
// RemoveBackend deletes every per-upstream series of a removed upstream so
// dead series do not accumulate as endpoints change.
func (m *Metrics) RemoveBackend(upstream string) {
m.backendHealth.DeleteLabelValues(upstream)
m.connErrors.DeleteLabelValues(upstream)
m.backendRTT.DeleteLabelValues(upstream)
m.backendTier.DeleteLabelValues(upstream)
}

// SetBackendRTT records the smoothed health check connect time of an upstream.
func (m *Metrics) SetBackendRTT(upstream string, rtt time.Duration) {
m.backendRTT.WithLabelValues(upstream).Set(rtt.Seconds())
}

// SetBackendTier records the latency tier of an upstream.
func (m *Metrics) SetBackendTier(upstream string, tier int) {
m.backendTier.WithLabelValues(upstream).Set(float64(tier))
}
36 changes: 36 additions & 0 deletions internal/metrics/metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"net/http"
"net/http/httptest"
"testing"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
Expand Down Expand Up @@ -124,3 +125,38 @@ func TestNew_IsolatedRegistries(t *testing.T) {
assert.Contains(t, body1, "extractedprism_connections_active 1")
assert.Contains(t, body2, "extractedprism_connections_active 0")
}

func TestSetBackendRTT_ExportsSeconds(t *testing.T) {
m := metrics.New()

m.SetBackendRTT("192.0.2.1:6443", 1500*time.Microsecond)

body := scrape(t, m)
assert.Contains(t, body, `extractedprism_upstream_rtt_seconds{upstream="192.0.2.1:6443"} 0.0015`)
}

func TestSetBackendTier_ExportsTierIndex(t *testing.T) {
m := metrics.New()

m.SetBackendTier("192.0.2.1:6443", 0)
m.SetBackendTier("192.0.2.2:6443", 2)

body := scrape(t, m)
assert.Contains(t, body, `extractedprism_upstream_latency_tier{upstream="192.0.2.1:6443"} 0`)
assert.Contains(t, body, `extractedprism_upstream_latency_tier{upstream="192.0.2.2:6443"} 2`)
}

func TestRemoveBackend_DeletesLatencySeries(t *testing.T) {
m := metrics.New()

m.SetBackendRTT("192.0.2.1:6443", time.Millisecond)
m.SetBackendTier("192.0.2.1:6443", 0)
m.SetBackendRTT("192.0.2.2:6443", time.Millisecond)
m.RemoveBackend("192.0.2.1:6443")

body := scrape(t, m)
assert.NotContains(t, body, `extractedprism_upstream_rtt_seconds{upstream="192.0.2.1:6443"}`)
assert.NotContains(t, body, `extractedprism_upstream_latency_tier{upstream="192.0.2.1:6443"}`)
assert.Contains(t, body, `extractedprism_upstream_rtt_seconds{upstream="192.0.2.2:6443"}`,
"other upstreams keep their series")
}
Loading
Loading