fix(distributed): converge model configs across frontends (#12558)

* fix(galleryop): announce model changes before the preload

After a gallery install or delete, the replica that ran it replaced its
config loader, then preloaded every installed model, and only then
published the models invalidation. The preload does remote lookups and
checksums for each model, so on a large models directory peers learned
about the change minutes after the originator listed it. When the
preload failed or the operation was cancelled, the event was never sent.

Publish the invalidation, and apply the delete lifecycle, as soon as
the loader holds the new set. The preload still runs afterwards with
its own error handling. Its failure is reported on the operation, but
it no longer rolls back a deletion that peers have already applied.

Assisted-by: Claude Code:claude-opus-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(distributed): resync model configs from the models directory

Frontends refresh their model configs only when a models invalidation
arrives on NATS. NATS keeps no history, so a frontend that is
disconnected when the message is published never applies the change.
It keeps serving the old config, for example an alias that points at
the previous model, until some later change happens to touch it.

Each frontend now reruns the peer reconcile against the shared models
directory after every NATS reconnect, and every
--model-config-resync-interval (default 30s) when a config file
changed. The pass names no model, so only models whose file changed
get a revision transition, and an unchanged directory costs one read
of the config files.

The reconcile replaced the whole loader with a parse of the models
directory, which dropped models loaded with --config-file and
published a deletion revision for them. Configs defined outside the
directory are now kept, both there and after a gallery install.

Assisted-by: Claude Code:claude-opus-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(nodes): stop stale frontends from retargeting alias rules

A scheduling rule keyed by an alias derives its target from the alias
mapping of the frontend that reads it. Every frontend keeps its own
copy of the model configs, so after an alias is repointed a frontend
that has not reloaded it still resolves the old target. Two frontends
then rewrote the rule's stored target_model against each other on
alternate reconciler ticks, and the outdated one scaled up the model
the alias used to point at.

The registry already records the accepted config revision of each
model. A frontend now derives a rule's target from its own alias
mapping only when its config revision for the rule's name matches
that record. Otherwise it keeps the stored target_model: it neither
writes the column nor reconciles replicas of the old target. The check
reads the database only for a rule whose stored and derived targets
differ. With no accepted revision on record, the old behaviour stays.

Assisted-by: Claude Code:claude-opus-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

---------

Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
This commit is contained in:
mudler-agentandEttore Di Giacinto authored and GitHub committed 2026-10-08 09:56:12 +02:00
1 parent 89a0955b5b
commit 895d50385f
16 files changed
+1005 -17

No files matched your search

+14
View File
@@ -286,6 +286,10 @@ func New(opts ...config.AppOption) (*Application, error) {
// revisionStore is built inside the distributed block below but used after
// the model configs are loaded, so it is declared out here.
var revisionStore modeladmin.RevisionStore
// modelConfigResync is built there too, and started only once the model
// configs are loaded: a pass against an empty loader would treat every
// model as new.
var modelConfigResync *modeladmin.DirectoryResync
distSvc, err := initDistributed(options, application.authDB, application.ModelConfigLoader(),
&failoverPinnedResolver{base: application.ModelConfigLoader(), fm: application.failoverManager})
@@ -399,6 +403,7 @@ func New(opts ...config.AppOption) (*Application, error) {
// Captured here, used after the model configs are loaded below: the
// resync reads the loader, which is still empty at this point.
revisionStore = modeladmin.NewRevisionStore(distSvc.Registry, modelRevisionLifecycle)
modelConfigResync = modeladmin.NewDirectoryResync(application.ModelConfigLoader(), sys.Model.ModelsPath, modelRevisionLifecycle, cfgLoaderOpts...)
gs.OnModelsChanged = func(evt messaging.CacheInvalidateEvent) {
// ApplyRemoteChange honors the op: a "delete" prunes the element
// (a reload-from-path is additive and cannot drop it), anything
@@ -509,6 +514,15 @@ func New(opts ...config.AppOption) (*Application, error) {
xlog.Error("error downloading models", "error", err)
}
// Catch up on model config changes whose invalidation this frontend
// missed. NATS keeps no history, so a change published while this
// frontend was disconnected is otherwise never applied here. Every
// frontend runs its own pass against the shared models directory; the
// pass is idempotent, so no leader is needed.
if modelConfigResync != nil && distSvc != nil {
modelConfigResync.Start(options.Context, options.Distributed.ModelConfigResyncIntervalOrDefault(), distSvc.Nats)
}
if options.PreloadJSONModels != "" {
if err := galleryop.ApplyGalleryFromString(options.SystemState, application.ModelLoader(), options.EnforcePredownloadScans, options.AutoloadBackendGalleries, options.Galleries, options.BackendGalleries, options.PreloadJSONModels, options.RequireBackendIntegrity, gallery.WithArtifactMaterializer(options.ModelArtifactMaterializer)); err != nil {
return nil, err
+8
View File
@@ -185,6 +185,7 @@ type RunCMD struct {
ModelLoadWait string `env:"LOCALAI_MODEL_LOAD_WAIT" help:"How long an inference request waits for a model that is still cold-loading onto a worker before it is answered with 503, a Retry-After header and live staging progress (default 60s). The request is served the moment the model becomes ready, so a model already most of the way staged needs no client retry. Set to 0 to wait as long as the load takes — only safe when no ingress or load balancer with an idle timeout sits in front." group:"distributed"`
StaleNodeThreshold string `env:"LOCALAI_STALE_NODE_THRESHOLD" help:"How long a worker node may go without a durable heartbeat before the health monitor marks it offline (default 5m). Because a beat that only carries a fresher timestamp is held back by --node-heartbeat-checkpoint, this must stay comfortably wider than that interval; raise both together. Dead-node detection through the per-model gRPC health check and through request-time failure is unaffected by this knob." group:"distributed"`
NodeHeartbeatCheckpoint string `env:"LOCALAI_NODE_HEARTBEAT_CHECKPOINT" help:"Minimum gap between durable heartbeat writes for a worker node (default 60s). A beat that only carries a fresher timestamp is dropped until this interval elapses; every field is compared against the value last written, so a node's first beat, a changed total VRAM/total disk/GPU vendor, and a free VRAM/RAM/disk reading that has moved more than 256 MiB from the written value all still write immediately, and a node that is not active is never suppressed. Set below the worker heartbeat interval to write on every beat." group:"distributed"`
ModelConfigResyncInterval string `env:"LOCALAI_MODEL_CONFIG_RESYNC_INTERVAL" help:"How often each frontend compares its model configs with the shared models directory (default 30s). Frontends normally learn about a model install, edit or removal from a message on NATS; this pass catches a change whose message a frontend missed, for example during a NATS reconnect, so an old config is served for at most this long. A pass reads only the config files and skips the parse when none of them changed." group:"distributed"`
NatsAccountSeed string `env:"LOCALAI_NATS_ACCOUNT_SEED" help:"NATS account signing seed (SU...) used to mint per-node worker JWTs at registration" group:"distributed"`
NatsServiceJWT string `env:"LOCALAI_NATS_SERVICE_JWT" help:"NATS user JWT for the frontend (and agent workers) to publish control-plane messages" group:"distributed"`
NatsServiceSeed string `env:"LOCALAI_NATS_SERVICE_SEED" help:"NATS user signing seed (SU...) paired with LOCALAI_NATS_SERVICE_JWT" group:"distributed"`
@@ -423,6 +424,13 @@ func (r *RunCMD) Run(ctx *cliContext.Context) error {
}
opts = append(opts, config.WithNodeHeartbeatCheckpoint(d))
}
if r.ModelConfigResyncInterval != "" {
d, err := parseDistributedDuration("LOCALAI_MODEL_CONFIG_RESYNC_INTERVAL", r.ModelConfigResyncInterval)
if err != nil {
return err
}
opts = append(opts, config.WithModelConfigResyncInterval(d))
}
if r.RegistrationToken != "" {
opts = append(opts, config.WithRegistrationToken(r.RegistrationToken))
}
+29
View File
@@ -64,6 +64,10 @@ type DistributedConfig struct {
HealthCheckInterval time.Duration // Health monitor check interval (default 15s)
StaleNodeThreshold time.Duration // Time before a node is considered stale (default 5m)
NodeHeartbeatCheckpoint time.Duration // Minimum gap between durable heartbeat writes (default 60s, 0 = every beat)
// ModelConfigResyncInterval is how often a frontend compares its model
// configs with the shared models directory, to catch a change whose
// invalidation event it missed (default 30s).
ModelConfigResyncInterval time.Duration
// DisablePerModelHealthCheck turns off the health monitor's per-model
// gRPC probe. When enabled (the default), the monitor pings each model's
// gRPC address and removes stale node_models rows whose backend has
@@ -181,6 +185,9 @@ func (c DistributedConfig) Validate() error {
return fmt.Errorf("%s must not be negative", name)
}
}
if c.ModelConfigResyncInterval < 0 {
return fmt.Errorf("%s must not be negative", FlagModelConfigResyncInterval)
}
return nil
}
@@ -338,6 +345,14 @@ func WithModelLoadWait(d time.Duration) AppOption {
}
}
// WithModelConfigResyncInterval sets how often a frontend resyncs its model
// configs from the shared models directory. Zero means the default.
func WithModelConfigResyncInterval(d time.Duration) AppOption {
return func(o *ApplicationConfig) {
o.Distributed.ModelConfigResyncInterval = d
}
}
// WithStaleNodeThreshold sets how long a node may go without a durable
// heartbeat before the health monitor marks it offline. It has to be raised
// alongside WithNodeHeartbeatCheckpoint: a checkpoint interval wider than this
@@ -431,6 +446,14 @@ const (
FlagDiskHeadroomCheck = "distributed-disk-headroom-check"
)
// FlagModelConfigResyncInterval names the model config resync interval.
const FlagModelConfigResyncInterval = "model-config-resync-interval"
// DefaultModelConfigResyncInterval bounds how long a frontend that missed a
// models invalidation serves an old model config. It matches the tick of the
// replica reconciler, which acts on those configs.
const DefaultModelConfigResyncInterval = 30 * time.Second
// Defaults for distributed timeouts.
const (
DefaultMCPToolTimeout = 360 * time.Second
@@ -544,6 +567,12 @@ func (c DistributedConfig) HealthCheckIntervalOrDefault() time.Duration {
}
// StaleNodeThresholdOrDefault returns the configured threshold or the default.
// ModelConfigResyncIntervalOrDefault returns the configured interval or the
// default.
func (c DistributedConfig) ModelConfigResyncIntervalOrDefault() time.Duration {
return cmp.Or(c.ModelConfigResyncInterval, DefaultModelConfigResyncInterval)
}
func (c DistributedConfig) StaleNodeThresholdOrDefault() time.Duration {
return cmp.Or(c.StaleNodeThreshold, DefaultStaleNodeThreshold)
}
+18
View File
@@ -88,6 +88,7 @@ var _ = Describe("DistributedConfig flag-name constants", func() {
Entry("backend install timeout", config.FlagBackendInstallTimeout, "backend-install-timeout"),
Entry("backend upgrade timeout", config.FlagBackendUpgradeTimeout, "backend-upgrade-timeout"),
Entry("model load timeout", config.FlagModelLoadTimeout, "model-load-timeout"),
Entry("model config resync interval", config.FlagModelConfigResyncInterval, "model-config-resync-interval"),
)
})
@@ -104,6 +105,23 @@ var _ = Describe("DistributedConfig.Validate negative-duration errors", func() {
Expect(err.Error()).To(ContainSubstring("must not be negative"))
})
It("rejects a negative ModelConfigResyncInterval with the flag name in the error", func() {
c := config.DistributedConfig{
Enabled: true,
NatsURL: "nats://localhost:4222",
ModelConfigResyncInterval: -1 * time.Second,
}
err := c.Validate()
Expect(err).To(HaveOccurred())
Expect(err.Error()).To(ContainSubstring(config.FlagModelConfigResyncInterval))
Expect(err.Error()).To(ContainSubstring("must not be negative"))
})
It("defaults the model config resync interval", func() {
Expect(config.DistributedConfig{}.ModelConfigResyncIntervalOrDefault()).To(Equal(config.DefaultModelConfigResyncInterval))
Expect(config.DistributedConfig{ModelConfigResyncInterval: time.Minute}.ModelConfigResyncIntervalOrDefault()).To(Equal(time.Minute))
})
It("rejects a negative BackendUpgradeTimeout with the flag name in the error", func() {
c := config.DistributedConfig{
Enabled: true,
+76
View File
@@ -0,0 +1,76 @@
package config
import (
"crypto/sha256"
"fmt"
"path/filepath"
)
// DefinedOutside reports whether this configuration was read from a file that
// is not directly inside dir, for example the list given with --config-file.
// A configuration with no recorded file counts as one from dir, so callers that
// replace dir-backed configurations keep the behaviour they had before.
func (c *ModelConfig) DefinedOutside(dir string) bool {
if c.modelConfigFile == "" {
return false
}
return !sameDir(filepath.Dir(c.modelConfigFile), dir)
}
func sameDir(a, b string) bool {
absA, errA := filepath.Abs(a)
absB, errB := filepath.Abs(b)
if errA != nil || errB != nil {
return filepath.Clean(a) == filepath.Clean(b)
}
return absA == absB
}
// MergeDirectorySnapshot returns the configuration set that results from
// replacing every configuration in current that was read from dir with
// snapshot, a fresh parse of dir.
//
// A configuration defined outside dir is kept, and it wins over a snapshot
// entry with the same name, as it does at startup, where --config-file is read
// after the models directory. A parse of dir cannot see these files, so
// replacing the whole set with it would drop them.
func MergeDirectorySnapshot(current, snapshot []ModelConfig, dir string) []ModelConfig {
outside := make(map[string]ModelConfig)
for _, cfg := range current {
if cfg.DefinedOutside(dir) {
outside[cfg.Name] = cfg
}
}
merged := make([]ModelConfig, 0, len(snapshot)+len(outside))
for _, cfg := range snapshot {
if _, shadowed := outside[cfg.Name]; !shadowed {
merged = append(merged, cfg)
}
}
for _, cfg := range outside {
merged = append(merged, cfg)
}
return merged
}
// DeletedModelConfigRevision is the config revision recorded for a model that
// no longer has a config. Every frontend derives the same value, so a deletion
// compares equal however it reached a frontend.
func DeletedModelConfigRevision(modelName string) string {
return fmt.Sprintf("%x", sha256.Sum256([]byte("deleted\x00"+modelName)))
}
// ConfigRevisionOf returns the revision of this loader's view of name: the
// revision stamped when its config was parsed, or DeletedModelConfigRevision
// when this loader has no config for name. It returns "" for a config that
// was never stamped, whose revision is unknown.
//
// Comparing it with the revision the cluster accepted for name tells whether
// this frontend's view of the model is current.
func (bcl *ModelConfigLoader) ConfigRevisionOf(name string) string {
cfg, ok := bcl.GetModelConfig(name)
if !ok {
return DeletedModelConfigRevision(name)
}
return cfg.PersistedConfigRevision()
}
@@ -0,0 +1,60 @@
package config_test
import (
"os"
"path/filepath"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"github.com/mudler/LocalAI/core/config"
)
var _ = Describe("MergeDirectorySnapshot", func() {
It("replaces directory configs and keeps configs defined outside the directory", func() {
dir := GinkgoT().TempDir()
Expect(os.WriteFile(filepath.Join(dir, "kept.yaml"), []byte("name: kept\nbackend: llama-cpp\n"), 0644)).To(Succeed())
Expect(os.WriteFile(filepath.Join(dir, "removed.yaml"), []byte("name: removed\nbackend: llama-cpp\n"), 0644)).To(Succeed())
cfgFile := filepath.Join(GinkgoT().TempDir(), "models.yaml")
Expect(os.WriteFile(cfgFile, []byte("- name: external\n backend: llama-cpp\n- name: kept\n backend: from-config-file\n"), 0644)).To(Succeed())
current := config.NewModelConfigLoader(dir)
Expect(current.LoadModelConfigsFromPath(dir)).To(Succeed())
Expect(current.LoadMultipleModelConfigsSingleFile(cfgFile)).To(Succeed())
Expect(os.Remove(filepath.Join(dir, "removed.yaml"))).To(Succeed())
Expect(os.WriteFile(filepath.Join(dir, "added.yaml"), []byte("name: added\nbackend: llama-cpp\n"), 0644)).To(Succeed())
snapshot := config.NewModelConfigLoader(dir)
Expect(snapshot.LoadModelConfigsFromPathStrict(dir)).To(Succeed())
merged := config.MergeDirectorySnapshot(current.GetAllModelsConfigs(), snapshot.GetAllModelsConfigs(), dir)
byName := map[string]string{}
for _, cfg := range merged {
byName[cfg.Name] = cfg.Backend
}
Expect(byName).To(Equal(map[string]string{
"added": "llama-cpp",
"external": "llama-cpp",
"kept": "from-config-file",
}))
})
It("treats a config with no recorded file as one from the directory", func() {
cfg := config.ModelConfig{Name: "synthesized"}
Expect(cfg.DefinedOutside(GinkgoT().TempDir())).To(BeFalse())
})
})
var _ = Describe("ModelConfigLoader.ConfigRevisionOf", func() {
It("returns the stamped revision of a loaded config and the deletion revision of an absent one", func() {
dir := GinkgoT().TempDir()
Expect(os.WriteFile(filepath.Join(dir, "present.yaml"), []byte("name: present\nbackend: llama-cpp\n"), 0644)).To(Succeed())
loader := config.NewModelConfigLoader(dir)
Expect(loader.LoadModelConfigsFromPath(dir)).To(Succeed())
revision, err := loader.RevisionForPath("present", dir)
Expect(err).ToNot(HaveOccurred())
Expect(loader.ConfigRevisionOf("present")).To(Equal(revision))
Expect(loader.ConfigRevisionOf("absent")).To(Equal(config.DeletedModelConfigRevision("absent")))
})
})
@@ -195,8 +195,6 @@ var _ = Describe("model deletion revision lifecycle", func() {
manager.afterDelete = func() error {
return os.WriteFile(filepath.Join(dir, "broken.yaml"), []byte("name: ["), 0644)
}
case "preload":
Expect(os.WriteFile(filepath.Join(dir, "survivor.yaml"), []byte("name: survivor\nbackend: transformers\nartifacts:\n - name: model\n target: model\n source: {type: huggingface, repo: owner/repo}\n"), 0644)).To(Succeed())
case "lifecycle":
lifecycle.err = errors.New("injected lifecycle failure")
}
@@ -225,7 +223,46 @@ var _ = Describe("model deletion revision lifecycle", func() {
Expect(ok).To(BeTrue())
},
Entry("when authoritative parsing fails", "parse"),
Entry("when preload fails", "preload"),
Entry("when lifecycle publication fails", "lifecycle"),
)
// The preload walks every remaining model, not the deleted one, so its
// failure says nothing about the deletion. Once the deletion is announced
// to peers it must stay applied here too, or this replica would list a
// model that every peer has already dropped.
It("keeps and announces a committed deletion when the preload fails", func() {
dir := GinkgoT().TempDir()
materializer := &rejectingDeleteMaterializer{}
appConfig := &config.ApplicationConfig{
SystemState: &system.SystemState{Model: system.Model{ModelsPath: dir}},
ModelArtifactMaterializer: materializer,
}
loader := config.NewModelConfigLoader(dir, config.WithArtifactMaterializer(materializer))
configPath := filepath.Join(dir, "doomed.yaml")
Expect(os.WriteFile(configPath, []byte("name: doomed\nbackend: llama-cpp\n"), 0640)).To(Succeed())
Expect(os.WriteFile(filepath.Join(dir, gallery.GalleryFileName("doomed")), []byte("files: []\n"), 0600)).To(Succeed())
Expect(os.WriteFile(filepath.Join(dir, "survivor.yaml"), []byte("name: survivor\nbackend: transformers\nartifacts:\n - name: model\n target: model\n source: {type: huggingface, repo: owner/repo}\n"), 0644)).To(Succeed())
Expect(loader.LoadModelConfigsFromPath(dir, appConfig.ToConfigLoaderOptions()...)).To(Succeed())
service := NewGalleryService(appConfig, nil)
bus := &countingMessagingClient{}
service.SetNATSClient(bus)
service.SetModelManager(&realDeletingManager{state: appConfig.SystemState})
lifecycle := &deleteRevisionLifecycle{}
service.SetModelRevisionLifecycle(lifecycle)
op := &ManagementOp[gallery.GalleryModel, gallery.ModelConfig]{
ID: "delete-operation", GalleryElementName: "doomed", Delete: true, Context: context.Background(),
}
err := service.modelHandler(op, loader, appConfig.SystemState)
Expect(err).To(MatchError(ContainSubstring("preloading model files failed")))
Expect(materializer.calls).To(BeNumerically(">", 0), "precondition: the preload ran and failed")
Expect(lifecycle.applied).To(BeTrue())
Expect(bus.subjects).To(ContainElement(messaging.SubjectCacheInvalidateModels))
Expect(configPath).NotTo(BeAnExistingFile())
_, ok := loader.GetModelConfig("doomed")
Expect(ok).To(BeFalse())
_, ok = loader.GetModelConfig("survivor")
Expect(ok).To(BeTrue())
})
})
+16 -7
View File
@@ -183,15 +183,11 @@ func (g *GalleryService) modelHandlerLocked(op *ManagementOp[gallery.GalleryMode
if err != nil {
return err
}
cl.ReplaceModelConfigs(authoritative.GetAllModelsConfigs())
err = cl.PreloadWithContext(operationCtx, systemState.Model.ModelsPath)
if err != nil {
return err
}
cl.ReplaceModelConfigs(config.MergeDirectorySnapshot(cl.GetAllModelsConfigs(), authoritative.GetAllModelsConfigs(), systemState.Model.ModelsPath))
// Lifecycle publication is the irreversible boundary. File mutation,
// authoritative parsing, loader replacement, and preload have all completed,
// so no later failure can roll local configuration back behind an accepted
// authoritative parsing and loader replacement have all completed, so no
// later failure can roll local configuration back behind an accepted
// registry revision.
if op.Delete && g.modelRevisionLifecycle != nil {
pending, lifecycleErr := g.modelRevisionLifecycle.ApplyConfigRevisions(operationCtx, []config.ModelConfigRevisionTransition{{
@@ -210,6 +206,12 @@ func (g *GalleryService) modelHandlerLocked(op *ManagementOp[gallery.GalleryMode
// authoritative replacement above already covered THIS replica; without
// this broadcast a chat completion routed by the load balancer to a peer
// would fail to find a model just installed.
//
// Publish before the preload below. This replica already lists the new
// set, and the preload walks every installed model (remote HEAD requests,
// checksums of existing files), which can take minutes on a large models
// directory. Publishing after it left peers behind for that long, and a
// failed or cancelled preload never published at all.
op2 := "install"
if op.Delete {
op2 = "delete"
@@ -220,6 +222,13 @@ func (g *GalleryService) modelHandlerLocked(op *ManagementOp[gallery.GalleryMode
ConfigRevision: configRevision,
})
// The configuration change is committed and announced, so a preload
// failure no longer rolls it back: it is reported on the operation as an
// error, and the model stays listed on every replica, as it does on this one.
if err := cl.PreloadWithContext(operationCtx, systemState.Model.ModelsPath); err != nil {
return fmt.Errorf("model configuration applied, but preloading model files failed: %w", err)
}
legacyCoalescer.Close()
g.UpdateStatus(op.ID,
&OpStatus{
@@ -0,0 +1,160 @@
package galleryop
import (
"context"
"errors"
"os"
"path/filepath"
"sync"
"time"
"github.com/mudler/LocalAI/core/config"
"github.com/mudler/LocalAI/core/gallery"
"github.com/mudler/LocalAI/core/services/messaging"
"github.com/mudler/LocalAI/pkg/modelartifacts"
"github.com/mudler/LocalAI/pkg/system"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
)
// blockingMaterializer holds the preload of an already-installed model open
// until the spec releases it, so a spec can look at both replicas while the
// originator is still inside the preload.
type blockingMaterializer struct {
entered chan struct{}
release chan struct{}
err error
once sync.Once
}
func (m *blockingMaterializer) Ensure(ctx context.Context, _ string, spec modelartifacts.Spec) (modelartifacts.Result, error) {
m.once.Do(func() { close(m.entered) })
select {
case <-m.release:
case <-ctx.Done():
return modelartifacts.Result{}, ctx.Err()
}
if m.err != nil {
return modelartifacts.Result{}, m.err
}
return modelartifacts.Result{Spec: spec}, nil
}
// writingInstallManager stands in for a gallery install: it writes one new
// config file into the shared models directory.
type writingInstallManager struct{ dir string }
func (m *writingInstallManager) InstallModel(context.Context, *ManagementOp[gallery.GalleryModel, gallery.ModelConfig], ProgressCallback) error {
return os.WriteFile(filepath.Join(m.dir, "fresh.yaml"), []byte("name: fresh\nbackend: llama-cpp\n"), 0644)
}
func (m *writingInstallManager) DeleteModel(string) error { return nil }
// peerReloadingBus reloads a peer replica's loader from the shared directory
// whenever the models invalidation is published, the way a subscribed peer
// does in distributed mode.
type peerReloadingBus struct {
mu sync.Mutex
published int
onPublish func()
}
func (b *peerReloadingBus) Publish(subject string, _ any) error {
if subject != messaging.SubjectCacheInvalidateModels {
return nil
}
b.mu.Lock()
b.published++
b.mu.Unlock()
if b.onPublish != nil {
b.onPublish()
}
return nil
}
func (b *peerReloadingBus) count() int { b.mu.Lock(); defer b.mu.Unlock(); return b.published }
func (b *peerReloadingBus) Subscribe(string, func([]byte)) (messaging.Subscription, error) {
return nil, nil
}
func (b *peerReloadingBus) QueueSubscribe(string, string, func([]byte)) (messaging.Subscription, error) {
return nil, nil
}
func (b *peerReloadingBus) QueueSubscribeReply(string, string, func([]byte, func([]byte))) (messaging.Subscription, error) {
return nil, nil
}
func (b *peerReloadingBus) SubscribeReply(string, func([]byte, func([]byte))) (messaging.Subscription, error) {
return nil, nil
}
func (b *peerReloadingBus) Request(string, []byte, time.Duration) ([]byte, error) { return nil, nil }
func (b *peerReloadingBus) IsConnected() bool { return true }
func (b *peerReloadingBus) Close() {}
var _ = Describe("gallery install as seen by a peer replica", func() {
var (
dir string
appConfig *config.ApplicationConfig
materializer *blockingMaterializer
originator *config.ModelConfigLoader
peer *config.ModelConfigLoader
bus *peerReloadingBus
service *GalleryService
op *ManagementOp[gallery.GalleryModel, gallery.ModelConfig]
)
BeforeEach(func() {
dir = GinkgoT().TempDir()
appConfig = &config.ApplicationConfig{SystemState: &system.SystemState{Model: system.Model{ModelsPath: dir}}}
// An already-installed model whose preload is slow: it stands in for
// the checksums and remote lookups the preload does for every model.
Expect(os.WriteFile(filepath.Join(dir, "existing.yaml"), []byte(
"name: existing\nbackend: llama-cpp\nartifacts:\n - name: model\n target: model\n source: {type: huggingface, repo: owner/repo}\n"), 0644)).To(Succeed())
materializer = &blockingMaterializer{entered: make(chan struct{}), release: make(chan struct{})}
originator = config.NewModelConfigLoader(dir, config.WithArtifactMaterializer(materializer))
Expect(originator.LoadModelConfigsFromPath(dir)).To(Succeed())
// The peer learns about changes only through the broadcast.
peer = config.NewModelConfigLoader(dir)
Expect(peer.LoadModelConfigsFromPath(dir)).To(Succeed())
bus = &peerReloadingBus{onPublish: func() {
defer GinkgoRecover()
Expect(peer.LoadModelConfigsFromPath(dir)).To(Succeed())
}}
service = NewGalleryService(appConfig, nil)
service.SetNATSClient(bus)
service.SetModelManager(&writingInstallManager{dir: dir})
op = &ManagementOp[gallery.GalleryModel, gallery.ModelConfig]{
ID: "install-op", GalleryElementName: "localai@fresh", Context: context.Background(),
}
})
It("tells peers about the new model before the slow preload finishes", func() {
done := make(chan error, 1)
go func() { done <- service.modelHandler(op, originator, appConfig.SystemState) }()
Eventually(materializer.entered, 5*time.Second).Should(BeClosed())
// The preload of an unrelated model is still in progress here.
_, onOriginator := originator.GetModelConfig("fresh")
_, onPeer := peer.GetModelConfig("fresh")
published := bus.count()
close(materializer.release)
Expect(<-done).To(Succeed())
Expect(onOriginator).To(BeTrue(), "precondition: the originator already lists the new model")
Expect(published).To(Equal(1), "the originator listed the model before it published the invalidation")
Expect(onPeer).To(BeTrue(), "the peer does not list a model the originator already lists")
})
It("still tells peers when the preload fails, and reports the failure", func() {
materializer.err = errors.New("injected preload failure")
close(materializer.release)
err := service.modelHandler(op, originator, appConfig.SystemState)
Expect(err).To(MatchError(ContainSubstring("injected preload failure")))
Expect(bus.count()).To(Equal(1))
_, onOriginator := originator.GetModelConfig("fresh")
_, onPeer := peer.GetModelConfig("fresh")
Expect(onOriginator).To(BeTrue())
Expect(onPeer).To(BeTrue())
})
})
@@ -0,0 +1,148 @@
package modeladmin
import (
"context"
"crypto/sha256"
"encoding/hex"
"fmt"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"time"
"github.com/mudler/LocalAI/core/config"
"github.com/mudler/LocalAI/core/services/messaging"
"github.com/mudler/xlog"
)
// DirectoryResync brings one frontend's in-memory model configs back in line
// with the shared models directory when no invalidation event told it to.
//
// The models invalidation on the message bus is the fast path, but it is fire
// and forget: NATS keeps no history, so a frontend that is disconnected, or
// whose handler failed, never hears about the change again. DirectoryResync is
// the catch-up path. It reruns the same reconcile as ApplyRemoteChange, with no
// named element, so only models whose file actually changed get a revision
// transition, and a pass over an unchanged directory changes nothing.
//
// It keeps no state that other frontends need: every frontend reads the same
// shared directory, and the only local memory is a fingerprint of the config
// files it last reconciled, used to skip a pass when nothing changed.
type DirectoryResync struct {
loader *config.ModelConfigLoader
modelsPath string
lifecycle ModelRevisionLifecycle
opts []config.ConfigLoaderOption
mu sync.Mutex
reconciled string // fingerprint of the last successful pass
lastFailure string // fingerprint of the last failed pass, to log it once
}
// NewDirectoryResync returns a resync for loader against modelsPath. lifecycle
// may be nil; opts are the loader options used everywhere else.
func NewDirectoryResync(loader *config.ModelConfigLoader, modelsPath string, lifecycle ModelRevisionLifecycle, opts ...config.ConfigLoaderOption) *DirectoryResync {
return &DirectoryResync{loader: loader, modelsPath: modelsPath, lifecycle: lifecycle, opts: opts}
}
// Resync reconciles the loader with the models directory. Unless force is set,
// it does nothing when the config files are byte-for-byte the ones the last
// successful pass reconciled.
func (r *DirectoryResync) Resync(ctx context.Context, force bool) error {
r.mu.Lock()
defer r.mu.Unlock()
fingerprint, err := configFilesFingerprint(r.modelsPath)
if err != nil {
return err
}
if !force && fingerprint == r.reconciled {
return nil
}
if err := ApplyRemoteChange(ctx, r.loader, r.modelsPath, messaging.CacheInvalidateEvent{}, r.lifecycle, r.opts...); err != nil {
r.reconciled = ""
if fingerprint != r.lastFailure {
r.lastFailure = fingerprint
return err
}
// Same directory content, same failure: already reported.
xlog.Debug("Model config resync still failing", "error", err)
return nil
}
r.reconciled = fingerprint
r.lastFailure = ""
return nil
}
// Start runs a pass every interval until ctx is done, and a forced pass after
// each reconnect of bus when bus reports reconnects (the NATS client does;
// a nil bus or one without OnReconnect gets only the periodic pass). The first
// periodic pass runs after one interval, because startup has just loaded the
// directory.
func (r *DirectoryResync) Start(ctx context.Context, interval time.Duration, bus any) {
if reconnecting, ok := bus.(interface{ OnReconnect(func()) }); ok {
reconnecting.OnReconnect(func() {
// Off the NATS callback goroutine: a pass reads the disk and the
// database.
go func() {
if err := r.Resync(ctx, true); err != nil {
xlog.Warn("Failed to resync model configs after a NATS reconnect", "error", err)
}
}()
})
}
go r.run(ctx, interval)
}
func (r *DirectoryResync) run(ctx context.Context, interval time.Duration) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := r.Resync(ctx, false); err != nil {
xlog.Warn("Failed to resync model configs from the models directory", "error", err)
}
}
}
}
// configFilesFingerprint hashes the names and contents of the config files the
// loader reads from dir: top-level .yaml and .yml files that are not dotfiles.
// Reading a few small text files is cheap; parsing them is not, because a parse
// also inspects the model files they point at.
func configFilesFingerprint(dir string) (string, error) {
entries, err := os.ReadDir(dir)
if err != nil {
return "", fmt.Errorf("read models directory: %w", err)
}
names := make([]string, 0, len(entries))
for _, entry := range entries {
name := entry.Name()
ext := strings.ToLower(filepath.Ext(name))
if entry.IsDir() || strings.HasPrefix(name, ".") || (ext != ".yaml" && ext != ".yml") {
continue
}
names = append(names, name)
}
sort.Strings(names)
h := sha256.New()
for _, name := range names {
data, err := os.ReadFile(filepath.Join(dir, name))
if err != nil {
if os.IsNotExist(err) {
continue // removed between the listing and the read
}
return "", fmt.Errorf("read model config %q: %w", name, err)
}
sum := sha256.Sum256(data)
h.Write([]byte(name))
h.Write([]byte{0})
h.Write(sum[:])
}
return hex.EncodeToString(h.Sum(nil)), nil
}
@@ -0,0 +1,193 @@
package modeladmin
import (
"context"
"os"
"path/filepath"
"sync"
"time"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"github.com/mudler/LocalAI/core/config"
"github.com/mudler/LocalAI/core/services/messaging"
"github.com/mudler/LocalAI/core/services/testutil"
)
// lockedRevisionLifecycle records transitions and is safe to read while a
// reconnect pass runs on another goroutine.
type lockedRevisionLifecycle struct {
mu sync.Mutex
transitions []ModelRevisionTransition
}
func (l *lockedRevisionLifecycle) ApplyConfigRevisions(_ context.Context, transitions []ModelRevisionTransition) (int, error) {
l.mu.Lock()
defer l.mu.Unlock()
l.transitions = append(l.transitions, transitions...)
return 0, nil
}
func (l *lockedRevisionLifecycle) names() []string {
l.mu.Lock()
defer l.mu.Unlock()
names := make([]string, 0, len(l.transitions))
for _, t := range l.transitions {
names = append(names, t.ModelName)
}
return names
}
func writeModelFile(dir, name, body string) {
Expect(os.WriteFile(filepath.Join(dir, name+".yaml"), []byte(body), 0644)).To(Succeed())
}
func aliasTarget(loader *config.ModelConfigLoader, name string) string {
cfg, ok := loader.GetModelConfig(name)
if !ok {
return ""
}
return cfg.Alias
}
var _ = Describe("ApplyRemoteChange with models defined outside the models directory", func() {
It("keeps a --config-file model and publishes no deletion for it", func() {
dir := GinkgoT().TempDir()
cfgFile := filepath.Join(GinkgoT().TempDir(), "models.yaml")
Expect(os.WriteFile(cfgFile, []byte("- name: from-config-file\n backend: llama-cpp\n"), 0644)).To(Succeed())
loader := config.NewModelConfigLoader(dir)
Expect(loader.LoadMultipleModelConfigsSingleFile(cfgFile)).To(Succeed())
writeModelFile(dir, "fresh", "name: fresh\nbackend: llama-cpp\n")
lifecycle := &lockedRevisionLifecycle{}
Expect(ApplyRemoteChange(context.Background(), loader, dir, messaging.CacheInvalidateEvent{Element: "fresh", Op: "install"}, lifecycle)).To(Succeed())
_, ok := loader.GetModelConfig("from-config-file")
Expect(ok).To(BeTrue())
_, ok = loader.GetModelConfig("fresh")
Expect(ok).To(BeTrue())
Expect(lifecycle.names()).To(ConsistOf("fresh"))
})
It("lets a --config-file model win over a models directory file of the same name", func() {
dir := GinkgoT().TempDir()
cfgFile := filepath.Join(GinkgoT().TempDir(), "models.yaml")
Expect(os.WriteFile(cfgFile, []byte("- name: shared\n backend: from-config-file\n"), 0644)).To(Succeed())
writeModelFile(dir, "shared", "name: shared\nbackend: from-directory\n")
loader := config.NewModelConfigLoader(dir)
Expect(loader.LoadModelConfigsFromPath(dir)).To(Succeed())
Expect(loader.LoadMultipleModelConfigsSingleFile(cfgFile)).To(Succeed())
lifecycle := &lockedRevisionLifecycle{}
Expect(ApplyRemoteChange(context.Background(), loader, dir, messaging.CacheInvalidateEvent{Element: "shared", Op: "install"}, lifecycle)).To(Succeed())
cfg, ok := loader.GetModelConfig("shared")
Expect(ok).To(BeTrue())
Expect(cfg.Backend).To(Equal("from-config-file"))
Expect(lifecycle.names()).To(BeEmpty())
})
})
var _ = Describe("DirectoryResync", func() {
var (
dir string
originator *config.ModelConfigLoader
peer *config.ModelConfigLoader
lifecycle *lockedRevisionLifecycle
resync *DirectoryResync
)
BeforeEach(func() {
dir = GinkgoT().TempDir()
writeModelFile(dir, "model-a", "name: model-a\nbackend: llama-cpp\n")
writeModelFile(dir, "model-b", "name: model-b\nbackend: llama-cpp\n")
writeModelFile(dir, "stable", "name: stable\nalias: model-a\n")
originator = config.NewModelConfigLoader(dir)
Expect(originator.LoadModelConfigsFromPath(dir)).To(Succeed())
peer = config.NewModelConfigLoader(dir)
Expect(peer.LoadModelConfigsFromPath(dir)).To(Succeed())
lifecycle = &lockedRevisionLifecycle{}
resync = NewDirectoryResync(peer, dir, lifecycle)
})
// repointAlias changes the alias on the originator and announces it on
// bus. The peer only hears it if it is subscribed.
repointAlias := func(bus *testutil.FakeBus) {
writeModelFile(dir, "stable", "name: stable\nalias: model-b\n")
Expect(originator.LoadModelConfigsFromPath(dir)).To(Succeed())
Expect(bus.Publish(messaging.SubjectCacheInvalidateModels, messaging.CacheInvalidateEvent{Element: "stable", Op: "install"})).To(Succeed())
}
It("catches up on a change whose invalidation the peer missed", func() {
bus := testutil.NewFakeBus()
// The peer is not subscribed: the message is lost, as it is when the
// peer is disconnected from NATS at publish time.
repointAlias(bus)
Expect(aliasTarget(originator, "stable")).To(Equal("model-b"))
Expect(aliasTarget(peer, "stable")).To(Equal("model-a"), "precondition: the peer missed the change")
Expect(resync.Resync(context.Background(), false)).To(Succeed())
Expect(aliasTarget(peer, "stable")).To(Equal("model-b"))
Expect(lifecycle.names()).To(ConsistOf("stable"), "only the changed model gets a revision transition")
})
It("changes nothing when the directory did not change", func() {
Expect(resync.Resync(context.Background(), false)).To(Succeed())
Expect(lifecycle.names()).To(BeEmpty())
// Even a forced pass over an unchanged directory is a no-op.
Expect(resync.Resync(context.Background(), true)).To(Succeed())
Expect(lifecycle.names()).To(BeEmpty())
_, ok := peer.GetModelConfig("model-a")
Expect(ok).To(BeTrue())
})
It("resyncs after a NATS reconnect", func() {
bus := testutil.NewFakeBus()
ctx, cancel := context.WithCancel(context.Background())
DeferCleanup(cancel)
// A long interval, so only the reconnect can trigger the pass.
resync.Start(ctx, time.Hour, bus)
repointAlias(bus)
Expect(aliasTarget(peer, "stable")).To(Equal("model-a"), "precondition: the peer missed the change")
bus.TriggerReconnect()
Eventually(func() string { return aliasTarget(peer, "stable") }, 5*time.Second, 10*time.Millisecond).Should(Equal("model-b"))
})
It("runs the pass periodically", func() {
ctx, cancel := context.WithCancel(context.Background())
DeferCleanup(cancel)
resync.Start(ctx, 20*time.Millisecond, nil)
repointAlias(testutil.NewFakeBus())
Eventually(func() string { return aliasTarget(peer, "stable") }, 5*time.Second, 10*time.Millisecond).Should(Equal("model-b"))
})
It("reports a failing directory once, and retries until it is fixed", func() {
writeModelFile(dir, "half-written", "name: [unterminated\n")
Expect(resync.Resync(context.Background(), false)).NotTo(Succeed())
Expect(resync.Resync(context.Background(), false)).To(Succeed(), "the same failure is not reported twice")
writeModelFile(dir, "half-written", "name: half-written\nbackend: llama-cpp\n")
Expect(resync.Resync(context.Background(), false)).To(Succeed())
_, ok := peer.GetModelConfig("half-written")
Expect(ok).To(BeTrue())
})
It("keeps --config-file models across passes", func() {
cfgFile := filepath.Join(GinkgoT().TempDir(), "models.yaml")
Expect(os.WriteFile(cfgFile, []byte("- name: from-config-file\n backend: llama-cpp\n"), 0644)).To(Succeed())
Expect(peer.LoadMultipleModelConfigsSingleFile(cfgFile)).To(Succeed())
repointAlias(testutil.NewFakeBus())
Expect(resync.Resync(context.Background(), true)).To(Succeed())
_, ok := peer.GetModelConfig("from-config-file")
Expect(ok).To(BeTrue())
Expect(lifecycle.names()).NotTo(ContainElement("from-config-file"))
})
})
+23 -5
View File
@@ -2,7 +2,6 @@ package modeladmin
import (
"context"
"crypto/sha256"
"fmt"
"sort"
@@ -33,10 +32,29 @@ func applyRemoteChange(ctx context.Context, cl *config.ModelConfigLoader, models
if err := authoritative.LoadModelConfigsFromPathStrict(modelsPath, opts...); err != nil {
return err
}
current := configsByName(cl.GetAllModelsConfigs())
currentConfigs := cl.GetAllModelsConfigs()
current := make(map[string]config.ModelConfig, len(currentConfigs))
outside := map[string]struct{}{}
for _, cfg := range currentConfigs {
// A config read from outside the models directory (--config-file) is
// invisible to the snapshot. Treating it as removed would drop it and
// publish a deletion revision for a model that still exists.
if cfg.DefinedOutside(modelsPath) {
outside[cfg.Name] = struct{}{}
continue
}
current[cfg.Name] = cfg
}
snapshotConfigs := authoritative.GetAllModelsConfigs()
snapshot := configsByName(snapshotConfigs)
changed, err := changedConfigNames(current, snapshot, evt.Element)
for name := range outside {
delete(snapshot, name)
}
named := evt.Element
if _, isOutside := outside[named]; isOutside {
named = ""
}
changed, err := changedConfigNames(current, snapshot, named)
if err != nil {
return err
}
@@ -63,7 +81,7 @@ func applyRemoteChange(ctx context.Context, cl *config.ModelConfigLoader, models
}
}
}
cl.ReplaceModelConfigs(snapshotConfigs)
cl.ReplaceModelConfigs(config.MergeDirectorySnapshot(currentConfigs, snapshotConfigs, modelsPath))
return nil
}
@@ -109,5 +127,5 @@ func changedConfigNames(current, snapshot map[string]config.ModelConfig, named s
// model. It lets every frontend derive the same authoritative state regardless
// of which reordered cache-invalidation event woke it up.
func DeletedModelConfigRevision(modelName string) string {
return fmt.Sprintf("%x", sha256.Sum256([]byte("deleted\x00"+modelName)))
return config.DeletedModelConfigRevision(modelName)
}
+75 -1
View File
@@ -6,6 +6,7 @@ import (
"fmt"
"github.com/mudler/xlog"
"gorm.io/gorm"
)
// AliasResolver maps a model name to the name of the model that actually
@@ -18,6 +19,13 @@ type AliasResolver interface {
ResolveAliasName(name string) (string, bool)
}
// configRevisionSource is implemented by an AliasResolver that can also report
// the config revision behind its view of a name. core/config.ModelConfigLoader
// implements it. A resolver without it is always treated as current.
type configRevisionSource interface {
ConfigRevisionOf(name string) string
}
// SetAliasResolver installs the resolver used to map a scheduling rule's model
// name onto the model the rule actually governs. Called once at startup before
// serving. Leaving it unset makes every rule govern its own name, which is the
@@ -36,6 +44,65 @@ func (r *NodeRegistry) resolveAlias(name string) (string, bool) {
return (*p).ResolveAliasName(name)
}
// localViewIsCurrent reports whether this frontend's config for name is the
// one the cluster accepted: its revision matches the model_config_states row
// for name. Every frontend keeps its own copy of the model configs, and a
// frontend that missed an update still resolves an alias the old way.
//
// It answers true when it cannot tell: no revision-aware resolver, a config
// with no stamped revision, or no accepted revision recorded for name. Those
// cases keep the behaviour from before this check existed.
func (r *NodeRegistry) localViewIsCurrent(ctx context.Context, name string) (bool, error) {
p := r.aliasResolver.Load()
if p == nil || *p == nil {
return true, nil
}
source, ok := (*p).(configRevisionSource)
if !ok {
return true, nil
}
local := source.ConfigRevisionOf(name)
if local == "" {
return true, nil
}
accepted, err := r.GetModelConfigRevision(ctx, name)
if errors.Is(err, gorm.ErrRecordNotFound) {
return true, nil
}
if err != nil {
return false, err
}
return accepted == local, nil
}
// applyCurrentTarget is applyTarget for the writers of shared scheduling state:
// the replica reconciler and RefreshSchedulingTargets. It derives the target
// from this frontend's alias mapping only when that mapping is current. When
// this frontend is behind the cluster, it keeps the target_model stored by a
// frontend that is current, so a stale frontend can neither rewrite the stored
// target back nor scale up the model the alias used to point at.
//
// It costs a database read only for a rule whose stored and derived targets
// differ, which is rare and short-lived.
func (r *NodeRegistry) applyCurrentTarget(ctx context.Context, cfg *ModelSchedulingConfig) {
stored := cfg.TargetModel
r.applyTarget(cfg)
if stored == "" || stored == cfg.TargetModel {
return
}
current, err := r.localViewIsCurrent(ctx, cfg.ModelName)
if err != nil {
xlog.Warn("Cannot tell whether this frontend's model config is current; keeping the stored scheduling target",
"rule", cfg.ModelName, "stored", stored, "local", cfg.TargetModel, "error", err)
} else if current {
return
} else {
xlog.Debug("This frontend's model config is behind the cluster; keeping the stored scheduling target",
"rule", cfg.ModelName, "stored", stored, "local", cfg.TargetModel)
}
cfg.TargetModel = stored
}
// applyTarget fills in the rule's derived TargetModel. Every read path runs a
// rule through this so callers can tell the rule's key (ModelName, the name
// the operator chose) apart from the model it governs (TargetModel).
@@ -207,6 +274,12 @@ func (r *NodeRegistry) ValidateSchedulingTarget(ctx context.Context, ruleName st
// alias therefore reaches that guard one reconciler tick later, which is early
// enough: until then the guard protects the previous target, and the reconciler
// is already reloading the new one.
//
// Only a frontend whose model config is current writes (see
// applyCurrentTarget). Every frontend can run the reconciler, and each one
// resolves aliases from its own copy of the configs, so without that check a
// frontend that missed an alias update and one that did not would rewrite the
// row against each other on alternate ticks.
func (r *NodeRegistry) RefreshSchedulingTargets(ctx context.Context) error {
var configs []ModelSchedulingConfig
if err := r.db.WithContext(ctx).Find(&configs).Error; err != nil {
@@ -214,7 +287,8 @@ func (r *NodeRegistry) RefreshSchedulingTargets(ctx context.Context) error {
}
for i := range configs {
stored := configs[i].TargetModel
live, _ := r.resolveAlias(configs[i].ModelName)
r.applyCurrentTarget(ctx, &configs[i])
live := configs[i].TargetModel
if stored == live {
continue
}
@@ -0,0 +1,109 @@
package nodes
import (
"context"
"runtime"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"github.com/mudler/LocalAI/core/services/testutil"
"gorm.io/gorm"
)
// revisionAliasResolver is one frontend's view of the model configs: where
// each alias points, and the config revision behind each name.
type revisionAliasResolver struct {
fakeAliasResolver
revisions map[string]string
}
func (f *revisionAliasResolver) ConfigRevisionOf(name string) string {
return f.revisions[name]
}
var _ = Describe("Alias-keyed scheduling rules across frontends", func() {
var (
ctx context.Context
db *gorm.DB
current *NodeRegistry // saw the alias repoint
stale *NodeRegistry // missed it
storedOf func() string
)
BeforeEach(func() {
if runtime.GOOS == "darwin" {
Skip("testcontainers requires Docker, not available on macOS CI")
}
ctx = context.Background()
db = testutil.SetupTestDB()
var err error
// Two frontends share one database, each with its own copy of the
// model configs.
current, err = NewNodeRegistry(db)
Expect(err).ToNot(HaveOccurred())
stale, err = NewNodeRegistry(db)
Expect(err).ToNot(HaveOccurred())
// The rule is created while "production" still points at qwen3.
stale.SetAliasResolver(&revisionAliasResolver{
fakeAliasResolver: fakeAliasResolver{aliases: map[string]string{"production": "qwen3"}},
revisions: map[string]string{"production": "rev-old"},
})
Expect(stale.SetModelScheduling(ctx, &ModelSchedulingConfig{ModelName: "production", MinReplicas: 1})).To(Succeed())
// The alias is repointed at llama4. The cluster accepts the new
// config revision; only one frontend has reloaded it.
current.SetAliasResolver(&revisionAliasResolver{
fakeAliasResolver: fakeAliasResolver{aliases: map[string]string{"production": "llama4"}},
revisions: map[string]string{"production": "rev-new"},
})
_, err = current.AdvanceModelConfigRevision(ctx, "production", "rev-new")
Expect(err).ToNot(HaveOccurred())
storedOf = func() string {
var row ModelSchedulingConfig
ExpectWithOffset(1, db.Where("model_name = ?", "production").First(&row).Error).To(Succeed())
return row.TargetModel
}
})
It("does not let a stale frontend rewrite the stored target back", func() {
// Alternate ticks between the two frontends, as the reconciler's
// lock hands them out.
for range 3 {
Expect(current.RefreshSchedulingTargets(ctx)).To(Succeed())
Expect(storedOf()).To(Equal("llama4"))
Expect(stale.RefreshSchedulingTargets(ctx)).To(Succeed())
Expect(storedOf()).To(Equal("llama4"))
}
})
It("makes a stale frontend's reconciler act on the stored target", func() {
Expect(current.RefreshSchedulingTargets(ctx)).To(Succeed())
configs, err := stale.ListAutoScalingConfigs(ctx)
Expect(err).ToNot(HaveOccurred())
Expect(configs).To(HaveLen(1))
Expect(configs[0].Target()).To(Equal("llama4"))
})
It("lets the stale frontend follow once it reloads the config", func() {
Expect(current.RefreshSchedulingTargets(ctx)).To(Succeed())
stale.SetAliasResolver(&revisionAliasResolver{
fakeAliasResolver: fakeAliasResolver{aliases: map[string]string{"production": "llama4"}},
revisions: map[string]string{"production": "rev-new"},
})
configs, err := stale.ListAutoScalingConfigs(ctx)
Expect(err).ToNot(HaveOccurred())
Expect(configs[0].Target()).To(Equal("llama4"))
})
It("keeps the previous behaviour when no revision was accepted for the rule", func() {
Expect(db.Where("model_name = ?", "production").Delete(&ModelConfigState{}).Error).To(Succeed())
Expect(current.RefreshSchedulingTargets(ctx)).To(Succeed())
Expect(storedOf()).To(Equal("llama4"))
})
})
+4 -1
View File
@@ -2410,11 +2410,14 @@ func (r *NodeRegistry) ListModelSchedulings(ctx context.Context) ([]ModelSchedul
}
// ListAutoScalingConfigs returns scheduling configs where auto-scaling is enabled.
// The replica reconciler acts on these, so a rule whose alias this frontend
// resolves from an outdated config keeps the stored target instead (see
// applyCurrentTarget).
func (r *NodeRegistry) ListAutoScalingConfigs(ctx context.Context) ([]ModelSchedulingConfig, error) {
var configs []ModelSchedulingConfig
err := r.db.WithContext(ctx).Where("min_replicas > 0 OR max_replicas > 0 OR spread_all = ?", true).Find(&configs).Error
for i := range configs {
r.applyTarget(&configs[i])
r.applyCurrentTarget(ctx, &configs[i])
}
return configs, err
}
+32
View File
@@ -78,6 +78,7 @@ The frontend is a standard LocalAI instance with distributed mode enabled. These
| *(env only)* | `LOCALAI_MODEL_LOAD_WAIT` | `60s` | How long an inference request waits for a model that is still cold-loading onto a worker before it is answered with `503`, a `Retry-After` header and live staging progress. The request is served the moment the model becomes ready, so a model already most of the way staged needs no client retry. Set to `0` to wait as long as the load takes — only safe when no ingress or load balancer with an idle timeout sits in front. See [Requests for a model that is still loading](#requests-for-a-model-that-is-still-loading). |
| `--node-heartbeat-checkpoint` | `LOCALAI_NODE_HEARTBEAT_CHECKPOINT` | `60s` | Minimum gap between **durable** heartbeat writes for a worker node. A beat that only carries a fresher timestamp is kept in memory until this interval elapses instead of being written to PostgreSQL; every reported field is compared against the value last written rather than merely tested for presence, so a node's first beat, a changed total VRAM / total disk / GPU vendor, and a free VRAM / RAM / disk reading that has moved more than 256 MiB from the written value all still write immediately, and a node that is not active is never suppressed. Set it below the worker's `--heartbeat-interval` to restore a write per beat. See [Heartbeat writes and stale-node detection](#heartbeat-writes-and-stale-node-detection). |
| `--stale-node-threshold` | `LOCALAI_STALE_NODE_THRESHOLD` | `5m` | How long a node may go without a **durable** heartbeat before the health monitor marks it `offline`. Because `--node-heartbeat-checkpoint` holds back a beat that only carries a fresher timestamp, this has to stay comfortably wider than that interval: raising the checkpoint without raising this marks healthy, beating nodes offline. Neither the per-model gRPC health check nor request-time failure reads `last_heartbeat`, so neither is affected by this knob. See [Heartbeat writes and stale-node detection](#heartbeat-writes-and-stale-node-detection). |
| `--model-config-resync-interval` | `LOCALAI_MODEL_CONFIG_RESYNC_INTERVAL` | `30s` | How often each frontend compares its model configs with the shared models directory, to apply a change whose NATS message it missed. A frontend that missed a message serves the old config for at most this long. See [Model configs across frontends](#model-configs-across-frontends). |
| `--expose-node-header` | `LOCALAI_EXPOSE_NODE_HEADER` | `false` | When enabled, inference responses carry an `X-LocalAI-Node` header with the ID of the worker node that served the request. Coverage spans the OpenAI-compatible endpoints (chat completions, completions, embeddings, audio transcriptions, audio speech / TTS, image generations, image inpainting), the Jina rerank endpoint (`/v1/rerank`), the VAD endpoints (`/v1/vad`, `/vad`), and the Anthropic Messages (`/v1/messages`) and Ollama (`/api/chat`, `/api/generate`, `/api/embed`) shims. Useful for debugging, observability and load-balancer attribution. Off by default: the node ID reveals internal cluster topology and should not be exposed on a public endpoint. Best-effort: under heavy concurrency for the same model across multiple replicas, the header may reflect a recent routing decision rather than this exact request's. Acceptable for observability and debugging. |
### The model load deadline scales with the checkpoint
@@ -649,6 +650,26 @@ Variant selection (`GET /api/models/variants/:id`) uses the same reading, and
judges backend compatibility against the union of the capabilities present in
the cluster, so a CUDA-only build is offered when any worker can run it.
### Model configs across frontends
Every frontend keeps its own in-memory copy of the model configs in the shared models directory. When a frontend installs, edits, toggles or deletes a model, it writes the change to the directory and publishes a message on NATS. The other frontends reload the directory when they receive it.
A gallery install or delete publishes this message as soon as the new config is in place, before the frontend preloads model files. The preload can take minutes on a large models directory, and other frontends do not wait for it. If the preload fails, the operation reports the error, but the config change stays applied on every frontend.
NATS keeps no history of these messages. A frontend that is disconnected when a message is published never receives it. To recover, each frontend also reloads the models directory:
- every `--model-config-resync-interval` (default `30s`), when a config file changed since its last pass, and
- after each NATS reconnect.
The pass is the same reconcile that a NATS message triggers, so it is idempotent. Only models whose file changed get a new [configuration revision](#model-configuration-revisions). A pass over an unchanged directory reads the config files and does nothing else.
Distributed-state mode: each frontend derives this state from the shared directory, so there is no leader and nothing to replicate. The only per-frontend memory is a hash of the config files from its last pass, which only saves work. The consequence is a bounded delay: a frontend that missed a message serves the previous config for at most one interval.
Two limits apply:
- Models loaded with `--config-file` exist only on the frontend that loaded them. A reload of the models directory keeps them, and a config-file model wins over a directory file with the same name, as it does at startup.
- The reload is strict: if any config file in the directory does not parse, the frontend keeps its current configs and logs the error once. It retries on each pass until the file is fixed. This keeps a half-written file from looking like a deleted model.
### Model configuration revisions
Distributed mode assigns a `config_revision` to each validated model configuration. It hashes the persisted semantic configuration, including fields such as `context_size` and parallel settings. YAML formatting, comments, and map order do not change it.
@@ -1192,6 +1213,17 @@ the slot, and the model filling it can change without rewriting the rule. The
WebUI lists aliases in the model picker on the **Scheduling** page, tagged with
the model each one resolves to.
Each frontend resolves the alias from its own copy of the model configs, and a
frontend that has not yet reloaded a repointed alias still resolves it the old
way (see [Model configs across frontends](#model-configs-across-frontends)).
The rule's stored target therefore follows the alias only through frontends
whose copy of the alias config matches the
[configuration revision](#model-configuration-revisions) the cluster accepted.
A frontend that is behind uses the stored target for the replica reconciler and
does not write it, so two frontends cannot overwrite the rule's target against
each other, and a frontend that is behind cannot reload the model the alias
used to point at.
Two constraints follow from replicas being shared. A single load of `llama3`
serves both `production` and any request that names `llama3` directly, so only
one rule can decide where it runs: a rule whose target is already governed by