diff --git a/go.mod b/go.mod index dc334abe799..fd4b3319cea 100644 --- a/go.mod +++ b/go.mod @@ -7,7 +7,7 @@ require ( github.com/alecthomas/units v0.0.0-20240927000941-0f3dac36c52b github.com/dustin/go-humanize v1.0.1 github.com/edsrzf/mmap-go v1.2.1-0.20241212181136-fad1cd13edbd - github.com/failsafe-go/failsafe-go v0.9.6 + github.com/failsafe-go/failsafe-go v0.9.8 github.com/go-kit/log v0.2.1 github.com/go-openapi/strfmt v0.27.2 github.com/go-openapi/swag v0.28.0 // indirect diff --git a/go.sum b/go.sum index 729cb70e186..5211d59c4f8 100644 --- a/go.sum +++ b/go.sum @@ -322,8 +322,8 @@ github.com/envoyproxy/protoc-gen-validate v1.3.3 h1:MVQghNeW+LZcmXe7SY1V36Z+WFMD github.com/envoyproxy/protoc-gen-validate v1.3.3/go.mod h1:TsndJ/ngyIdQRhMcVVGDDHINPLWB7C82oDArY51KfB0= github.com/facette/natsort v0.0.0-20181210072756-2cd4dd1e2dcb h1:IT4JYU7k4ikYg1SCxNI1/Tieq/NFvh6dzLdgi7eu0tM= github.com/facette/natsort v0.0.0-20181210072756-2cd4dd1e2dcb/go.mod h1:bH6Xx7IW64qjjJq8M2u4dxNaBiDfKK+z/3eGDpXEQhc= -github.com/failsafe-go/failsafe-go v0.9.6 h1:vPSH2cry0Ee5cnR9wc9qshCDO6jdrMA9elBJNwyo4Uk= -github.com/failsafe-go/failsafe-go v0.9.6/go.mod h1:IeRpglkcwzKagjDMh90ZhN2l4Ovt3+jemQBUbThag54= +github.com/failsafe-go/failsafe-go v0.9.8 h1:NzahTEc+vg6FqCJV0Sy9mgwcaCemB81c8iSHWTF/urU= +github.com/failsafe-go/failsafe-go v0.9.8/go.mod h1:wMKUFRbxxjWwvwhJttSiim9BaQssrlZV4lk0IwPkwUc= github.com/fatih/color v1.7.0/go.mod h1:Zm6kSWBoL9eyXnKyktHP6abPY2pDugNf5KwzbycvMj4= github.com/fatih/color v1.9.0/go.mod h1:eQcE1qtQxscV5RaZvpXrrb8Drkc3/DdQ+uUYCNjL+zU= github.com/fatih/color v1.10.0/go.mod h1:ELkj/draVOlAH/xkhN6mQ50Qd0MPOk5AAr3maGEBuJM= diff --git a/vendor/github.com/failsafe-go/failsafe-go/CHANGELOG.md b/vendor/github.com/failsafe-go/failsafe-go/CHANGELOG.md index 51f75336015..acaa18b22ad 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/CHANGELOG.md +++ b/vendor/github.com/failsafe-go/failsafe-go/CHANGELOG.md @@ -1,4 +1,25 @@ -## Upcoming Release +## 0.9.8 + +### Improvements + +- Added #142 - Random and jitter delay support added for circuit breakers. +- Added #143 - `OnAcquired` and `OnReleased` listeners to `Bulkhead` for tracking executions that hold a permit. + +### Changes + +- Upgraded Go dependency to 1.22. + +## 0.9.7 + +### Bug Fixes + +- Fixed #136 - `failsafegrpc.NewUnaryClientInterceptor` should preserve per request contexts. +- Fixed #139 - Budgets should correctly threshold against min concurrency. +- Fixed #141 - An exceeded budget should not cause hedges to fail. + +### Improvements + +- Added dynamic delay support for hedge policies. ## 0.9.6 diff --git a/vendor/github.com/failsafe-go/failsafe-go/README.md b/vendor/github.com/failsafe-go/failsafe-go/README.md index 4f40e0af99a..cb194f87002 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/README.md +++ b/vendor/github.com/failsafe-go/failsafe-go/README.md @@ -1,7 +1,6 @@ # Failsafe-go [![Build Status](https://img.shields.io/github/actions/workflow/status/failsafe-go/failsafe-go/test.yml)](https://github.com/failsafe-go/failsafe-go/actions/workflows/test.yml) -[![Go Report Card](https://goreportcard.com/badge/github.com/failsafe-go/failsafe-go)](https://goreportcard.com/report/github.com/failsafe-go/failsafe-go) [![codecov](https://codecov.io/gh/failsafe-go/failsafe-go/graph/badge.svg?token=UC2BU7NTJ7)](https://codecov.io/gh/failsafe-go/failsafe-go) [![License](http://img.shields.io/:license-mit-brightgreen.svg)](https://opensource.org/licenses/MIT) [![Godoc](https://pkg.go.dev/badge/github.com/failsafe-go/failsafe-go)](https://pkg.go.dev/github.com/failsafe-go/failsafe-go) diff --git a/vendor/github.com/failsafe-go/failsafe-go/adaptivelimiter/adaptivelimiter.go b/vendor/github.com/failsafe-go/failsafe-go/adaptivelimiter/adaptivelimiter.go index fb1f923e131..94308faf4fa 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/adaptivelimiter/adaptivelimiter.go +++ b/vendor/github.com/failsafe-go/failsafe-go/adaptivelimiter/adaptivelimiter.go @@ -171,7 +171,7 @@ type Builder[R any] interface { // WithRecentQuantile configures the recentQuantile of recent execution times to consider when adjusting the concurrency limit. // // Defaults to 0.9 which uses p90 samples. - // Panics if recentQuantile is negative. + // Panics if recentQuantile is not > 0 and < 1. WithRecentQuantile(quantile float64) Builder[R] // WithBaselineWindow configures how the baseline execution times are maintained and updated. The baseline represents @@ -337,7 +337,7 @@ func (c *config[R]) WithRecentWindow(minDuration time.Duration, maxDuration time } func (c *config[R]) WithRecentQuantile(quantile float64) Builder[R] { - util.Assert(quantile >= 0, "recentQuantile must be >= 0") + util.Assert(quantile > 0 && quantile < 1, "recentQuantile must be between 0 and 1 exclusive") c.recentQuantile = quantile return c } @@ -387,9 +387,9 @@ func (c *config[R]) Build() AdaptiveLimiter[R] { semaphore: util.NewDynamicSemaphore(int(c.initialLimit)), limit: float64(c.initialLimit), recentRTT: tdigestSample{TDigest: tdigest.NewWithCompression(100)}, - medianFilter: util.NewMedianFilter(smoothedSamples), - smoothedRecentRTT: util.NewEwma(smoothedSamples, warmupSamples), - baselineRTT: util.NewEwma(c.baselineWindowAge, warmupSamples), + medianFilter: util.NewMovingMedian(smoothedSamples), + smoothedRecentRTT: util.NewMovingAverage(smoothedSamples, warmupSamples), + baselineRTT: util.NewMovingAverage(c.baselineWindowAge, warmupSamples), nextUpdateTime: time.Now(), rttCorrelation: util.NewCorrelationWindow(c.correlationWindowSize, warmupSamples), throughputCorrelation: util.NewCorrelationWindow(c.correlationWindowSize, warmupSamples), @@ -465,9 +465,9 @@ type adaptiveLimiter[R any] struct { maxInflightWindow util.MaxWindow // Tracks the max inflight over a stabilization window recentRTT tdigestSample // Recent execution times lastMaxInflight int // The max inflight requests for the last sampling period - medianFilter util.MedianFilter - smoothedRecentRTT util.Ewma - baselineRTT util.Ewma // Tracks baseline execution time + medianFilter util.MovingMedian + smoothedRecentRTT util.MovingAverage + baselineRTT util.MovingAverage // Tracks baseline execution time nextUpdateTime time.Time // Tracks when the limit can next be updated throughputCorrelation util.CorrelationWindow // Tracks the correlation between concurrency and throughput rttCorrelation util.CorrelationWindow // Tracks the correlation between concurrency and round trip times (RTT) diff --git a/vendor/github.com/failsafe-go/failsafe-go/adaptivelimiter/queueinglimiter.go b/vendor/github.com/failsafe-go/failsafe-go/adaptivelimiter/queueinglimiter.go index 38f39d44e31..180fd806ff1 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/adaptivelimiter/queueinglimiter.go +++ b/vendor/github.com/failsafe-go/failsafe-go/adaptivelimiter/queueinglimiter.go @@ -2,7 +2,7 @@ package adaptivelimiter import ( "context" - "math/rand" + "math/rand/v2" "time" "github.com/failsafe-go/failsafe-go/policy" diff --git a/vendor/github.com/failsafe-go/failsafe-go/budget/budget.go b/vendor/github.com/failsafe-go/failsafe-go/budget/budget.go index d42fa6fee7b..83efac56c6a 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/budget/budget.go +++ b/vendor/github.com/failsafe-go/failsafe-go/budget/budget.go @@ -14,11 +14,8 @@ var ErrExceeded = errors.New("budget exceeded") // // This type is concurrency safe. type Budget interface { - // RetryRate returns the current rate of retries relative to total executions, from 0 to 1. - RetryRate() float64 - - // HedgeRate returns the current rate of hedges relative to total executions, from 0 to 1. - HedgeRate() float64 + // Rate returns the current combined rate of retries and hedges relative to total inflight executions, from 0 to 1. + Rate() float64 } // Builder builds Budget instances. @@ -28,8 +25,9 @@ type Builder interface { // WithMaxRate configures the max rate of inflight executions that can be retries and/or hedges. WithMaxRate(maxRate float64) Builder - // WithMinConcurrency configures the min number of budgeted retries and/or hedges that can be executed, regardless of - // the total number of inflight executions. + // WithMinConcurrency configures a minimum execution budget. At least minConcurrency retries and/or hedges may execute + // concurrently, even when the budget computed from WithMaxRate and the current inflight execution count falls below + // minConcurrency. WithMinConcurrency(minConcurrency uint) Builder // OnBudgetExceeded registers the listener to be called when the budget is exceeded. @@ -108,46 +106,39 @@ type budget struct { config executions atomic.Int32 - retries atomic.Int32 - hedges atomic.Int32 + inflight atomic.Int32 // current count of inflight retries and hedges } -func (b *budget) TryAcquireRetryPermit() bool { - if b.RetryRate() > b.maxRate { - return false - } - - b.retries.Add(1) +func (b *budget) RecordExecution() { b.executions.Add(1) - return true } -func (b *budget) TryAcquireHedgePermit() bool { - if b.HedgeRate() > b.maxRate { +func (b *budget) ReleaseExecution() { + b.executions.Add(-1) +} + +func (b *budget) TryAcquirePermit() bool { + budgeted := b.inflight.Load() + if budgeted >= int32(b.minConcurrency) && b.Rate() > b.maxRate { return false } - b.hedges.Add(1) - b.executions.Add(1) + b.inflight.Add(1) + b.RecordExecution() return true } -func (b *budget) ReleaseRetryPermit() { - b.retries.Add(-1) - b.executions.Add(-1) -} - -func (b *budget) ReleaseHedgePermit() { - b.hedges.Add(-1) - b.executions.Add(-1) -} - -func (b *budget) RetryRate() float64 { - return float64(b.retries.Load()) / float64(b.executions.Load()) +func (b *budget) ReleasePermit() { + b.inflight.Add(-1) + b.ReleaseExecution() } -func (b *budget) HedgeRate() float64 { - return float64(b.hedges.Load()) / float64(b.executions.Load()) +func (b *budget) Rate() float64 { + executions := b.executions.Load() + if executions == 0 { + return 0 + } + return float64(b.inflight.Load()) / float64(executions) } func (b *budget) OnBudgetExceeded(executionType ExecutionType, info failsafe.ExecutionInfo) { diff --git a/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitbreaker.go b/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitbreaker.go index e9ff465aec4..e983b80faa9 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitbreaker.go +++ b/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitbreaker.go @@ -3,6 +3,7 @@ package circuitbreaker import ( "context" "errors" + "math/rand/v2" "sync" "time" @@ -46,9 +47,10 @@ CircuitBreaker is a policy that temporarily blocks execution when a configured n breakers have three states: closed, open, and half-open. When a circuit breaker is in the ClosedState (default), executions are allowed. If a configurable number of failures occur, optionally over some time period, the circuit breaker transitions to OpenState. In the OpenState a circuit breaker will fail executions with ErrOpen. -After a configurable delay, the circuit breaker will transition to HalfOpenState. In the HalfOpenState a configurable -number of trial executions will be allowed, after which the circuit breaker will transition to either ClosedState or -OpenState depending on how many were successful. +After a configurable delay, the circuit breaker will transition to HalfOpenState. The delay can be fixed, computed via a +DelayFunc, or randomized, and can be jittered to desynchronize the half-open probes of breakers that opened together. In +the HalfOpenState a configurable number of trial executions will be allowed, after which the circuit breaker will +transition to either ClosedState or OpenState depending on how many were successful. A circuit breaker can be count based or time based: @@ -299,11 +301,7 @@ func (cb *circuitBreaker[R]) transitionTo(newState State, exec failsafe.Executio case ClosedState: cb.state = newClosedState(cb) case OpenState: - delay := cb.ComputeDelay(exec) - if delay == -1 { - delay = cb.Delay - } - cb.state = newOpenState(cb, cb.state, delay) + cb.state = newOpenState(cb, cb.state, cb.getDelay(exec)) case HalfOpenState: cb.state = newHalfOpenState(cb) } @@ -361,6 +359,24 @@ func (cb *circuitBreaker[R]) tryAcquirePermit() bool { return cb.state.tryAcquirePermit() } +// getDelay returns the delay to wait in the OpenState before transitioning to HalfOpenState. It resolves the base +// delay from the DelayFunc, else the fixed delay, else a configured random delay range, and then applies any configured +// jitter. Jitter desynchronizes the half-open probes of many breakers that opened at the same time. +func (cb *circuitBreaker[R]) getDelay(exec failsafe.Execution[R]) time.Duration { + var delay time.Duration + if computed := cb.ComputeDelay(exec); computed != -1 { + delay = computed + } else if cb.Delay != 0 { + delay = cb.Delay + } else if cb.delayMin != 0 && cb.delayMax != 0 { + delay = time.Duration(util.RandomDelayInRange(cb.delayMin.Nanoseconds(), cb.delayMax.Nanoseconds(), rand.Float64())) + } + if delay != 0 { + delay = util.ApplyJitter(delay, cb.jitter, cb.jitterFactor) + } + return max(0, delay) +} + // Opens the circuit breaker and considers the execution when computing the delay before the circuit breaker // will transition to half open. // diff --git a/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitbreakerbuilder.go b/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitbreakerbuilder.go index 7547f5656ce..1111aade40b 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitbreakerbuilder.go +++ b/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitbreakerbuilder.go @@ -73,6 +73,23 @@ type Builder[R any] interface { // WithDelayFunc configures a function that provides the delay to wait in OpenState before transitioning to HalfOpenState. WithDelayFunc(delayFunc failsafe.DelayFunc[R]) Builder[R] + // WithRandomDelay configures a random delay to wait in OpenState before transitioning to HalfOpenState, chosen + // randomly between the delayMin and delayMax. Replaces any previously configured fixed delay. + WithRandomDelay(delayMin time.Duration, delayMax time.Duration) Builder[R] + + // WithJitter configures the jitter to randomly vary the OpenState delay by. For each delay, a random portion of the + // jitter will be added or subtracted to the delay. For example: a jitter of 100 milliseconds will randomly add + // between -100 and 100 milliseconds to each delay. Jittering the delay prevents multiple circuit breakers that open + // at the same time from transitioning to HalfOpenState in lockstep and re-probing an unhealthy dependency together. + // Replaces any previously configured jitter factor. + WithJitter(jitter time.Duration) Builder[R] + + // WithJitterFactor configures the jitterFactor to randomly vary the OpenState delay by. For each delay, a random + // portion of the delay multiplied by the jitterFactor will be added or subtracted to the delay. For example: a delay + // of 100 milliseconds and a jitterFactor of .25 will result in a random delay between 75 and 125 milliseconds. + // Replaces any previously configured jitter duration. + WithJitterFactor(jitterFactor float64) Builder[R] + // WithSuccessThreshold configures count based success thresholding by setting the number of consecutive successful // executions that must occur when in a HalfOpenState in order to close the circuit, else the circuit is re-opened when a // failure occurs. @@ -90,7 +107,14 @@ type Builder[R any] interface { type config[R any] struct { policy.BaseFailurePolicy[R] policy.BaseDelayablePolicy[R] - clock util.Clock + clock util.Clock + + // Delay config that augments BaseDelayablePolicy.Delay + delayMin time.Duration + delayMax time.Duration + jitter time.Duration + jitterFactor float64 + stateChangedListener func(StateChangedEvent) openListener func(StateChangedEvent) halfOpenListener func(StateChangedEvent) @@ -206,6 +230,25 @@ func (c *config[R]) WithDelayFunc(delayFunc failsafe.DelayFunc[R]) Builder[R] { return c } +func (c *config[R]) WithRandomDelay(delayMin time.Duration, delayMax time.Duration) Builder[R] { + c.delayMin = delayMin + c.delayMax = delayMax + + // Clear fixed delay + c.Delay = 0 + return c +} + +func (c *config[R]) WithJitter(jitter time.Duration) Builder[R] { + c.jitter = jitter + return c +} + +func (c *config[R]) WithJitterFactor(jitterFactor float64) Builder[R] { + c.jitterFactor = jitterFactor + return c +} + func (c *config[R]) OnStateChanged(listener func(event StateChangedEvent)) Builder[R] { c.stateChangedListener = listener return c diff --git a/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitstates.go b/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitstates.go index c70bef9be77..40e865da575 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitstates.go +++ b/vendor/github.com/failsafe-go/failsafe-go/circuitbreaker/circuitstates.go @@ -73,7 +73,7 @@ func (s *closedState[R]) checkThresholdAndReleasePermit(exec failsafe.Execution[ type openState[R any] struct { breaker *circuitBreaker[R] util.ExecutionStats - startTime int64 + startTime time.Time delay time.Duration } @@ -81,7 +81,7 @@ func newOpenState[R any](breaker *circuitBreaker[R], previousState circuitState[ return &openState[R]{ breaker: breaker, ExecutionStats: previousState, - startTime: breaker.clock.Now().UnixNano(), + startTime: breaker.clock.Now(), delay: delay, } } @@ -91,12 +91,11 @@ func (s *openState[R]) state() State { } func (s *openState[R]) remainingDelay() time.Duration { - elapsedTime := s.breaker.clock.Now().UnixNano() - s.startTime - return max(0, s.delay-time.Duration(elapsedTime)) + return max(0, s.delay-s.breaker.clock.Now().Sub(s.startTime)) } func (s *openState[R]) tryAcquirePermit() bool { - if s.breaker.clock.Now().UnixNano()-s.startTime >= s.delay.Nanoseconds() { + if s.breaker.clock.Now().Sub(s.startTime) >= s.delay { s.breaker.halfOpen() return s.breaker.tryAcquirePermit() } diff --git a/vendor/github.com/failsafe-go/failsafe-go/internal/budget.go b/vendor/github.com/failsafe-go/failsafe-go/internal/budget.go index 55b5a128807..9c3b7b0fcd0 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/internal/budget.go +++ b/vendor/github.com/failsafe-go/failsafe-go/internal/budget.go @@ -8,17 +8,17 @@ import ( type Budget interface { budget.Budget - // TryAcquireRetryPermit acquires a permit to retry an execution, else returns false if the budget is exceeded. - TryAcquireRetryPermit() bool + // RecordExecution records that a primary execution has started. + RecordExecution() - // TryAcquireHedgePermit acquires a permit to perform a hedged execution, else returns false if the budget is exceeded. - TryAcquireHedgePermit() bool + // ReleaseExecution records that a primary execution has ended. + ReleaseExecution() - // ReleaseRetryPermit releases a previously acquired retry permit back to the budget. - ReleaseRetryPermit() + // TryAcquirePermit acquires a permit for a retry or hedge execution, returning false if the budget is exceeded. + TryAcquirePermit() bool - // ReleaseHedgePermit releases a previously acquired hedge permit back to the budget. - ReleaseHedgePermit() + // ReleasePermit releases a previously acquired retry or hedge permit. + ReleasePermit() // OnBudgetExceeded calls the OnBudgetExceeded event listener, if one is configured. OnBudgetExceeded(executionType budget.ExecutionType, info failsafe.ExecutionInfo) diff --git a/vendor/github.com/failsafe-go/failsafe-go/internal/util/ewma.go b/vendor/github.com/failsafe-go/failsafe-go/internal/util/average.go similarity index 56% rename from vendor/github.com/failsafe-go/failsafe-go/internal/util/ewma.go rename to vendor/github.com/failsafe-go/failsafe-go/internal/util/average.go index 5e93c7ad70c..fe7b81654a8 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/internal/util/ewma.go +++ b/vendor/github.com/failsafe-go/failsafe-go/internal/util/average.go @@ -1,9 +1,9 @@ package util -// Ewma is an exponentially weighted moving average. +// MovingAverage is an exponentially weighted moving average. // // This type is not concurrency safe. -type Ewma struct { +type MovingAverage struct { warmupSamples uint8 smoothingFactor float64 @@ -13,21 +13,21 @@ type Ewma struct { sum float64 } -// NewEwma creates a new Ewma for the given age and warmupSamples. The age controls how far back in time the -// Ewma effectively "remembers" - smaller ages adapt faster to recent changes, while larger ages provide -// more stability by retaining influence from older samples. The warmupSamples parameter controls how many +// NewMovingAverage creates a new MovingAverage for the given age and warmupSamples. The age controls how far back in +// time the MovingAverage effectively "remembers" - smaller ages adapt faster to recent changes, while larger ages +// provide more stability by retaining influence from older samples. The warmupSamples parameter controls how many // samples must be recorded before exponential decay begins, during which a simple average is used instead. -func NewEwma(age uint, warmupSamples uint8) Ewma { - return Ewma{ +func NewMovingAverage(age uint, warmupSamples uint8) MovingAverage { + return MovingAverage{ warmupSamples: warmupSamples, smoothingFactor: 2 / (float64(age) + 1), } } -// Add adds a value to the series and updates the moving average. Add decays the Ewma value via: +// Add adds a value to the series and updates the moving average. Add decays the MovingAverage value via: // // (oldValue * (1 - smoothingFactor)) + (newValue * smoothingFactor) -func (e *Ewma) Add(newValue float64) float64 { +func (e *MovingAverage) Add(newValue float64) float64 { switch { case e.count < e.warmupSamples: e.count++ @@ -40,12 +40,12 @@ func (e *Ewma) Add(newValue float64) float64 { } // Value gets the current value of the moving average. -func (e *Ewma) Value() float64 { +func (e *MovingAverage) Value() float64 { return e.value } // Reset resets the value of the moving average and requires a new warmup if one was configured. -func (e *Ewma) Reset() { +func (e *MovingAverage) Reset() { e.count = 0 e.value = 0 e.sum = 0 diff --git a/vendor/github.com/failsafe-go/failsafe-go/internal/util/median.go b/vendor/github.com/failsafe-go/failsafe-go/internal/util/median.go index d4f066859bd..47f031791be 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/internal/util/median.go +++ b/vendor/github.com/failsafe-go/failsafe-go/internal/util/median.go @@ -4,52 +4,52 @@ import ( "slices" ) -// MedianFilter provides the median value over a rolling window. +// MovingMedian provides the median value over a moving window. // // This type is not concurrency safe. -type MedianFilter struct { +type MovingMedian struct { values []float64 sorted []float64 index int size int } -func NewMedianFilter(size int) MedianFilter { - return MedianFilter{ +func NewMovingMedian(size int) MovingMedian { + return MovingMedian{ values: make([]float64, size), sorted: make([]float64, size), } } -// Add adds a value to the filter, sorts the values, and returns the current median. -func (f *MedianFilter) Add(value float64) float64 { - f.values[f.index] = value - f.index = (f.index + 1) % len(f.values) +// Add adds a value to the window, sorts the values, and returns the current median. +func (m *MovingMedian) Add(value float64) float64 { + m.values[m.index] = value + m.index = (m.index + 1) % len(m.values) - if f.size < len(f.values)-1 { - f.size++ + if m.size < len(m.values)-1 { + m.size++ return value } - copy(f.sorted, f.values) - slices.Sort(f.sorted) - return f.Median() + copy(m.sorted, m.values) + slices.Sort(m.sorted) + return m.Median() } -// Median returns the current median, else 0 if the filter isn't full yet. -func (f *MedianFilter) Median() float64 { - if f.size < len(f.values)-1 { +// Median returns the current median, else 0 if the window isn't full yet. +func (m *MovingMedian) Median() float64 { + if m.size < len(m.values)-1 { return 0 } - return f.sorted[len(f.sorted)/2] + return m.sorted[len(m.sorted)/2] } -// Reset resets the filter to its initial value. -func (f *MedianFilter) Reset() { - for i := range f.values { - f.values[i] = 0 - f.sorted[i] = 0 +// Reset resets the window to its initial value. +func (m *MovingMedian) Reset() { + for i := range m.values { + m.values[i] = 0 + m.sorted[i] = 0 } - f.index = 0 - f.size = 0 + m.index = 0 + m.size = 0 } diff --git a/vendor/github.com/failsafe-go/failsafe-go/internal/util/quantile.go b/vendor/github.com/failsafe-go/failsafe-go/internal/util/quantile.go new file mode 100644 index 00000000000..b799448de2c --- /dev/null +++ b/vendor/github.com/failsafe-go/failsafe-go/internal/util/quantile.go @@ -0,0 +1,77 @@ +package util + +import "math" + +// MovingQuantile estimates a streaming quantile using the Windowless Moving Percentile algorithm (Martin Jambon). +// This provides O(1) time and space quantile estimation that adapts to distribution changes. +// +// This type is not concurrency safe. +type MovingQuantile struct { + quantile float64 + r float64 + alpha float64 + + // Mutable state + count int + value float64 + mean float64 + variance float64 +} + +// NewMovingQuantile creates a new MovingQuantile for the given quantile (0-1), step ratio r, and age. The age controls +// how far back in time the estimate effectively "remembers" - smaller ages adapt faster to recent changes, while larger +// ages provide more stability by retaining influence from older samples. +func NewMovingQuantile(quantile float64, r float64, age uint) MovingQuantile { + return MovingQuantile{ + quantile: quantile, + r: r, + alpha: 2 / (float64(age) + 1), + } +} + +// Add adds a sample and returns the updated quantile estimate. +func (q *MovingQuantile) Add(sample float64) float64 { + q.count++ + if q.count == 1 { + q.value = sample + q.mean = sample + return q.value + } + + // Update EMA mean and variance + oldMean := q.mean + q.mean = Smooth(q.mean, sample, q.alpha) + q.variance = Smooth(q.variance, (sample-oldMean)*(sample-q.mean), q.alpha) + + // Compute step size + delta := math.Sqrt(q.variance) * q.r + if delta == 0 { + return q.value + } + + // Adjust estimate + if sample < q.value { + q.value -= delta / q.quantile + } else if sample > q.value { + q.value += delta / (1 - q.quantile) + } + return q.value +} + +// Value returns the current quantile estimate. +func (q *MovingQuantile) Value() float64 { + return q.value +} + +// Count returns the number of samples added. +func (q *MovingQuantile) Count() int { + return q.count +} + +// Reset resets the quantile estimate. +func (q *MovingQuantile) Reset() { + q.value = 0 + q.mean = 0 + q.variance = 0 + q.count = 0 +} diff --git a/vendor/github.com/failsafe-go/failsafe-go/internal/util/util.go b/vendor/github.com/failsafe-go/failsafe-go/internal/util/util.go index 7124fea8c22..93a5614862e 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/internal/util/util.go +++ b/vendor/github.com/failsafe-go/failsafe-go/internal/util/util.go @@ -3,6 +3,7 @@ package util import ( "context" "math" + "math/rand/v2" "reflect" "time" ) @@ -13,7 +14,7 @@ type number interface { func noop(_ error) {} -var errorType = reflect.TypeOf((*error)(nil)).Elem() +var errorType = reflect.TypeFor[error]() // ErrorTypesMatch indicates whether the err or any unwrapped causes of the err are assignable to the target type. This is // similar to the test that errors.As performs, but does not actually assign a value and allows a non-pointer target. @@ -161,6 +162,15 @@ func RandomDelayFactor[T number](delay T, jitterFactor float64, random float64) return T(float64(delay) * randomFactor) } +func ApplyJitter[T number](delay T, jitter T, jitterFactor float64) T { + if jitter != 0 { + delay = RandomDelay(delay, jitter, rand.Float64()) + } else if jitterFactor != 0 { + delay = RandomDelayFactor(delay, jitterFactor, rand.Float64()) + } + return delay +} + // Smooth returns a value that is decreased by some portion of the oldValue, and increased by some portion of the // newValue, based on the factor. func Smooth(oldValue, newValue, factor float64) float64 { @@ -170,7 +180,7 @@ func Smooth(oldValue, newValue, factor float64) float64 { var log10Values []int func init() { - for i := 0; i < 100; i++ { + for range 100 { log10Values = append(log10Values, 1) } for i := 100; i < 1000; i++ { diff --git a/vendor/github.com/failsafe-go/failsafe-go/internal/util/windows.go b/vendor/github.com/failsafe-go/failsafe-go/internal/util/windows.go index 579275e52e1..d901d0450a4 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/internal/util/windows.go +++ b/vendor/github.com/failsafe-go/failsafe-go/internal/util/windows.go @@ -5,27 +5,27 @@ import ( "time" ) -// RollingSum maintains a sum over a rolling window. +// MovingSum maintains a sum over a moving window. // // This type is not concurrency safe. -type RollingSum struct { +type MovingSum struct { // For variation and covariance samples []float64 size int index int - // Rolling sum fields + // Moving sum fields sumY float64 // Y values are the samples sumSquares float64 } -func NewRollingSum(capacity uint) RollingSum { - return RollingSum{samples: make([]float64, capacity)} +func NewMovingSum(capacity uint) MovingSum { + return MovingSum{samples: make([]float64, capacity)} } // Add adds the value to the window if it's non-zero, updates the sums, and returns the old value along with whether the // window is full. -func (r *RollingSum) Add(value float64) (oldValue float64, full bool) { +func (r *MovingSum) Add(value float64) (oldValue float64, full bool) { if value != 0 { if r.size == len(r.samples) { full = true @@ -54,7 +54,7 @@ func (r *RollingSum) Add(value float64) (oldValue float64, full bool) { // CalculateCV calculates the coefficient of variation (relative variance), mean, and variance for the sum. Returns NaN // values if there are < 2 samples, the variance is < 0, or the mean is 0. -func (r *RollingSum) CalculateCV() (cv, mean, variance float64) { +func (r *MovingSum) CalculateCV() (cv, mean, variance float64) { if r.size < 2 { return math.NaN(), math.NaN(), math.NaN() } @@ -70,7 +70,7 @@ func (r *RollingSum) CalculateCV() (cv, mean, variance float64) { } // Reset resets the sum to its initial state. -func (r *RollingSum) Reset() { +func (r *MovingSum) Reset() { for i := range r.samples { r.samples[i] = 0 } @@ -87,16 +87,16 @@ type CorrelationWindow struct { warmupSamples uint8 // Mutable state - xSamples RollingSum - ySamples RollingSum + xSamples MovingSum + ySamples MovingSum corrSumXY float64 } func NewCorrelationWindow(capacity uint, warmupSamples uint8) CorrelationWindow { return CorrelationWindow{ warmupSamples: warmupSamples, - xSamples: NewRollingSum(capacity), - ySamples: NewRollingSum(capacity), + xSamples: NewMovingSum(capacity), + ySamples: NewMovingSum(capacity), } } @@ -235,7 +235,7 @@ func (w *BucketedWindow[T]) ExpireBuckets() *T { if newHead > w.HeadTime { bucketsToMove := min(w.BucketCount, newHead-w.HeadTime) - for i := int64(0); i < bucketsToMove; i++ { + for i := range bucketsToMove { bucket := &w.Buckets[(w.HeadTime+i+1)%w.BucketCount] w.RemoveFn(&w.Summary, bucket) w.ResetFn(bucket) diff --git a/vendor/github.com/failsafe-go/failsafe-go/priority/priority.go b/vendor/github.com/failsafe-go/failsafe-go/priority/priority.go index db0bee4b44c..5935c5b084d 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/priority/priority.go +++ b/vendor/github.com/failsafe-go/failsafe-go/priority/priority.go @@ -3,7 +3,7 @@ package priority import ( "context" "math" - "math/rand" + "math/rand/v2" "sync" ) @@ -23,7 +23,7 @@ const totalLevels = 500 // RandomLevel returns a random level for the Priority. func (p Priority) RandomLevel() int { r := priorityLevelRanges[p] - return rand.Intn(r.upper-r.lower+1) + r.lower + return rand.IntN(r.upper-r.lower+1) + r.lower } // AddTo returns the ctx with the priority added to it as a value with the PriorityKey. @@ -161,14 +161,11 @@ func (lt *windowedLevelTracker) GetLevel(quantile float64) int { if currentSize > 0 { // Determine how many recorded levels we need to find to match the quantile - targetLevels := int(math.Ceil(float64(currentSize) * quantile)) - if targetLevels < 1 { - targetLevels = 1 - } + targetLevels := max(1, int(math.Ceil(float64(currentSize)*quantile))) // Count the levels until we hit the desired quantile countedLevels := 0 - for level := 0; level < totalLevels; level++ { + for level := range totalLevels { countedLevels += lt.levelCounts[level] if countedLevels >= targetLevels { return level diff --git a/vendor/github.com/failsafe-go/failsafe-go/priority/usage.go b/vendor/github.com/failsafe-go/failsafe-go/priority/usage.go index 2632e6978c1..a7ea717aa33 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/priority/usage.go +++ b/vendor/github.com/failsafe-go/failsafe-go/priority/usage.go @@ -3,6 +3,7 @@ package priority import ( "container/list" "context" + "slices" "sort" "sync" "time" @@ -163,9 +164,7 @@ func (ut *usageTracker) Calibrate() { } } - sort.Slice(usages, func(i, j int) bool { - return usages[i] < usages[j] - }) + slices.Sort(usages) // Update percentiles for all active users for _, entry := range ut.users { diff --git a/vendor/github.com/failsafe-go/failsafe-go/retrypolicy/retryexecutor.go b/vendor/github.com/failsafe-go/failsafe-go/retrypolicy/retryexecutor.go index 38259ea63a8..392db06eb7f 100644 --- a/vendor/github.com/failsafe-go/failsafe-go/retrypolicy/retryexecutor.go +++ b/vendor/github.com/failsafe-go/failsafe-go/retrypolicy/retryexecutor.go @@ -1,7 +1,7 @@ package retrypolicy import ( - "math/rand" + "math/rand/v2" "time" "github.com/failsafe-go/failsafe-go" @@ -30,12 +30,17 @@ func (e *executor[R]) Apply(innerFn func(failsafe.Execution[R]) *common.PolicyRe execInternal := exec.(policy.ExecutionInternal[R]) isRetry := false + if e.budget != nil { + e.budget.RecordExecution() + defer e.budget.ReleaseExecution() + } + for { // Perform the execution result := innerFn(exec) if isRetry && e.budget != nil { - e.budget.ReleaseRetryPermit() + e.budget.ReleasePermit() } // Check for cancellation during execution @@ -77,7 +82,7 @@ func (e *executor[R]) Apply(innerFn func(failsafe.Execution[R]) *common.PolicyRe } // Check the retry budget, if any - if e.budget != nil && !e.budget.TryAcquireRetryPermit() { + if e.budget != nil && !e.budget.TryAcquirePermit() { e.budget.OnBudgetExceeded(budget.RetryExecution, exec) return internal.FailureResult[R](budget.ErrExceeded) } @@ -130,7 +135,7 @@ func (e *executor[R]) getDelay(exec failsafe.ExecutionAttempt[R]) time.Duration delay = e.getFixedOrRandomDelay(exec) } if delay != 0 { - delay = e.adjustForJitter(delay) + delay = util.ApplyJitter(delay, e.jitter, e.jitterFactor) } delay = e.adjustForMaxDuration(delay, exec.ElapsedTime()) return delay @@ -153,15 +158,6 @@ func (e *executor[R]) getFixedOrRandomDelay(exec failsafe.ExecutionAttempt[R]) t return 0 } -func (e *executor[R]) adjustForJitter(delay time.Duration) time.Duration { - if e.jitter != 0 { - delay = util.RandomDelay(delay, e.jitter, rand.Float64()) - } else if e.jitterFactor != 0 { - delay = util.RandomDelayFactor(delay, e.jitterFactor, rand.Float64()) - } - return delay -} - func (e *executor[R]) adjustForMaxDuration(delay time.Duration, elapsed time.Duration) time.Duration { if e.maxDuration != 0 { delay = min(delay, e.maxDuration-elapsed) diff --git a/vendor/modules.txt b/vendor/modules.txt index 12ebeca0382..15ff5b105c3 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -420,8 +420,8 @@ github.com/envoyproxy/protoc-gen-validate/validate # github.com/facette/natsort v0.0.0-20181210072756-2cd4dd1e2dcb ## explicit github.com/facette/natsort -# github.com/failsafe-go/failsafe-go v0.9.6 -## explicit; go 1.21 +# github.com/failsafe-go/failsafe-go v0.9.8 +## explicit; go 1.22 github.com/failsafe-go/failsafe-go github.com/failsafe-go/failsafe-go/adaptivelimiter github.com/failsafe-go/failsafe-go/budget