From 882f51fd62ce3c9ab429ad16b2dc55755ad71dbe Mon Sep 17 00:00:00 2001 From: Ettore Di Giacinto Date: Sat, 26 Sep 2026 23:20:31 +0000 Subject: [PATCH] fix(failover): trip a rate-limited or exhausted target instead of skipping it Treating ResourceExhausted as a capability gap skipped the target without counting a failure, so a target that stays rate limited or out of memory kept its traffic. It is now an ordinary retryable failure: the request moves to the next target and the exhausted one trips. Assisted-by: Claude:claude-opus-5-5 Signed-off-by: Ettore Di Giacinto --- backend/go/localai-proxy/client.go | 12 +++++++----- backend/go/localai-proxy/text_test.go | 2 +- core/http/middleware/failover_test.go | 4 ++-- core/services/failover/classify.go | 25 ++++++++++--------------- core/services/failover/classify_test.go | 4 ++-- core/services/failover/manager.go | 5 ++--- core/services/failover/manager_test.go | 4 ++-- 7 files changed, 26 insertions(+), 30 deletions(-) diff --git a/backend/go/localai-proxy/client.go b/backend/go/localai-proxy/client.go index 5d07562ec..100a931dd 100644 --- a/backend/go/localai-proxy/client.go +++ b/backend/go/localai-proxy/client.go @@ -189,9 +189,10 @@ func transportError(path string, err error) error { // statusError maps a non-2xx upstream reply to a gRPC status. 5xx means the // upstream is unhealthy (Unavailable, so failover retries elsewhere); 4xx // means the request itself is wrong (InvalidArgument, so failover does not -// trip a healthy target over a client error), except 429. 501 is the upstream saying it -// cannot serve this kind of request, which failover treats as a capability -// gap, like our own Unimplemented methods. +// trip a healthy target over a client error). 429 is the exception, see +// below. 501 is the upstream saying it cannot serve this kind of request, +// which failover treats as a capability gap, like our own Unimplemented +// methods. func statusError(path string, resp *http.Response) error { raw, _ := io.ReadAll(io.LimitReader(resp.Body, maxErrorBody+1)) msg := strings.TrimSpace(string(raw)) @@ -209,8 +210,9 @@ func statusError(path string, resp *http.Response) error { case resp.StatusCode >= 500: code = codes.Unavailable case resp.StatusCode == http.StatusTooManyRequests: - // Rate limited: the upstream is healthy but out of capacity, so - // failover skips to the next target without tripping this one. + // Rate limited: the request is fine but the upstream is out of + // capacity. Failover retries it on the next target and trips this + // one, so traffic moves off it for a while. code = codes.ResourceExhausted case resp.StatusCode >= 400: code = codes.InvalidArgument diff --git a/backend/go/localai-proxy/text_test.go b/backend/go/localai-proxy/text_test.go index b47051b4c..53ad73c9b 100644 --- a/backend/go/localai-proxy/text_test.go +++ b/backend/go/localai-proxy/text_test.go @@ -176,7 +176,7 @@ var _ = Describe("localai-proxy", func() { Expect(len(status.Convert(err).Message())).To(BeNumerically("<", 700)) }) - It("maps a 429 upstream to ResourceExhausted so failover skips without tripping", func() { + It("maps a 429 upstream to ResourceExhausted so failover moves to the next target", func() { p := loadProxy(up, nil) up.script("/v1/completions", scriptedResponse{Status: http.StatusTooManyRequests, Body: "slow down"}) diff --git a/core/http/middleware/failover_test.go b/core/http/middleware/failover_test.go index 00439abcc..173468d73 100644 --- a/core/http/middleware/failover_test.go +++ b/core/http/middleware/failover_test.go @@ -293,7 +293,7 @@ var _ = Describe("failover chains in the request pipeline", func() { Expect(st.Targets[0].State).To(Equal(failover.StateHealthy)) }) - It("spills a rate-limited target (gRPC ResourceExhausted) to the next target without tripping it", func() { + It("fails a rate-limited target (gRPC ResourceExhausted) over to the next target and trips it", func() { behavior["a"] = func(echo.Context) error { return grpcstatus.Error(codes.ResourceExhausted, "localai-proxy: upstream /v1/chat/completions returned 429: slow down") } @@ -302,7 +302,7 @@ var _ = Describe("failover chains in the request pipeline", func() { Expect(rec.Body.String()).To(ContainSubstring(`"served":"b"`)) Expect(calls).To(Equal([]string{"a", "b"})) st, _ := fm.ChainStatus("chain") - Expect(st.Targets[0].State).To(Equal(failover.StateHealthy)) + Expect(st.Targets[0].State).To(Equal(failover.StateDown)) }) It("skips a disabled target without tripping it", func() { diff --git a/core/services/failover/classify.go b/core/services/failover/classify.go index 3a50299b6..ac45c3c20 100644 --- a/core/services/failover/classify.go +++ b/core/services/failover/classify.go @@ -46,7 +46,11 @@ func IsRetryable(err error, status int) bool { } if st, ok := grpcstatus.FromError(err); ok { switch st.Code() { - case codes.Unavailable, codes.Internal, codes.DeadlineExceeded, codes.Unknown: + // ResourceExhausted (a rate-limited upstream, what localai-proxy + // returns for a 429, or a backend out of memory) is retried + // elsewhere and trips the target: a gap would skip a chronically + // exhausted target forever without moving traffic off it. + case codes.Unavailable, codes.Internal, codes.DeadlineExceeded, codes.Unknown, codes.ResourceExhausted: return !isRequestError(st.Message()) default: return false @@ -61,25 +65,16 @@ func IsRetryable(err error, status int) bool { return !isRequestError(msg) } -// IsCapabilityGap reports a target that cannot serve this request right now -// for a reason that says nothing about its health: it cannot serve this kind -// of request at all (gRPC Unimplemented), or it is out of capacity, such as a -// rate-limited upstream (gRPC ResourceExhausted, what localai-proxy returns -// for an upstream 429). Matched anywhere in the error chain. The next target -// may serve it, so callers must skip this one without tripping it. +// IsCapabilityGap reports a target that cannot serve this kind of request at +// all (gRPC Unimplemented, anywhere in the error chain). The next target may +// serve it, and this target is not broken: the failure carries no signal +// about its health, so callers must skip it without tripping. func IsCapabilityGap(err error) bool { if err == nil { return false } st, ok := grpcstatus.FromError(err) - if !ok { - return false - } - switch st.Code() { - case codes.Unimplemented, codes.ResourceExhausted: - return true - } - return false + return ok && st.Code() == codes.Unimplemented } func retryableStatus(code int) bool { diff --git a/core/services/failover/classify_test.go b/core/services/failover/classify_test.go index baafff337..c19fd0aba 100644 --- a/core/services/failover/classify_test.go +++ b/core/services/failover/classify_test.go @@ -32,6 +32,7 @@ var _ = DescribeTable("IsRetryable", Entry("grpc deadline", grpcstatus.Error(codes.DeadlineExceeded, "x"), 0, true), Entry("grpc unknown", grpcstatus.Error(codes.Unknown, "x"), 0, true), Entry("grpc invalid argument", grpcstatus.Error(codes.InvalidArgument, "x"), 0, false), + Entry("grpc resource exhausted (rate limit, OOM)", grpcstatus.Error(codes.ResourceExhausted, "x"), 0, true), Entry("cloud-proxy upstream 503", errors.New("cloud-proxy: upstream 503: no healthy nodes"), 0, true), Entry("cloud-proxy upstream 429 stays 4xx", errors.New("cloud-proxy: upstream 429: slow down"), 0, false), Entry("context overflow", errors.New("the request exceeds the available context size"), 0, false), @@ -45,8 +46,7 @@ var _ = DescribeTable("IsCapabilityGap", Entry("nil error", nil, false), Entry("grpc unimplemented", grpcstatus.Error(codes.Unimplemented, "x"), true), Entry("wrapped grpc unimplemented", fmt.Errorf("call: %w", grpcstatus.Error(codes.Unimplemented, "x")), true), - Entry("grpc resource exhausted", grpcstatus.Error(codes.ResourceExhausted, "x"), true), - Entry("wrapped grpc resource exhausted", fmt.Errorf("call: %w", grpcstatus.Error(codes.ResourceExhausted, "x")), true), + Entry("grpc resource exhausted is a failure, not a gap", grpcstatus.Error(codes.ResourceExhausted, "x"), false), Entry("grpc unavailable", grpcstatus.Error(codes.Unavailable, "x"), false), Entry("plain error", errors.New("boom"), false), ) diff --git a/core/services/failover/manager.go b/core/services/failover/manager.go index 147b95b34..071e9ab0c 100644 --- a/core/services/failover/manager.go +++ b/core/services/failover/manager.go @@ -728,9 +728,8 @@ func (m *Manager) Do(ctx context.Context, chain string, fn func(ctx context.Cont att.Succeed() return nil case IsCapabilityGap(err) && !committed.Load(): - // This target cannot serve this kind of request, or is out of - // capacity; it is not broken, so move on without counting a - // failure. + // This target cannot serve this kind of request at all; it is + // not broken, so move on without counting a failure. if !att.Skip() { return err } diff --git a/core/services/failover/manager_test.go b/core/services/failover/manager_test.go index 948f5aef0..75a773f04 100644 --- a/core/services/failover/manager_test.go +++ b/core/services/failover/manager_test.go @@ -306,7 +306,7 @@ var _ = Describe("Manager", func() { Expect(st.Targets[0].State).To(Equal(StateHealthy)) }) - It("skips a rate-limited (ResourceExhausted) target without tripping it", func() { + It("fails a rate-limited (ResourceExhausted) target over and trips it", func() { var tried []string err := m.Do(context.Background(), "chain", func(_ context.Context, target string, _ func()) error { tried = append(tried, target) @@ -318,7 +318,7 @@ var _ = Describe("Manager", func() { Expect(err).ToNot(HaveOccurred()) Expect(tried).To(Equal([]string{"a", "b"})) st, _ := m.ChainStatus("chain") - Expect(st.Targets[0].State).To(Equal(StateHealthy)) + Expect(st.Targets[0].State).To(Equal(StateDown)) }) }) })