diff --git a/core/http/endpoints/localai/backend_monitor.go b/core/http/endpoints/localai/backend_monitor.go index d138c5bee..c60e92a92 100644 --- a/core/http/endpoints/localai/backend_monitor.go +++ b/core/http/endpoints/localai/backend_monitor.go @@ -1,9 +1,15 @@ package localai import ( + "errors" + "fmt" + "net/http" + "github.com/labstack/echo/v4" "github.com/mudler/LocalAI/core/schema" "github.com/mudler/LocalAI/core/services/monitoring" + "github.com/mudler/LocalAI/pkg/model" + "github.com/mudler/xlog" ) // BackendMonitorEndpoint returns the status of the specified backend @@ -48,7 +54,25 @@ func BackendShutdownEndpoint(bm *monitoring.BackendMonitorService) echo.HandlerF if err := c.Bind(input); err != nil { return err } + if input.Model == "" { + return echo.NewHTTPError(http.StatusBadRequest, "model is required") + } - return bm.ShutdownModel(input.Model) + if err := bm.ShutdownModel(input.Model); err != nil { + // "Not loaded" is a client-side condition, not a server fault, so + // it should not surface as a 500. In distributed mode this branch + // used to be reached for models that were running on a worker — + // only this replica's local store was empty. The loader now + // consults the node registry first, so reaching it means the + // model is loaded neither here nor anywhere in the cluster. + if errors.Is(err, model.ErrModelNotFound) { + xlog.Info("Shutdown requested for a model that is not loaded", "model", input.Model) + return echo.NewHTTPError(http.StatusNotFound, + fmt.Sprintf("model %q is not loaded on this instance or on any worker node", input.Model)) + } + xlog.Error("Failed to shut down model", "model", input.Model, "error", err) + return err + } + return c.JSON(http.StatusOK, map[string]string{"message": "model stopped"}) } } diff --git a/core/services/nodes/unloader.go b/core/services/nodes/unloader.go index 869996817..f6f1db211 100644 --- a/core/services/nodes/unloader.go +++ b/core/services/nodes/unloader.go @@ -12,6 +12,7 @@ import ( "github.com/mudler/LocalAI/core/services/galleryop" "github.com/mudler/LocalAI/core/services/messaging" + "github.com/mudler/LocalAI/pkg/model" "github.com/mudler/xlog" ) @@ -91,8 +92,12 @@ func (a *RemoteUnloaderAdapter) UnloadRemoteModelContext(ctx context.Context, mo return fmt.Errorf("finding nodes with model %q: %w", modelName, err) } if len(nodes) == 0 { + // Distinguish "no node has it" from "stopped successfully" so the + // caller can report the truth to an operator. Returning nil here made + // a no-op stop indistinguishable from a real one, and left the loader + // unable to tell a cluster-wide miss from a completed unload. xlog.Debug("No remote nodes found with model", "model", modelName) - return nil + return model.ErrRemoteModelNotLoaded } var unloadErr error diff --git a/core/services/nodes/unloader_test.go b/core/services/nodes/unloader_test.go index dc2c7e8b1..b502af144 100644 --- a/core/services/nodes/unloader_test.go +++ b/core/services/nodes/unloader_test.go @@ -14,6 +14,7 @@ import ( "github.com/mudler/LocalAI/core/services/galleryop" "github.com/mudler/LocalAI/core/services/messaging" + "github.com/mudler/LocalAI/pkg/model" ) // --- Fakes --- @@ -127,9 +128,13 @@ var _ = Describe("RemoteUnloaderAdapter", func() { }) Describe("UnloadRemoteModel", func() { - It("with no nodes returns nil", func() { + It("with no nodes reports the model is not loaded", func() { + // Previously returned nil, making "stopped it" and "there was + // nothing to stop" indistinguishable. ShutdownModel relies on + // this distinction to choose between success and a genuine 404, + // so a cluster-wide miss must be reported, not swallowed. locator.nodes = nil - Expect(adapter.UnloadRemoteModel("my-model")).To(Succeed()) + Expect(adapter.UnloadRemoteModel("my-model")).To(MatchError(model.ErrRemoteModelNotLoaded)) Expect(mc.published).To(BeEmpty()) }) diff --git a/core/services/worker/backend_delete_test.go b/core/services/worker/backend_delete_test.go new file mode 100644 index 000000000..1fb0c52ec --- /dev/null +++ b/core/services/worker/backend_delete_test.go @@ -0,0 +1,208 @@ +package worker + +import ( + "github.com/mudler/LocalAI/core/gallery" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +// Production incident (Jetson/Thor worker): deleting a backend returned HTTP +// 200 but left its gRPC process alive for ~40 minutes with its directory +// removed from disk. A later model load was routed to that orphan and failed +// with a certifi path pointing into the deleted directory. +// +// Root cause: s.processes is keyed by `modelID#replicaIndex` +// (buildProcessKey), so a delete keyed on the *backend* name resolved to no +// keys at all and the stop silently no-op'd. These specs pin the two lookups +// that must work by backend name. +var _ = Describe("Backend deletion reaps the backend's processes", func() { + const ( + concrete = "cuda13-nvidia-l4t-arm64-longcat-video" + development = "cuda13-nvidia-l4t-arm64-longcat-video-development" + alias = "longcat-video" + ) + + Describe("resolveProcessKeysForBackend", func() { + It("finds processes started for a backend even though they are keyed by modelID", func() { + s := &backendSupervisor{ + processes: map[string]*backendProcess{ + // Started by a model load: key is modelID#replica, and the + // backend name survives only in backendName. + "LongCat-Video#0": {addr: "127.0.0.1:30232", backendName: concrete}, + "LongCat-Video#1": {addr: "127.0.0.1:30234", backendName: concrete}, + "other-model#0": {addr: "127.0.0.1:30235", backendName: "llama-cpp"}, + }, + } + + Expect(s.resolveProcessKeysForBackend(setOf(concrete))).To( + ConsistOf("LongCat-Video#0", "LongCat-Video#1"), + "a delete keyed on the backend name must reach every process started for it") + }) + + It("matches through the alias the model config used at install time", func() { + // The model config referenced `backend: longcat-video` (the alias), + // so the process recorded the alias; the delete request carries the + // concrete directory name. Without alias resolution the orphan + // survives exactly as it did in production. + s := &backendSupervisor{ + processes: map[string]*backendProcess{ + "LongCat-Video#0": {addr: "127.0.0.1:30232", backendName: alias}, + }, + } + + Expect(s.resolveProcessKeysForBackend(setOf(concrete, alias))).To( + ConsistOf("LongCat-Video#0"), + "alias and concrete name must resolve to the same process") + }) + + It("does not reap processes belonging to a different backend", func() { + s := &backendSupervisor{ + processes: map[string]*backendProcess{ + "LongCat-Video#0": {addr: "127.0.0.1:30233", backendName: development}, + }, + } + + Expect(s.resolveProcessKeysForBackend(setOf(concrete))).To(BeEmpty(), + "deleting the non-development backend must not kill the -development process") + }) + + It("still resolves legacy entries that predate backendName via the modelID prefix", func() { + // Installs with an empty modelID key the map by the backend name + // itself (buildProcessKey falls back). Those must keep working. + s := &backendSupervisor{ + processes: map[string]*backendProcess{ + concrete + "#0": {addr: "127.0.0.1:30232"}, + }, + } + + Expect(s.resolveProcessKeysForBackend(setOf(concrete))).To(ConsistOf(concrete + "#0")) + }) + }) + + Describe("resolveStopTargets", func() { + // backend.stop is ambiguous by design: the admin backends UI publishes + // a BACKEND name, UnloadRemoteModel publishes a MODEL name, and the + // router's abandoned-load reap (#10948) publishes an exact + // modelID#replica key. Resolving only one of the three silently + // strands the others. + newSupervisor := func() *backendSupervisor { + return &backendSupervisor{ + processes: map[string]*backendProcess{ + "LongCat-Video#0": {addr: "127.0.0.1:30232", backendName: concrete}, + }, + } + } + + It("stops a process addressed by model ID", func() { + Expect(newSupervisor().resolveStopTargets("LongCat-Video")).To(ConsistOf("LongCat-Video#0"), + "model-name stop (UnloadRemoteModel) must still reach the process") + }) + + It("stops a process addressed by backend name", func() { + Expect(newSupervisor().resolveStopTargets(concrete)).To(ConsistOf("LongCat-Video#0"), + "backend-name stop must reach the process too") + }) + + It("stops a process addressed by an exact replica key", func() { + Expect(newSupervisor().resolveStopTargets("LongCat-Video#0")).To(ConsistOf("LongCat-Video#0"), + "the router reaps an abandoned load by exact modelID#replica key") + }) + + It("does not stop unrelated processes", func() { + Expect(newSupervisor().resolveStopTargets("some-other-model")).To(BeEmpty()) + }) + }) + + Describe("processMatchesBackend", func() { + It("rejects reusing a process that was started for a different backend", func() { + // The install fast path returns the address of any live process + // under this processKey. After the concrete backend was deleted and + // the -development variant installed, the same model+replica got + // the deleted backend's port handed back to it. + s := &backendSupervisor{ + processes: map[string]*backendProcess{ + "LongCat-Video#0": {addr: "127.0.0.1:30232", backendName: concrete}, + }, + } + + Expect(s.processMatchesBackend("LongCat-Video#0", development)).To(BeFalse(), + "a process started for the deleted backend must not be reused for another backend") + }) + + It("accepts reusing a process started for the same backend", func() { + s := &backendSupervisor{ + processes: map[string]*backendProcess{ + "LongCat-Video#0": {addr: "127.0.0.1:30232", backendName: concrete}, + }, + } + + Expect(s.processMatchesBackend("LongCat-Video#0", concrete)).To(BeTrue()) + }) + + It("accepts a process recorded under the alias of the requested backend", func() { + s := &backendSupervisor{ + processes: map[string]*backendProcess{ + "LongCat-Video#0": {addr: "127.0.0.1:30232", backendName: alias}, + }, + } + + Expect(s.processMatchesBackend("LongCat-Video#0", alias)).To(BeTrue()) + }) + + It("accepts legacy entries with no recorded backend name", func() { + // Pre-upgrade processes carry no backendName. Treating them as a + // mismatch would restart every running backend once on rollout. + s := &backendSupervisor{ + processes: map[string]*backendProcess{ + "LongCat-Video#0": {addr: "127.0.0.1:30232"}, + }, + } + + Expect(s.processMatchesBackend("LongCat-Video#0", concrete)).To(BeTrue()) + }) + }) + + Describe("backendIdentitySet", func() { + It("maps an alias to its concrete backend and back", func() { + backends := gallery.SystemBackends{ + concrete: {Name: concrete, Metadata: &gallery.BackendMetadata{Name: concrete, Alias: alias}}, + alias: {Name: alias, Metadata: &gallery.BackendMetadata{Name: concrete, Alias: alias}}, + } + + Expect(backendIdentitySet(backends, concrete)).To(HaveKey(alias)) + Expect(backendIdentitySet(backends, alias)).To(HaveKey(concrete)) + }) + + It("falls back to the bare name when the backend is unknown", func() { + // A delete must never fail (or over-reach) because the gallery + // listing is unavailable or the entry is already gone from disk. + Expect(backendIdentitySet(gallery.SystemBackends{}, concrete)).To( + Equal(map[string]struct{}{concrete: {}})) + }) + + It("does not conflate two concrete backends sharing an alias", func() { + // Both variants declare alias "longcat-video"; ListSystemBackends + // picks one for the alias row. Deleting the non-chosen concrete + // must not pull in the chosen one's identity. + backends := gallery.SystemBackends{ + development: {Name: development, Metadata: &gallery.BackendMetadata{Name: development, Alias: alias}}, + alias: {Name: alias, Metadata: &gallery.BackendMetadata{Name: development, Alias: alias}}, + } + + set := backendIdentitySet(backends, development) + Expect(set).To(HaveKey(development)) + Expect(set).ToNot(HaveKey(concrete)) + }) + }) +}) + +// setOf builds the identity set the resolver consumes, keeping the specs +// readable without repeating the map literal. +func setOf(names ...string) map[string]struct{} { + set := make(map[string]struct{}, len(names)) + for _, n := range names { + set[n] = struct{}{} + } + return set +} diff --git a/core/services/worker/free_timeout_test.go b/core/services/worker/free_timeout_test.go new file mode 100644 index 000000000..eee10b28e --- /dev/null +++ b/core/services/worker/free_timeout_test.go @@ -0,0 +1,114 @@ +package worker + +import ( + "context" + "net" + + process "github.com/mudler/go-processmanager" + gogrpc "google.golang.org/grpc" + + pb "github.com/mudler/LocalAI/pkg/grpc/proto" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +// hangingBackend is a real gRPC backend server whose Free handler never +// returns. It models the production failure mode: 37 Python backends default to +// PYTHON_GRPC_MAX_WORKERS=1, so while a LoadModel handler is wedged (e.g. +// inside torch.load) the single gRPC worker thread is occupied and a Free RPC +// queues behind it forever. +// +// The server must be real rather than a socket that merely accepts: without a +// completed HTTP/2 handshake the call is bounded by gRPC's own ~20s connect +// timeout, which would mask an unbounded Free instead of exposing it. Here the +// connection reaches READY, so nothing but the caller's own deadline can end +// the call. +type hangingBackend struct { + pb.UnimplementedBackendServer + entered chan struct{} + release chan struct{} +} + +func (h *hangingBackend) Free(_ context.Context, _ *pb.HealthMessage) (*pb.Result, error) { + close(h.entered) + // Deliberately ignores the server-side context: a backend whose only + // worker thread is blocked cannot observe cancellation either. + <-h.release + return &pb.Result{Success: true}, nil +} + +var _ = Describe("Stopping a backend whose Free never returns", func() { + var ( + proc *process.Process + backend *hangingBackend + server *gogrpc.Server + s *backendSupervisor + ) + + BeforeEach(func() { + lis, err := net.Listen("tcp", "127.0.0.1:0") + Expect(err).ToNot(HaveOccurred()) + + backend = &hangingBackend{ + entered: make(chan struct{}), + release: make(chan struct{}), + } + server = gogrpc.NewServer() + pb.RegisterBackendServer(server, backend) + go func() { _ = server.Serve(lis) }() + + // The process is deliberately never Run(): go-processmanager v0.1.1 + // writes Process.pid from readPID() with no synchronization, so a live + // process races its own monitor goroutine under -race (reproducible + // with a bare Run()+Stop(), independent of this spec). Since + // scripts/model-lifecycle-conformance.sh runs this package with -race + // and is fail-closed, starting one here would turn that gate red on an + // upstream defect. An unstarted process still drives the branch this + // spec is about: Stop() returns promptly, which is all that is needed + // to prove the stop was reached at all. Whether SIGTERM/SIGKILL then + // lands is go-processmanager's contract, not this supervisor's. + proc = process.New( + process.WithTemporaryStateDir(), + process.WithName("/bin/sleep"), + process.WithArgs("300"), + ) + + s = &backendSupervisor{ + cfg: &Config{}, + processes: map[string]*backendProcess{ + "wedged-model#0": { + proc: proc, + addr: lis.Addr().String(), + port: lis.Addr().(*net.TCPAddr).Port, + backendName: "longcat-video", + }, + }, + } + }) + + AfterEach(func() { + close(backend.release) + server.Stop() + }) + + It("reaches the stop even though Free never returns", func() { + done := make(chan error, 1) + go func() { done <- s.stopBackendExact("wedged-model#0", false) }() + + Eventually(backend.entered, "20s", "100ms").Should(BeClosed(), + "Free must be attempted — it is a courtesy call we still want to make") + + // The whole point: a Free that never answers must not hold the stop + // hostage. Unbounded, this never receives, the process is never + // signalled, and the router's abandoned-load reap (#10948) is inert + // against a wedged single-worker Python backend. + Eventually(done, "60s", "200ms").Should(Receive(BeNil()), + "stopBackendExact must fall through to the process stop despite a hung Free") + + s.mu.Lock() + defer s.mu.Unlock() + Expect(s.processes).ToNot(HaveKey("wedged-model#0"), + "the supervisor must release the slot so the port can be recycled") + }) +}) diff --git a/core/services/worker/install.go b/core/services/worker/install.go index 19056ea7a..8efb3664e 100644 --- a/core/services/worker/install.go +++ b/core/services/worker/install.go @@ -72,30 +72,46 @@ func (s *backendSupervisor) installBackend(req messaging.BackendInstallRequest, // replica failed, retries the install, the supervisor says "already // running" again, and the cluster loops on a dead replica forever. if addr := s.getAddr(processKey); addr != "" { - if s.isRunning(processKey) { + switch { + case !s.processMatchesBackend(processKey, req.Backend): + // The slot is held by a process started from a DIFFERENT + // backend — typically this model's previous backend was + // deleted (or superseded by a -development variant) while its + // process stayed up. Reusing that address would serve the load + // from a backend directory that may no longer exist on disk. + xlog.Warn("Process for this model replica belongs to another backend; restarting it", + "backend", req.Backend, "model", req.ModelID, "replica", req.ReplicaIndex, "addr", addr) + if err := s.stopBackendExact(processKey, false); err != nil { + return "", fmt.Errorf("stopping mismatched backend process before reinstall: %w", err) + } + case s.isRunning(processKey): xlog.Info("Backend already running for model replica", "backend", req.Backend, "model", req.ModelID, "replica", req.ReplicaIndex, "addr", addr) return addr, nil + default: + xlog.Warn("Stale process entry for backend (dead process); cleaning up before reinstall", + "backend", req.Backend, "model", req.ModelID, "replica", req.ReplicaIndex, "addr", addr) + if err := s.stopBackendExact(processKey, false); err != nil { + xlog.Warn("Failed to clean up stale process entry", "processKey", processKey, "error", err) + } } - xlog.Warn("Stale process entry for backend (dead process); cleaning up before reinstall", - "backend", req.Backend, "model", req.ModelID, "replica", req.ReplicaIndex, "addr", addr) - s.stopBackendExact(processKey, false) } } else { // Upgrade path: stop every live process that shares this backend so the // gallery install can overwrite the on-disk artifact and the restarted - // process picks up the new binary. resolveProcessKeys catches peer - // replicas of the same backend (whisper#0, whisper#1, ...) on workers - // configured with MaxReplicasPerModel>1. We also stop the exact - // processKey from the request tuple — keys created with an explicit - // modelID don't share the bare-name prefix the resolver matches, but - // they're still using the old binary and need to come down. Both calls - // are no-ops on missing keys. - toStop := s.resolveProcessKeys(req.Backend) + // process picks up the new binary. resolveProcessKeysForBackend finds + // them by recorded backend name (and alias), which also covers peer + // replicas (whisper#0, whisper#1, ...) on workers configured with + // MaxReplicasPerModel>1. We also stop the exact processKey from the + // request tuple, whose recorded name may be absent on legacy entries. + // Both are no-ops on missing keys. + toStop := s.resolveProcessKeysForBackend(s.backendIdentity(req.Backend)) toStop = append(toStop, processKey) for _, key := range toStop { xlog.Info("Force install: stopping running backend before reinstall", "backend", req.Backend, "processKey", key) - s.stopBackendExact(key, true) + if err := s.stopBackendExact(key, true); err != nil { + return "", fmt.Errorf("stopping running backend before reinstall: %w", err) + } } } @@ -160,8 +176,10 @@ func (s *backendSupervisor) installBackend(req messaging.BackendInstallRequest, xlog.Info("Found backend binary", "path", backendPath, "processKey", processKey) - // Start the gRPC process on a new port (keyed by model, not just backend) - return s.startBackend(processKey, backendPath) + // Start the gRPC process on a new port (keyed by model, not just backend). + // req.Backend is recorded on the process so a later backend.delete/stop can + // find it by backend name. + return s.startBackend(processKey, req.Backend, backendPath) } // upgradeBackend stops every running process for `backend`, force-reinstalls @@ -173,12 +191,14 @@ func (s *backendSupervisor) installBackend(req messaging.BackendInstallRequest, func (s *backendSupervisor) upgradeBackend(req messaging.BackendUpgradeRequest) error { // Stop every live process for this backend (peer replicas + the bare // processKey). Same logic as the force branch in installBackend. - toStop := s.resolveProcessKeys(req.Backend) + toStop := s.resolveProcessKeysForBackend(s.backendIdentity(req.Backend)) toStop = append(toStop, buildProcessKey("", req.Backend, int(req.ReplicaIndex))) for _, key := range toStop { xlog.Info("Upgrade: stopping running backend before reinstall", "backend", req.Backend, "processKey", key) - s.stopBackendExact(key, true) + if err := s.stopBackendExact(key, true); err != nil { + return fmt.Errorf("stopping running backend before upgrade: %w", err) + } } galleries := s.galleries diff --git a/core/services/worker/lifecycle.go b/core/services/worker/lifecycle.go index 0348e41ec..53e44c1eb 100644 --- a/core/services/worker/lifecycle.go +++ b/core/services/worker/lifecycle.go @@ -4,7 +4,9 @@ import ( "context" "encoding/json" "fmt" + "maps" "net" + "slices" "syscall" "github.com/mudler/LocalAI/core/gallery" @@ -137,7 +139,14 @@ func (s *backendSupervisor) handleBackendStop(data []byte) { return } xlog.Info("Received NATS backend.stop event", "backend", req.Backend, "force", req.Force) - s.stopBackend(req.Backend, req.Force) + // The identifier may be a backend name, a model name, or an exact + // modelID#replica key depending on the publisher; resolveStopTargets + // handles all three. stopBackend alone resolves only the model meanings. + for _, key := range s.resolveStopTargets(req.Backend) { + if err := s.stopBackendExact(key, req.Force); err != nil { + xlog.Error("Failed to stop backend process", "backend", req.Backend, "processKey", key, "error", err) + } + } } func decodeBackendStopRequest(data []byte) (messaging.BackendStopRequest, bool, error) { @@ -162,9 +171,36 @@ func (s *backendSupervisor) handleBackendDelete(data []byte, reply func([]byte)) } xlog.Info("Received NATS backend.delete event", "backend", req.Backend) - // Stop if running this backend - if s.isRunning(req.Backend) { - s.stopBackend(req.Backend, false) + // Resolve the backend's identity (concrete name + alias) BEFORE touching + // the filesystem: DeleteBackendFromSystem removes the metadata.json that + // carries the alias, and a model loaded via the alias records the alias as + // its process's backend name. + identity := s.backendIdentity(req.Backend) + + // Stop every process started for this backend. Processes are keyed by + // modelID#replica, so the lookup must match the recorded backend name — a + // lookup by backend name alone resolved to nothing and left the process + // running with its directory deleted underneath it. + keys := s.resolveProcessKeysForBackend(identity) + if len(keys) == 0 { + // Not an error: deleting a backend that was never loaded is routine. + // But log it — silence here is what made the orphan case invisible. + xlog.Info("Deleting backend with no matching running process", + "backend", req.Backend, "identity", slices.Sorted(maps.Keys(identity))) + } + for _, key := range keys { + if err := s.stopBackendExact(key, false); err != nil { + // We knew about this process and could not kill it. Replying + // success would repeat the original defect: the operator is told + // "backend deleted" while the process keeps serving requests. + xlog.Error("Failed to stop backend process during delete; aborting delete", + "backend", req.Backend, "processKey", key, "error", err) + replyJSON(reply, messaging.BackendDeleteReply{ + Success: false, + Error: fmt.Sprintf("could not stop running process %s: %v", key, err), + }) + return + } } // Delete the backend files diff --git a/core/services/worker/supervisor.go b/core/services/worker/supervisor.go index 702071b74..12df8e79c 100644 --- a/core/services/worker/supervisor.go +++ b/core/services/worker/supervisor.go @@ -11,6 +11,7 @@ import ( "time" "github.com/mudler/LocalAI/core/config" + "github.com/mudler/LocalAI/core/gallery" "github.com/mudler/LocalAI/core/services/messaging" grpc "github.com/mudler/LocalAI/pkg/grpc" "github.com/mudler/LocalAI/pkg/model" @@ -25,10 +26,58 @@ type backendProcess struct { addr string // gRPC address (host:port) port int stopping bool + // backendName is the gallery backend this process was started for (e.g. + // "cuda13-nvidia-l4t-arm64-longcat-video"). It is NOT derivable from the + // map key: keys are `modelID#replicaIndex`, so without it a backend.delete + // keyed on the backend name resolves to nothing and the process is + // orphaned while its files are removed from disk. + backendName string } const workerBackendFreeTimeout = 5 * time.Second +// backendIdentitySet returns every name that refers to the same on-disk +// backend as `name`: the concrete directory name plus its alias. Deletes +// arrive with the concrete name while installs arrive with whatever the model +// config declared (often the alias), so both must land on one identity. +// +// An unknown name yields just itself. A delete must never over-reach (or fail) +// because the gallery listing is unavailable or the entry is already gone. +func backendIdentitySet(backends gallery.SystemBackends, name string) map[string]struct{} { + set := map[string]struct{}{name: {}} + b, ok := backends[name] + if !ok || b.Metadata == nil { + return set + } + // Only this entry's own metadata is consulted. Two concrete backends can + // share one alias (foo and foo-development); walking every candidate of + // that alias would let deleting one reap the other's process. + if b.Metadata.Name != "" { + set[b.Metadata.Name] = struct{}{} + } + if b.Metadata.Alias != "" { + set[b.Metadata.Alias] = struct{}{} + } + return set +} + +// backendIdentity resolves `name` against this worker's installed backends. +// Callers must invoke it BEFORE deleting files: DeleteBackendFromSystem +// removes the metadata.json that carries the alias mapping. +func (s *backendSupervisor) backendIdentity(name string) map[string]struct{} { + if s.systemState == nil { + return map[string]struct{}{name: {}} + } + backends, err := gallery.ListSystemBackends(s.systemState) + if err != nil { + // Degrade to name-only matching rather than failing the operation. + xlog.Warn("Could not list backends for alias resolution; matching on name only", + "backend", name, "error", err) + return map[string]struct{}{name: {}} + } + return backendIdentitySet(backends, name) +} + // backendSupervisor manages multiple backend gRPC processes on different ports. // Each backend type (e.g., llama-cpp, bert-embeddings) gets its own process and port. type backendSupervisor struct { @@ -56,7 +105,11 @@ type backendSupervisor struct { // startBackend starts a gRPC backend process on a dynamically allocated port. // Returns the gRPC address. -func (s *backendSupervisor) startBackend(backend, backendPath string) (string, error) { +// +// `backend` is the process key (modelID#replica); `backendName` is the gallery +// backend it was started from, recorded so delete/stop can find this process +// by backend name later. +func (s *backendSupervisor) startBackend(backend, backendName, backendPath string) (string, error) { s.mu.Lock() // Already running? @@ -94,9 +147,10 @@ func (s *backendSupervisor) startBackend(backend, backendPath string) (string, e } s.processes[backend] = &backendProcess{ - proc: proc, - addr: clientAddr, - port: port, + proc: proc, + addr: clientAddr, + port: port, + backendName: backendName, } xlog.Info("Backend process started", "backend", backend, "addr", clientAddr) @@ -219,12 +273,119 @@ func (s *backendSupervisor) resolveProcessKeys(id string) []string { return keys } +// resolveProcessKeysForBackend returns the process keys of every process +// started for any name in `names` (the backend's identity set). +// +// resolveProcessKeys alone is not enough here: it resolves a bare *modelID* +// via the `id#N` prefix, but backend.delete/backend.stop carry a *backend* +// name, which never prefixes a `modelID#replica` key. That mismatch is what +// let a deleted backend's process survive with its files removed. Matching on +// the recorded backendName is the primary path; the prefix resolution is +// unioned in so legacy entries keyed by the backend name itself (installs with +// an empty modelID) keep resolving. +func (s *backendSupervisor) resolveProcessKeysForBackend(names map[string]struct{}) []string { + seen := make(map[string]struct{}) + var keys []string + add := func(k string) { + if _, dup := seen[k]; dup { + return + } + seen[k] = struct{}{} + keys = append(keys, k) + } + + s.mu.Lock() + for key, bp := range s.processes { + if bp.backendName == "" { + continue + } + if _, ok := names[bp.backendName]; ok { + add(key) + } + } + s.mu.Unlock() + + // Legacy fallback for entries predating backendName: an install with an + // empty modelID keys the map by the backend name itself. Restricted to + // entries with no recorded name so it cannot reap a process that belongs + // to a different backend but whose modelID happens to equal this backend's + // name or alias. + for name := range names { + for _, k := range s.resolveProcessKeys(name) { + s.mu.Lock() + bp, ok := s.processes[k] + unnamed := ok && bp.backendName == "" + s.mu.Unlock() + if unnamed { + add(k) + } + } + } + return keys +} + +// resolveStopTargets resolves the ambiguous identifier carried on backend.stop. +// +// The payload field is named "backend", but it is published with both +// meanings: the admin backends UI sends a BACKEND name, while +// UnloadRemoteModel and the router's abandoned-load reap send a MODEL name (or +// an exact modelID#replica key). Resolving only one meaning silently strands +// the other — matching on backend name alone stops nothing for a model unload, +// and matching on the modelID prefix alone is the bug that orphaned deleted +// backends' processes. Stop is best-effort and idempotent, so accepting both is +// safe here; delete stays strict because its identifier is unambiguously a +// backend name. +func (s *backendSupervisor) resolveStopTargets(id string) []string { + keys := s.resolveProcessKeysForBackend(s.backendIdentity(id)) + seen := make(map[string]struct{}, len(keys)) + for _, k := range keys { + seen[k] = struct{}{} + } + for _, k := range s.resolveProcessKeys(id) { + if _, dup := seen[k]; !dup { + seen[k] = struct{}{} + keys = append(keys, k) + } + } + return keys +} + +// processMatchesBackend reports whether the process under `key` may be reused +// to serve `backend`. The install fast path returns any live process for a +// (model, replica) slot; without this check a slot whose backend was deleted +// and replaced handed the new install the deleted backend's port, so the load +// ran against a directory that no longer exists. +// +// A process with no recorded backendName predates this field and is accepted: +// treating it as a mismatch would restart every running backend once on +// rollout. +func (s *backendSupervisor) processMatchesBackend(key, backend string) bool { + s.mu.Lock() + bp, ok := s.processes[key] + if !ok { + s.mu.Unlock() + return false + } + recorded := bp.backendName + s.mu.Unlock() + + if recorded == "" || recorded == backend { + return true + } + // Names differ textually — they may still be alias and concrete for the + // same on-disk backend, which is a legitimate reuse. + _, match := s.backendIdentity(backend)[recorded] + return match +} + // stopBackend stops the backend process(es) matching the given identifier. // Accepts a bare modelID (stops every replica) or a full processKey // (stops just that replica). func (s *backendSupervisor) stopBackend(id string, force bool) { for _, key := range s.resolveProcessKeys(id) { - s.stopBackendExact(key, force) + if err := s.stopBackendExact(key, force); err != nil { + xlog.Error("Failed to stop backend process", "processKey", key, "error", err) + } } } @@ -232,10 +393,15 @@ func (s *backendSupervisor) stopBackend(id string, force bool) { // entry as stopping under the lock, then runs Free() and proc.Stop() after // release so network I/O cannot block other supervisor operations. The entry // and port remain reserved until process termination completes. -func (s *backendSupervisor) stopBackendExact(key string, force bool) { +// +// Returns an error only when a process was found and is still alive after the +// stop attempt — a missing key is a no-op, not a failure. Callers that report +// success to an operator (backend.delete) must not claim success while the +// process keeps serving requests from a directory they just removed. +func (s *backendSupervisor) stopBackendExact(key string, force bool) error { bp := s.beginBackendStop(key) if bp == nil { - return + return nil } if !force { @@ -248,12 +414,12 @@ func (s *backendSupervisor) stopBackendExact(key string, force bool) { cancel() } - xlog.Info("Stopping backend process", "backend", key, "addr", bp.addr, "force", force) + xlog.Info("Stopping backend process", "backend", key, "addr", bp.addr, "force", force, "backendName", bp.backendName) stopErr := bp.proc.Stop() if stopErr != nil { xlog.Error("Error stopping backend process", "backend", key, "error", stopErr) } - s.finishBackendStop(key, bp, stopErr) + return s.finishBackendStop(key, bp, stopErr) } // beginBackendStop reserves both the process entry and its port while network @@ -269,25 +435,28 @@ func (s *backendSupervisor) beginBackendStop(key string) *backendProcess { return bp } -func (s *backendSupervisor) finishBackendStop(key string, bp *backendProcess, stopErr error) { +// finishBackendStop returns a non-nil error when the process survived the stop +// attempt, so callers can distinguish a real reap from a failed one. +func (s *backendSupervisor) finishBackendStop(key string, bp *backendProcess, stopErr error) error { // Keep the process and port reserved until termination completes. Recycling // the port before this point can start a second backend on an address still // owned by the stuck process. s.mu.Lock() defer s.mu.Unlock() if current, exists := s.processes[key]; !exists || current != bp { - return + return nil } if stopErr != nil && bp.proc.IsAlive() { bp.stopping = false - return + return fmt.Errorf("stopping backend process %s: %w", key, stopErr) } delete(s.processes, key) if bp.port <= 0 { xlog.Error("Cannot recycle backend port: process has invalid recorded port", "backend", key, "addr", bp.addr, "port", bp.port) - return + return nil } s.freePorts = append(s.freePorts, bp.port) + return nil } // stopAllBackends stops all running backend processes. @@ -303,8 +472,12 @@ func (s *backendSupervisor) stopAllBackends(force bool) { // isRunning returns whether at least one backend process matching the given // identifier is currently running. Accepts a bare modelID (matches any -// replica) or a full processKey (exact match). Callers like the -// backend.delete pre-check rely on the bare-name path. +// replica) or a full processKey (exact match). +// +// It does NOT accept a backend name: backend names never prefix a +// modelID#replica key. Callers holding a backend name must go through +// resolveProcessKeysForBackend — assuming otherwise here is what left deleted +// backends' processes running. func (s *backendSupervisor) isRunning(id string) bool { keys := s.resolveProcessKeys(id) if len(keys) == 0 { diff --git a/pkg/model/process.go b/pkg/model/process.go index 33eb6377e..524d675a2 100644 --- a/pkg/model/process.go +++ b/pkg/model/process.go @@ -20,10 +20,18 @@ import ( var forceBackendShutdown bool = os.Getenv("LOCALAI_FORCE_BACKEND_SHUTDOWN") == "true" var ( - modelNotFoundErr = errors.New("model not found") + // ErrModelNotFound reports that a model is not loaded. Exported so HTTP + // handlers can map it to 404 instead of a blanket 500. + ErrModelNotFound = errors.New("model not found") // ErrModelBusy indicates that a graceful shutdown context ended while // requests were still in flight. ErrModelBusy = errors.New("model is still busy") + + // ErrRemoteModelNotLoaded is returned by a RemoteModelUnloader when no + // node in the cluster has the model loaded. It exists so the local store + // miss and the cluster-wide miss stay distinguishable: only when BOTH are + // empty may we tell the operator the model is not loaded. + ErrRemoteModelNotLoaded = errors.New("model not loaded on any node") ) const ( @@ -33,6 +41,16 @@ const ( busyPollInterval = 100 * time.Millisecond ) +// unloadRemote asks the remote unloader to stop `s` on whichever node holds +// it, preferring the context-aware extension so a forced shutdown and the +// caller's deadline both survive the distributed boundary. +func unloadRemote(ctx context.Context, u RemoteModelUnloader, s string, force bool) error { + if contextUnloader, ok := u.(RemoteModelContextUnloader); ok { + return contextUnloader.UnloadRemoteModelContext(ctx, s, force) + } + return u.UnloadRemoteModel(s) +} + // deleteProcess stops and removes a backend. The force flag trades a graceful // shutdown for a prompt one and is meant for the watchdog's busy-killer: a // backend that has been busy past the watchdog timeout may be stuck in an @@ -58,8 +76,26 @@ func (ml *ModelLoader) deleteProcess(ctx context.Context, s string, force bool) model, ok := store.Get(s) if !ok { + // A local-store miss is not proof the model is unloaded. In + // distributed mode the model runs on a worker and the authoritative + // record is the shared node registry: any frontend replica that did + // not itself serve the model (load balancer picked a peer, or this + // replica restarted) has no local entry. Returning "model not found" + // here reported a model that was demonstrably running as absent, and + // left its backend process untouched. + if remoteUnloader != nil { + xlog.Debug("Model not in local store; asking the remote unloader", "model", s) + if err := unloadRemote(ctx, remoteUnloader, s, force); err != nil { + if errors.Is(err, ErrRemoteModelNotLoaded) { + // Absent locally AND cluster-wide: genuinely not loaded. + return ErrModelNotFound + } + return err + } + return nil + } xlog.Debug("Model not found", "model", s) - return modelNotFoundErr + return ErrModelNotFound } if !force { @@ -114,11 +150,7 @@ func (ml *ModelLoader) deleteProcess(ctx context.Context, s string, force bool) var unloadErr error if remoteUnloader != nil { xlog.Debug("Delegating model unload to remote unloader", "model", s) - if contextUnloader, ok := remoteUnloader.(RemoteModelContextUnloader); ok { - unloadErr = contextUnloader.UnloadRemoteModelContext(ctx, s, force) - } else { - unloadErr = remoteUnloader.UnloadRemoteModel(s) - } + unloadErr = unloadRemote(ctx, remoteUnloader, s, force) if unloadErr != nil { xlog.Warn("Remote model unload failed", "model", s, "error", unloadErr) } diff --git a/pkg/model/remote_shutdown_test.go b/pkg/model/remote_shutdown_test.go new file mode 100644 index 000000000..b857982b6 --- /dev/null +++ b/pkg/model/remote_shutdown_test.go @@ -0,0 +1,72 @@ +package model_test + +import ( + "github.com/mudler/LocalAI/pkg/model" + "github.com/mudler/LocalAI/pkg/system" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +// fakeRemoteUnloader records the models it was asked to unload so the specs +// can assert the remote path was actually taken (not merely that no error +// surfaced). +type fakeRemoteUnloader struct { + called []string + err error +} + +func (f *fakeRemoteUnloader) UnloadRemoteModel(modelName string) error { + f.called = append(f.called, modelName) + return f.err +} + +// In distributed mode the authoritative record of "is this model loaded" is +// the shared node registry, not this replica's in-memory store. A frontend +// replica that never served the model itself (load balancer picked another +// replica, or this one restarted) has no local entry, so ShutdownModel +// short-circuited on a local-store miss and reported "model not found" for a +// model that was demonstrably running on a worker — while the remote unload +// path it documents was never reached. +var _ = Describe("ShutdownModel in distributed mode", func() { + var ( + modelLoader *model.ModelLoader + unloader *fakeRemoteUnloader + ) + + BeforeEach(func() { + systemState, err := system.GetSystemState(system.WithModelPath(GinkgoT().TempDir())) + Expect(err).ToNot(HaveOccurred()) + modelLoader = model.NewModelLoader(systemState) + unloader = &fakeRemoteUnloader{} + }) + + It("delegates to the remote unloader when the model is not in the local store", func() { + modelLoader.SetRemoteUnloader(unloader) + + err := modelLoader.ShutdownModel("longcat-video-avatar-1.5") + + Expect(unloader.called).To(ConsistOf("longcat-video-avatar-1.5"), + "a model absent locally may still be loaded on a worker; the remote unloader must be consulted") + Expect(err).ToNot(HaveOccurred(), + "stopping a model that is running on a worker must succeed, not report 'model not found'") + }) + + It("reports not-found only after the remote unloader confirms no node has it", func() { + unloader.err = model.ErrRemoteModelNotLoaded + modelLoader.SetRemoteUnloader(unloader) + + err := modelLoader.ShutdownModel("never-loaded") + + Expect(unloader.called).To(ConsistOf("never-loaded"), + "the registry must be consulted before declaring a model not found") + Expect(err).To(HaveOccurred(), + "a genuinely absent model must still error — silently succeeding would hide a failed stop") + }) + + It("still reports not-found when no remote unloader is configured", func() { + // Single-node behavior must be unchanged. + err := modelLoader.ShutdownModel("never-loaded") + Expect(err).To(HaveOccurred()) + }) +}) diff --git a/pkg/model/watchdog.go b/pkg/model/watchdog.go index efe400e41..3a604668a 100644 --- a/pkg/model/watchdog.go +++ b/pkg/model/watchdog.go @@ -1,6 +1,7 @@ package model import ( + "errors" "slices" "sync" "time" @@ -899,7 +900,7 @@ func (wd *WatchDog) evictLRUModel() { if wasBusy { shutdown = wd.pm.ShutdownModelForce } - if err := shutdown(lruModel.model); err != nil && err != modelNotFoundErr { + if err := shutdown(lruModel.model); 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