Files
LocalAI/core/services/modeladmin/remote_sync.go
T
mudler-agentandEttore Di Giacinto 895d50385f 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>
2026-10-08 09:56:12 +02:00

132 lines
4.5 KiB
Go

package modeladmin
import (
"context"
"fmt"
"sort"
"github.com/mudler/LocalAI/core/config"
"github.com/mudler/LocalAI/core/services/messaging"
)
// ApplyRemoteChange refreshes this replica's in-memory model state from a peer
// replica's model-config change broadcast (messaging.CacheInvalidateEvent on
// SubjectCacheInvalidateModels). It is the subscriber-side counterpart to
// GalleryService.BroadcastModelsChanged.
//
// The event is only a wake-up signal. Its operation and revision may be stale
// or reordered, so named changes are always reconciled against the current
// shared filesystem state.
//
// Revision-aware events apply the same idempotent lifecycle transition as the
// originating frontend. modelsPath and opts are forwarded to
// LoadModelConfigsFromPath.
func ApplyRemoteChange(ctx context.Context, cl *config.ModelConfigLoader, modelsPath string, evt messaging.CacheInvalidateEvent, lifecycle ModelRevisionLifecycle, opts ...config.ConfigLoaderOption) error {
return cl.WithModelConfigMutation(func() error {
return applyRemoteChange(ctx, cl, modelsPath, evt, lifecycle, opts...)
})
}
func applyRemoteChange(ctx context.Context, cl *config.ModelConfigLoader, modelsPath string, evt messaging.CacheInvalidateEvent, lifecycle ModelRevisionLifecycle, opts ...config.ConfigLoaderOption) error {
authoritative := config.NewModelConfigLoader(modelsPath)
if err := authoritative.LoadModelConfigsFromPathStrict(modelsPath, opts...); err != nil {
return err
}
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)
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
}
if lifecycle != nil {
transitions := make([]ModelRevisionTransition, 0, len(changed))
for _, name := range changed {
cfg, exists := snapshot[name]
revision := DeletedModelConfigRevision(name)
disabled := true
if exists {
var err error
revision, err = authoritative.RevisionForPath(name, modelsPath, opts...)
if err != nil {
return fmt.Errorf("resolve authoritative model config revision for %q: %w", name, err)
}
disabled = cfg.IsDisabled()
}
transitions = append(transitions, ModelRevisionTransition{ModelName: name, ConfigRevision: revision, Disabled: disabled})
}
if len(transitions) > 0 {
if _, err := lifecycle.ApplyConfigRevisions(ctx, transitions); err != nil {
return err
}
}
}
cl.ReplaceModelConfigs(config.MergeDirectorySnapshot(currentConfigs, snapshotConfigs, modelsPath))
return nil
}
func configsByName(configs []config.ModelConfig) map[string]config.ModelConfig {
result := make(map[string]config.ModelConfig, len(configs))
for _, cfg := range configs {
result[cfg.Name] = cfg
}
return result
}
func changedConfigNames(current, snapshot map[string]config.ModelConfig, named string) ([]string, error) {
changed := map[string]struct{}{}
for name, cfg := range snapshot {
previous, exists := current[name]
if !exists {
changed[name] = struct{}{}
continue
}
// Both sides come from a loader, so both carry the revision stamped
// when their file was parsed. Comparing the stamps compares the files.
if previous.PersistedConfigRevision() != cfg.PersistedConfigRevision() {
changed[name] = struct{}{}
}
}
for name := range current {
if _, exists := snapshot[name]; !exists {
changed[name] = struct{}{}
}
}
if named != "" {
changed[named] = struct{}{}
}
names := make([]string, 0, len(changed))
for name := range changed {
names = append(names, name)
}
sort.Strings(names)
return names, nil
}
// DeletedModelConfigRevision is a stable tombstone generation for an absent
// 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 config.DeletedModelConfigRevision(modelName)
}