diff --git a/docs/content/operations/backend-monitor.md b/docs/content/operations/backend-monitor.md index 1fa1750bf..359bab8f2 100644 --- a/docs/content/operations/backend-monitor.md +++ b/docs/content/operations/backend-monitor.md @@ -123,6 +123,8 @@ curl -X POST http://localhost:8080/backend/shutdown \ Returns `200 OK` with the shutdown confirmation message on success. +Stopping a backend removes its watchdog timers and eviction state. A timeout from a stopped backend does not shut down a replacement at a different address. + ## Error Responses | Status Code | Description | diff --git a/pkg/model/loader.go b/pkg/model/loader.go index 9a3a9fd78..85e39f3b5 100644 --- a/pkg/model/loader.go +++ b/pkg/model/loader.go @@ -614,6 +614,32 @@ func (ml *ModelLoader) ShutdownModelForce(modelName string) error { return ml.shutdownModel(ctx, modelName, true) } +// ShutdownModelAtAddress ignores stale watchdog evictions after a reload. +// Address validation and teardown share the same lifecycle lock as loading. +func (ml *ModelLoader) ShutdownModelAtAddress(modelName, address string, force bool) error { + ctx, cancel := context.WithTimeout(context.Background(), gracefulShutdownTimeout) + defer cancel() + release, err := ml.operations.acquireContext(ctx, modelName, true) + if err != nil { + return fmt.Errorf("waiting to shut down model %q: %w", modelName, err) + } + defer release() + ml.mu.Lock() + store := ml.store + ml.mu.Unlock() + m, ok := store.Get(modelName) + if !ok || m.address != address { + return nil + } + err = ml.deleteProcess(ctx, modelName, force) + if errors.Is(err, ErrModelBusy) && !force && forceBackendShutdown { + forceCtx, forceCancel := context.WithTimeout(context.Background(), forcedShutdownTimeout) + defer forceCancel() + return ml.deleteProcess(forceCtx, modelName, true) + } + return err +} + // ShutdownModelContext is the cancellation-aware lifecycle primitive used by // both graceful and forced shutdown wrappers. func (ml *ModelLoader) ShutdownModelContext(ctx context.Context, modelName string, force bool) error { diff --git a/pkg/model/process.go b/pkg/model/process.go index 9ec13bb51..36bff8c80 100644 --- a/pkg/model/process.go +++ b/pkg/model/process.go @@ -164,6 +164,9 @@ func (ml *ModelLoader) deleteProcess(ctx context.Context, s string, force bool) // at a known-unreachable worker, while the distributed registry remains // the source of truth for anything that is still running remotely. store.Delete(s) + if wd != nil { + wd.Untrack(model.address) + } return unloadErr } @@ -310,6 +313,7 @@ func (ml *ModelLoader) startProcess(grpcProcess, id string, serverAddress string } if err := grpcControlProcess.Run(); err != nil { + ml.untrackProcess(grpcControlProcess) runtime.cleanup() return grpcControlProcess, err } @@ -366,6 +370,7 @@ func (ml *ModelLoader) startProcess(grpcProcess, id string, serverAddress string // whether the child is alive. go func() { <-grpcControlProcess.Done() + ml.untrackProcess(grpcControlProcess) // LoadAndDelete both reads the intentional-stop marker and frees the // map entry so it doesn't accumulate across the process's lifetime. _, intentional := ml.stoppingProcs.LoadAndDelete(grpcControlProcess) @@ -403,6 +408,7 @@ func (ml *ModelLoader) cleanupProcessRuntime(process *process.Process) { if process == nil { return } + ml.untrackProcess(process) value, ok := ml.processRuntimes.LoadAndDelete(process) if !ok { return @@ -414,6 +420,18 @@ func (ml *ModelLoader) cleanupProcessRuntime(process *process.Process) { }() } +// Use the current watchdog because settings updates can replace it while a +// backend is running. Match the process identity so a late exit notification +// cannot remove a replacement that happens to reuse the same address. +func (ml *ModelLoader) untrackProcess(p *process.Process) { + ml.mu.Lock() + wd := ml.wd + ml.mu.Unlock() + if wd != nil { + wd.untrackProcess(p) + } +} + // CleanupProcessRuntime releases state and scratch owned by a process started // through StartProcess. Callers that supervise processes outside ModelLoader's // model store must invoke it after they have consumed exit diagnostics. diff --git a/pkg/model/watchdog.go b/pkg/model/watchdog.go index 3a604668a..977eca1ff 100644 --- a/pkg/model/watchdog.go +++ b/pkg/model/watchdog.go @@ -468,6 +468,7 @@ type modelUsageInfo struct { // opted into forceEvictionWhenBusy — either way, waiting for the graceful // deadline would defeat prompt eviction. type evictionTarget struct { + address string model string wasBusy bool } @@ -575,7 +576,7 @@ func (wd *WatchDog) collectEvictionsLocked(candidates []modelUsageInfo, maxToEvi continue } xlog.Info("[WatchDog] evicting model", "model", m.model, "busy", isBusy) - evicted = append(evicted, evictionTarget{model: m.model, wasBusy: isBusy}) + evicted = append(evicted, evictionTarget{address: m.address, model: m.model, wasBusy: isBusy}) wd.untrack(m.address) } return evicted, skippedBusy @@ -586,12 +587,7 @@ func (wd *WatchDog) collectEvictionsLocked(candidates []modelUsageInfo, maxToEvi // the graceful shutdown deadline. func (wd *WatchDog) shutdownEvicted(targets []evictionTarget, label string) { for _, t := range targets { - var err error - if t.wasBusy { - err = wd.pm.ShutdownModelForce(t.model) - } else { - err = wd.pm.ShutdownModel(t.model) - } + err := wd.shutdownTarget(t) if err != nil { xlog.Error("[WatchDog] error shutting down model during "+label, "error", err, "model", t.model, "busy", t.wasBusy) } @@ -599,6 +595,21 @@ func (wd *WatchDog) shutdownEvicted(targets []evictionTarget, label string) { } } +// shutdownTarget retains the address selected under the watchdog lock. The +// loader checks it under its per-model lifecycle lock so a queued eviction +// cannot stop a backend loaded after the original target exited. +func (wd *WatchDog) shutdownTarget(t evictionTarget) error { + if pm, ok := wd.pm.(interface { + ShutdownModelAtAddress(string, string, bool) error + }); ok { + return pm.ShutdownModelAtAddress(t.model, t.address, t.wasBusy) + } + if t.wasBusy { + return wd.pm.ShutdownModelForce(t.model) + } + return wd.pm.ShutdownModel(t.model) +} + // EnforceGroupExclusivity evicts every loaded model that shares at least one // concurrency group with the requested model. The pinned/busy/retry semantics // match EnforceLRULimit so the loader's retry loop can stay generic. @@ -713,7 +724,7 @@ func (wd *WatchDog) checkIdle() { xlog.Debug("[WatchDog] Watchdog checks for idle connections") // Collect models to shutdown while holding the lock - var modelsToShutdown []string + var modelsToShutdown []evictionTarget for address, t := range wd.idleTime { xlog.Debug("[WatchDog] idle connection", "address", address) if time.Since(t) > wd.idletimeout { @@ -723,8 +734,8 @@ func (wd *WatchDog) checkIdle() { xlog.Debug("[WatchDog] Skipping idle eviction for pinned model", "model", model) continue } - xlog.Warn("[WatchDog] Address is idle for too long, killing it", "address", address) - modelsToShutdown = append(modelsToShutdown, model) + xlog.Warn("[WatchDog] Address is idle for too long, killing it", "address", address, "model", model) + modelsToShutdown = append(modelsToShutdown, evictionTarget{model: model, address: address}) } else { xlog.Warn("[WatchDog] Address unresolvable", "address", address) } @@ -733,13 +744,7 @@ func (wd *WatchDog) checkIdle() { } wd.Unlock() - // Now shutdown models without holding the watchdog lock to prevent deadlock - for _, model := range modelsToShutdown { - if err := wd.pm.ShutdownModel(model); err != nil { - xlog.Error("[watchdog] error shutting down model", "error", err, "model", model) - } - xlog.Debug("[WatchDog] model shut down", "model", model) - } + wd.shutdownEvicted(modelsToShutdown, "idle timeout") } func (wd *WatchDog) checkBusy() { @@ -747,7 +752,7 @@ func (wd *WatchDog) checkBusy() { xlog.Debug("[WatchDog] Watchdog checks for busy connections") // Collect models to shutdown while holding the lock - var modelsToShutdown []string + var modelsToShutdown []evictionTarget for address, t := range wd.busyTime { xlog.Debug("[WatchDog] active connection", "address", address) @@ -755,7 +760,7 @@ func (wd *WatchDog) checkBusy() { model, ok := wd.addressModelMap[address] if ok { xlog.Warn("[WatchDog] Model is busy for too long, killing it", "model", model) - modelsToShutdown = append(modelsToShutdown, model) + modelsToShutdown = append(modelsToShutdown, evictionTarget{model: model, address: address, wasBusy: true}) } else { xlog.Warn("[WatchDog] Address unresolvable", "address", address) } @@ -764,16 +769,7 @@ func (wd *WatchDog) checkBusy() { } wd.Unlock() - // The busy-killer targets backends whose in-flight gRPC call has been - // stuck past the busy timeout. Use the force path so the loader stops - // the process FIRST (dropping the stuck call's gRPC connection) instead - // of waiting for the graceful shutdown deadline. - for _, model := range modelsToShutdown { - if err := wd.pm.ShutdownModelForce(model); err != nil { - xlog.Error("[watchdog] error shutting down model", "error", err, "model", model) - } - xlog.Debug("[WatchDog] busy model shut down", "model", model) - } + wd.shutdownEvicted(modelsToShutdown, "busy timeout") } // checkMemory monitors memory usage (GPU VRAM if available, otherwise RAM) and evicts backends when usage exceeds threshold @@ -896,11 +892,7 @@ func (wd *WatchDog) evictLRUModel() { wd.Unlock() // Shutdown the model - shutdown := wd.pm.ShutdownModel - if wasBusy { - shutdown = wd.pm.ShutdownModelForce - } - if err := shutdown(lruModel.model); err != nil && !errors.Is(err, ErrModelNotFound) { + if err := wd.shutdownTarget(evictionTarget{model: lruModel.model, address: lruModel.address, wasBusy: wasBusy}); err != nil && !errors.Is(err, ErrModelNotFound) { xlog.Error("[WatchDog] error shutting down model during memory reclamation", "error", err, "model", lruModel.model) } else { // Untrack the model @@ -913,7 +905,18 @@ func (wd *WatchDog) evictLRUModel() { func (wd *WatchDog) untrack(address string) { if modelID, ok := wd.addressModelMap[address]; ok { - delete(wd.modelSizes, modelID) + // A stale address can coexist with the replacement's registration. + // Its cleanup must not erase the replacement's size estimate. + otherAddress := false + for addr, name := range wd.addressModelMap { + if addr != address && name == modelID { + otherAddress = true + break + } + } + if !otherAddress { + delete(wd.modelSizes, modelID) + } } delete(wd.busyTime, address) delete(wd.inFlight, address) @@ -924,3 +927,20 @@ func (wd *WatchDog) untrack(address string) { delete(wd.addressModelMap, address) delete(wd.addressMap, address) } + +// Untrack removes request and eviction state after a backend is removed. +func (wd *WatchDog) Untrack(address string) { + wd.Lock() + defer wd.Unlock() + wd.untrack(address) +} + +func (wd *WatchDog) untrackProcess(p *process.Process) { + wd.Lock() + defer wd.Unlock() + for address, tracked := range wd.addressMap { + if tracked == p { + wd.untrack(address) + } + } +} diff --git a/pkg/model/watchdog_lifecycle_test.go b/pkg/model/watchdog_lifecycle_test.go new file mode 100644 index 000000000..50e4e84d3 --- /dev/null +++ b/pkg/model/watchdog_lifecycle_test.go @@ -0,0 +1,157 @@ +// SPDX-License-Identifier: MIT + +package model + +import ( + "context" + "os" + "path/filepath" + "time" + + grpc "github.com/mudler/LocalAI/pkg/grpc" + "github.com/mudler/LocalAI/pkg/system" + process "github.com/mudler/go-processmanager" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +type watchdogLifecycleBackend struct{ grpc.Backend } + +func (*watchdogLifecycleBackend) IsBusy() bool { return false } +func (*watchdogLifecycleBackend) Free(context.Context) error { return nil } + +var _ = Describe("Watchdog backend lifecycle", func() { + var loader *ModelLoader + var wd *WatchDog + BeforeEach(func() { + loader = NewModelLoader(&system.SystemState{Model: system.Model{ModelsPath: GinkgoT().TempDir()}}) + wd = NewWatchDog(WithProcessManager(loader), WithIdleTimeout(time.Second), WithBusyTimeout(time.Second), WithLRULimit(1)) + loader.SetWatchDog(wd) + }) + + It("removes shutdown tracking and ignores late request completions", func() { + loader.store.Set("model", NewModelWithClient("model", "old", &watchdogLifecycleBackend{})) + wd.AddAddressModelMap("old", "model") + finish := wd.TrackRequest("old") + Expect(loader.ShutdownModelForce("model")).To(Succeed()) + finish() + state := wd.GetState() + Expect(state.AddressModelMap).To(BeEmpty()) + Expect(state.BusyTime).To(BeEmpty()) + Expect(state.IdleTime).To(BeEmpty()) + Expect(state.InFlight).To(BeEmpty()) + Expect(state.LastUsed).To(BeEmpty()) + }) + + DescribeTable("does not evict a replacement backend from a stale address", + func(evict func(*WatchDog)) { + replacement := NewModelWithClient("model", "new", &watchdogLifecycleBackend{}) + loader.store.Set("model", replacement) + wd.AddAddressModelMap("old", "model") + wd.AddAddressModelMap("new", "model") + wd.RegisterModelSize("model", 123) + wd.lastUsed["old"] = time.Now().Add(-time.Hour) + wd.lastUsed["new"] = time.Now() + wd.idleTime["old"] = time.Now().Add(-time.Hour) + evict(wd) + resident, ok := loader.store.Get("model") + Expect(ok).To(BeTrue()) + Expect(resident).To(BeIdenticalTo(replacement)) + Expect(wd.GetState().AddressModelMap).To(HaveKeyWithValue("new", "model")) + Expect(wd.modelSizes).To(HaveKeyWithValue("model", int64(123))) + }, + Entry("idle timeout", func(w *WatchDog) { w.checkIdle() }), + Entry("busy timeout", func(w *WatchDog) { w.busyTime["old"] = time.Now().Add(-time.Hour); w.checkBusy() }), + Entry("memory eviction", func(w *WatchDog) { w.evictLRUModel() }), + ) + + DescribeTable("does not stop a replacement after eviction selection", + func(force bool) { + wd.AddAddressModelMap("old", "model") + if force { + wd.TrackRequest("old") + } + targets, _ := wd.collectEvictionsLocked([]modelUsageInfo{{model: "model", address: "old"}}, 1, force) + replacement := NewModelWithClient("model", "new", &watchdogLifecycleBackend{}) + loader.store.Set("model", replacement) + wd.AddAddressModelMap("new", "model") + wd.RegisterModelSize("model", 123) + wd.shutdownEvicted(targets, "test") + resident, ok := loader.store.Get("model") + Expect(ok).To(BeTrue()) + Expect(resident).To(BeIdenticalTo(replacement)) + Expect(wd.modelSizes).To(HaveKeyWithValue("model", int64(123))) + }, + Entry("graceful", false), + Entry("forced", true), + ) + + DescribeTable("still stops the matching backend", + func(force bool) { + loader.store.Set("model", NewModelWithClient("model", "current", &watchdogLifecycleBackend{})) + wd.AddAddressModelMap("current", "model") + Expect(wd.shutdownTarget(evictionTarget{model: "model", address: "current", wasBusy: force})).To(Succeed()) + _, ok := loader.store.Get("model") + Expect(ok).To(BeFalse()) + Expect(wd.GetState().AddressModelMap).To(BeEmpty()) + }, + Entry("graceful", false), + Entry("forced", true), + ) + + It("keeps replacement tracking when an old process cleanup arrives late", func() { + oldProcess := process.New() + newProcess := process.New() + wd.Add("same-address", oldProcess) + wd.AddAddressModelMap("same-address", "model") + wd.Add("same-address", newProcess) + wd.RegisterModelSize("model", 123) + finish := wd.TrackRequest("same-address") + loader.cleanupProcessRuntime(oldProcess) + Expect(wd.GetState().AddressMap).To(HaveKeyWithValue("same-address", newProcess)) + Expect(wd.modelSizes).To(HaveKeyWithValue("model", int64(123))) + finish() + Expect(wd.GetState().IdleTime).To(HaveKey("same-address")) + }) + + It("removes local process tracking synchronously during shutdown", func() { + p := process.New() + m := NewModelWithClient("model", "old", &watchdogLifecycleBackend{}) + m.process = p + loader.store.Set("model", m) + wd.Add("old", p) + wd.AddAddressModelMap("old", "model") + wd.TrackRequest("old")() + Expect(loader.ShutdownModelForce("model")).To(Succeed()) + Expect(wd.GetState().AddressMap).To(BeEmpty()) + Expect(wd.GetState().IdleTime).To(BeEmpty()) + }) + + It("cleans the current watchdog after its configuration is replaced", func() { + p := process.New() + wd.Add("old", p) + wd.AddAddressModelMap("old", "model") + replacement := NewWatchDog(WithProcessManager(loader)) + replacement.RestoreState(wd.GetState()) + loader.SetWatchDog(replacement) + loader.cleanupProcessRuntime(p) + Expect(replacement.GetState().AddressModelMap).To(BeEmpty()) + }) + + It("untracks a backend that fails to start", func() { + _, err := loader.startProcess(filepath.Join(GinkgoT().TempDir(), "missing"), "model", "old", nil) + Expect(err).To(HaveOccurred()) + Expect(wd.GetState().AddressModelMap).To(BeEmpty()) + Expect(wd.GetState().AddressMap).To(BeEmpty()) + }) + + It("untracks a backend after an unexpected exit", func() { + backend := filepath.Join(GinkgoT().TempDir(), "backend") + Expect(os.WriteFile(backend, []byte("#!/bin/sh\nexit 42\n"), 0o700)).To(Succeed()) + p, err := loader.startProcess(backend, "model", "old", nil) + Expect(err).NotTo(HaveOccurred()) + DeferCleanup(func() { loader.cleanupProcessRuntime(p) }) + Eventually(p.Done()).Should(BeClosed()) + Eventually(func() map[string]string { return wd.GetState().AddressModelMap }).Should(BeEmpty()) + }) +})