From c541dbeef40786ac004d48bb993b79f13a945ea7 Mon Sep 17 00:00:00 2001 From: Ettore Di Giacinto Date: Sun, 23 Aug 2026 20:44:43 +0000 Subject: [PATCH] fix(distributed): check a node answers before scheduling onto it A node's status comes from its HTTP heartbeat. Backend installs travel over NATS. The two are independent, so a worker that dies stops answering on the bus at once but stays healthy in the database until its heartbeat ages out. Inside that window the scheduler picked a node it could not reach, and the request failed with "no responders available" rather than moving to a node that was up. The scheduler now probes the node it selected and, when nothing answers, marks it unhealthy and selects again. The demotion is what makes the retry terminate: the next selection reads only healthy nodes. It also tells the other frontends what this one learned, so the cluster does not rediscover a dead worker one failed request at a time. Only nats.ErrNoResponders counts as absent. A worker that answers slowly stays eligible, because dropping it would cost capacity that is really there. The probe reuses the models.running subject: a new subject would go unanswered by workers that have not been upgraded, and every one of them would then look dead. Signed-off-by: Ettore Di Giacinto Assisted-by: Claude Code:claude-opus-5 [golangci-lint] --- core/services/nodes/interfaces.go | 1 + core/services/nodes/model_router_test.go | 4 + core/services/nodes/router.go | 56 +++++---- core/services/nodes/router_liveness.go | 60 +++++++++ .../nodes/router_nats_liveness_test.go | 118 ++++++++++++++++++ core/services/nodes/router_test.go | 31 ++++- core/services/nodes/unloader.go | 25 ++++ docs/content/features/distributed-mode.md | 6 + 8 files changed, 277 insertions(+), 24 deletions(-) create mode 100644 core/services/nodes/router_liveness.go create mode 100644 core/services/nodes/router_nats_liveness_test.go diff --git a/core/services/nodes/interfaces.go b/core/services/nodes/interfaces.go index c204752de..93399bc14 100644 --- a/core/services/nodes/interfaces.go +++ b/core/services/nodes/interfaces.go @@ -60,6 +60,7 @@ type ModelRouter interface { GetNodeLabels(ctx context.Context, nodeID string) ([]NodeLabel, error) FindNodesWithModel(ctx context.Context, modelName string) ([]BackendNode, error) LoadedReplicaStats(ctx context.Context, modelName string, candidateNodeIDs []string) ([]ReplicaCandidate, error) + MarkUnhealthy(ctx context.Context, nodeID string) error LoadJobStore } diff --git a/core/services/nodes/model_router_test.go b/core/services/nodes/model_router_test.go index 43002006a..13e36b12d 100644 --- a/core/services/nodes/model_router_test.go +++ b/core/services/nodes/model_router_test.go @@ -250,3 +250,7 @@ var _ = Describe("ModelRouterAdapter", func() { }) }) }) + +func (f *fakeModelRouterForSmartRouter) MarkUnhealthy(_ context.Context, _ string) error { + return nil +} diff --git a/core/services/nodes/router.go b/core/services/nodes/router.go index 465b2d44e..14b19092e 100644 --- a/core/services/nodes/router.go +++ b/core/services/nodes/router.go @@ -1089,34 +1089,44 @@ func (r *SmartRouter) scheduleNewModel(ctx context.Context, backendType, modelID // If freeSlotNodes is empty (everyone full), candidateNodeIDs is whatever // it was — we'll fall through to eviction below. - var node *BackendNode - - if estimatedVRAM > 0 { - if candidateNodeIDs != nil { - node, err = r.registry.FindNodeWithVRAMFromSet(ctx, estimatedVRAM, candidateNodeIDs) - } else { - node, err = r.registry.FindNodeWithVRAM(ctx, estimatedVRAM) - } - if err != nil { - xlog.Warn("No nodes with enough VRAM, falling back to standard scheduling", - "required_vram", vram.FormatBytes(estimatedVRAM), "error", err) - } - } - - if node == nil { - if candidateNodeIDs != nil { - node, err = r.registry.FindIdleNodeFromSet(ctx, candidateNodeIDs) - if err != nil { - node, err = r.registry.FindLeastLoadedNodeFromSet(ctx, candidateNodeIDs) + // Node choice is wrapped in a liveness check: a node's stored status comes + // from its HTTP heartbeat, which is a different channel from the bus that + // carries the install. A worker that has died stops answering on the bus at + // once but stays healthy in the database until its heartbeat ages out, so + // without this the scheduler could commit to a node it cannot reach. + selectNode := func() *BackendNode { + var candidate *BackendNode + var selErr error + if estimatedVRAM > 0 { + if candidateNodeIDs != nil { + candidate, selErr = r.registry.FindNodeWithVRAMFromSet(ctx, estimatedVRAM, candidateNodeIDs) + } else { + candidate, selErr = r.registry.FindNodeWithVRAM(ctx, estimatedVRAM) } - } else { - node, err = r.registry.FindIdleNode(ctx) - if err != nil { - node, err = r.registry.FindLeastLoadedNode(ctx) + if selErr != nil { + xlog.Warn("No nodes with enough VRAM, falling back to standard scheduling", + "required_vram", vram.FormatBytes(estimatedVRAM), "error", selErr) } } + + if candidate == nil { + if candidateNodeIDs != nil { + candidate, selErr = r.registry.FindIdleNodeFromSet(ctx, candidateNodeIDs) + if selErr != nil { + candidate, _ = r.registry.FindLeastLoadedNodeFromSet(ctx, candidateNodeIDs) + } + } else { + candidate, selErr = r.registry.FindIdleNode(ctx) + if selErr != nil { + candidate, _ = r.registry.FindLeastLoadedNode(ctx) + } + } + } + return candidate } + node := r.pickReachableNode(ctx, selectNode) + // 4. Preemptive eviction: if no suitable node found, evict the LRU model with zero in-flight if node == nil { evictedNode, evictErr := r.evictLRUAndFreeNode(ctx) diff --git a/core/services/nodes/router_liveness.go b/core/services/nodes/router_liveness.go new file mode 100644 index 000000000..88646162f --- /dev/null +++ b/core/services/nodes/router_liveness.go @@ -0,0 +1,60 @@ +package nodes + +import ( + "context" + "errors" + + "github.com/mudler/xlog" + "github.com/nats-io/nats.go" +) + +// maxNodeLivenessRetries bounds how many unreachable nodes a single scheduling +// attempt discards before giving up. Each discarded node is marked unhealthy, +// so the bound only has to cover one burst of dead workers rather than the +// whole fleet. +const maxNodeLivenessRetries = 3 + +// nodeAnswersOnBus reports whether a node still has a live subscription. +// +// Only nats.ErrNoResponders means "absent". Any other outcome, a timeout or a +// transport hiccup, leaves the node eligible: wrongly excluding a node that is +// merely slow costs real capacity, while the install that follows already +// reports its own failure. When no command sender is configured there is no bus +// to consult and every node is treated as reachable, which preserves the +// behaviour of deployments that do not run one. +func (r *SmartRouter) nodeAnswersOnBus(node *BackendNode) bool { + if r.unloader == nil || node == nil { + return true + } + err := r.unloader.PingNode(node.ID) + return !errors.Is(err, nats.ErrNoResponders) +} + +// pickReachableNode calls selectNode until it yields a node that still answers +// on the bus, and returns nil when it cannot find one. +// +// A node that does not answer is marked unhealthy before the next attempt. That +// both removes it from the next selection, which queries only healthy nodes, +// and tells every other scheduler in the cluster what this one just learned, so +// the discovery is not repeated one failed request at a time. +func (r *SmartRouter) pickReachableNode(ctx context.Context, selectNode func() *BackendNode) *BackendNode { + for range maxNodeLivenessRetries { + node := selectNode() + if node == nil { + return nil + } + if r.nodeAnswersOnBus(node) { + return node + } + xlog.Warn("Scheduled node is not answering on the bus, marking unhealthy and re-scheduling", + "node", node.Name, "nodeID", node.ID) + if err := r.registry.MarkUnhealthy(ctx, node.ID); err != nil { + // Without the demotion the next selection would hand back the same + // node, so stop rather than spin. + xlog.Warn("Failed to mark unreachable node unhealthy", + "node", node.Name, "nodeID", node.ID, "error", err) + return nil + } + } + return nil +} diff --git a/core/services/nodes/router_nats_liveness_test.go b/core/services/nodes/router_nats_liveness_test.go new file mode 100644 index 000000000..ec4820c9f --- /dev/null +++ b/core/services/nodes/router_nats_liveness_test.go @@ -0,0 +1,118 @@ +package nodes + +import ( + "context" + "errors" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +// A node's stored status comes from its HTTP heartbeat, but work is dispatched +// over NATS. A worker that dies stops answering on the bus immediately and +// keeps its healthy status until the heartbeat ages out, so the scheduler could +// commit to a node it could not reach. The request then failed outright with +// "no responders available" rather than moving to a node that was actually up. +var _ = Describe("Scheduling past a node that left the bus", func() { + var ( + reg *fakeModelRouter + fake *fakeUnloader + router *SmartRouter + ) + + newNode := func(id string) *BackendNode { + return &BackendNode{ID: id, Name: id, Address: id + ":50051"} + } + + // selectorReturning hands back each node in turn, mimicking a scheduler + // that re-picks after the previous choice was demoted. + selectorReturning := func(nodes ...*BackendNode) func() *BackendNode { + i := 0 + return func() *BackendNode { + if i >= len(nodes) { + return nil + } + n := nodes[i] + i++ + return n + } + } + + BeforeEach(func() { + reg = &fakeModelRouter{} + fake = &fakeUnloader{deadNodes: map[string]bool{}} + router = NewSmartRouter(reg, SmartRouterOptions{Unloader: fake}) + }) + + It("passes over a node that no longer answers and takes one that does", func() { + dead, alive := newNode("dead-node"), newNode("alive-node") + fake.deadNodes["dead-node"] = true + + picked := router.pickReachableNode(context.Background(), selectorReturning(dead, alive)) + + Expect(picked).ToNot(BeNil()) + Expect(picked.ID).To(Equal("alive-node")) + Expect(fake.pingCalls).To(Equal([]string{"dead-node", "alive-node"})) + }) + + It("demotes the absent node so other schedulers stop choosing it", func() { + dead, alive := newNode("dead-node"), newNode("alive-node") + fake.deadNodes["dead-node"] = true + + router.pickReachableNode(context.Background(), selectorReturning(dead, alive)) + + Expect(reg.markedUnhealthy).To(Equal([]string{"dead-node"})) + }) + + It("takes the first node when it answers, without probing further", func() { + first, second := newNode("first"), newNode("second") + + picked := router.pickReachableNode(context.Background(), selectorReturning(first, second)) + + Expect(picked.ID).To(Equal("first")) + Expect(fake.pingCalls).To(Equal([]string{"first"})) + }) + + It("gives up rather than spinning when every node is gone", func() { + a, b, c, d := newNode("a"), newNode("b"), newNode("c"), newNode("d") + for _, id := range []string{"a", "b", "c", "d"} { + fake.deadNodes[id] = true + } + + picked := router.pickReachableNode(context.Background(), selectorReturning(a, b, c, d)) + + Expect(picked).To(BeNil()) + Expect(len(fake.pingCalls)).To(BeNumerically("<=", maxNodeLivenessRetries)) + }) + + It("stops when the demotion itself fails, so it cannot loop on one node", func() { + dead := newNode("dead-node") + fake.deadNodes["dead-node"] = true + reg.markUnhealthyErr = errors.New("database is down") + + picked := router.pickReachableNode(context.Background(), selectorReturning(dead, dead, dead)) + + Expect(picked).To(BeNil()) + Expect(fake.pingCalls).To(Equal([]string{"dead-node"})) + }) + + // Only a no-responders answer proves absence. Excluding a node that is + // merely slow would cost real capacity. + It("keeps a node that answers slowly or errors for another reason", func() { + slow := newNode("slow-node") + fake.pingErr = errors.New("timeout waiting for reply") + + picked := router.pickReachableNode(context.Background(), selectorReturning(slow)) + + Expect(picked).ToNot(BeNil()) + Expect(picked.ID).To(Equal("slow-node")) + Expect(reg.markedUnhealthy).To(BeEmpty()) + }) + + It("treats every node as reachable when no command sender is configured", func() { + plain := NewSmartRouter(reg, SmartRouterOptions{}) + node := newNode("only-node") + + Expect(plain.pickReachableNode(context.Background(), selectorReturning(node))).To(Equal(node)) + }) +}) diff --git a/core/services/nodes/router_test.go b/core/services/nodes/router_test.go index 96db9b93f..0f0938335 100644 --- a/core/services/nodes/router_test.go +++ b/core/services/nodes/router_test.go @@ -17,6 +17,7 @@ import ( "github.com/mudler/LocalAI/pkg/distributedhdr" grpc "github.com/mudler/LocalAI/pkg/grpc" pb "github.com/mudler/LocalAI/pkg/grpc/proto" + "github.com/nats-io/nats.go" ggrpc "google.golang.org/grpc" "google.golang.org/protobuf/proto" "gorm.io/gorm" @@ -61,6 +62,10 @@ func (f *fakeFileStager) ListRemoteDir(_ context.Context, _, _ string) ([]string // fakeModelRouter implements ModelRouter with configurable return values. type fakeModelRouter struct { + // markedUnhealthy records nodes demoted by the scheduler's liveness check. + markedUnhealthy []string + markUnhealthyErr error + fakeLoadJobStore // FindAndLockNodeWithModel returns @@ -474,7 +479,15 @@ type fakeUnloader struct { stopCalls []string // "nodeID:model" stopErr error unloadCalls []string - unloadErr error + + // deadNodes names the nodes PingNode reports as absent from the bus, and + // pingCalls records every node it was asked about, in order. + deadNodes map[string]bool + pingCalls []string + // pingErr is returned for nodes not in deadNodes, so a spec can model a + // node that is reachable but answering badly. + pingErr error + unloadErr error } // installCall captures the args we care about when asserting that the @@ -532,6 +545,22 @@ func (f *fakeUnloader) UnloadModelOnNode(nodeID, modelName string) error { return f.unloadErr } +func (f *fakeModelRouter) MarkUnhealthy(_ context.Context, nodeID string) error { + f.markedUnhealthy = append(f.markedUnhealthy, nodeID) + return f.markUnhealthyErr +} + +func (f *fakeUnloader) PingNode(nodeID string) error { + f.mu.Lock() + f.pingCalls = append(f.pingCalls, nodeID) + dead := f.deadNodes[nodeID] + f.mu.Unlock() + if dead { + return nats.ErrNoResponders + } + return f.pingErr +} + // --------------------------------------------------------------------------- // Tests // --------------------------------------------------------------------------- diff --git a/core/services/nodes/unloader.go b/core/services/nodes/unloader.go index 8d47d71a6..3b1cd15b8 100644 --- a/core/services/nodes/unloader.go +++ b/core/services/nodes/unloader.go @@ -36,6 +36,10 @@ type NodeCommandSender interface { ListBackends(nodeID string) (*messaging.BackendListReply, error) StopBackend(nodeID, backend string) error UnloadModelOnNode(nodeID, modelName string) error + // PingNode reports whether the node is still subscribed on the bus. It + // returns nats.ErrNoResponders when nothing answers for the node, which is + // the only condition callers may read as "this node cannot be given work". + PingNode(nodeID string) error } // RemoteUnloaderAdapter implements NodeCommandSender and model.RemoteModelUnloader @@ -360,6 +364,27 @@ func (a *RemoteUnloaderAdapter) ListBackends(nodeID string) (*messaging.BackendL return messaging.RequestJSON[messaging.BackendListRequest, messaging.BackendListReply](a.nats, subject, messaging.BackendListRequest{}, 30*time.Second) } +// PingNode checks that a worker still has a live subscription on the bus. +// +// A node's status in the database comes from its HTTP heartbeat, which is a +// separate channel from NATS. A worker that has died stops answering on NATS +// at once but keeps its healthy status until the heartbeat ages out, so the +// scheduler could pick a node that could not be given work and the request +// failed with "no responders available". +// +// It reuses the models.running subject rather than a dedicated ping subject on +// purpose: a new subject would go unanswered by any worker that has not been +// upgraded yet, and this check would then report every one of them as dead. +// The worker answers out of its in-memory process table, so a live node +// replies immediately, and NATS reports no-responders without waiting out the +// timeout. +func (a *RemoteUnloaderAdapter) PingNode(nodeID string) error { + subject := messaging.SubjectNodeModelsRunning(nodeID) + _, err := messaging.RequestJSON[messaging.ModelsRunningRequest, messaging.ModelsRunningReply]( + a.nats, subject, messaging.ModelsRunningRequest{}, 5*time.Second) + return err +} + // ListRunningModels asks a worker node which model backend processes it // currently has running, via NATS request-reply. // diff --git a/docs/content/features/distributed-mode.md b/docs/content/features/distributed-mode.md index 54c3abbab..667408fbb 100644 --- a/docs/content/features/distributed-mode.md +++ b/docs/content/features/distributed-mode.md @@ -1020,6 +1020,12 @@ Notes: - Upgrade the worker when it does not support the exact model-stop request. - Stop and restart the stale backend only as an operational recovery action. LocalAI keeps it non-routable while durable cleanup is pending. +**A request fails with `nats: no responders available for request`:** +- The chosen worker was not subscribed on the bus when the frontend tried to install the backend on it. A node's status comes from its HTTP heartbeat, which is a separate channel: a worker that stops stays `healthy` until that heartbeat ages out. +- The scheduler now checks that a node still answers on the bus before it commits to it, marks one that does not as unhealthy, and picks another. A request should therefore see this only when no reachable node is left. +- Only a no-responders answer counts as absent. A worker that answers slowly stays eligible, because excluding it would cost capacity that is really there. +- Check the worker process is running and its NATS connection is up. `Scheduled node is not answering on the bus` in the frontend log names each node demoted this way. + **A worker fills its own disk over time:** - A request that carries a file (an image, an audio clip, a video) stages that file to the worker under `/../staging/ephemeral/`. The worker deletes these 6 hours after the request that needed them, and sweeps every 30 minutes plus once at startup, so a worker that crashed mid-request still reclaims the space. - Releases before this sweep existed kept every staged input for the lifetime of the worker. Delete `/../staging/ephemeral/` on an affected worker once, as the user the worker runs as; the sweep keeps it bounded from then on.