mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
atenet/router: address self-review feedback
- Adopt the ParkedRequest* vocabulary for parking flags and config (bowei's suggestion): --parked-request-budget / --parked-request-max, matching fields and default consts. - Make the parked-retry backoff configurable: --parked-request-retry-interval/-factor/-jitter, validated at startup (factor >= 1, jitter in [0,1)); the backoff still has no cap and no attempt limit, so the budget alone bounds the wait. - Resolve the effective parking config once in Run() so the resumer's retry loop and the Envoy ext_proc timeout always agree, even when the budget flag is set non-positive. - Drop timeline-relative wording from docs and identifiers (failFastResumeBudget, fail-fast behavior). - Guard the parking-lot counter against going negative, loudly. - Document exactly when the wait-duration metric is recorded and what each outcome label means.
This commit is contained in:
@@ -60,8 +60,11 @@ func NewRouterCmd() *cobra.Command {
|
||||
cmd.Flags().StringVar(&cfg.Auth.AteapiServerName, "ateapi-server-name", "", "SNI / hostname expected on the ateapi server cert. Optional.")
|
||||
cmd.Flags().BoolVar(&cfg.Auth.AteapiUseTokenAuth, "ateapi-use-token-auth", false, "Authenticate to ateapi with the Bearer token from --ateapi-token-file instead of the client certificate from --ateapi-client-cert.")
|
||||
cmd.Flags().StringVar(&cfg.Auth.AteapiTokenFile, "ateapi-token-file", "", "Projected SA token file used as Bearer credential. Required with --ateapi-use-token-auth, ignored otherwise.")
|
||||
cmd.Flags().DurationVar(&cfg.ParkingMaxWait, "parking-max-wait", defaultParkingMaxWait, "Maximum time a request may be parked (held and retried) waiting for its actor to become routable")
|
||||
cmd.Flags().IntVar(&cfg.ParkingMaxParked, "parking-max-parked", defaultParkingMaxParked, "Maximum number of requests that may be parked simultaneously; excess requests are shed with 503. 0 disables parking (requests fail fast on worker-pool saturation)")
|
||||
cmd.Flags().DurationVar(&cfg.ParkedRequestBudget, "parked-request-budget", defaultParkedRequestBudget, "Maximum time a request may be parked (held and retried) waiting for its actor to become routable")
|
||||
cmd.Flags().IntVar(&cfg.ParkedRequestMax, "parked-request-max", defaultParkedRequestMax, "Maximum number of requests that may be parked simultaneously; excess requests are shed with 503. 0 disables parking (requests fail fast on worker-pool saturation)")
|
||||
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")
|
||||
|
||||
return cmd
|
||||
}
|
||||
|
||||
@@ -58,7 +58,10 @@ type routerConfig struct {
|
||||
|
||||
// Request parking: hold and retry requests whose actor cannot be served
|
||||
// immediately due to transient worker-pool saturation, instead of failing
|
||||
// fast. A non-positive ParkingMaxParked disables parking. See parkingConfig.
|
||||
ParkingMaxWait time.Duration
|
||||
ParkingMaxParked int
|
||||
// fast. A non-positive ParkedRequestMax disables parking. See parkingConfig.
|
||||
ParkedRequestBudget time.Duration
|
||||
ParkedRequestMax int
|
||||
ParkedRequestRetryInterval time.Duration
|
||||
ParkedRequestRetryFactor float64
|
||||
ParkedRequestRetryJitter float64
|
||||
}
|
||||
|
||||
@@ -51,7 +51,7 @@ func NewExtProcServer(port int, apiClient ateapipb.ControlClient, routeDuration
|
||||
port: port,
|
||||
apiClient: apiClient,
|
||||
recorder: NewQueryRecorder(100),
|
||||
resumer: NewActorResumer(apiClient, withParking(parkCfg.enabled(), parkCfg.maxWait)),
|
||||
resumer: NewActorResumer(apiClient, withParking(parkCfg)),
|
||||
routeDuration: routeDuration,
|
||||
parking: newParkingLot(parkCfg, parkMetrics),
|
||||
}
|
||||
|
||||
@@ -271,7 +271,7 @@ func TestExtProc_ParkingLotFull(t *testing.T) {
|
||||
|
||||
// A 1-slot lot with the slot already occupied deterministically simulates a
|
||||
// full lot without needing a concurrent in-flight request.
|
||||
s := NewExtProcServer(50051, clientMock, nil, parkingConfig{maxWait: time.Second, maxParked: 1}, nil)
|
||||
s := NewExtProcServer(50051, clientMock, nil, parkingConfig{budget: time.Second, maxParked: 1}, nil)
|
||||
occupy, ok := s.parking.enter(context.Background())
|
||||
if !ok {
|
||||
t.Fatal("priming enter should be admitted")
|
||||
|
||||
@@ -17,15 +17,23 @@ package router
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Default request-parking parameters. See parkingConfig for the meaning of each
|
||||
// field; these are also the flag defaults wired up in NewCmd.
|
||||
// field; these are also the flag defaults wired up in NewRouterCmd.
|
||||
const (
|
||||
defaultParkingMaxWait = 30 * time.Second
|
||||
defaultParkingMaxParked = 2048
|
||||
defaultParkedRequestBudget = 5 * time.Second
|
||||
defaultParkedRequestMax = 2048
|
||||
|
||||
// Retry cadence between resume attempts while a request is parked: a gentle
|
||||
// exponential backoff.
|
||||
defaultParkedRequestRetryInterval = 100 * time.Millisecond
|
||||
defaultParkedRequestRetryFactor = 1.1
|
||||
defaultParkedRequestRetryJitter = 0.1
|
||||
)
|
||||
|
||||
// parkOutcome is the terminal disposition of a parked request. It is recorded
|
||||
@@ -47,26 +55,66 @@ const (
|
||||
// control plane before routing. If the worker pool is momentarily saturated the
|
||||
// control plane returns FailedPrecondition ("no free workers available"). With
|
||||
// parking enabled the router holds ("parks") such a request and keeps retrying
|
||||
// the resume until the actor becomes routable or maxWait elapses, instead of
|
||||
// the resume until the actor becomes routable or budget elapses, instead of
|
||||
// failing the request immediately. maxParked bounds how many requests may be
|
||||
// parked at once so the router sheds load rather than queueing without bound;
|
||||
// a non-positive maxParked disables parking entirely.
|
||||
//
|
||||
// retryInterval/retryFactor/retryJitter shape the backoff between resume
|
||||
// attempts while a request is parked. The backoff deliberately has no cap and
|
||||
// no step limit: the budget alone bounds the wait.
|
||||
type parkingConfig struct {
|
||||
maxWait time.Duration
|
||||
budget time.Duration
|
||||
maxParked int
|
||||
|
||||
retryInterval time.Duration
|
||||
retryFactor float64
|
||||
retryJitter float64
|
||||
}
|
||||
|
||||
// enabled reports whether request parking is active. Parking has no separate
|
||||
// on/off switch: setting maxParked to 0 disables it, preserving the legacy
|
||||
// fail-fast behavior (no admission cap, no retry on pool saturation).
|
||||
// on/off switch: setting maxParked to 0 disables it, applying a fail-fast
|
||||
// behavior (no admission cap, no retry on pool saturation).
|
||||
func (c parkingConfig) enabled() bool { return c.maxParked > 0 }
|
||||
|
||||
// normalized returns the config with non-positive budget and retry parameters
|
||||
// replaced by their defaults, so every consumer (the resumer's retry loop and
|
||||
// the Envoy ext_proc timeout) sees the same effective values.
|
||||
func (c parkingConfig) normalized() parkingConfig {
|
||||
if c.budget <= 0 {
|
||||
c.budget = defaultParkedRequestBudget
|
||||
}
|
||||
if c.retryInterval <= 0 {
|
||||
c.retryInterval = defaultParkedRequestRetryInterval
|
||||
}
|
||||
if c.retryFactor == 0 {
|
||||
c.retryFactor = defaultParkedRequestRetryFactor
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// validate rejects retry parameters that would make parking misbehave rather
|
||||
// than merely differ: a factor below 1 shrinks delays toward zero and turns
|
||||
// the parked retry loop into a hot loop against the control plane.
|
||||
func (c parkingConfig) validate() error {
|
||||
if c.retryFactor != 0 && c.retryFactor < 1.0 {
|
||||
return fmt.Errorf("parked-request retry factor must be >= 1.0, got %v", c.retryFactor)
|
||||
}
|
||||
if c.retryJitter < 0 || c.retryJitter >= 1 {
|
||||
return fmt.Errorf("parked-request retry jitter must be in [0, 1), got %v", c.retryJitter)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// defaultParkingConfig returns the built-in parking configuration (matching the
|
||||
// NewCmd flag defaults).
|
||||
// NewRouterCmd flag defaults).
|
||||
func defaultParkingConfig() parkingConfig {
|
||||
return parkingConfig{
|
||||
maxWait: defaultParkingMaxWait,
|
||||
maxParked: defaultParkingMaxParked,
|
||||
budget: defaultParkedRequestBudget,
|
||||
maxParked: defaultParkedRequestMax,
|
||||
retryInterval: defaultParkedRequestRetryInterval,
|
||||
retryFactor: defaultParkedRequestRetryFactor,
|
||||
retryJitter: defaultParkedRequestRetryJitter,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -76,7 +124,7 @@ func defaultParkingConfig() parkingConfig {
|
||||
// router applies backpressure instead of accumulating waiters without bound.
|
||||
//
|
||||
// With parking disabled (maxParked <= 0) enter always admits and performs no
|
||||
// accounting, preserving the router's legacy behavior.
|
||||
// accounting, applying the router's fail-fast behavior.
|
||||
type parkingLot struct {
|
||||
cfg parkingConfig
|
||||
metrics *parkingMetrics
|
||||
@@ -116,7 +164,14 @@ func (l *parkingLot) enter(ctx context.Context) (release func(outcome parkOutcom
|
||||
return func(outcome parkOutcome) {
|
||||
once.Do(func() {
|
||||
l.mu.Lock()
|
||||
l.active--
|
||||
// The counter cannot go negative today (a release only exists after a
|
||||
// successful enter, and it is Once-guarded), so a violation means an
|
||||
// accounting bug elsewhere: clamp, but say so loudly.
|
||||
if l.active > 0 {
|
||||
l.active--
|
||||
} else {
|
||||
slog.Error("parking lot slot released more times than acquired")
|
||||
}
|
||||
l.mu.Unlock()
|
||||
l.metrics.addActive(ctx, -1)
|
||||
l.metrics.recordWait(ctx, time.Since(start), outcome)
|
||||
@@ -137,7 +192,7 @@ func (l *parkingLot) status() ParkingStatus {
|
||||
Enabled: l.cfg.enabled(),
|
||||
Active: l.activeCount(),
|
||||
MaxParked: l.cfg.maxParked,
|
||||
MaxWait: l.cfg.maxWait.String(),
|
||||
MaxWait: l.cfg.budget.String(),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -22,7 +22,7 @@ import (
|
||||
)
|
||||
|
||||
func TestParkingLot_CapacityAndRelease(t *testing.T) {
|
||||
lot := newParkingLot(parkingConfig{maxWait: time.Second, maxParked: 2}, nil)
|
||||
lot := newParkingLot(parkingConfig{budget: time.Second, maxParked: 2}, nil)
|
||||
ctx := context.Background()
|
||||
|
||||
r1, ok := lot.enter(ctx)
|
||||
@@ -60,7 +60,7 @@ func TestParkingLot_CapacityAndRelease(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestParkingLot_ReleaseIsIdempotent(t *testing.T) {
|
||||
lot := newParkingLot(parkingConfig{maxWait: time.Second, maxParked: 1}, nil)
|
||||
lot := newParkingLot(parkingConfig{budget: time.Second, maxParked: 1}, nil)
|
||||
|
||||
release, ok := lot.enter(context.Background())
|
||||
if !ok {
|
||||
@@ -96,7 +96,7 @@ func TestParkingLot_DisabledAlwaysAdmits(t *testing.T) {
|
||||
func TestParkingLot_ConcurrentEntryRespectsCapacity(t *testing.T) {
|
||||
const capacity = 8
|
||||
const goroutines = 100
|
||||
lot := newParkingLot(parkingConfig{maxWait: time.Second, maxParked: capacity}, nil)
|
||||
lot := newParkingLot(parkingConfig{budget: time.Second, maxParked: capacity}, nil)
|
||||
|
||||
var admitted int64
|
||||
var mu sync.Mutex
|
||||
|
||||
@@ -31,32 +31,27 @@ import (
|
||||
"k8s.io/apimachinery/pkg/util/wait"
|
||||
)
|
||||
|
||||
// legacyResumeBudget is the total time the resumer spends retrying a resume when
|
||||
// request parking is disabled. It preserves the historical fail-fast-on-capacity
|
||||
// behavior (only concurrent-update conflicts are retried).
|
||||
const legacyResumeBudget = 15 * time.Second
|
||||
// failFastResumeBudget is the total time the resumer spends retrying a resume
|
||||
// when request parking is disabled. In that mode only concurrent-update
|
||||
// conflicts are retried; capacity errors fail immediately.
|
||||
const failFastResumeBudget = 15 * time.Second
|
||||
|
||||
// Retry cadence between resume attempts: a gentle exponential backoff.
|
||||
const (
|
||||
resumeBackoffBase = 500 * time.Millisecond
|
||||
resumeBackoffFactor = 1.1
|
||||
resumeBackoffJitter = 0.1
|
||||
)
|
||||
|
||||
// resumeBackoff is the backoff between resume attempts while a request is parked.
|
||||
// resumeBackoff builds the backoff between resume attempts while a request is
|
||||
// parked, from the configured retry parameters.
|
||||
//
|
||||
// It intentionally sets NO Cap. wait.Backoff's delay() zeroes Steps the moment
|
||||
// the delay reaches Cap, which would end retries long before the parking budget
|
||||
// (a Cap of 2s stops the loop in ~7 steps / ~5s regardless of the budget). A
|
||||
// gentle Factor keeps the gap small on its own — from 500ms it only grows to
|
||||
// ~3.5s over a 30s budget — while Steps is set high so the budget context passed
|
||||
// to ExponentialBackoffWithContext, not the step count, bounds the wait.
|
||||
func resumeBackoff() wait.Backoff {
|
||||
// (a Cap of 2s stops the loop in ~7 steps regardless of the budget). A gentle
|
||||
// Factor keeps the gap small on its own — from 100ms at the default 1.1 the gap
|
||||
// only grows to ~0.5s over a 5s budget — while Steps is set high so the budget
|
||||
// context passed to ExponentialBackoffWithContext, not the step count, bounds
|
||||
// the wait.
|
||||
func resumeBackoff(interval time.Duration, factor, jitter float64) wait.Backoff {
|
||||
return wait.Backoff{
|
||||
Steps: math.MaxInt32,
|
||||
Duration: resumeBackoffBase,
|
||||
Factor: resumeBackoffFactor,
|
||||
Jitter: resumeBackoffJitter,
|
||||
Duration: interval,
|
||||
Factor: factor,
|
||||
Jitter: jitter,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -82,28 +77,34 @@ type ActorResumer struct {
|
||||
// budget bounds the total time a single resume operation retries before the
|
||||
// underlying error is returned.
|
||||
budget time.Duration
|
||||
// backoff paces the retries within the budget.
|
||||
backoff wait.Backoff
|
||||
}
|
||||
|
||||
// resumerOption configures an ActorResumer.
|
||||
type resumerOption func(*ActorResumer)
|
||||
|
||||
// withParking configures parking behavior. When enabled, FailedPrecondition
|
||||
// ("no free workers available") becomes retryable and the resume is retried for
|
||||
// up to maxWait; a non-positive maxWait keeps the default budget. When disabled,
|
||||
// the resumer preserves its legacy fail-fast-on-capacity behavior.
|
||||
func withParking(enabled bool, maxWait time.Duration) resumerOption {
|
||||
// withParking configures parking behavior from cfg. When parking is enabled,
|
||||
// FailedPrecondition ("no free workers available") becomes retryable and the
|
||||
// resume is retried, at cfg's retry cadence, for up to cfg's budget. When
|
||||
// disabled, the resumer applies fail-fast-on-capacity behavior.
|
||||
func withParking(cfg parkingConfig) resumerOption {
|
||||
cfg = cfg.normalized()
|
||||
return func(r *ActorResumer) {
|
||||
r.parkEnabled = enabled
|
||||
if maxWait > 0 {
|
||||
r.budget = maxWait
|
||||
r.parkEnabled = cfg.enabled()
|
||||
if r.parkEnabled {
|
||||
r.budget = cfg.budget
|
||||
}
|
||||
r.backoff = resumeBackoff(cfg.retryInterval, cfg.retryFactor, cfg.retryJitter)
|
||||
}
|
||||
}
|
||||
|
||||
func NewActorResumer(apiClient ateapipb.ControlClient, opts ...resumerOption) *ActorResumer {
|
||||
r := &ActorResumer{
|
||||
apiClient: apiClient,
|
||||
budget: legacyResumeBudget,
|
||||
budget: failFastResumeBudget,
|
||||
backoff: resumeBackoff(defaultParkedRequestRetryInterval,
|
||||
defaultParkedRequestRetryFactor, defaultParkedRequestRetryJitter),
|
||||
}
|
||||
for _, opt := range opts {
|
||||
opt(r)
|
||||
@@ -144,7 +145,7 @@ func (r *ActorResumer) ResumeActor(ctx context.Context, actorRef resources.Actor
|
||||
bgCtx, bgCancel := context.WithTimeout(context.Background(), r.budget)
|
||||
defer bgCancel()
|
||||
|
||||
backoff := resumeBackoff()
|
||||
backoff := r.backoff
|
||||
|
||||
var resumeResp *ateapipb.ResumeActorResponse
|
||||
var lastRetryErr error
|
||||
|
||||
@@ -197,7 +197,7 @@ func TestActorResumer_Parking(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
resumer := NewActorResumer(mock, withParking(true, 5*time.Second))
|
||||
resumer := NewActorResumer(mock, withParking(parkingConfig{maxParked: 1, budget: 5 * time.Second}))
|
||||
actor, err := resumer.ResumeActor(context.Background(), testAtespace, testActorName)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
@@ -224,9 +224,9 @@ func TestActorResumer_Parking(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
// Budget large enough for a few ~500ms-spaced retries before it elapses;
|
||||
// Budget large enough for a few ~100ms-spaced retries before it elapses;
|
||||
// the pool never frees up.
|
||||
resumer := NewActorResumer(mock, withParking(true, 1500*time.Millisecond))
|
||||
resumer := NewActorResumer(mock, withParking(parkingConfig{maxParked: 1, budget: 1500 * time.Millisecond}))
|
||||
_, err := resumer.ResumeActor(context.Background(), testAtespace, testActorName)
|
||||
// The client must see the meaningful capacity error, not a generic
|
||||
// timeout: status.Code must unwrap through the budget-exhaustion marker.
|
||||
@@ -256,7 +256,7 @@ func TestActorResumer_Parking(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
// Default constructor => parking disabled => legacy fail-fast.
|
||||
// Default constructor => parking disabled => fail-fast.
|
||||
resumer := NewActorResumer(mock)
|
||||
_, err := resumer.ResumeActor(context.Background(), testAtespace, testActorName)
|
||||
if got := status.Code(err); got != codes.FailedPrecondition {
|
||||
@@ -275,7 +275,7 @@ func TestResumeBackoffHasNoCap(t *testing.T) {
|
||||
// Steps the moment the delay reaches Cap, which would end parking retries far
|
||||
// short of the budget (a 2s Cap stops the loop in ~7 steps / ~5s). The budget
|
||||
// context — not the step count or a cap — must bound how long a request parks.
|
||||
b := resumeBackoff()
|
||||
b := resumeBackoff(defaultParkedRequestRetryInterval, defaultParkedRequestRetryFactor, defaultParkedRequestRetryJitter)
|
||||
if b.Cap != 0 {
|
||||
t.Errorf("resume backoff must not set Cap (it would stop retries at the cap); got %v", b.Cap)
|
||||
}
|
||||
|
||||
@@ -189,10 +189,26 @@ func (s *RouterServer) Run(ctx context.Context) error {
|
||||
if err := xdsSrv.SetOtlpCollector(s.cfg.OtlpCollectorAddress); err != nil {
|
||||
return fmt.Errorf("configure OTLP collector: %w", err)
|
||||
}
|
||||
if s.cfg.ParkingMaxParked > 0 && s.cfg.ParkingMaxWait > 0 {
|
||||
|
||||
// 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()
|
||||
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.
|
||||
xdsSrv.SetExtProcMessageTimeout(s.cfg.ParkingMaxWait + 5*time.Second)
|
||||
xdsSrv.SetExtProcMessageTimeout(parkCfg.budget + 5*time.Second)
|
||||
}
|
||||
|
||||
xdsSrv.SetTlsConfig(s.cfg.HttpsPort, s.cfg.EnvoyCertPath)
|
||||
@@ -205,10 +221,6 @@ func (s *RouterServer) Run(ctx context.Context) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to create parking metrics: %w", err)
|
||||
}
|
||||
parkCfg := parkingConfig{
|
||||
maxWait: s.cfg.ParkingMaxWait,
|
||||
maxParked: s.cfg.ParkingMaxParked,
|
||||
}
|
||||
s.extprocSrv = NewExtProcServer(s.cfg.ExtprocPort, s.apiClient, routeDuration, parkCfg, parkMetrics)
|
||||
}
|
||||
ctrl := NewController(s.k8sClient, s.clientset, s.cfg, xdsSrv, s.extprocSrv)
|
||||
|
||||
@@ -163,8 +163,8 @@ func TestStatuszEndpoint(t *testing.T) {
|
||||
if !dashboard.Parking.Enabled {
|
||||
t.Errorf("expected parking reported as enabled in status JSON")
|
||||
}
|
||||
if dashboard.Parking.MaxParked != defaultParkingMaxParked {
|
||||
t.Errorf("expected parking max_parked %d, got %d", defaultParkingMaxParked, dashboard.Parking.MaxParked)
|
||||
if dashboard.Parking.MaxParked != defaultParkedRequestMax {
|
||||
t.Errorf("expected parking max_parked %d, got %d", defaultParkedRequestMax, dashboard.Parking.MaxParked)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -60,8 +60,8 @@ kubectl port-forward -n ate-system svc/atenet-router 8000:80
|
||||
|
||||
## How to Use
|
||||
|
||||
Parking is **on by default** (`--parking-max-wait=30s`,
|
||||
`--parking-max-parked=2048`), so the cluster you just deployed already parks.
|
||||
Parking is **on by default** (`--parked-request-budget=5s`,
|
||||
`--parked-request-max=2048`), so the cluster you just deployed already parks.
|
||||
|
||||
### A. Watch a 503 become a served request
|
||||
|
||||
@@ -84,7 +84,7 @@ curl -s -w '\n-> HTTP %{http_code} in %{time_total}s\n' \
|
||||
```
|
||||
|
||||
While that is hanging, in a **second terminal** free a worker by suspending p1
|
||||
(within the 30s park budget):
|
||||
(within the 5s park budget):
|
||||
|
||||
```bash
|
||||
kubectl ate suspend actor p1 --atespace parking
|
||||
@@ -120,7 +120,7 @@ non-zero `active`):
|
||||
```bash
|
||||
kubectl -n ate-system port-forward deployment/atenet-router 4040:4040
|
||||
curl -s 'http://localhost:4040/statusz?format=json' | jq .parking
|
||||
# { "enabled": true, "active": 3, "max_parked": 2048, "max_wait": "30s" }
|
||||
# { "enabled": true, "active": 3, "max_parked": 2048, "max_wait": "5s" }
|
||||
```
|
||||
|
||||
The parking metrics are also exported on the router's metrics endpoint
|
||||
@@ -135,7 +135,7 @@ container's args:
|
||||
|
||||
```bash
|
||||
kubectl -n ate-system patch deployment atenet-router --type=json \
|
||||
-p='[{"op":"add","path":"/spec/template/spec/containers/0/args/-","value":"--parking-max-parked=0"}]'
|
||||
-p='[{"op":"add","path":"/spec/template/spec/containers/0/args/-","value":"--parked-request-max=0"}]'
|
||||
kubectl -n ate-system rollout status deployment/atenet-router
|
||||
```
|
||||
|
||||
@@ -154,8 +154,8 @@ kubectl -n ate-system rollout undo deployment/atenet-router
|
||||
```
|
||||
|
||||
> [!TIP]
|
||||
> You can tune parking instead of disabling it: add `--parking-max-wait=10s` or
|
||||
> `--parking-max-parked=512` to the same args list.
|
||||
> You can tune parking instead of disabling it: add `--parked-request-budget=10s` or
|
||||
> `--parked-request-max=512` to the same args list.
|
||||
|
||||
## How to Uninstall
|
||||
|
||||
|
||||
@@ -25,7 +25,7 @@
|
||||
#
|
||||
# Compare two runs:
|
||||
# * parking ON (default router config) -> ~all 200, some elevated latency
|
||||
# * parking OFF (--parking-max-parked=0) -> a burst of 503s
|
||||
# * parking OFF (--parked-request-max=0) -> a burst of 503s
|
||||
# See README.md for how to flip the router flag.
|
||||
#
|
||||
# Usage:
|
||||
@@ -143,7 +143,7 @@ printf " 200 latency : avg %ss, slowest %ss <- parked requests sit here\n
|
||||
echo
|
||||
if [[ "${busy}" -eq 0 ]]; then
|
||||
echo " => 0 failures under saturation: parking absorbed the contention."
|
||||
echo " Re-run with the router started --parking-max-parked=0 to see 503s."
|
||||
echo " Re-run with the router started --parked-request-max=0 to see 503s."
|
||||
else
|
||||
echo " => ${busy} requests were shed with 503 (parking off, or lot full / budget exceeded)."
|
||||
fi
|
||||
|
||||
+30
-16
@@ -33,16 +33,16 @@ user-visible error.
|
||||
With parking enabled (the default), the router treats `FailedPrecondition` from
|
||||
`ResumeActor` as a **retryable** condition (alongside the existing `Aborted`
|
||||
concurrent-resume conflict). The request is *parked*: the resumer keeps retrying
|
||||
with capped exponential backoff until either
|
||||
with exponential backoff until either
|
||||
|
||||
- the resume succeeds (the actor is `RUNNING` and has a worker IP) — the request
|
||||
is then routed normally; or
|
||||
- the **park budget** (`--parking-max-wait`, default `30s`) elapses — the
|
||||
- the **park budget** (`--parked-request-budget`, default `5s`) elapses — the
|
||||
underlying capacity error is returned, surfacing as `503 "actor <id>
|
||||
unavailable: no free workers available"`.
|
||||
|
||||
To bound resource use and provide backpressure, the router admits requests to a
|
||||
**parking lot** of fixed capacity (`--parking-max-parked`, default `2048`). Each
|
||||
**parking lot** of fixed capacity (`--parked-request-max`, default `2048`). Each
|
||||
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.
|
||||
@@ -56,7 +56,7 @@ control-plane RPC.
|
||||
|
||||
Only transient capacity (`FailedPrecondition`) and concurrency (`Aborted`)
|
||||
conditions are parked. Errors that will not resolve by waiting are returned
|
||||
immediately, preserving prior semantics:
|
||||
immediately (fail fast):
|
||||
|
||||
| Resume result | Behavior |
|
||||
| ------------------------------------- | --------------------------------- |
|
||||
@@ -68,17 +68,22 @@ immediately, preserving prior semantics:
|
||||
| `DeadlineExceeded` | Fail fast → `504` |
|
||||
| `PermissionDenied` / `Unauthenticated`| Fail fast → `403` / `401` |
|
||||
|
||||
When parking is **disabled** (`--parking-max-parked=0`), the router preserves
|
||||
its legacy fail-fast behavior: `FailedPrecondition` fails fast, there is no
|
||||
admission cap, and only `Aborted` (concurrent-resume) conflicts are retried,
|
||||
within the historical `15s` budget.
|
||||
When parking is **disabled** (`--parked-request-max=0`), the router fails fast:
|
||||
`FailedPrecondition` is returned immediately, there is no admission cap, and
|
||||
only `Aborted` (concurrent-resume) conflicts are retried, within a `15s` budget.
|
||||
|
||||
## Configuration
|
||||
|
||||
| Flag | Default | Meaning |
|
||||
| ---------------------- | ------- | ------------------------------------------------------------------ |
|
||||
| `--parking-max-wait` | `30s` | Max time a single request may stay parked awaiting resume. |
|
||||
| `--parking-max-parked` | `2048` | Max concurrent parked/in-flight resume requests; excess shed (503). `0` disables parking. |
|
||||
| Flag | Default | Meaning |
|
||||
| -------------------------------- | ------- | ------------------------------------------------------------------ |
|
||||
| `--parked-request-budget` | `5s` | Max time a single request may stay parked awaiting resume. |
|
||||
| `--parked-request-max` | `2048` | 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. |
|
||||
|
||||
The retry backoff deliberately has no cap and no attempt limit: the budget alone
|
||||
bounds the wait.
|
||||
|
||||
## Observability
|
||||
|
||||
@@ -86,10 +91,19 @@ within the historical `15s` budget.
|
||||
|
||||
- `atenet.router.parking.active` — up/down counter: requests currently parked.
|
||||
- `atenet.router.parking.wait.duration` — histogram (seconds) of time spent
|
||||
parked, labeled `outcome` ∈ {`served`, `budget_exhausted`, `timeout`,
|
||||
`canceled`, `error`}. `budget_exhausted` means the full park budget elapsed
|
||||
while the pool stayed saturated — the signal that capacity, not a fault, is
|
||||
the bottleneck.
|
||||
parked. Recorded **exactly once per admitted request**, at the moment its
|
||||
resume attempt completes; never recorded for shed requests (those only
|
||||
increment `parking.rejected`) nor when parking is disabled. The `outcome`
|
||||
label says how the park ended:
|
||||
|
||||
| `outcome` | When it is set |
|
||||
| ------------------ | --------------------------------------------------------------------------- |
|
||||
| `served` | The resume succeeded and the request was routed to its worker. |
|
||||
| `budget_exhausted` | The park budget elapsed while the resume was still blocked on a retryable condition (pool saturated, or a concurrent operation holding the actor) — the signal that capacity, not a fault, is the bottleneck. |
|
||||
| `canceled` | The client disconnected while parked (request context canceled). |
|
||||
| `timeout` | The request's own deadline expired while parked (distinct from the park budget). |
|
||||
| `error` | The resume failed with a non-retryable error (`NotFound`, `Unavailable`, ...). |
|
||||
|
||||
- `atenet.router.parking.rejected` — counter: requests shed because the lot was
|
||||
full.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user