mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
atenet/router: make the ext_proc circuit breaker an explicit, validated flag
Add --extproc-max-requests (default 2048) and set circuit_breakers.max_requests on the ext_proc cluster from it, replacing the implicit Envoy default and the hand-maintained doc coupling with a guarantee. Every request's header exchange occupies one slot briefly and every parked request holds one for its entire wait, so startup validation enforces extproc-max-requests >= parked-request-max — a breaker below the lot silently truncates it with Envoy-generated 503s that bypass parking.rejected. The default leaves the lot's worth of fast-path headroom (1024 lot / 2048 breaker), so a saturated lot cannot starve requests to already-running actors.
This commit is contained in:
@@ -65,6 +65,7 @@ func NewRouterCmd() *cobra.Command {
|
||||
cmd.Flags().DurationVar(&cfg.ParkedRequestRetryInterval, "parked-request-retry-interval", defaultParkedRequestRetryInterval, "Delay before a parked request's first resume retry")
|
||||
cmd.Flags().Float64Var(&cfg.ParkedRequestRetryFactor, "parked-request-retry-factor", defaultParkedRequestRetryFactor, "Multiplier applied to the retry delay after each attempt; must be >= 1")
|
||||
cmd.Flags().Float64Var(&cfg.ParkedRequestRetryJitter, "parked-request-retry-jitter", defaultParkedRequestRetryJitter, "Random fraction in [0, 1) added to each retry delay to de-synchronize parked requests")
|
||||
cmd.Flags().IntVar(&cfg.ExtProcMaxRequests, "extproc-max-requests", defaultExtProcMaxRequests, "Circuit-breaker max_requests for Envoy's ext_proc cluster; every parked request holds one slot for its full wait, so this must be >= --parked-request-max (the excess is fast-path headroom)")
|
||||
|
||||
return cmd
|
||||
}
|
||||
|
||||
@@ -15,6 +15,7 @@
|
||||
package router
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
@@ -64,4 +65,23 @@ type routerConfig struct {
|
||||
ParkedRequestRetryInterval time.Duration
|
||||
ParkedRequestRetryFactor float64
|
||||
ParkedRequestRetryJitter float64
|
||||
|
||||
// ExtProcMaxRequests is the circuit-breaker max_requests Envoy applies to
|
||||
// the ext_proc cluster. Every parked request holds one slot for its entire
|
||||
// wait, so this must be >= ParkedRequestMax (validated at startup); the
|
||||
// excess is fast-path headroom for requests to already-running actors.
|
||||
ExtProcMaxRequests int
|
||||
}
|
||||
|
||||
// validate rejects flag combinations that would make the router misbehave
|
||||
// rather than merely differ.
|
||||
func (c routerConfig) validate() error {
|
||||
if c.ExtProcMaxRequests <= 0 {
|
||||
return fmt.Errorf("--extproc-max-requests must be positive, got %d", c.ExtProcMaxRequests)
|
||||
}
|
||||
if c.ParkedRequestMax > 0 && c.ExtProcMaxRequests < c.ParkedRequestMax {
|
||||
return fmt.Errorf("--extproc-max-requests (%d) must be >= --parked-request-max (%d): a circuit breaker below the parking lot silently truncates it with Envoy-generated 503s",
|
||||
c.ExtProcMaxRequests, c.ParkedRequestMax)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,65 @@
|
||||
// Copyright 2026 Google LLC
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package router
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestRouterConfigValidate(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
cfg routerConfig
|
||||
wantErr string // substring; empty means valid
|
||||
}{
|
||||
{
|
||||
name: "defaults are valid",
|
||||
cfg: routerConfig{ExtProcMaxRequests: defaultExtProcMaxRequests, ParkedRequestMax: defaultParkedRequestMax},
|
||||
},
|
||||
{
|
||||
name: "zero extproc-max-requests rejected",
|
||||
cfg: routerConfig{ExtProcMaxRequests: 0, ParkedRequestMax: 0},
|
||||
wantErr: "must be positive",
|
||||
},
|
||||
{
|
||||
name: "breaker below the lot rejected",
|
||||
cfg: routerConfig{ExtProcMaxRequests: 512, ParkedRequestMax: 1024},
|
||||
wantErr: "must be >= --parked-request-max",
|
||||
},
|
||||
{
|
||||
name: "breaker equal to the lot accepted",
|
||||
cfg: routerConfig{ExtProcMaxRequests: 1024, ParkedRequestMax: 1024},
|
||||
},
|
||||
{
|
||||
name: "parking disabled ignores the relation",
|
||||
cfg: routerConfig{ExtProcMaxRequests: 8, ParkedRequestMax: 0},
|
||||
},
|
||||
}
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
err := tc.cfg.validate()
|
||||
if tc.wantErr == "" {
|
||||
if err != nil {
|
||||
t.Fatalf("expected valid, got %v", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
if err == nil || !strings.Contains(err.Error(), tc.wantErr) {
|
||||
t.Fatalf("expected error containing %q, got %v", tc.wantErr, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -28,14 +28,12 @@ import (
|
||||
const (
|
||||
defaultParkedRequestBudget = 5 * time.Second
|
||||
|
||||
// defaultParkedRequestMax is sized to Envoy's default per-cluster circuit
|
||||
// breaker (max_requests = 1024): each parked request holds one ext_proc
|
||||
// stream, i.e. one active request against the ext_proc cluster, and that
|
||||
// cluster (buildCluster in xds.go) sets no explicit circuit_breakers. A lot
|
||||
// larger than the circuit breaker is unreachable — Envoy rejects the
|
||||
// overflow itself with 503s that never reach the lot (and so never count in
|
||||
// parking.rejected). Raise --parked-request-max beyond 1024 only together
|
||||
// with an explicit circuit_breakers.max_requests on the ext_proc cluster.
|
||||
// defaultParkedRequestMax is sized together with the ext_proc cluster's
|
||||
// circuit breaker (--extproc-max-requests, default 2048): each parked
|
||||
// request holds one ext_proc stream, i.e. one active request against that
|
||||
// cluster, for its entire wait. Startup validation keeps the breaker >= the
|
||||
// lot, and the default pair (1024 lot / 2048 breaker) leaves equal headroom
|
||||
// for the fast path. See buildCluster in xds.go.
|
||||
defaultParkedRequestMax = 1024
|
||||
|
||||
// Retry cadence between resume attempts while a request is parked: a gentle
|
||||
|
||||
@@ -126,6 +126,26 @@ func (s *RouterServer) Run(ctx context.Context) error {
|
||||
cancel()
|
||||
}()
|
||||
|
||||
// Validate the configuration before doing any other work, so a bad flag
|
||||
// combination fails fast — no tracing, metrics, or connections are set up
|
||||
// for a router that is about to refuse to start. The parking config is
|
||||
// resolved once here so every consumer — the resumer's retry loop and the
|
||||
// Envoy ext_proc timeout — sees the same effective values.
|
||||
parkCfg := parkingConfig{
|
||||
budget: s.cfg.ParkedRequestBudget,
|
||||
maxParked: s.cfg.ParkedRequestMax,
|
||||
retryInterval: s.cfg.ParkedRequestRetryInterval,
|
||||
retryFactor: s.cfg.ParkedRequestRetryFactor,
|
||||
retryJitter: s.cfg.ParkedRequestRetryJitter,
|
||||
}
|
||||
if err := parkCfg.validate(); err != nil {
|
||||
return fmt.Errorf("invalid parking configuration: %w", err)
|
||||
}
|
||||
if err := s.cfg.validate(); err != nil {
|
||||
return fmt.Errorf("invalid router configuration: %w", err)
|
||||
}
|
||||
parkCfg = parkCfg.normalized()
|
||||
|
||||
var level slog.Level
|
||||
switch strings.ToLower(s.cfg.LogLevel) {
|
||||
case "debug":
|
||||
@@ -190,21 +210,7 @@ func (s *RouterServer) Run(ctx context.Context) error {
|
||||
return fmt.Errorf("configure OTLP collector: %w", err)
|
||||
}
|
||||
|
||||
// Resolve the parking configuration once so every consumer — the resumer's
|
||||
// retry loop and the Envoy ext_proc timeout below — sees the same effective
|
||||
// values (a non-positive budget falls back to the default rather than
|
||||
// leaving Envoy on its short parking-off timeout).
|
||||
parkCfg := parkingConfig{
|
||||
budget: s.cfg.ParkedRequestBudget,
|
||||
maxParked: s.cfg.ParkedRequestMax,
|
||||
retryInterval: s.cfg.ParkedRequestRetryInterval,
|
||||
retryFactor: s.cfg.ParkedRequestRetryFactor,
|
||||
retryJitter: s.cfg.ParkedRequestRetryJitter,
|
||||
}
|
||||
if err := parkCfg.validate(); err != nil {
|
||||
return fmt.Errorf("invalid parking configuration: %w", err)
|
||||
}
|
||||
parkCfg = parkCfg.normalized()
|
||||
xdsSrv.SetExtProcMaxRequests(s.cfg.ExtProcMaxRequests)
|
||||
if parkCfg.enabled() {
|
||||
// Envoy must keep a parked request open at least as long as the router
|
||||
// will hold it; add a margin so the router surfaces its own 503 first.
|
||||
|
||||
@@ -59,14 +59,15 @@ func TestStatuszEndpoint(t *testing.T) {
|
||||
caPath, clientCertPath := writeTestTLSMaterial(t)
|
||||
|
||||
cfg := routerConfig{
|
||||
Standalone: true,
|
||||
Namespace: "default",
|
||||
StatusPort: httpPort,
|
||||
HttpPort: 8080,
|
||||
XdsPort: 18000,
|
||||
ExtprocPort: 50051,
|
||||
TemplatesFile: tmpFile.Name(),
|
||||
MetricsAddr: "127.0.0.1:0",
|
||||
Standalone: true,
|
||||
Namespace: "default",
|
||||
StatusPort: httpPort,
|
||||
HttpPort: 8080,
|
||||
XdsPort: 18000,
|
||||
ExtprocPort: 50051,
|
||||
ExtProcMaxRequests: defaultExtProcMaxRequests,
|
||||
TemplatesFile: tmpFile.Name(),
|
||||
MetricsAddr: "127.0.0.1:0",
|
||||
Auth: authConfig{
|
||||
AteapiCAFile: caPath,
|
||||
AteapiClientCertPath: clientCertPath,
|
||||
|
||||
@@ -79,6 +79,12 @@ const (
|
||||
// otherwise Envoy abandons a parked request (500) long before the router does.
|
||||
const defaultExtProcMessageTimeout = 5 * time.Second
|
||||
|
||||
// defaultExtProcMaxRequests is the circuit-breaker max_requests set on the
|
||||
// ext_proc cluster: defaultParkedRequestMax plus equal fast-path headroom, so a
|
||||
// full parking lot cannot starve the millisecond-scale header exchanges of
|
||||
// requests to already-running actors. See buildCluster.
|
||||
const defaultExtProcMaxRequests = 2048
|
||||
|
||||
// XdsServer implements an aggregated discovery service server for dynamic Envoy router nodes.
|
||||
type XdsServer struct {
|
||||
xdsPort int
|
||||
@@ -100,6 +106,12 @@ type XdsServer struct {
|
||||
// extProcMessageTimeout bounds how long Envoy waits for the router's ext_proc
|
||||
// response. Must be >= the parking budget so parked requests aren't cut short.
|
||||
extProcMessageTimeout time.Duration
|
||||
|
||||
// extProcMaxRequests is the circuit-breaker max_requests on the ext_proc
|
||||
// cluster — the hard ceiling on concurrent requests held open against the
|
||||
// router's processing server, parked requests included. Must be >= the
|
||||
// parking lot size (enforced at startup in Run).
|
||||
extProcMaxRequests uint32
|
||||
}
|
||||
|
||||
func NewXdsServer(xdsPort int) *XdsServer {
|
||||
@@ -114,6 +126,7 @@ func NewXdsServer(xdsPort int) *XdsServer {
|
||||
extprocAddr: "127.0.0.1",
|
||||
ingressPort: 8080,
|
||||
extProcMessageTimeout: defaultExtProcMessageTimeout,
|
||||
extProcMaxRequests: defaultExtProcMaxRequests,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -137,6 +150,17 @@ func (x *XdsServer) SetExtProcMessageTimeout(d time.Duration) {
|
||||
}
|
||||
}
|
||||
|
||||
// SetExtProcMaxRequests sets the circuit-breaker max_requests on the ext_proc
|
||||
// cluster. Size it to the parking lot plus fast-path headroom (validated in
|
||||
// Run()); a non-positive value leaves the default unchanged.
|
||||
func (x *XdsServer) SetExtProcMaxRequests(n int) {
|
||||
x.mu.Lock()
|
||||
defer x.mu.Unlock()
|
||||
if n > 0 {
|
||||
x.extProcMaxRequests = uint32(n)
|
||||
}
|
||||
}
|
||||
|
||||
func (x *XdsServer) SetTlsConfig(httpsPort int, certPath string) {
|
||||
x.mu.Lock()
|
||||
defer x.mu.Unlock()
|
||||
@@ -248,17 +272,6 @@ func (x *XdsServer) Serve(ctx context.Context, lis net.Listener) error {
|
||||
}
|
||||
}
|
||||
|
||||
// buildCluster builds the ext_proc cluster Envoy uses to reach the router's
|
||||
// processing server.
|
||||
//
|
||||
// It sets no explicit circuit_breakers, so Envoy's defaults apply — notably
|
||||
// max_requests = 1024 concurrent requests against this cluster. Each PARKED
|
||||
// request holds one ext_proc stream, i.e. one active request here, which makes
|
||||
// that circuit breaker the true upper bound on concurrent parked requests:
|
||||
// defaultParkedRequestMax (parking.go) is deliberately sized to it. If
|
||||
// --parked-request-max is ever raised beyond 1024, add an explicit
|
||||
// circuit_breakers.max_requests >= that value here, or the overflow is
|
||||
// rejected by Envoy itself (503s that bypass the lot and parking.rejected).
|
||||
func (x *XdsServer) buildCluster() *clusterv3.Cluster {
|
||||
h2Opts, _ := anypb.New(&httpv3.HttpProtocolOptions{
|
||||
UpstreamProtocolOptions: &httpv3.HttpProtocolOptions_ExplicitHttpConfig_{
|
||||
@@ -275,6 +288,12 @@ func (x *XdsServer) buildCluster() *clusterv3.Cluster {
|
||||
Type: clusterv3.Cluster_STATIC,
|
||||
},
|
||||
LbPolicy: clusterv3.Cluster_ROUND_ROBIN,
|
||||
CircuitBreakers: &clusterv3.CircuitBreakers{
|
||||
Thresholds: []*clusterv3.CircuitBreakers_Thresholds{{
|
||||
Priority: corev3.RoutingPriority_DEFAULT,
|
||||
MaxRequests: wrapperspb.UInt32(x.extProcMaxRequests),
|
||||
}},
|
||||
},
|
||||
LoadAssignment: &endpointv3.ClusterLoadAssignment{
|
||||
ClusterName: ClusterName,
|
||||
Endpoints: []*endpointv3.LocalityLbEndpoints{
|
||||
|
||||
@@ -259,3 +259,34 @@ func TestDynamicForwardProxyCluster_DisablesConnectionReuse(t *testing.T) {
|
||||
t.Errorf("Expected max_requests_per_connection 1, got %d", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestXdsServer_ExtProcCircuitBreaker(t *testing.T) {
|
||||
t.Run("DefaultCoversLotPlusHeadroom", func(t *testing.T) {
|
||||
x := NewXdsServer(0)
|
||||
got := x.buildCluster().GetCircuitBreakers().GetThresholds()[0].GetMaxRequests().GetValue()
|
||||
if got != uint32(defaultExtProcMaxRequests) {
|
||||
t.Errorf("default max_requests = %d, want %d", got, defaultExtProcMaxRequests)
|
||||
}
|
||||
if got < uint32(defaultParkedRequestMax) {
|
||||
t.Errorf("default breaker (%d) below the default lot (%d): a full lot would be truncated by Envoy", got, defaultParkedRequestMax)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("SetterOverrides", func(t *testing.T) {
|
||||
x := NewXdsServer(0)
|
||||
x.SetExtProcMaxRequests(4096)
|
||||
got := x.buildCluster().GetCircuitBreakers().GetThresholds()[0].GetMaxRequests().GetValue()
|
||||
if got != 4096 {
|
||||
t.Errorf("max_requests after SetExtProcMaxRequests(4096) = %d, want 4096", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("NonPositiveKeepsDefault", func(t *testing.T) {
|
||||
x := NewXdsServer(0)
|
||||
x.SetExtProcMaxRequests(0)
|
||||
got := x.buildCluster().GetCircuitBreakers().GetThresholds()[0].GetMaxRequests().GetValue()
|
||||
if got != uint32(defaultExtProcMaxRequests) {
|
||||
t.Errorf("max_requests after SetExtProcMaxRequests(0) = %d, want default %d", got, defaultExtProcMaxRequests)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
+13
-9
@@ -49,14 +49,17 @@ in-flight resume occupies one slot. When the lot is full, further requests are
|
||||
shed immediately with `503 "actor <id> unavailable: router at capacity"` rather
|
||||
than queueing without bound.
|
||||
|
||||
The default lot size is deliberately equal to Envoy's default per-cluster
|
||||
circuit breaker (`max_requests = 1024`) on the ext_proc cluster: every parked
|
||||
request holds one ext_proc stream, i.e. one active request against that
|
||||
cluster, so the circuit breaker is the true upper bound on concurrent parked
|
||||
requests. Raising `--parked-request-max` beyond `1024` requires an explicit
|
||||
`circuit_breakers.max_requests` on the ext_proc cluster to match — otherwise
|
||||
the overflow is rejected by Envoy itself, with 503s that never reach the lot
|
||||
(and never count in `parking.rejected`).
|
||||
Every parked request holds one ext_proc stream — one active request against
|
||||
Envoy's ext_proc cluster — for its entire wait, while ordinary requests hold
|
||||
one only for a millisecond-scale header exchange. The cluster's circuit breaker
|
||||
is therefore the hard ceiling on concurrent parked requests, and the router
|
||||
sets it explicitly from `--extproc-max-requests` (default `2048`). Startup
|
||||
validation enforces `--extproc-max-requests >= --parked-request-max`; the
|
||||
excess is **fast-path headroom**, so a saturated lot cannot starve requests to
|
||||
already-running actors (the default pair is a `1024` lot with `1024` of
|
||||
headroom). A breaker below the lot would silently truncate it — Envoy would
|
||||
reject the overflow itself, with 503s that never reach the lot and never count
|
||||
in `parking.rejected`.
|
||||
|
||||
Concurrent requests for the *same* actor are de-duplicated by the resumer's
|
||||
`singleflight` group: they share a single in-flight `ResumeActor` call and all
|
||||
@@ -98,10 +101,11 @@ within a `15s` budget.
|
||||
| Flag | Default | Meaning |
|
||||
| -------------------------------- | ------- | ------------------------------------------------------------------ |
|
||||
| `--parked-request-budget` | `5s` | Park budget per resume *flight*; requests de-duplicated onto an in-flight resume share its remaining budget (see Behavior). |
|
||||
| `--parked-request-max` | `1024` | Max concurrent parked/in-flight resume requests; excess shed (503). `0` disables parking. Sized to Envoy's ext_proc circuit breaker (see above). |
|
||||
| `--parked-request-max` | `1024` | Max concurrent parked/in-flight resume requests; excess shed (503). `0` disables parking. |
|
||||
| `--parked-request-retry-interval` | `100ms` | Delay before a parked request's first resume retry. |
|
||||
| `--parked-request-retry-factor` | `1.1` | Multiplier applied to the retry delay after each attempt (>= 1). |
|
||||
| `--parked-request-retry-jitter` | `0.1` | Random fraction in `[0, 1)` added per retry to de-synchronize parked requests. |
|
||||
| `--extproc-max-requests` | `2048` | Envoy circuit-breaker `max_requests` for the ext_proc cluster. Must be `>= --parked-request-max` (enforced at startup); the excess is fast-path headroom (see Behavior). |
|
||||
|
||||
The retry backoff deliberately has no cap and no attempt limit: the budget alone
|
||||
bounds the wait.
|
||||
|
||||
Reference in New Issue
Block a user