diff --git a/core/cli/worker.go b/core/cli/worker.go index affde4b08..186fe298e 100644 --- a/core/cli/worker.go +++ b/core/cli/worker.go @@ -738,6 +738,9 @@ func (s *backendSupervisor) subscribeLifecycleEvents() { if b.Metadata != nil { info.InstalledAt = b.Metadata.InstalledAt info.GalleryURL = b.Metadata.GalleryURL + info.Version = b.Metadata.Version + info.URI = b.Metadata.URI + info.Digest = b.Metadata.Digest } infos = append(infos, info) } diff --git a/core/gallery/backends.go b/core/gallery/backends.go index c2622c272..ee3ca906d 100644 --- a/core/gallery/backends.go +++ b/core/gallery/backends.go @@ -394,6 +394,23 @@ type SystemBackend struct { Metadata *BackendMetadata UpgradeAvailable bool `json:"upgrade_available,omitempty"` AvailableVersion string `json:"available_version,omitempty"` + // Nodes holds per-node attribution in distributed mode. Empty in single-node. + // Each entry describes a node that has this backend installed, with the + // version/digest it reports. Lets the UI surface drift and per-node status. + Nodes []NodeBackendRef `json:"nodes,omitempty"` +} + +// NodeBackendRef describes one node's view of an installed backend. Used both +// for per-node attribution in the UI and for drift detection during upgrade +// checks (a cluster with mismatched versions/digests is flagged upgradeable). +type NodeBackendRef struct { + NodeID string `json:"node_id"` + NodeName string `json:"node_name"` + NodeStatus string `json:"node_status"` // healthy | unhealthy | offline | draining | pending + Version string `json:"version,omitempty"` + Digest string `json:"digest,omitempty"` + URI string `json:"uri,omitempty"` + InstalledAt string `json:"installed_at,omitempty"` } type SystemBackends map[string]SystemBackend diff --git a/core/gallery/upgrade.go b/core/gallery/upgrade.go index dde33300f..d0671617e 100644 --- a/core/gallery/upgrade.go +++ b/core/gallery/upgrade.go @@ -23,22 +23,45 @@ type UpgradeInfo struct { AvailableVersion string `json:"available_version"` InstalledDigest string `json:"installed_digest,omitempty"` AvailableDigest string `json:"available_digest,omitempty"` + // NodeDrift lists nodes whose installed version or digest differs from + // the cluster majority. Non-empty means the cluster has diverged and an + // upgrade will realign it. Empty in single-node mode. + NodeDrift []NodeDriftInfo `json:"node_drift,omitempty"` } -// CheckBackendUpgrades compares installed backends against gallery entries -// and returns a map of backend names to UpgradeInfo for those that have -// newer versions or different OCI digests available. +// NodeDriftInfo describes one node that disagrees with the cluster majority +// on which version/digest of a backend is installed. +type NodeDriftInfo struct { + NodeID string `json:"node_id"` + NodeName string `json:"node_name"` + Version string `json:"version,omitempty"` + Digest string `json:"digest,omitempty"` +} + +// CheckBackendUpgrades is the single-node entrypoint. Distributed callers use +// CheckUpgradesAgainst directly with their aggregated SystemBackends. func CheckBackendUpgrades(ctx context.Context, galleries []config.Gallery, systemState *system.SystemState) (map[string]UpgradeInfo, error) { + installed, err := ListSystemBackends(systemState) + if err != nil { + return nil, fmt.Errorf("failed to list installed backends: %w", err) + } + return CheckUpgradesAgainst(ctx, galleries, systemState, installed) +} + +// CheckUpgradesAgainst compares a caller-supplied SystemBackends set against +// the gallery. Fixes the distributed-mode bug where the old code passed the +// frontend's (empty) local filesystem through ListSystemBackends and so never +// surfaced any upgrades. +// +// Cluster drift policy: if a backend's per-node versions/digests disagree, the +// row is flagged upgradeable regardless of whether any node matches the gallery +// — next Upgrade All realigns the cluster. NodeDrift lists the outliers. +func CheckUpgradesAgainst(ctx context.Context, galleries []config.Gallery, systemState *system.SystemState, installedBackends SystemBackends) (map[string]UpgradeInfo, error) { galleryBackends, err := AvailableBackends(galleries, systemState) if err != nil { return nil, fmt.Errorf("failed to list available backends: %w", err) } - installedBackends, err := ListSystemBackends(systemState) - if err != nil { - return nil, fmt.Errorf("failed to list installed backends: %w", err) - } - result := make(map[string]UpgradeInfo) for _, installed := range installedBackends { @@ -57,34 +80,48 @@ func CheckBackendUpgrades(ctx context.Context, galleries []config.Gallery, syste } installedVersion := installed.Metadata.Version + installedDigest := installed.Metadata.Digest galleryVersion := galleryEntry.Version - // If both sides have versions, compare them + // Detect cluster drift: does every node report the same version+digest? + // In single-node mode this stays empty (Nodes is nil). + majority, drift := summarizeNodeDrift(installed.Nodes) + if majority.version != "" { + installedVersion = majority.version + } + if majority.digest != "" { + installedDigest = majority.digest + } + + makeInfo := func(info UpgradeInfo) UpgradeInfo { + info.NodeDrift = drift + return info + } + + // If versions are available on both sides, they're the source of truth. if galleryVersion != "" && installedVersion != "" { - if galleryVersion != installedVersion { - result[installed.Metadata.Name] = UpgradeInfo{ + if galleryVersion != installedVersion || len(drift) > 0 { + result[installed.Metadata.Name] = makeInfo(UpgradeInfo{ BackendName: installed.Metadata.Name, InstalledVersion: installedVersion, AvailableVersion: galleryVersion, - } + }) } - // Versions match — no upgrade needed continue } - // Gallery has a version but installed doesn't — this happens for backends - // installed before version tracking was added. Flag as upgradeable so - // users can re-install to pick up version metadata. + // Gallery has a version but installed doesn't — backends installed before + // version tracking was added. Flag as upgradeable to pick up metadata. if galleryVersion != "" && installedVersion == "" { - result[installed.Metadata.Name] = UpgradeInfo{ + result[installed.Metadata.Name] = makeInfo(UpgradeInfo{ BackendName: installed.Metadata.Name, InstalledVersion: "", AvailableVersion: galleryVersion, - } + }) continue } - // Fall back to OCI digest comparison when versions are unavailable + // Fall back to OCI digest comparison when versions are unavailable. if downloader.URI(galleryEntry.URI).LooksLikeOCI() { remoteDigest, err := oci.GetImageDigest(galleryEntry.URI, "", nil, nil) if err != nil { @@ -92,21 +129,68 @@ func CheckBackendUpgrades(ctx context.Context, galleries []config.Gallery, syste continue } // If we have a stored digest, compare; otherwise any remote digest - // means we can't confirm we're up to date — flag as upgradeable - if installed.Metadata.Digest == "" || remoteDigest != installed.Metadata.Digest { - result[installed.Metadata.Name] = UpgradeInfo{ + // means we can't confirm we're up to date — flag as upgradeable. + if installedDigest == "" || remoteDigest != installedDigest || len(drift) > 0 { + result[installed.Metadata.Name] = makeInfo(UpgradeInfo{ BackendName: installed.Metadata.Name, - InstalledDigest: installed.Metadata.Digest, + InstalledDigest: installedDigest, AvailableDigest: remoteDigest, - } + }) } + } else if len(drift) > 0 { + // No version/digest path but nodes disagree — still worth flagging. + result[installed.Metadata.Name] = makeInfo(UpgradeInfo{ + BackendName: installed.Metadata.Name, + InstalledVersion: installedVersion, + InstalledDigest: installedDigest, + }) } - // No version info and non-OCI URI — cannot determine, skip } return result, nil } +// summarizeNodeDrift collapses per-node version/digest tuples to a majority +// pair and returns the outliers. In single-node mode (empty nodes slice) this +// returns zero values and a nil drift list. +func summarizeNodeDrift(nodes []NodeBackendRef) (majority struct{ version, digest string }, drift []NodeDriftInfo) { + if len(nodes) == 0 { + return majority, nil + } + + type key struct{ version, digest string } + counts := map[key]int{} + var topKey key + var topCount int + for _, n := range nodes { + k := key{n.Version, n.Digest} + counts[k]++ + if counts[k] > topCount { + topCount = counts[k] + topKey = k + } + } + + majority.version = topKey.version + majority.digest = topKey.digest + + if len(counts) == 1 { + return majority, nil // unanimous — no drift + } + for _, n := range nodes { + if n.Version == majority.version && n.Digest == majority.digest { + continue + } + drift = append(drift, NodeDriftInfo{ + NodeID: n.NodeID, + NodeName: n.NodeName, + Version: n.Version, + Digest: n.Digest, + }) + } + return majority, drift +} + // UpgradeBackend upgrades a single backend to the latest gallery version using // an atomic swap with backup-based rollback on failure. func UpgradeBackend(ctx context.Context, systemState *system.SystemState, modelLoader *model.ModelLoader, galleries []config.Gallery, backendName string, downloadStatus func(string, string, string, float64)) error { diff --git a/core/gallery/upgrade_test.go b/core/gallery/upgrade_test.go index f65b4276b..6fd386b2e 100644 --- a/core/gallery/upgrade_test.go +++ b/core/gallery/upgrade_test.go @@ -144,6 +144,97 @@ var _ = Describe("Upgrade Detection and Execution", func() { }) }) + // CheckUpgradesAgainst is the entry point used by DistributedBackendManager. + // It takes installed backends directly — typically aggregated from workers — + // instead of reading the frontend filesystem. These tests exercise drift + // detection, which is the feature the distributed path relies on. + Describe("CheckUpgradesAgainst (distributed)", func() { + It("flags upgrade when cluster nodes disagree on version, even if gallery matches majority", func() { + writeGalleryYAML([]GalleryBackend{ + { + Metadata: Metadata{Name: "my-backend"}, + URI: filepath.Join(tempDir, "some-source"), + Version: "2.0.0", + }, + }) + + installed := SystemBackends{ + "my-backend": SystemBackend{ + Name: "my-backend", + Metadata: &BackendMetadata{Name: "my-backend", Version: "2.0.0"}, + Nodes: []NodeBackendRef{ + {NodeID: "a", NodeName: "worker-1", Version: "2.0.0"}, + {NodeID: "b", NodeName: "worker-2", Version: "2.0.0"}, + {NodeID: "c", NodeName: "worker-3", Version: "1.0.0"}, // drift + }, + }, + } + + upgrades, err := CheckUpgradesAgainst(context.Background(), galleries, systemState, installed) + Expect(err).NotTo(HaveOccurred()) + Expect(upgrades).To(HaveKey("my-backend")) + info := upgrades["my-backend"] + Expect(info.AvailableVersion).To(Equal("2.0.0")) + Expect(info.NodeDrift).To(HaveLen(1)) + Expect(info.NodeDrift[0].NodeName).To(Equal("worker-3")) + Expect(info.NodeDrift[0].Version).To(Equal("1.0.0")) + }) + + It("does not flag upgrade when all nodes agree and match gallery", func() { + writeGalleryYAML([]GalleryBackend{ + { + Metadata: Metadata{Name: "my-backend"}, + URI: filepath.Join(tempDir, "some-source"), + Version: "2.0.0", + }, + }) + + installed := SystemBackends{ + "my-backend": SystemBackend{ + Name: "my-backend", + Metadata: &BackendMetadata{Name: "my-backend", Version: "2.0.0"}, + Nodes: []NodeBackendRef{ + {NodeID: "a", NodeName: "worker-1", Version: "2.0.0"}, + {NodeID: "b", NodeName: "worker-2", Version: "2.0.0"}, + }, + }, + } + + upgrades, err := CheckUpgradesAgainst(context.Background(), galleries, systemState, installed) + Expect(err).NotTo(HaveOccurred()) + Expect(upgrades).To(BeEmpty()) + }) + + It("surfaces empty-installed-version path the old distributed code silently missed", func() { + // Simulates the real-world bug: worker has a backend, its version + // is empty (pre-tracking or OCI-pinned-to-latest), gallery has a + // version. Pre-fix CheckUpgrades returned nothing; now it surfaces. + writeGalleryYAML([]GalleryBackend{ + { + Metadata: Metadata{Name: "my-backend"}, + URI: filepath.Join(tempDir, "some-source"), + Version: "2.0.0", + }, + }) + + installed := SystemBackends{ + "my-backend": SystemBackend{ + Name: "my-backend", + Metadata: &BackendMetadata{Name: "my-backend"}, + Nodes: []NodeBackendRef{ + {NodeID: "a", NodeName: "worker-1"}, + }, + }, + } + + upgrades, err := CheckUpgradesAgainst(context.Background(), galleries, systemState, installed) + Expect(err).NotTo(HaveOccurred()) + Expect(upgrades).To(HaveKey("my-backend")) + Expect(upgrades["my-backend"].InstalledVersion).To(BeEmpty()) + Expect(upgrades["my-backend"].AvailableVersion).To(Equal("2.0.0")) + }) + }) + Describe("UpgradeBackend", func() { It("should replace backend directory and update metadata", func() { // Install v1 diff --git a/core/services/messaging/subjects.go b/core/services/messaging/subjects.go index 397f63ff0..3e9af53a9 100644 --- a/core/services/messaging/subjects.go +++ b/core/services/messaging/subjects.go @@ -157,6 +157,12 @@ type NodeBackendInfo struct { IsMeta bool `json:"is_meta"` InstalledAt string `json:"installed_at,omitempty"` GalleryURL string `json:"gallery_url,omitempty"` + // Version, URI and Digest enable cluster-wide upgrade detection — + // without them, the frontend cannot tell whether the installed OCI + // image matches the gallery entry, and upgrades silently never surface. + Version string `json:"version,omitempty"` + URI string `json:"uri,omitempty"` + Digest string `json:"digest,omitempty"` } // SubjectNodeBackendStop tells a worker node to stop its gRPC backend process. diff --git a/core/services/nodes/managers_distributed.go b/core/services/nodes/managers_distributed.go index 62cb32552..0934a6282 100644 --- a/core/services/nodes/managers_distributed.go +++ b/core/services/nodes/managers_distributed.go @@ -10,6 +10,7 @@ import ( "github.com/mudler/LocalAI/core/gallery" "github.com/mudler/LocalAI/core/services/galleryop" "github.com/mudler/LocalAI/pkg/model" + "github.com/mudler/LocalAI/pkg/system" "github.com/mudler/xlog" "github.com/nats-io/nats.go" ) @@ -53,6 +54,7 @@ type DistributedBackendManager struct { adapter *RemoteUnloaderAdapter registry *NodeRegistry backendGalleries []config.Gallery + systemState *system.SystemState } // NewDistributedBackendManager creates a DistributedBackendManager. @@ -62,6 +64,7 @@ func NewDistributedBackendManager(appConfig *config.ApplicationConfig, ml *model adapter: adapter, registry: registry, backendGalleries: appConfig.BackendGalleries, + systemState: appConfig.SystemState, } } @@ -101,7 +104,14 @@ func (d *DistributedBackendManager) DeleteBackend(name string) error { return errors.Join(errs...) } -// ListBackends aggregates installed backends from all healthy worker nodes. +// ListBackends aggregates installed backends from all worker nodes, preserving +// per-node attribution. Each SystemBackend.Nodes entry records which node has +// the backend and the version/digest it reports. The top-level Metadata is +// populated from the first node seen so single-node-minded callers still work. +// +// Pending/offline/draining nodes are skipped because they aren't expected to +// answer NATS requests; unhealthy nodes are still queried — ErrNoResponders +// then marks them unhealthy and the loop continues. func (d *DistributedBackendManager) ListBackends() (gallery.SystemBackends, error) { result := make(gallery.SystemBackends) allNodes, err := d.registry.List(context.Background()) @@ -110,7 +120,7 @@ func (d *DistributedBackendManager) ListBackends() (gallery.SystemBackends, erro } for _, node := range allNodes { - if node.Status != StatusHealthy { + if node.Status == StatusPending || node.Status == StatusOffline || node.Status == StatusDraining { continue } reply, err := d.adapter.ListBackends(node.ID) @@ -128,17 +138,33 @@ func (d *DistributedBackendManager) ListBackends() (gallery.SystemBackends, erro continue } for _, b := range reply.Backends { - if _, exists := result[b.Name]; !exists { - result[b.Name] = gallery.SystemBackend{ + ref := gallery.NodeBackendRef{ + NodeID: node.ID, + NodeName: node.Name, + NodeStatus: node.Status, + Version: b.Version, + Digest: b.Digest, + URI: b.URI, + InstalledAt: b.InstalledAt, + } + entry, exists := result[b.Name] + if !exists { + entry = gallery.SystemBackend{ Name: b.Name, IsSystem: b.IsSystem, IsMeta: b.IsMeta, Metadata: &gallery.BackendMetadata{ + Name: b.Name, InstalledAt: b.InstalledAt, GalleryURL: b.GalleryURL, + Version: b.Version, + URI: b.URI, + Digest: b.Digest, }, } } + entry.Nodes = append(entry.Nodes, ref) + result[b.Name] = entry } } return result, nil @@ -209,8 +235,21 @@ func (d *DistributedBackendManager) UpgradeBackend(ctx context.Context, name str return errors.Join(errs...) } -// CheckUpgrades checks for available backend upgrades. -// Gallery comparison is global (not per-node), so we delegate to the local manager. +// CheckUpgrades checks for available backend upgrades across the cluster. +// +// The previous implementation delegated to d.local, which called +// ListSystemBackends on the frontend — but in distributed mode the frontend +// has no backends installed locally, so the upgrade loop never ran and the UI +// never surfaced any upgrades. We now feed the cluster-wide aggregation +// (including per-node versions/digests) into gallery.CheckUpgradesAgainst so +// digest-based detection actually works and cluster drift is visible. func (d *DistributedBackendManager) CheckUpgrades(ctx context.Context) (map[string]gallery.UpgradeInfo, error) { - return d.local.CheckUpgrades(ctx) + installed, err := d.ListBackends() + if err != nil { + return nil, err + } + // systemState is used by AvailableBackends (gallery paths + meta-backend + // resolution). The `installed` argument is what the old code got wrong — + // it used to come from the empty frontend filesystem. + return gallery.CheckUpgradesAgainst(ctx, d.backendGalleries, d.systemState, installed) }