diff --git a/core/services/failover/manager.go b/core/services/failover/manager.go index 3c647a2ce..ae997a5ad 100644 --- a/core/services/failover/manager.go +++ b/core/services/failover/manager.go @@ -10,6 +10,7 @@ import ( "sync/atomic" "time" + "github.com/google/uuid" "github.com/mudler/LocalAI/core/config" "github.com/mudler/xlog" ) @@ -43,7 +44,10 @@ func WithOnWarmChanged(fn func(warm []string)) Option { return func(m *Manager) // Manager tracks health per target and the active target per chain. type Manager struct { - mu sync.Mutex + mu sync.Mutex + // id tells this manager's own target publishes apart when the sync + // layer echoes them back. + id string src ConfigSource clock Clock prober Prober @@ -120,6 +124,7 @@ type chainState struct { func New(src ConfigSource, opts ...Option) *Manager { m := &Manager{ + id: uuid.NewString(), src: src, clock: realClock{}, targets: map[string]*targetState{}, diff --git a/core/services/failover/statesync.go b/core/services/failover/statesync.go index 49515d394..14e9de431 100644 --- a/core/services/failover/statesync.go +++ b/core/services/failover/statesync.go @@ -16,6 +16,8 @@ type TargetSnapshot struct { Error string `json:"error,omitempty"` ConsecutiveOK int `json:"consecutive_ok"` Since time.Time `json:"since"` + // Origin is the publishing manager, so it can drop its own echoes. + Origin string `json:"origin,omitempty"` } // ChainSnapshot is one chain's active target as decided by the leader. @@ -129,7 +131,7 @@ func (m *Manager) targetSnapshotLocked(ts *targetState) TargetSnapshot { } return TargetSnapshot{ Target: ts.name, State: ts.state, Reason: ts.reason, Error: ts.lastError, - ConsecutiveOK: ts.consecutiveOK, Since: since, + ConsecutiveOK: ts.consecutiveOK, Since: since, Origin: m.id, } } @@ -158,11 +160,17 @@ func (m *Manager) queuePublishChainLocked(ch *chainState, reason Reason) { m.pending = append(m.pending, func() { s.PublishChain(snap) }) } -// ApplyTarget takes a target state published by any frontend, this one -// included. The echo of an own publish finds the same state and does nothing. +// ApplyTarget takes a target state published by another frontend. The echo +// of an own publish is dropped: this manager already holds that state or a +// newer one, and publishes are snapshotted when queued, so an echo can arrive +// after a later local transition (down -> recovering -> healthy under one +// lock queues two) and would roll the target back. func (m *Manager) ApplyTarget(s TargetSnapshot) { m.mu.Lock() defer m.unlockAndFlush() + if s.Origin != "" && s.Origin == m.id { + return + } ts := m.targetLocked(s.Target) if ts == nil || ts.state == StateMissing || s.State == StateMissing { return diff --git a/core/services/failover/statesync_test.go b/core/services/failover/statesync_test.go index 7e45da87c..6021636de 100644 --- a/core/services/failover/statesync_test.go +++ b/core/services/failover/statesync_test.go @@ -5,6 +5,8 @@ import ( "sync" "time" + "github.com/mudler/LocalAI/core/config" + . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" ) @@ -106,6 +108,28 @@ var _ = Describe("Manager state sync", func() { Expect(n).To(Equal(1)) }) + It("does not let the echo of its own earlier publish undo a newer local state", func() { + // One recovery probe: a success moves x down -> recovering -> healthy + // under one lock, which queues two publishes. Their echoes arrive + // after x is already healthy here. + src.Put(chainCfg("chain", &config.FailoverConfig{Recovery: config.FailoverRecovery{Probes: 1}}, t("x"), t("y"))) + a.Sync() + b.Sync() + a.ReportFailure("x", errBoom) + st, _ := a.ChainStatus("chain") + Expect(st.Active).To(Equal("y")) + + clock.Advance(10 * time.Minute) // past any dwell + a.ReportSuccess("x") + + st, _ = a.ChainStatus("chain") + Expect(st.Targets[0].State).To(Equal(StateHealthy)) + Expect(st.Active).To(Equal("x"), "the leader must stay failed back to x") + st, _ = b.ChainStatus("chain") + Expect(st.Targets[0].State).To(Equal(StateHealthy)) + Expect(st.Active).To(Equal("x")) + }) + It("a trip on one frontend is skipped by the other's plan", func() { b.ReportFailure("x", errBoom) att, err := a.Plan("chain")