mirror of
https://github.com/mudler/LocalAI.git
synced 2026-09-29 17:44:30 -04:00
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 <mudler@localai.io>
This commit is contained in:
1 parent
6df1767133
commit
882f51fd62
7 files changed
+26
-30
No files matched your search
@@ -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
|
||||
|
||||
@@ -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"})
|
||||
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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),
|
||||
)
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
})
|
||||
})
|
||||
})
|
||||
Reference in new issue
Block a user