diff --git a/core/http/middleware/failover_test.go b/core/http/middleware/failover_test.go index a16337545..81fff5e02 100644 --- a/core/http/middleware/failover_test.go +++ b/core/http/middleware/failover_test.go @@ -87,6 +87,9 @@ var _ = Describe("failover chains in the request pipeline", func() { Expect(mcl.LoadModelConfigsFromPath(dir)).To(Succeed()) re = NewRequestExtractor(mcl, model.NewModelLoader(ss), appConfig) fm = failover.New(mcl) + // The scheduler's first tick syncs in the application; HasChains + // answers from the last sync. + fm.Sync() re.SetFailoverManager(fm) app = echo.New() diff --git a/core/services/failover/fakes_test.go b/core/services/failover/fakes_test.go index cad6b3552..232b7b4df 100644 --- a/core/services/failover/fakes_test.go +++ b/core/services/failover/fakes_test.go @@ -26,8 +26,9 @@ func (c *fakeClock) Advance(d time.Duration) { } type fakeSource struct { - mu sync.Mutex - cfgs map[string]config.ModelConfig + mu sync.Mutex + cfgs map[string]config.ModelConfig + scans int // GetAllModelsConfigs calls } func newFakeSource(cfgs ...config.ModelConfig) *fakeSource { @@ -39,6 +40,7 @@ func newFakeSource(cfgs ...config.ModelConfig) *fakeSource { } func (s *fakeSource) Put(c config.ModelConfig) { s.mu.Lock(); s.cfgs[c.Name] = c; s.mu.Unlock() } func (s *fakeSource) Delete(name string) { s.mu.Lock(); delete(s.cfgs, name); s.mu.Unlock() } +func (s *fakeSource) Scans() int { s.mu.Lock(); defer s.mu.Unlock(); return s.scans } func (s *fakeSource) GetModelConfig(n string) (config.ModelConfig, bool) { s.mu.Lock() defer s.mu.Unlock() @@ -48,6 +50,7 @@ func (s *fakeSource) GetModelConfig(n string) (config.ModelConfig, bool) { func (s *fakeSource) GetAllModelsConfigs() []config.ModelConfig { s.mu.Lock() defer s.mu.Unlock() + s.scans++ out := make([]config.ModelConfig, 0, len(s.cfgs)) for _, c := range s.cfgs { out = append(out, c) diff --git a/core/services/failover/manager.go b/core/services/failover/manager.go index ab130d1ee..ea31ee999 100644 --- a/core/services/failover/manager.go +++ b/core/services/failover/manager.go @@ -55,6 +55,9 @@ type Manager struct { warm []string warmPending bool closed bool + // hasChains mirrors len(chains) > 0 as of the last sync, so the request + // path can check it without the lock or a config-source scan. + hasChains atomic.Bool } type targetState struct { @@ -168,6 +171,7 @@ func (m *Manager) syncLocked() { delete(m.targets, name) } } + m.hasChains.Store(len(m.chains) > 0) for _, ch := range m.chains { m.recomputeLocked(ch, "") } @@ -201,26 +205,18 @@ func (m *Manager) lookupTarget(name string) (config.ModelConfig, bool) { return c, ok } -// HasChains reports whether any failover chain is configured. The request -// path uses it to skip chain bookkeeping on installations without chains. -// Chains are synced lazily, so with none known yet the config source is -// consulted, which catches a chain added since the last sync. +// HasChains reports whether any failover chain was configured at the last +// sync. The request path calls it on every request to skip chain bookkeeping +// on installations without chains, so it reads a flag instead of scanning +// the config source (which takes the loader's lock and copies every config). +// A chain added since the last sync is still served, because Plan syncs on a +// miss; only in-request retry is missing for it until the scheduler's next +// tick, at most one second later. func (m *Manager) HasChains() bool { if m == nil { return false } - m.mu.Lock() - known := len(m.chains) > 0 - m.mu.Unlock() - if known { - return true - } - for _, c := range m.src.GetAllModelsConfigs() { - if c.IsFailover() { - return true - } - } - return false + return m.hasChains.Load() } func (m *Manager) chainLocked(name string) *chainState { diff --git a/core/services/failover/manager_test.go b/core/services/failover/manager_test.go index d0736cf19..0d2746331 100644 --- a/core/services/failover/manager_test.go +++ b/core/services/failover/manager_test.go @@ -60,18 +60,34 @@ var _ = Describe("Manager", func() { Expect(st.Targets[0].State).To(Equal(StateHealthy)) }) - It("reports whether any chain is configured, including one added since the last sync", func() { + It("reports whether any chain was configured at the last sync", func() { + m.Sync() + Expect(m.HasChains()).To(BeTrue()) + var nilManager *Manager + Expect(nilManager.HasChains()).To(BeFalse()) + }) + + It("answers HasChains from the last sync without scanning the config source", func() { empty := New(newFakeSource(remote("a")), WithClock(clock)) Expect(empty.HasChains()).To(BeFalse()) lateSrc := newFakeSource(remote("a"), local("b")) late := New(lateSrc, WithClock(clock)) late.Sync() - Expect(late.HasChains()).To(BeFalse()) + scans := lateSrc.Scans() + for range 100 { + Expect(late.HasChains()).To(BeFalse()) + } + Expect(lateSrc.Scans()).To(Equal(scans)) + // A chain added since the last sync is seen at the next sync, or + // sooner by Plan, which syncs on a miss. lateSrc.Put(chainCfg("chain", nil, t("a"), t("b"))) + Expect(late.HasChains()).To(BeFalse()) + _, err := late.Plan("chain") + Expect(err).ToNot(HaveOccurred()) Expect(late.HasChains()).To(BeTrue()) - Expect(m.HasChains()).To(BeTrue()) - var nilManager *Manager - Expect(nilManager.HasChains()).To(BeFalse()) + lateSrc.Delete("chain") + late.Sync() + Expect(late.HasChains()).To(BeFalse()) }) It("returns ErrChainNotFound for an unknown chain", func() {