Skip to content

Commit 4843d57

Browse files
committed
fix(distributed): drop replica rows when a worker stops a backend
A worker returns a stopped backend's gRPC port to its allocator as soon as the process is confirmed dead, and hands it to the next backend that starts. The controller's NodeModel row for the old address survives, and both SmartRouter.probeHealth and the HealthMonitor per-model probe verify liveness, not identity, so once an unrelated backend binds the recycled port the stale row passes every check and the request is served by the wrong backend instead of failing. backend.delete is newly able to trigger this: before #10956 a delete never actually stopped a process, so it never recycled a port. backend.upgrade has the identical gap and always did — upgradeBackend force-stops every process using the binary and starts none back up, while DistributedBackendManager.UpgradeBackend never removes rows. model.unload is the one path that gets this right today: it calls RemoveAllNodeModelReplicas straight after StopBackend. Report the process keys the worker terminated on the delete and upgrade replies, and drop the matching rows in RemoteUnloaderAdapter, which already holds a ModelLocator with RemoveNodeModel. All three call sites funnel through that adapter, so no new interface, DB migration, or proto change is needed. A key is reported only once its process is confirmed gone, so the list stays trustworthy on the partial-failure replies too. Old workers never populate the new fields. ReportsStoppedProcesses tells "stopped nothing" apart from "does not report", so an old worker's silence falls back to the pre-existing probe-based staleness recovery instead of being mistaken for a completed cleanup. Quarantine released ports for a short window as an interlock covering the NATS round-trip between the worker freeing the port and the controller dropping the row. It is deliberately not derived from HealthCheckInterval: that cadence is operator-tunable and the per-model reaper can be disabled outright, so coupling a worker-local constant to it would be silently wrong on some clusters. Eager row removal is the fix; the delay only closes the handoff gap. Identity verification in probeHealth was considered and rejected: Health and Status carry no backend identity, so it needs a proto change plus an implementation in 36 Python and 4 C++ Health servicers, it is fail-open for any backend not yet rebuilt, and the probeCache short-circuit means it would not even execute during the 30s window where the misroute happens. Fixes #10952 Refs #10954, #10956 Assisted-by: Claude:claude-opus-4-8 golangci-lint Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
1 parent 2b7a049 commit 4843d57

11 files changed

Lines changed: 613 additions & 39 deletions

File tree

core/services/messaging/subjects.go

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -226,6 +226,14 @@ type BackendUpgradeRequest struct {
226226
type BackendUpgradeReply struct {
227227
Success bool `json:"success"`
228228
Error string `json:"error,omitempty"`
229+
230+
// StoppedProcessKeys / ReportsStoppedProcesses carry the same
231+
// stale-row-invalidation contract as on BackendDeleteReply; an upgrade
232+
// force-stops every process using the binary and starts none back up, so it
233+
// recycles ports exactly the way a delete does. See that type for why the
234+
// boolean is not redundant with an empty list.
235+
StoppedProcessKeys []string `json:"stopped_process_keys,omitempty"`
236+
ReportsStoppedProcesses bool `json:"reports_stopped_processes,omitempty"`
229237
}
230238

231239
// SubjectNodeBackendList queries a worker node for its installed backends.
@@ -290,6 +298,23 @@ type BackendDeleteRequest struct {
290298
type BackendDeleteReply struct {
291299
Success bool `json:"success"`
292300
Error string `json:"error,omitempty"`
301+
302+
// StoppedProcessKeys names every `modelID#replica` process the worker
303+
// terminated while serving this delete. Stopping a process returns its gRPC
304+
// port to the worker's allocator, so any NodeModel row still pointing at
305+
// that address becomes a live misroute the moment an unrelated backend
306+
// binds the recycled port: probeHealth verifies liveness, not identity, so
307+
// the request is served by the wrong backend rather than failing. The
308+
// controller uses these keys to drop the rows eagerly.
309+
StoppedProcessKeys []string `json:"stopped_process_keys,omitempty"`
310+
311+
// ReportsStoppedProcesses distinguishes "this worker enumerates what it
312+
// stopped and stopped nothing" from "this worker predates the field". Both
313+
// send an empty list, and only the first is authoritative. Without this
314+
// flag a controller cannot tell them apart and would eventually be tempted
315+
// to read silence as a completed cleanup, which is precisely the wrong
316+
// conclusion against an older worker.
317+
ReportsStoppedProcesses bool `json:"reports_stopped_processes,omitempty"`
293318
}
294319

295320
// SubjectNodeModelUnload tells a worker node to unload a model (gRPC Free) without killing the backend.

core/services/nodes/unloader.go

Lines changed: 52 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -270,6 +270,9 @@ func (a *RemoteUnloaderAdapter) UpgradeBackend(nodeID, backendType, galleriesJSO
270270
return nil, fmt.Errorf("%w (subject=%s nodeID=%s backend=%s): %v",
271271
galleryop.ErrWorkerStillInstalling, subject, nodeID, backendType, err)
272272
}
273+
if err == nil {
274+
a.dropStoppedReplicaRows(nodeID, "backend.upgrade", backendType, reply.StoppedProcessKeys, reply.ReportsStoppedProcesses)
275+
}
273276
return reply, err
274277
}
275278

@@ -337,7 +340,55 @@ func (a *RemoteUnloaderAdapter) DeleteBackend(nodeID, backendName string) (*mess
337340
subject := messaging.SubjectNodeBackendDelete(nodeID)
338341
xlog.Info("Sending NATS backend.delete", "nodeID", nodeID, "backend", backendName)
339342

340-
return messaging.RequestJSON[messaging.BackendDeleteRequest, messaging.BackendDeleteReply](a.nats, subject, messaging.BackendDeleteRequest{Backend: backendName}, 2*time.Minute)
343+
reply, err := messaging.RequestJSON[messaging.BackendDeleteRequest, messaging.BackendDeleteReply](a.nats, subject, messaging.BackendDeleteRequest{Backend: backendName}, 2*time.Minute)
344+
if err != nil {
345+
return reply, err
346+
}
347+
a.dropStoppedReplicaRows(nodeID, "backend.delete", backendName, reply.StoppedProcessKeys, reply.ReportsStoppedProcesses)
348+
return reply, nil
349+
}
350+
351+
// dropStoppedReplicaRows removes the NodeModel rows addressing processes a
352+
// worker just terminated.
353+
//
354+
// Why eagerly, rather than leaving it to the existing health checks: stopping a
355+
// process returns its gRPC port to the worker's allocator, and the next backend
356+
// started there can be handed that same port. Until the row is gone it names a
357+
// live address, so both SmartRouter.probeHealth and the HealthMonitor per-model
358+
// probe — which verify liveness, not identity — pass against whatever now
359+
// occupies the port, and the request is served by the wrong backend rather than
360+
// failing. Nothing else on the delete/upgrade path tells the controller the
361+
// address just became invalid, unlike model.unload which drops its rows itself.
362+
//
363+
// reported=false means the worker predates this reply field. Its empty list is
364+
// then indistinguishable from "stopped nothing", so it must NOT be read as a
365+
// completed cleanup: leave the rows alone and fall back to the probe-based
366+
// staleness recovery that was the only mechanism before this change.
367+
func (a *RemoteUnloaderAdapter) dropStoppedReplicaRows(nodeID, op, backendName string, processKeys []string, reported bool) {
368+
if !reported {
369+
xlog.Debug("Worker did not report stopped processes; relying on probe-based staleness recovery",
370+
"nodeID", nodeID, "op", op, "backend", backendName)
371+
return
372+
}
373+
ctx := context.Background()
374+
for _, key := range processKeys {
375+
modelName, replicaIndex, ok := model.ParseBackendProcessKey(key)
376+
if !ok {
377+
// Acting on a guess could evict the row of a healthy sibling replica.
378+
xlog.Warn("Ignoring unparseable process key reported by worker",
379+
"nodeID", nodeID, "op", op, "backend", backendName, "processKey", key)
380+
continue
381+
}
382+
xlog.Info("Dropping replica row for a process the worker stopped",
383+
"nodeID", nodeID, "op", op, "backend", backendName, "model", modelName, "replica", replicaIndex)
384+
if err := a.registry.RemoveNodeModel(ctx, nodeID, modelName, replicaIndex); err != nil {
385+
// Best-effort: probe-based recovery remains the backstop, and failing
386+
// the operator's delete over a bookkeeping error would be worse than
387+
// the stale row this prevents.
388+
xlog.Warn("Failed to drop replica row for a stopped process",
389+
"nodeID", nodeID, "op", op, "model", modelName, "replica", replicaIndex, "error", err)
390+
}
391+
}
341392
}
342393

343394
// UnloadModelOnNode sends a model.unload request to a specific node.
Lines changed: 186 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,186 @@
1+
package nodes
2+
3+
import (
4+
"encoding/json"
5+
"time"
6+
7+
. "github.com/onsi/ginkgo/v2"
8+
. "github.com/onsi/gomega"
9+
10+
"github.com/mudler/LocalAI/core/services/messaging"
11+
)
12+
13+
// Replies are handed to the adapter as raw JSON rather than as marshalled
14+
// structs on purpose: these specs pin the behavior a controller must exhibit
15+
// when a worker on a *different* build answers it. Marshalling the current
16+
// struct would only ever produce the current field set and could not express
17+
// an old worker's reply at all.
18+
var _ = Describe("RemoteUnloaderAdapter stale replica rows", func() {
19+
var (
20+
locator *fakeModelLocator
21+
mc *fakeMessagingClient
22+
adapter *RemoteUnloaderAdapter
23+
)
24+
25+
BeforeEach(func() {
26+
locator = &fakeModelLocator{}
27+
mc = &fakeMessagingClient{}
28+
adapter = NewRemoteUnloaderAdapter(locator, mc, 3*time.Minute, 15*time.Minute)
29+
})
30+
31+
Describe("DeleteBackend", func() {
32+
It("drops the replica rows for every process the worker stopped", func() {
33+
// The worker has just terminated these processes and returned their
34+
// gRPC ports to its allocator. Any row still pointing at those
35+
// addresses will pass probeHealth as soon as an unrelated backend
36+
// binds the recycled port, and the request is then silently served
37+
// by the wrong backend.
38+
mc.requestReply = []byte(`{
39+
"success": true,
40+
"reports_stopped_processes": true,
41+
"stopped_process_keys": ["qwen3-0.6b#0", "qwen3-0.6b#2"]
42+
}`)
43+
44+
reply, err := adapter.DeleteBackend("node-1", "llama-cpp")
45+
Expect(err).NotTo(HaveOccurred())
46+
Expect(reply.Success).To(BeTrue())
47+
48+
Expect(locator.removedReplicas).To(ConsistOf(
49+
modelReplicaRef{"node-1", "qwen3-0.6b", 0},
50+
modelReplicaRef{"node-1", "qwen3-0.6b", 2},
51+
))
52+
})
53+
54+
It("keeps model IDs that contain the replica separator intact", func() {
55+
// Process keys are `modelID#replica` and model IDs are free to
56+
// contain '#' themselves, so the split must be anchored at the last
57+
// separator. Splitting at the first one addresses a row that does
58+
// not exist and leaves the real stale row in place.
59+
mc.requestReply = []byte(`{
60+
"success": true,
61+
"reports_stopped_processes": true,
62+
"stopped_process_keys": ["weird#name#3"]
63+
}`)
64+
65+
_, err := adapter.DeleteBackend("node-1", "llama-cpp")
66+
Expect(err).NotTo(HaveOccurred())
67+
Expect(locator.removedReplicas).To(ConsistOf(modelReplicaRef{"node-1", "weird#name", 3}))
68+
})
69+
70+
It("removes nothing when a worker that reports stopped processes stopped none", func() {
71+
// Deleting a backend that was never loaded is routine. The reply is
72+
// authoritative here, so "no keys" genuinely means "no rows".
73+
mc.requestReply = []byte(`{"success": true, "reports_stopped_processes": true}`)
74+
75+
_, err := adapter.DeleteBackend("node-1", "llama-cpp")
76+
Expect(err).NotTo(HaveOccurred())
77+
Expect(locator.removedReplicas).To(BeEmpty())
78+
})
79+
80+
It("does not read an old worker's silence as a completed cleanup", func() {
81+
// A worker predating this change never sets either field. Its empty
82+
// list is indistinguishable from "stopped nothing", so the
83+
// controller must not conclude the node is clean, and equally must
84+
// not guess at rows to delete. It falls back to the probe-based
85+
// self-heal in SmartRouter.probeHealth, which is exactly the
86+
// pre-change behavior.
87+
mc.requestReply = []byte(`{"success": true}`)
88+
89+
reply, err := adapter.DeleteBackend("node-1", "llama-cpp")
90+
Expect(err).NotTo(HaveOccurred())
91+
Expect(reply.Success).To(BeTrue())
92+
Expect(reply.ReportsStoppedProcesses).To(BeFalse())
93+
Expect(locator.removedReplicas).To(BeEmpty())
94+
Expect(locator.removedPairs).To(BeEmpty())
95+
})
96+
97+
It("leaves rows alone when the worker could not stop the process", func() {
98+
// The worker aborts the delete without listing the key it failed to
99+
// kill: that process survived, so its address is still correct and
100+
// dropping the row would force a needless reload of a live replica.
101+
mc.requestReply = []byte(`{
102+
"success": false,
103+
"error": "could not stop running process qwen3-0.6b#0",
104+
"reports_stopped_processes": true
105+
}`)
106+
107+
reply, err := adapter.DeleteBackend("node-1", "llama-cpp")
108+
Expect(err).NotTo(HaveOccurred())
109+
Expect(reply.Success).To(BeFalse())
110+
Expect(locator.removedReplicas).To(BeEmpty())
111+
})
112+
113+
It("drops rows for processes that died before a later step failed", func() {
114+
// The worker reports a key only once that process is confirmed gone,
115+
// so a failure further along (removing files, re-registering) does
116+
// not make the already-recycled ports any less dangerous. Gating
117+
// removal on overall success would strand exactly those rows.
118+
mc.requestReply = []byte(`{
119+
"success": false,
120+
"error": "failed to delete backend files",
121+
"reports_stopped_processes": true,
122+
"stopped_process_keys": ["qwen3-0.6b#0"]
123+
}`)
124+
125+
_, err := adapter.DeleteBackend("node-1", "llama-cpp")
126+
Expect(err).NotTo(HaveOccurred())
127+
Expect(locator.removedReplicas).To(ConsistOf(modelReplicaRef{"node-1", "qwen3-0.6b", 0}))
128+
})
129+
130+
It("ignores malformed process keys instead of removing a wrong row", func() {
131+
mc.requestReply = []byte(`{
132+
"success": true,
133+
"reports_stopped_processes": true,
134+
"stopped_process_keys": ["no-replica-suffix", "qwen3-0.6b#notanumber", "good#1"]
135+
}`)
136+
137+
_, err := adapter.DeleteBackend("node-1", "llama-cpp")
138+
Expect(err).NotTo(HaveOccurred())
139+
Expect(locator.removedReplicas).To(ConsistOf(modelReplicaRef{"node-1", "good", 1}))
140+
})
141+
})
142+
143+
Describe("UpgradeBackend", func() {
144+
It("drops the replica rows for every process the worker stopped", func() {
145+
// An upgrade force-stops every process using the binary and starts
146+
// none of them back up, so it recycles ports exactly as delete does
147+
// while leaving the same rows behind.
148+
mc.requestReply = []byte(`{
149+
"success": true,
150+
"reports_stopped_processes": true,
151+
"stopped_process_keys": ["whisper#0", "whisper#1"]
152+
}`)
153+
154+
reply, err := adapter.UpgradeBackend("node-1", "whisper", "", "", "", "", 0, "", nil)
155+
Expect(err).NotTo(HaveOccurred())
156+
Expect(reply.Success).To(BeTrue())
157+
158+
Expect(locator.removedReplicas).To(ConsistOf(
159+
modelReplicaRef{"node-1", "whisper", 0},
160+
modelReplicaRef{"node-1", "whisper", 1},
161+
))
162+
})
163+
164+
It("does not read an old worker's silence as a completed cleanup", func() {
165+
mc.requestReply = []byte(`{"success": true}`)
166+
167+
reply, err := adapter.UpgradeBackend("node-1", "whisper", "", "", "", "", 0, "", nil)
168+
Expect(err).NotTo(HaveOccurred())
169+
Expect(reply.ReportsStoppedProcesses).To(BeFalse())
170+
Expect(locator.removedReplicas).To(BeEmpty())
171+
})
172+
})
173+
174+
Describe("reply wire compatibility", func() {
175+
It("decodes an old worker's delete reply without the new fields", func() {
176+
// Guards the rolling-upgrade direction that matters: a new
177+
// controller must keep working against every worker already
178+
// deployed, not just ones rebuilt from this commit.
179+
var reply messaging.BackendDeleteReply
180+
Expect(json.Unmarshal([]byte(`{"success": true}`), &reply)).To(Succeed())
181+
Expect(reply.Success).To(BeTrue())
182+
Expect(reply.ReportsStoppedProcesses).To(BeFalse())
183+
Expect(reply.StoppedProcessKeys).To(BeEmpty())
184+
})
185+
})
186+
})

core/services/nodes/unloader_test.go

Lines changed: 17 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -20,22 +20,35 @@ import (
2020

2121
// fakeModelLocator implements ModelLocator with configurable node lists.
2222
type fakeModelLocator struct {
23-
nodes []BackendNode
24-
findErr error
25-
removedPairs []modelNodePair // records RemoveNodeModel calls
23+
nodes []BackendNode
24+
findErr error
25+
removedPairs []modelNodePair // records RemoveNodeModel calls
26+
removedReplicas []modelReplicaRef // records RemoveNodeModel calls including the replica index
2627
}
2728

2829
type modelNodePair struct {
2930
nodeID string
3031
modelName string
3132
}
3233

34+
// modelReplicaRef records a row removal at full replica granularity.
35+
// modelNodePair drops the index because RemoveAllNodeModelReplicas has none;
36+
// the backend delete/upgrade paths address exactly one replica row, so an
37+
// assertion that ignored the index could not tell a correct removal from one
38+
// that wiped a sibling replica still serving traffic.
39+
type modelReplicaRef struct {
40+
nodeID string
41+
modelName string
42+
replicaIndex int
43+
}
44+
3345
func (f *fakeModelLocator) FindNodesWithModel(_ context.Context, _ string) ([]BackendNode, error) {
3446
return f.nodes, f.findErr
3547
}
3648

37-
func (f *fakeModelLocator) RemoveNodeModel(_ context.Context, nodeID, modelName string, _ int) error {
49+
func (f *fakeModelLocator) RemoveNodeModel(_ context.Context, nodeID, modelName string, replicaIndex int) error {
3850
f.removedPairs = append(f.removedPairs, modelNodePair{nodeID, modelName})
51+
f.removedReplicas = append(f.removedReplicas, modelReplicaRef{nodeID, modelName, replicaIndex})
3952
return nil
4053
}
4154

0 commit comments

Comments
 (0)