From a7a72e193443be94f3d4b97a15809ca074e96e46 Mon Sep 17 00:00:00 2001 From: Tarcisio Ferraz Date: Tue, 25 Aug 2026 11:30:50 -0300 Subject: [PATCH 1/5] implement last-known-good fallback for dynamic limit settings --- pkg/settings/limits/bound.go | 23 ++- pkg/settings/limits/default_fallback_test.go | 200 +++++++++++++++++++ pkg/settings/limits/gate.go | 23 ++- pkg/settings/limits/range.go | 23 ++- pkg/settings/limits/time.go | 40 ++-- pkg/settings/limits/updater.go | 38 +++- 6 files changed, 308 insertions(+), 39 deletions(-) create mode 100644 pkg/settings/limits/default_fallback_test.go diff --git a/pkg/settings/limits/bound.go b/pkg/settings/limits/bound.go index a9e2670ce8..39ddc66c4c 100644 --- a/pkg/settings/limits/bound.go +++ b/pkg/settings/limits/bound.go @@ -213,17 +213,18 @@ func (b *boundLimiter[N]) Limit(ctx context.Context) (N, error) { defer b.wg.Done() tenant, bound, err := b.get(ctx) - if err != nil { - return zero, err + if err != nil && tenant == "" && b.scope != settings.ScopeGlobal { + return zero, err // no tenant, so get() never read a value at all } - if tenant == "" && b.scope != settings.ScopeGlobal { + if tenant == "" && b.scope != settings.ScopeGlobal && err == nil { return zero, nil // fail open } - return bound, nil + return bound, err // bound is always usable; err is advisory } func (b *boundLimiter[N]) get(ctx context.Context) (tenant string, bound N, err error) { + u := b.updater if b.scope != settings.ScopeGlobal { tenant = b.scope.Value(ctx) if tenant == "" { @@ -236,10 +237,11 @@ func (b *boundLimiter[N]) get(ctx context.Context) (tenant string, bound N, err return } - u := newUpdater(b.lggr, b.getLimitFn, b.subFn) - actual, loaded := b.updaters.LoadOrStore(tenant, u) + newU := newUpdater(b.lggr, b.getLimitFn, b.subFn) + actual, loaded := b.updaters.LoadOrStore(tenant, newU) creCtx := contexts.WithCRE(ctx, b.scope.RoundCRE(contexts.CREValue(ctx))) if !loaded { + u = newU go u.updateLoop(creCtx) } else { u = actual.(*updater[N]) @@ -249,7 +251,14 @@ func (b *boundLimiter[N]) get(ctx context.Context) (tenant string, bound N, err bound, err = b.getLimitFn(ctx) if err != nil { - b.lggr.Errorw("Failed to get limit. Using default value", "default", bound, "err", err) + if last, ok := u.lastGood(); ok { + b.lggr.Errorw("Failed to get limit. Using last known value", "value", last, "err", err) + bound = last + } else { + b.lggr.Errorw("Failed to get limit. Using compiled default", "default", bound, "err", err) + } + } else { + u.setLast(bound) } b.recordBound(ctx, bound, withScope(ctx, b.scope)) return diff --git a/pkg/settings/limits/default_fallback_test.go b/pkg/settings/limits/default_fallback_test.go new file mode 100644 index 0000000000..efed4d923c --- /dev/null +++ b/pkg/settings/limits/default_fallback_test.go @@ -0,0 +1,200 @@ +package limits + +import ( + "context" + "errors" + "strconv" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/smartcontractkit/chainlink-common/pkg/config" + "github.com/smartcontractkit/chainlink-common/pkg/settings" +) + +// toggleGetter is a settings.Getter that can be switched between a fixed value and +// failing, to drive the last-known-good fallback path deterministically. +type toggleGetter struct { + mu sync.Mutex + value string + failing bool +} + +func (g *toggleGetter) succeedWith(value string) { + g.mu.Lock() + defer g.mu.Unlock() + g.value = value + g.failing = false +} + +func (g *toggleGetter) fail() { + g.mu.Lock() + defer g.mu.Unlock() + g.failing = true +} + +var errGetterUnavailable = errors.New("settings getter unavailable") + +func (g *toggleGetter) GetScoped(context.Context, settings.Scope, string) (string, error) { + g.mu.Lock() + defer g.mu.Unlock() + if g.failing { + return "", errGetterUnavailable + } + return g.value, nil +} + +func TestBoundLimiter_LastKnownGoodOnReadFailure(t *testing.T) { + t.Parallel() + + setting := settings.Size(1 * config.GByte) // compiled default + setting.Key = "test.bound" + setting.Scope = settings.ScopeGlobal + + getter := &toggleGetter{} + bl, err := MakeUpperBoundLimiter(Factory{Settings: getter}, setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, bl.Close()) }) + + ctx := t.Context() + + getter.succeedWith("20gb") + v, err := bl.Limit(ctx) + require.NoError(t, err) + assert.Equal(t, 20*config.GByte, v) + + getter.fail() + v, err = bl.Limit(ctx) + require.ErrorIs(t, err, errGetterUnavailable) + assert.Equal(t, 20*config.GByte, v, "should fall back to the last resolved value, not the compiled default") +} + +func TestBoundLimiter_CompiledDefaultWhenNeverResolved(t *testing.T) { + t.Parallel() + + setting := settings.Size(1 * config.GByte) + setting.Key = "test.bound.never-resolved" + setting.Scope = settings.ScopeGlobal + + getter := &toggleGetter{} + getter.fail() + bl, err := MakeUpperBoundLimiter(Factory{Settings: getter}, setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, bl.Close()) }) + + v, err := bl.Limit(t.Context()) + require.ErrorIs(t, err, errGetterUnavailable) + assert.Equal(t, 1*config.GByte, v, "with no last-known-good value yet, must fall back to the compiled default") +} + +func TestTimeLimiter_LastKnownGoodOnReadFailure(t *testing.T) { + t.Parallel() + + setting := settings.Duration(1 * time.Minute) // compiled default + setting.Key = "test.time" + setting.Scope = settings.ScopeGlobal + + getter := &toggleGetter{} + tl, err := Factory{Settings: getter}.MakeTimeLimiter(setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, tl.Close()) }) + + ctx := t.Context() + + getter.succeedWith("10s") + d, err := tl.Limit(ctx) + require.NoError(t, err) + assert.Equal(t, 10*time.Second, d) + + getter.fail() + d, err = tl.Limit(ctx) + require.ErrorIs(t, err, errGetterUnavailable) + assert.Equal(t, 10*time.Second, d, "should fall back to the last resolved value, not the compiled default") +} + +// TestTimeLimiter_WithTimeout_UsableOnReadFailure: WithTimeout used to return +// (nil, nil, err) on a read failure, causing callers to drop the unit of work instead +// of running it with a real timeout. +func TestTimeLimiter_WithTimeout_UsableOnReadFailure(t *testing.T) { + t.Parallel() + + setting := settings.Duration(1 * time.Minute) + setting.Key = "test.time.with-timeout" + setting.Scope = settings.ScopeGlobal + + getter := &toggleGetter{} + getter.succeedWith("10s") + tl, err := Factory{Settings: getter}.MakeTimeLimiter(setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, tl.Close()) }) + + // Warm the last-known-good value via a successful call, then fail the getter. + _, err = tl.Limit(t.Context()) + require.NoError(t, err) + getter.fail() + + before := time.Now() + ctx, done, withTimeoutErr := tl.WithTimeout(t.Context()) + require.Error(t, withTimeoutErr, "err is advisory, not fatal - the read did fail") + require.NotNil(t, ctx, "ctx must be usable despite the read failure") + require.NotNil(t, done) + defer done() + + deadline, ok := ctx.Deadline() + require.True(t, ok, "ctx must carry a real deadline, not be unbounded") + assert.InDelta(t, 10*time.Second, deadline.Sub(before), float64(2*time.Second), + "deadline should be based on the last known good timeout (10s), not hang open or fire instantly") +} + +func TestGateLimiter_LastKnownGoodOnReadFailure(t *testing.T) { + t.Parallel() + + setting := settings.Bool(false) // compiled default + setting.Key = "test.gate" + setting.Scope = settings.ScopeGlobal + + getter := &toggleGetter{} + gl, err := MakeGateLimiter(Factory{Settings: getter}, setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, gl.Close()) }) + + ctx := t.Context() + + getter.succeedWith("true") + open, err := gl.Limit(ctx) + require.NoError(t, err) + assert.True(t, open) + + getter.fail() + open, err = gl.Limit(ctx) + require.ErrorIs(t, err, errGetterUnavailable) + assert.True(t, open, "should fall back to the last resolved value, not the compiled default") +} + +func TestRangeLimiter_LastKnownGoodOnReadFailure(t *testing.T) { + t.Parallel() + + setting := settings.NewSetting(settings.Range[int]{Lower: 0, Upper: 5}, settings.ParseRangeFn(strconv.Atoi)) + setting.Key = "test.range" + setting.Scope = settings.ScopeGlobal + + getter := &toggleGetter{} + rl, err := MakeRangeLimiter[int](Factory{Settings: getter}, setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, rl.Close()) }) + + ctx := t.Context() + + getter.succeedWith("[1,50]") + got, err := rl.Limit(ctx) + require.NoError(t, err) + assert.Equal(t, settings.Range[int]{Lower: 1, Upper: 50}, got) + + getter.fail() + got, err = rl.Limit(ctx) + require.ErrorIs(t, err, errGetterUnavailable) + assert.Equal(t, settings.Range[int]{Lower: 1, Upper: 50}, got, "should fall back to the last resolved value, not the compiled default") +} diff --git a/pkg/settings/limits/gate.go b/pkg/settings/limits/gate.go index 611362ebf6..59523a52f7 100644 --- a/pkg/settings/limits/gate.go +++ b/pkg/settings/limits/gate.go @@ -174,12 +174,12 @@ func (g *gateLimiter) Limit(ctx context.Context) (bool, error) { } defer g.wg.Done() - _, limit, err := g.get(ctx) - if err != nil { - return false, err + tenant, limit, err := g.get(ctx) + if err != nil && tenant == "" && g.scope != settings.ScopeGlobal { + return false, err // no tenant, so get() never read a value at all } - return limit, nil + return limit, err // limit is always usable; err is advisory } func (g *gateLimiter) AllowErr(ctx context.Context) error { @@ -200,6 +200,7 @@ func (g *gateLimiter) AllowErr(ctx context.Context) error { } func (g *gateLimiter) get(ctx context.Context) (tenant string, open bool, err error) { + u := g.updater if g.scope != settings.ScopeGlobal { tenant = g.scope.Value(ctx) if tenant == "" { @@ -212,10 +213,11 @@ func (g *gateLimiter) get(ctx context.Context) (tenant string, open bool, err er return } - u := newUpdater(g.lggr, g.getLimitFn, g.subFn) - actual, loaded := g.updaters.LoadOrStore(tenant, u) + newU := newUpdater(g.lggr, g.getLimitFn, g.subFn) + actual, loaded := g.updaters.LoadOrStore(tenant, newU) creCtx := contexts.WithCRE(ctx, g.scope.RoundCRE(contexts.CREValue(ctx))) if !loaded { + u = newU go u.updateLoop(creCtx) } else { u = actual.(*updater[bool]) @@ -225,7 +227,14 @@ func (g *gateLimiter) get(ctx context.Context) (tenant string, open bool, err er open, err = g.getLimitFn(ctx) if err != nil { - g.lggr.Errorw("Failed to get status. Using default value", "default", open, "err", err) + if last, ok := u.lastGood(); ok { + g.lggr.Errorw("Failed to get status. Using last known value", "value", last, "err", err) + open = last + } else { + g.lggr.Errorw("Failed to get status. Using compiled default", "default", open, "err", err) + } + } else { + u.setLast(open) } // TODO: include map key in attributes g.recordStatus(ctx, open, withScope(ctx, g.scope)) diff --git a/pkg/settings/limits/range.go b/pkg/settings/limits/range.go index 136ff31f65..44c13eaf17 100644 --- a/pkg/settings/limits/range.go +++ b/pkg/settings/limits/range.go @@ -202,17 +202,18 @@ func (b *rangeLimiter[N]) Limit(ctx context.Context) (settings.Range[N], error) defer b.wg.Done() tenant, bound, err := b.get(ctx) - if err != nil { - return zero, err + if err != nil && tenant == "" && b.scope != settings.ScopeGlobal { + return zero, err // no tenant, so get() never read a value at all } - if tenant == "" && b.scope != settings.ScopeGlobal { + if tenant == "" && b.scope != settings.ScopeGlobal && err == nil { return zero, nil // fail open } - return bound, nil + return bound, err // bound is always usable; err is advisory } func (b *rangeLimiter[N]) get(ctx context.Context) (tenant string, bound settings.Range[N], err error) { + u := b.updater if b.scope != settings.ScopeGlobal { tenant = b.scope.Value(ctx) if tenant == "" { @@ -225,10 +226,11 @@ func (b *rangeLimiter[N]) get(ctx context.Context) (tenant string, bound setting return } - u := newUpdater(b.lggr, b.getLimitFn, b.subFn) - actual, loaded := b.updaters.LoadOrStore(tenant, u) + newU := newUpdater(b.lggr, b.getLimitFn, b.subFn) + actual, loaded := b.updaters.LoadOrStore(tenant, newU) creCtx := contexts.WithCRE(ctx, b.scope.RoundCRE(contexts.CREValue(ctx))) if !loaded { + u = newU go u.updateLoop(creCtx) } else { u = actual.(*updater[settings.Range[N]]) @@ -238,7 +240,14 @@ func (b *rangeLimiter[N]) get(ctx context.Context) (tenant string, bound setting bound, err = b.getLimitFn(ctx) if err != nil { - b.lggr.Errorw("Failed to get limit. Using default value", "default", bound, "err", err) + if last, ok := u.lastGood(); ok { + b.lggr.Errorw("Failed to get limit. Using last known value", "value", last, "err", err) + bound = last + } else { + b.lggr.Errorw("Failed to get limit. Using compiled default", "default", bound, "err", err) + } + } else { + u.setLast(bound) } b.recordBound(ctx, bound, withScope(ctx, b.scope)) return diff --git a/pkg/settings/limits/time.go b/pkg/settings/limits/time.go index ce2e23b477..bee273617d 100644 --- a/pkg/settings/limits/time.go +++ b/pkg/settings/limits/time.go @@ -191,27 +191,28 @@ func (l *timeLimiter) WithTimeout(ctx context.Context) (context.Context, func(), defer l.wg.Done() tenant, timeout, err := l.get(ctx) - if err != nil { - return nil, nil, err + if err != nil && tenant == "" && l.scope != settings.ScopeGlobal { + return nil, nil, err // no tenant, so get() never read a value at all } - if tenant == "" && l.scope != settings.ScopeGlobal { + if tenant == "" && l.scope != settings.ScopeGlobal && err == nil { return ctx, func() {}, nil // fail open } + // timeout is always usable; err is advisory. Still build a real deadline from it. countTimeout := func() { l.countTimeout(ctx) } // constructing this first to reference the original ctx - ctx, cancel := context.WithTimeoutCause(ctx, timeout, ErrorTimeLimited{Key: l.key, Scope: l.scope, Tenant: tenant, Timeout: timeout}) - stop := context.AfterFunc(ctx, countTimeout) + timeoutCtx, cancel := context.WithTimeoutCause(ctx, timeout, ErrorTimeLimited{Key: l.key, Scope: l.scope, Tenant: tenant, Timeout: timeout}) + stop := context.AfterFunc(timeoutCtx, countTimeout) start := time.Now() - return ctx, func() { + return timeoutCtx, func() { elapsed := time.Since(start) - l.recordRuntime(ctx, elapsed) + l.recordRuntime(timeoutCtx, elapsed) if stop() { - l.countSuccess(ctx) + l.countSuccess(timeoutCtx) } cancel() - }, nil + }, err } func (l *timeLimiter) Limit(ctx context.Context) (time.Duration, error) { @@ -221,17 +222,18 @@ func (l *timeLimiter) Limit(ctx context.Context) (time.Duration, error) { defer l.wg.Done() tenant, timeout, err := l.get(ctx) - if err != nil { + if err != nil && tenant == "" && l.scope != settings.ScopeGlobal { return -1, err } - if tenant == "" && l.scope != settings.ScopeGlobal { + if tenant == "" && l.scope != settings.ScopeGlobal && err == nil { return -1, nil // fail open } - return timeout, nil + return timeout, err } func (l *timeLimiter) get(ctx context.Context) (tenant string, timeout time.Duration, err error) { + u := l.updater if l.scope != settings.ScopeGlobal { tenant = l.scope.Value(ctx) if tenant == "" { @@ -244,10 +246,11 @@ func (l *timeLimiter) get(ctx context.Context) (tenant string, timeout time.Dura return } - u := newUpdater(l.lggr, l.getLimitFn, l.subFn) - actual, loaded := l.updaters.LoadOrStore(tenant, u) + newU := newUpdater(l.lggr, l.getLimitFn, l.subFn) + actual, loaded := l.updaters.LoadOrStore(tenant, newU) creCtx := contexts.WithCRE(ctx, l.scope.RoundCRE(contexts.CREValue(ctx))) if !loaded { + u = newU go u.updateLoop(creCtx) } else { u = actual.(*updater[time.Duration]) @@ -257,7 +260,14 @@ func (l *timeLimiter) get(ctx context.Context) (tenant string, timeout time.Dura timeout, err = l.getLimitFn(ctx) if err != nil { - l.lggr.Errorw("Failed to get limit. Using default value", "default", timeout, "err", err) + if last, ok := u.lastGood(); ok { + l.lggr.Errorw("Failed to get limit. Using last known value", "value", last, "err", err) + timeout = last + } else { + l.lggr.Errorw("Failed to get limit. Using compiled default", "default", timeout, "err", err) + } + } else { + u.setLast(timeout) } l.recordTimeout(ctx, timeout) diff --git a/pkg/settings/limits/updater.go b/pkg/settings/limits/updater.go index 5956770bf3..2de946650b 100644 --- a/pkg/settings/limits/updater.go +++ b/pkg/settings/limits/updater.go @@ -3,6 +3,7 @@ package limits import ( "context" "sync" + "sync/atomic" "time" "github.com/smartcontractkit/chainlink-common/pkg/logger" @@ -25,6 +26,9 @@ type updater[N any] struct { stopCh services.StopChan done chan struct{} cancelSub func() // optional + + // lastGoodValue is the fallback used on a read failure, instead of Setting.DefaultValue. + lastGoodValue atomic.Pointer[N] } // newUpdater returns a new updater. lggr and subFn are optional, but getLimitFn is required. @@ -62,6 +66,19 @@ func (u *updater[N]) updateCtx(ctx context.Context) { } } +func (u *updater[N]) setLast(n N) { + u.lastGoodValue.Store(&n) +} + +// lastGood returns false if no value has ever resolved successfully. +func (u *updater[N]) lastGood() (N, bool) { + if v := u.lastGoodValue.Load(); v != nil { + return *v, true + } + var zero N + return zero, false +} + // updateLoop updates the limit either by subscribing via subFn or polling if subFn is not set. It also processes // contexts.CRE updates. Stopped by Close. // opt: reap after period of non-use @@ -90,7 +107,14 @@ func (u *updater[N]) updateLoop(ctx context.Context) { case <-c: limit, err := u.getLimitFn(ctx) if err != nil { - u.lggr.Errorw("Failed to get limit. Using default value", "default", limit, "err", err) + if last, ok := u.lastGood(); ok { + u.lggr.Errorw("Failed to get limit. Using last known value", "value", last, "err", err) + limit = last + } else { + u.lggr.Errorw("Failed to get limit. Using compiled default", "default", limit, "err", err) + } + } else { + u.setLast(limit) } u.recordLimit(ctx, limit) if u.onLimitUpdate != nil { @@ -98,10 +122,18 @@ func (u *updater[N]) updateLoop(ctx context.Context) { } case update := <-updates: + limit := update.Value if update.Err != nil { - u.lggr.Errorw("Failed to update limit. Using default value", "default", update.Value, "err", update.Err) + if last, ok := u.lastGood(); ok { + u.lggr.Errorw("Failed to update limit. Using last known value", "value", last, "err", update.Err) + limit = last + } else { + u.lggr.Errorw("Failed to update limit. Using compiled default", "default", update.Value, "err", update.Err) + } + } else { + u.setLast(limit) } - u.recordLimit(ctx, update.Value) + u.recordLimit(ctx, limit) if u.onLimitUpdate != nil { u.onLimitUpdate(ctx) } From 6886fa6a30cafd111172a66363a55a6e3ba03c7e Mon Sep 17 00:00:00 2001 From: Tarcisio Ferraz Date: Tue, 25 Aug 2026 13:17:28 -0300 Subject: [PATCH 2/5] use compiled default on error --- pkg/settings/limits/bound.go | 24 +-- pkg/settings/limits/default_fallback_test.go | 200 ------------------- pkg/settings/limits/default_on_error_test.go | 109 ++++++++++ pkg/settings/limits/gate.go | 21 +- pkg/settings/limits/range.go | 24 +-- pkg/settings/limits/time.go | 45 ++--- pkg/settings/limits/updater.go | 38 +--- 7 files changed, 146 insertions(+), 315 deletions(-) delete mode 100644 pkg/settings/limits/default_fallback_test.go create mode 100644 pkg/settings/limits/default_on_error_test.go diff --git a/pkg/settings/limits/bound.go b/pkg/settings/limits/bound.go index 39ddc66c4c..58d5d8f2aa 100644 --- a/pkg/settings/limits/bound.go +++ b/pkg/settings/limits/bound.go @@ -213,18 +213,14 @@ func (b *boundLimiter[N]) Limit(ctx context.Context) (N, error) { defer b.wg.Done() tenant, bound, err := b.get(ctx) - if err != nil && tenant == "" && b.scope != settings.ScopeGlobal { - return zero, err // no tenant, so get() never read a value at all - } - if tenant == "" && b.scope != settings.ScopeGlobal && err == nil { - return zero, nil // fail open + if tenant == "" && b.scope != settings.ScopeGlobal { + return zero, err // no tenant: get() never resolved a value } - return bound, err // bound is always usable; err is advisory + return bound, err // bound is get()'s resolved (or default) value; err is advisory } func (b *boundLimiter[N]) get(ctx context.Context) (tenant string, bound N, err error) { - u := b.updater if b.scope != settings.ScopeGlobal { tenant = b.scope.Value(ctx) if tenant == "" { @@ -237,11 +233,10 @@ func (b *boundLimiter[N]) get(ctx context.Context) (tenant string, bound N, err return } - newU := newUpdater(b.lggr, b.getLimitFn, b.subFn) - actual, loaded := b.updaters.LoadOrStore(tenant, newU) + u := newUpdater(b.lggr, b.getLimitFn, b.subFn) + actual, loaded := b.updaters.LoadOrStore(tenant, u) creCtx := contexts.WithCRE(ctx, b.scope.RoundCRE(contexts.CREValue(ctx))) if !loaded { - u = newU go u.updateLoop(creCtx) } else { u = actual.(*updater[N]) @@ -251,14 +246,7 @@ func (b *boundLimiter[N]) get(ctx context.Context) (tenant string, bound N, err bound, err = b.getLimitFn(ctx) if err != nil { - if last, ok := u.lastGood(); ok { - b.lggr.Errorw("Failed to get limit. Using last known value", "value", last, "err", err) - bound = last - } else { - b.lggr.Errorw("Failed to get limit. Using compiled default", "default", bound, "err", err) - } - } else { - u.setLast(bound) + b.lggr.Errorw("Failed to get limit. Using default value", "default", bound, "err", err) } b.recordBound(ctx, bound, withScope(ctx, b.scope)) return diff --git a/pkg/settings/limits/default_fallback_test.go b/pkg/settings/limits/default_fallback_test.go deleted file mode 100644 index efed4d923c..0000000000 --- a/pkg/settings/limits/default_fallback_test.go +++ /dev/null @@ -1,200 +0,0 @@ -package limits - -import ( - "context" - "errors" - "strconv" - "sync" - "testing" - "time" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - - "github.com/smartcontractkit/chainlink-common/pkg/config" - "github.com/smartcontractkit/chainlink-common/pkg/settings" -) - -// toggleGetter is a settings.Getter that can be switched between a fixed value and -// failing, to drive the last-known-good fallback path deterministically. -type toggleGetter struct { - mu sync.Mutex - value string - failing bool -} - -func (g *toggleGetter) succeedWith(value string) { - g.mu.Lock() - defer g.mu.Unlock() - g.value = value - g.failing = false -} - -func (g *toggleGetter) fail() { - g.mu.Lock() - defer g.mu.Unlock() - g.failing = true -} - -var errGetterUnavailable = errors.New("settings getter unavailable") - -func (g *toggleGetter) GetScoped(context.Context, settings.Scope, string) (string, error) { - g.mu.Lock() - defer g.mu.Unlock() - if g.failing { - return "", errGetterUnavailable - } - return g.value, nil -} - -func TestBoundLimiter_LastKnownGoodOnReadFailure(t *testing.T) { - t.Parallel() - - setting := settings.Size(1 * config.GByte) // compiled default - setting.Key = "test.bound" - setting.Scope = settings.ScopeGlobal - - getter := &toggleGetter{} - bl, err := MakeUpperBoundLimiter(Factory{Settings: getter}, setting) - require.NoError(t, err) - t.Cleanup(func() { assert.NoError(t, bl.Close()) }) - - ctx := t.Context() - - getter.succeedWith("20gb") - v, err := bl.Limit(ctx) - require.NoError(t, err) - assert.Equal(t, 20*config.GByte, v) - - getter.fail() - v, err = bl.Limit(ctx) - require.ErrorIs(t, err, errGetterUnavailable) - assert.Equal(t, 20*config.GByte, v, "should fall back to the last resolved value, not the compiled default") -} - -func TestBoundLimiter_CompiledDefaultWhenNeverResolved(t *testing.T) { - t.Parallel() - - setting := settings.Size(1 * config.GByte) - setting.Key = "test.bound.never-resolved" - setting.Scope = settings.ScopeGlobal - - getter := &toggleGetter{} - getter.fail() - bl, err := MakeUpperBoundLimiter(Factory{Settings: getter}, setting) - require.NoError(t, err) - t.Cleanup(func() { assert.NoError(t, bl.Close()) }) - - v, err := bl.Limit(t.Context()) - require.ErrorIs(t, err, errGetterUnavailable) - assert.Equal(t, 1*config.GByte, v, "with no last-known-good value yet, must fall back to the compiled default") -} - -func TestTimeLimiter_LastKnownGoodOnReadFailure(t *testing.T) { - t.Parallel() - - setting := settings.Duration(1 * time.Minute) // compiled default - setting.Key = "test.time" - setting.Scope = settings.ScopeGlobal - - getter := &toggleGetter{} - tl, err := Factory{Settings: getter}.MakeTimeLimiter(setting) - require.NoError(t, err) - t.Cleanup(func() { assert.NoError(t, tl.Close()) }) - - ctx := t.Context() - - getter.succeedWith("10s") - d, err := tl.Limit(ctx) - require.NoError(t, err) - assert.Equal(t, 10*time.Second, d) - - getter.fail() - d, err = tl.Limit(ctx) - require.ErrorIs(t, err, errGetterUnavailable) - assert.Equal(t, 10*time.Second, d, "should fall back to the last resolved value, not the compiled default") -} - -// TestTimeLimiter_WithTimeout_UsableOnReadFailure: WithTimeout used to return -// (nil, nil, err) on a read failure, causing callers to drop the unit of work instead -// of running it with a real timeout. -func TestTimeLimiter_WithTimeout_UsableOnReadFailure(t *testing.T) { - t.Parallel() - - setting := settings.Duration(1 * time.Minute) - setting.Key = "test.time.with-timeout" - setting.Scope = settings.ScopeGlobal - - getter := &toggleGetter{} - getter.succeedWith("10s") - tl, err := Factory{Settings: getter}.MakeTimeLimiter(setting) - require.NoError(t, err) - t.Cleanup(func() { assert.NoError(t, tl.Close()) }) - - // Warm the last-known-good value via a successful call, then fail the getter. - _, err = tl.Limit(t.Context()) - require.NoError(t, err) - getter.fail() - - before := time.Now() - ctx, done, withTimeoutErr := tl.WithTimeout(t.Context()) - require.Error(t, withTimeoutErr, "err is advisory, not fatal - the read did fail") - require.NotNil(t, ctx, "ctx must be usable despite the read failure") - require.NotNil(t, done) - defer done() - - deadline, ok := ctx.Deadline() - require.True(t, ok, "ctx must carry a real deadline, not be unbounded") - assert.InDelta(t, 10*time.Second, deadline.Sub(before), float64(2*time.Second), - "deadline should be based on the last known good timeout (10s), not hang open or fire instantly") -} - -func TestGateLimiter_LastKnownGoodOnReadFailure(t *testing.T) { - t.Parallel() - - setting := settings.Bool(false) // compiled default - setting.Key = "test.gate" - setting.Scope = settings.ScopeGlobal - - getter := &toggleGetter{} - gl, err := MakeGateLimiter(Factory{Settings: getter}, setting) - require.NoError(t, err) - t.Cleanup(func() { assert.NoError(t, gl.Close()) }) - - ctx := t.Context() - - getter.succeedWith("true") - open, err := gl.Limit(ctx) - require.NoError(t, err) - assert.True(t, open) - - getter.fail() - open, err = gl.Limit(ctx) - require.ErrorIs(t, err, errGetterUnavailable) - assert.True(t, open, "should fall back to the last resolved value, not the compiled default") -} - -func TestRangeLimiter_LastKnownGoodOnReadFailure(t *testing.T) { - t.Parallel() - - setting := settings.NewSetting(settings.Range[int]{Lower: 0, Upper: 5}, settings.ParseRangeFn(strconv.Atoi)) - setting.Key = "test.range" - setting.Scope = settings.ScopeGlobal - - getter := &toggleGetter{} - rl, err := MakeRangeLimiter[int](Factory{Settings: getter}, setting) - require.NoError(t, err) - t.Cleanup(func() { assert.NoError(t, rl.Close()) }) - - ctx := t.Context() - - getter.succeedWith("[1,50]") - got, err := rl.Limit(ctx) - require.NoError(t, err) - assert.Equal(t, settings.Range[int]{Lower: 1, Upper: 50}, got) - - getter.fail() - got, err = rl.Limit(ctx) - require.ErrorIs(t, err, errGetterUnavailable) - assert.Equal(t, settings.Range[int]{Lower: 1, Upper: 50}, got, "should fall back to the last resolved value, not the compiled default") -} diff --git a/pkg/settings/limits/default_on_error_test.go b/pkg/settings/limits/default_on_error_test.go new file mode 100644 index 0000000000..7d455626e5 --- /dev/null +++ b/pkg/settings/limits/default_on_error_test.go @@ -0,0 +1,109 @@ +package limits + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/smartcontractkit/chainlink-common/pkg/config" + "github.com/smartcontractkit/chainlink-common/pkg/settings" +) + +var errGetterUnavailable = errors.New("settings getter unavailable") + +// failingGetter always fails GetScoped, standing in for a settings-service outage. +type failingGetter struct{} + +func (failingGetter) GetScoped(context.Context, settings.Scope, string) (string, error) { + return "", errGetterUnavailable +} + +// TestLimiter_Limit_ReturnsDefaultOnReadFailure is the parity fix: Limit() must return +// the value get() already resolved (the compiled default, on a read failure) alongside +// the error, instead of discarding it. +func TestLimiter_Limit_ReturnsDefaultOnReadFailure(t *testing.T) { + t.Parallel() + + t.Run("bound", func(t *testing.T) { + t.Parallel() + setting := settings.Size(1 * config.GByte) + setting.Key, setting.Scope = "test.bound", settings.ScopeGlobal + bl, err := MakeUpperBoundLimiter(Factory{Settings: failingGetter{}}, setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, bl.Close()) }) + + v, err := bl.Limit(t.Context()) + require.ErrorIs(t, err, errGetterUnavailable) + assert.Equal(t, 1*config.GByte, v) + }) + + t.Run("time", func(t *testing.T) { + t.Parallel() + setting := settings.Duration(1 * time.Minute) + setting.Key, setting.Scope = "test.time", settings.ScopeGlobal + tl, err := Factory{Settings: failingGetter{}}.MakeTimeLimiter(setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, tl.Close()) }) + + d, err := tl.Limit(t.Context()) + require.ErrorIs(t, err, errGetterUnavailable) + assert.Equal(t, 1*time.Minute, d) + }) + + t.Run("gate", func(t *testing.T) { + t.Parallel() + setting := settings.Bool(true) + setting.Key, setting.Scope = "test.gate", settings.ScopeGlobal + gl, err := MakeGateLimiter(Factory{Settings: failingGetter{}}, setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, gl.Close()) }) + + open, err := gl.Limit(t.Context()) + require.ErrorIs(t, err, errGetterUnavailable) + assert.True(t, open) + }) + + t.Run("range", func(t *testing.T) { + t.Parallel() + setting := settings.NewSetting(settings.Range[int]{Lower: 1, Upper: 5}, settings.ParseRangeFn(func(s string) (int, error) { + return 0, nil + })) + setting.Key, setting.Scope = "test.range", settings.ScopeGlobal + rl, err := MakeRangeLimiter[int](Factory{Settings: failingGetter{}}, setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, rl.Close()) }) + + got, err := rl.Limit(t.Context()) + require.ErrorIs(t, err, errGetterUnavailable) + assert.Equal(t, settings.Range[int]{Lower: 1, Upper: 5}, got) + }) +} + +// TestTimeLimiter_WithTimeout_UsableOnReadFailure is the fix for the actual production +// incident: WithTimeout used to return (nil, nil, err) on a read failure, causing callers +// to drop the unit of work instead of running it with the compiled default timeout. +func TestTimeLimiter_WithTimeout_UsableOnReadFailure(t *testing.T) { + t.Parallel() + + setting := settings.Duration(10 * time.Second) + setting.Key, setting.Scope = "test.time.with-timeout", settings.ScopeGlobal + tl, err := Factory{Settings: failingGetter{}}.MakeTimeLimiter(setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, tl.Close()) }) + + before := time.Now() + ctx, done, withTimeoutErr := tl.WithTimeout(t.Context()) + require.Error(t, withTimeoutErr, "err is advisory, not fatal - the read did fail") + require.NotNil(t, ctx, "ctx must be usable despite the read failure") + require.NotNil(t, done) + defer done() + + deadline, ok := ctx.Deadline() + require.True(t, ok, "ctx must carry a real deadline, not be unbounded") + assert.InDelta(t, 10*time.Second, deadline.Sub(before), float64(2*time.Second), + "deadline should be based on the compiled default (10s), not hang open or fire instantly") +} diff --git a/pkg/settings/limits/gate.go b/pkg/settings/limits/gate.go index 59523a52f7..2ca06fede3 100644 --- a/pkg/settings/limits/gate.go +++ b/pkg/settings/limits/gate.go @@ -175,11 +175,11 @@ func (g *gateLimiter) Limit(ctx context.Context) (bool, error) { defer g.wg.Done() tenant, limit, err := g.get(ctx) - if err != nil && tenant == "" && g.scope != settings.ScopeGlobal { - return false, err // no tenant, so get() never read a value at all + if tenant == "" && g.scope != settings.ScopeGlobal { + return false, err // no tenant: get() never resolved a value } - return limit, err // limit is always usable; err is advisory + return limit, err // limit is get()'s resolved (or default) value; err is advisory } func (g *gateLimiter) AllowErr(ctx context.Context) error { @@ -200,7 +200,6 @@ func (g *gateLimiter) AllowErr(ctx context.Context) error { } func (g *gateLimiter) get(ctx context.Context) (tenant string, open bool, err error) { - u := g.updater if g.scope != settings.ScopeGlobal { tenant = g.scope.Value(ctx) if tenant == "" { @@ -213,11 +212,10 @@ func (g *gateLimiter) get(ctx context.Context) (tenant string, open bool, err er return } - newU := newUpdater(g.lggr, g.getLimitFn, g.subFn) - actual, loaded := g.updaters.LoadOrStore(tenant, newU) + u := newUpdater(g.lggr, g.getLimitFn, g.subFn) + actual, loaded := g.updaters.LoadOrStore(tenant, u) creCtx := contexts.WithCRE(ctx, g.scope.RoundCRE(contexts.CREValue(ctx))) if !loaded { - u = newU go u.updateLoop(creCtx) } else { u = actual.(*updater[bool]) @@ -227,14 +225,7 @@ func (g *gateLimiter) get(ctx context.Context) (tenant string, open bool, err er open, err = g.getLimitFn(ctx) if err != nil { - if last, ok := u.lastGood(); ok { - g.lggr.Errorw("Failed to get status. Using last known value", "value", last, "err", err) - open = last - } else { - g.lggr.Errorw("Failed to get status. Using compiled default", "default", open, "err", err) - } - } else { - u.setLast(open) + g.lggr.Errorw("Failed to get status. Using default value", "default", open, "err", err) } // TODO: include map key in attributes g.recordStatus(ctx, open, withScope(ctx, g.scope)) diff --git a/pkg/settings/limits/range.go b/pkg/settings/limits/range.go index 44c13eaf17..3d8cbad73a 100644 --- a/pkg/settings/limits/range.go +++ b/pkg/settings/limits/range.go @@ -202,18 +202,14 @@ func (b *rangeLimiter[N]) Limit(ctx context.Context) (settings.Range[N], error) defer b.wg.Done() tenant, bound, err := b.get(ctx) - if err != nil && tenant == "" && b.scope != settings.ScopeGlobal { - return zero, err // no tenant, so get() never read a value at all - } - if tenant == "" && b.scope != settings.ScopeGlobal && err == nil { - return zero, nil // fail open + if tenant == "" && b.scope != settings.ScopeGlobal { + return zero, err // no tenant: get() never resolved a value } - return bound, err // bound is always usable; err is advisory + return bound, err // bound is get()'s resolved (or default) value; err is advisory } func (b *rangeLimiter[N]) get(ctx context.Context) (tenant string, bound settings.Range[N], err error) { - u := b.updater if b.scope != settings.ScopeGlobal { tenant = b.scope.Value(ctx) if tenant == "" { @@ -226,11 +222,10 @@ func (b *rangeLimiter[N]) get(ctx context.Context) (tenant string, bound setting return } - newU := newUpdater(b.lggr, b.getLimitFn, b.subFn) - actual, loaded := b.updaters.LoadOrStore(tenant, newU) + u := newUpdater(b.lggr, b.getLimitFn, b.subFn) + actual, loaded := b.updaters.LoadOrStore(tenant, u) creCtx := contexts.WithCRE(ctx, b.scope.RoundCRE(contexts.CREValue(ctx))) if !loaded { - u = newU go u.updateLoop(creCtx) } else { u = actual.(*updater[settings.Range[N]]) @@ -240,14 +235,7 @@ func (b *rangeLimiter[N]) get(ctx context.Context) (tenant string, bound setting bound, err = b.getLimitFn(ctx) if err != nil { - if last, ok := u.lastGood(); ok { - b.lggr.Errorw("Failed to get limit. Using last known value", "value", last, "err", err) - bound = last - } else { - b.lggr.Errorw("Failed to get limit. Using compiled default", "default", bound, "err", err) - } - } else { - u.setLast(bound) + b.lggr.Errorw("Failed to get limit. Using default value", "default", bound, "err", err) } b.recordBound(ctx, bound, withScope(ctx, b.scope)) return diff --git a/pkg/settings/limits/time.go b/pkg/settings/limits/time.go index bee273617d..9f25d5fd3d 100644 --- a/pkg/settings/limits/time.go +++ b/pkg/settings/limits/time.go @@ -191,28 +191,27 @@ func (l *timeLimiter) WithTimeout(ctx context.Context) (context.Context, func(), defer l.wg.Done() tenant, timeout, err := l.get(ctx) - if err != nil && tenant == "" && l.scope != settings.ScopeGlobal { - return nil, nil, err // no tenant, so get() never read a value at all - } - if tenant == "" && l.scope != settings.ScopeGlobal && err == nil { + if tenant == "" && l.scope != settings.ScopeGlobal { + if err != nil { + return nil, nil, err // no tenant: get() never resolved a value + } return ctx, func() {}, nil // fail open } - // timeout is always usable; err is advisory. Still build a real deadline from it. countTimeout := func() { l.countTimeout(ctx) } // constructing this first to reference the original ctx - timeoutCtx, cancel := context.WithTimeoutCause(ctx, timeout, ErrorTimeLimited{Key: l.key, Scope: l.scope, Tenant: tenant, Timeout: timeout}) - stop := context.AfterFunc(timeoutCtx, countTimeout) + ctx, cancel := context.WithTimeoutCause(ctx, timeout, ErrorTimeLimited{Key: l.key, Scope: l.scope, Tenant: tenant, Timeout: timeout}) + stop := context.AfterFunc(ctx, countTimeout) start := time.Now() - return timeoutCtx, func() { + return ctx, func() { elapsed := time.Since(start) - l.recordRuntime(timeoutCtx, elapsed) + l.recordRuntime(ctx, elapsed) if stop() { - l.countSuccess(timeoutCtx) + l.countSuccess(ctx) } cancel() - }, err + }, err // timeout is get()'s resolved (or default) value; err is advisory } func (l *timeLimiter) Limit(ctx context.Context) (time.Duration, error) { @@ -222,18 +221,14 @@ func (l *timeLimiter) Limit(ctx context.Context) (time.Duration, error) { defer l.wg.Done() tenant, timeout, err := l.get(ctx) - if err != nil && tenant == "" && l.scope != settings.ScopeGlobal { - return -1, err - } - if tenant == "" && l.scope != settings.ScopeGlobal && err == nil { - return -1, nil // fail open + if tenant == "" && l.scope != settings.ScopeGlobal { + return -1, err // no tenant: get() never resolved a value } - return timeout, err + return timeout, err // timeout is get()'s resolved (or default) value; err is advisory } func (l *timeLimiter) get(ctx context.Context) (tenant string, timeout time.Duration, err error) { - u := l.updater if l.scope != settings.ScopeGlobal { tenant = l.scope.Value(ctx) if tenant == "" { @@ -246,11 +241,10 @@ func (l *timeLimiter) get(ctx context.Context) (tenant string, timeout time.Dura return } - newU := newUpdater(l.lggr, l.getLimitFn, l.subFn) - actual, loaded := l.updaters.LoadOrStore(tenant, newU) + u := newUpdater(l.lggr, l.getLimitFn, l.subFn) + actual, loaded := l.updaters.LoadOrStore(tenant, u) creCtx := contexts.WithCRE(ctx, l.scope.RoundCRE(contexts.CREValue(ctx))) if !loaded { - u = newU go u.updateLoop(creCtx) } else { u = actual.(*updater[time.Duration]) @@ -260,14 +254,7 @@ func (l *timeLimiter) get(ctx context.Context) (tenant string, timeout time.Dura timeout, err = l.getLimitFn(ctx) if err != nil { - if last, ok := u.lastGood(); ok { - l.lggr.Errorw("Failed to get limit. Using last known value", "value", last, "err", err) - timeout = last - } else { - l.lggr.Errorw("Failed to get limit. Using compiled default", "default", timeout, "err", err) - } - } else { - u.setLast(timeout) + l.lggr.Errorw("Failed to get limit. Using default value", "default", timeout, "err", err) } l.recordTimeout(ctx, timeout) diff --git a/pkg/settings/limits/updater.go b/pkg/settings/limits/updater.go index 2de946650b..5956770bf3 100644 --- a/pkg/settings/limits/updater.go +++ b/pkg/settings/limits/updater.go @@ -3,7 +3,6 @@ package limits import ( "context" "sync" - "sync/atomic" "time" "github.com/smartcontractkit/chainlink-common/pkg/logger" @@ -26,9 +25,6 @@ type updater[N any] struct { stopCh services.StopChan done chan struct{} cancelSub func() // optional - - // lastGoodValue is the fallback used on a read failure, instead of Setting.DefaultValue. - lastGoodValue atomic.Pointer[N] } // newUpdater returns a new updater. lggr and subFn are optional, but getLimitFn is required. @@ -66,19 +62,6 @@ func (u *updater[N]) updateCtx(ctx context.Context) { } } -func (u *updater[N]) setLast(n N) { - u.lastGoodValue.Store(&n) -} - -// lastGood returns false if no value has ever resolved successfully. -func (u *updater[N]) lastGood() (N, bool) { - if v := u.lastGoodValue.Load(); v != nil { - return *v, true - } - var zero N - return zero, false -} - // updateLoop updates the limit either by subscribing via subFn or polling if subFn is not set. It also processes // contexts.CRE updates. Stopped by Close. // opt: reap after period of non-use @@ -107,14 +90,7 @@ func (u *updater[N]) updateLoop(ctx context.Context) { case <-c: limit, err := u.getLimitFn(ctx) if err != nil { - if last, ok := u.lastGood(); ok { - u.lggr.Errorw("Failed to get limit. Using last known value", "value", last, "err", err) - limit = last - } else { - u.lggr.Errorw("Failed to get limit. Using compiled default", "default", limit, "err", err) - } - } else { - u.setLast(limit) + u.lggr.Errorw("Failed to get limit. Using default value", "default", limit, "err", err) } u.recordLimit(ctx, limit) if u.onLimitUpdate != nil { @@ -122,18 +98,10 @@ func (u *updater[N]) updateLoop(ctx context.Context) { } case update := <-updates: - limit := update.Value if update.Err != nil { - if last, ok := u.lastGood(); ok { - u.lggr.Errorw("Failed to update limit. Using last known value", "value", last, "err", update.Err) - limit = last - } else { - u.lggr.Errorw("Failed to update limit. Using compiled default", "default", update.Value, "err", update.Err) - } - } else { - u.setLast(limit) + u.lggr.Errorw("Failed to update limit. Using default value", "default", update.Value, "err", update.Err) } - u.recordLimit(ctx, limit) + u.recordLimit(ctx, update.Value) if u.onLimitUpdate != nil { u.onLimitUpdate(ctx) } From 9d7b12537b0c637498610970bc1471e58ae64549 Mon Sep 17 00:00:00 2001 From: Tarcisio Ferraz Date: Tue, 25 Aug 2026 17:05:04 -0300 Subject: [PATCH 3/5] change return behavior --- pkg/settings/limits/bound.go | 10 +++------- pkg/settings/limits/gate.go | 8 ++------ pkg/settings/limits/limits.go | 5 ++++- pkg/settings/limits/range.go | 10 +++------- 4 files changed, 12 insertions(+), 21 deletions(-) diff --git a/pkg/settings/limits/bound.go b/pkg/settings/limits/bound.go index 58d5d8f2aa..b566ce0d6d 100644 --- a/pkg/settings/limits/bound.go +++ b/pkg/settings/limits/bound.go @@ -206,18 +206,14 @@ func (b *boundLimiter[N]) Check(ctx context.Context, amount N) error { } func (b *boundLimiter[N]) Limit(ctx context.Context) (N, error) { - var zero N if err := b.wg.TryAdd(1); err != nil { + var zero N return zero, err } defer b.wg.Done() - tenant, bound, err := b.get(ctx) - if tenant == "" && b.scope != settings.ScopeGlobal { - return zero, err // no tenant: get() never resolved a value - } - - return bound, err // bound is get()'s resolved (or default) value; err is advisory + _, bound, err := b.get(ctx) + return bound, err // bound is get()'s resolved value; zero if no tenant, or default on error } func (b *boundLimiter[N]) get(ctx context.Context) (tenant string, bound N, err error) { diff --git a/pkg/settings/limits/gate.go b/pkg/settings/limits/gate.go index 2ca06fede3..51a7e9a683 100644 --- a/pkg/settings/limits/gate.go +++ b/pkg/settings/limits/gate.go @@ -174,12 +174,8 @@ func (g *gateLimiter) Limit(ctx context.Context) (bool, error) { } defer g.wg.Done() - tenant, limit, err := g.get(ctx) - if tenant == "" && g.scope != settings.ScopeGlobal { - return false, err // no tenant: get() never resolved a value - } - - return limit, err // limit is get()'s resolved (or default) value; err is advisory + _, limit, err := g.get(ctx) + return limit, err // limit is get()'s resolved value; false if no tenant, or default on error } func (g *gateLimiter) AllowErr(ctx context.Context) error { diff --git a/pkg/settings/limits/limits.go b/pkg/settings/limits/limits.go index ccc1e22a99..201e8608d0 100644 --- a/pkg/settings/limits/limits.go +++ b/pkg/settings/limits/limits.go @@ -29,7 +29,10 @@ type Number interface { type Limiter[N any] interface { io.Closer // Limiters spawn background goroutines and must be closed. - // Limit returns the current limit. + // Limit returns the current limit. The value is always usable, even when err is + // non-nil: on a read failure it's the compiled default (Setting.DefaultValue), so err + // is advisory, not a signal to discard the value. The exception is a required tenant + // missing from ctx, where no read is attempted and the zero value is returned instead. Limit(context.Context) (N, error) } diff --git a/pkg/settings/limits/range.go b/pkg/settings/limits/range.go index 3d8cbad73a..550111a9ee 100644 --- a/pkg/settings/limits/range.go +++ b/pkg/settings/limits/range.go @@ -195,18 +195,14 @@ func (b *rangeLimiter[N]) Check(ctx context.Context, amount N) error { } func (b *rangeLimiter[N]) Limit(ctx context.Context) (settings.Range[N], error) { - var zero settings.Range[N] if err := b.wg.TryAdd(1); err != nil { + var zero settings.Range[N] return zero, err } defer b.wg.Done() - tenant, bound, err := b.get(ctx) - if tenant == "" && b.scope != settings.ScopeGlobal { - return zero, err // no tenant: get() never resolved a value - } - - return bound, err // bound is get()'s resolved (or default) value; err is advisory + _, bound, err := b.get(ctx) + return bound, err // bound is get()'s resolved value; zero if no tenant, or default on error } func (b *rangeLimiter[N]) get(ctx context.Context) (tenant string, bound settings.Range[N], err error) { From 6a03afe28cae694b0a46068071390d61f3fd594a Mon Sep 17 00:00:00 2001 From: Tarcisio Ferraz Date: Wed, 26 Aug 2026 15:21:02 -0300 Subject: [PATCH 4/5] cleanup comments --- pkg/settings/limits/limits.go | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/pkg/settings/limits/limits.go b/pkg/settings/limits/limits.go index 201e8608d0..ffd7c9bb21 100644 --- a/pkg/settings/limits/limits.go +++ b/pkg/settings/limits/limits.go @@ -29,10 +29,7 @@ type Number interface { type Limiter[N any] interface { io.Closer // Limiters spawn background goroutines and must be closed. - // Limit returns the current limit. The value is always usable, even when err is - // non-nil: on a read failure it's the compiled default (Setting.DefaultValue), so err - // is advisory, not a signal to discard the value. The exception is a required tenant - // missing from ctx, where no read is attempted and the zero value is returned instead. + // Limit returns the current limit, or an error along with a usable fallback value. Limit(context.Context) (N, error) } From 3f19d9553cd629ea396c3b0ba9426100d08e772d Mon Sep 17 00:00:00 2001 From: Tarcisio Ferraz Date: Thu, 3 Sep 2026 17:56:23 -0300 Subject: [PATCH 5/5] fallback to compiled default timeout when tenant is missing --- pkg/settings/limits/default_on_error_test.go | 41 ++++++++++++++++++++ pkg/settings/limits/time.go | 23 ++++------- 2 files changed, 49 insertions(+), 15 deletions(-) diff --git a/pkg/settings/limits/default_on_error_test.go b/pkg/settings/limits/default_on_error_test.go index 7d455626e5..f98f5d6c45 100644 --- a/pkg/settings/limits/default_on_error_test.go +++ b/pkg/settings/limits/default_on_error_test.go @@ -107,3 +107,44 @@ func TestTimeLimiter_WithTimeout_UsableOnReadFailure(t *testing.T) { assert.InDelta(t, 10*time.Second, deadline.Sub(before), float64(2*time.Second), "deadline should be based on the compiled default (10s), not hang open or fire instantly") } + +// TestTimeLimiter_WithTimeout_UsableWithoutTenant covers the missing-tenant path. The timeout +// doesn't depend on the tenant, so the context must still be bounded by the compiled default: +// a zero timeout would expire immediately, and a nil context would make callers drop the work. +// A tenant that was required but missing still surfaces an error, so the bug stays visible. +func TestTimeLimiter_WithTimeout_UsableWithoutTenant(t *testing.T) { + t.Parallel() + + for _, tt := range []struct { + scope settings.Scope + expectErr bool + }{ + {settings.ScopeOrg, false}, // tenant not required + {settings.ScopeWorkflow, true}, // tenant required + } { + t.Run(tt.scope.String(), func(t *testing.T) { + t.Parallel() + setting := settings.Duration(10 * time.Second) + setting.Key, setting.Scope = "test.time.no-tenant", tt.scope + tl, err := Factory{}.MakeTimeLimiter(setting) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, tl.Close()) }) + + before := time.Now() + ctx, done, err := tl.WithTimeout(t.Context()) // no contexts.WithCRE, so no tenant + if tt.expectErr { + require.Error(t, err, "a required but missing tenant must stay visible") + } else { + require.NoError(t, err) + } + require.NotNil(t, ctx, "ctx must be usable without a tenant") + require.NotNil(t, done) + defer done() + + deadline, ok := ctx.Deadline() + require.True(t, ok, "ctx must carry a real deadline, not be unbounded") + assert.InDelta(t, 10*time.Second, deadline.Sub(before), float64(2*time.Second), + "deadline should be based on the compiled default (10s), not fire instantly") + }) + } +} diff --git a/pkg/settings/limits/time.go b/pkg/settings/limits/time.go index 9f25d5fd3d..9a9fa5ab8b 100644 --- a/pkg/settings/limits/time.go +++ b/pkg/settings/limits/time.go @@ -191,12 +191,6 @@ func (l *timeLimiter) WithTimeout(ctx context.Context) (context.Context, func(), defer l.wg.Done() tenant, timeout, err := l.get(ctx) - if tenant == "" && l.scope != settings.ScopeGlobal { - if err != nil { - return nil, nil, err // no tenant: get() never resolved a value - } - return ctx, func() {}, nil // fail open - } countTimeout := func() { l.countTimeout(ctx) } // constructing this first to reference the original ctx ctx, cancel := context.WithTimeoutCause(ctx, timeout, ErrorTimeLimited{Key: l.key, Scope: l.scope, Tenant: tenant, Timeout: timeout}) @@ -211,30 +205,29 @@ func (l *timeLimiter) WithTimeout(ctx context.Context) (context.Context, func(), l.countSuccess(ctx) } cancel() - }, err // timeout is get()'s resolved (or default) value; err is advisory + }, err // timeout is get()'s resolved value, or the compiled default; err is advisory } func (l *timeLimiter) Limit(ctx context.Context) (time.Duration, error) { if err := l.wg.TryAdd(1); err != nil { - return -1, err + return 0, err } defer l.wg.Done() - tenant, timeout, err := l.get(ctx) - if tenant == "" && l.scope != settings.ScopeGlobal { - return -1, err // no tenant: get() never resolved a value - } - - return timeout, err // timeout is get()'s resolved (or default) value; err is advisory + _, timeout, err := l.get(ctx) + return timeout, err // timeout is get()'s resolved value, or the compiled default; err is advisory } func (l *timeLimiter) get(ctx context.Context) (tenant string, timeout time.Duration, err error) { if l.scope != settings.ScopeGlobal { tenant = l.scope.Value(ctx) if tenant == "" { + // The timeout doesn't depend on the tenant, so fall back to the compiled default + // rather than leaving it at zero, which would mean an already-expired context. + timeout = l.defaultTimeout if !l.scope.IsTenantRequired() { kvs := contexts.CREValue(ctx).LoggerKVs() - l.lggr.Errorw("Unable to get scoped time limit due to missing tenant: failing open", append([]any{"scope", l.scope}, kvs...)...) + l.lggr.Errorw("Unable to get scoped time limit due to missing tenant: using default value", append([]any{"scope", l.scope, "default", timeout}, kvs...)...) return } err = fmt.Errorf("unable to get scoped time limit due to missing tenant for scope: %s", l.scope)