diff --git a/core/services/failover/manager.go b/core/services/failover/manager.go index 1412d7edf..da5ba6079 100644 --- a/core/services/failover/manager.go +++ b/core/services/failover/manager.go @@ -43,11 +43,16 @@ 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 - src ConfigSource - clock Clock - prober Prober - onWarm func([]string) + mu sync.Mutex + src ConfigSource + clock Clock + prober Prober + onWarm func([]string) + // pins holds every known pin by chain name, including pins for chains + // or targets this frontend's config does not have yet: with a sync + // layer the pin is delivered once, and a chain that appears (or is + // rebuilt) later must still pick it up. + pins map[string]string targets map[string]*targetState chains map[string]*chainState subs map[int]chan Event @@ -119,6 +124,7 @@ func New(src ConfigSource, opts ...Option) *Manager { clock: realClock{}, targets: map[string]*targetState{}, chains: map[string]*chainState{}, + pins: map[string]string{}, subs: map[int]chan Event{}, } for _, o := range opts { @@ -155,9 +161,14 @@ func (m *Manager) syncLocked() { } ch := m.chains[c.Name] if ch == nil || !slices.Equal(ch.targets, names) { - pinned := "" - if ch != nil && slices.Contains(names, ch.pinned) { - pinned = ch.pinned + pinned := m.pins[c.Name] + if !slices.Contains(names, pinned) { + if m.sync == nil { + // Standalone, the pin lives with the chain: a target + // dropped from it takes the pin along. + delete(m.pins, c.Name) + } + pinned = "" } ch = &chainState{name: c.Name, targets: names, activeSince: now, state: ChainPrimary, pinned: pinned} m.chains[c.Name] = ch @@ -191,6 +202,11 @@ func (m *Manager) syncLocked() { for name := range m.chains { if !seenChains[name] { delete(m.chains, name) + if m.sync == nil { + // With a sync layer the shared pin outlives a chain this + // frontend has not (re)loaded yet; standalone it does not. + delete(m.pins, name) + } } } for name := range m.targets { @@ -572,6 +588,7 @@ func (m *Manager) Pin(chain, target string) error { return fmt.Errorf("%w: %q", ErrTargetNotInChain, target) } ch.pinned = target + m.pins[chain] = target m.recomputeLocked(ch, ReasonManual) s := m.sync m.unlockAndFlush() @@ -589,6 +606,7 @@ func (m *Manager) Unpin(chain string) error { return fmt.Errorf("%w: %q", ErrChainNotFound, chain) } ch.pinned = "" + delete(m.pins, chain) m.recomputeLocked(ch, ReasonManual) s := m.sync m.unlockAndFlush() diff --git a/core/services/failover/statesync.go b/core/services/failover/statesync.go index d63a01b96..bcab901ef 100644 --- a/core/services/failover/statesync.go +++ b/core/services/failover/statesync.go @@ -176,16 +176,23 @@ func (m *Manager) ApplyChain(s ChainSnapshot) { } } -// ApplyPin sets (target != "") or clears a pin set on any frontend. +// ApplyPin sets (target != "") or clears a pin set on any frontend. A pin for +// a chain or target this frontend does not know yet is kept and applied by +// syncLocked once the config catches up. func (m *Manager) ApplyPin(chain, target string) { m.mu.Lock() defer m.unlockAndFlush() + if target == "" { + delete(m.pins, chain) + } else { + m.pins[chain] = target + } ch := m.chainLocked(chain) if ch == nil { return } if target != "" && !slices.Contains(ch.targets, target) { - xlog.Debug("failover: ignoring pin to a target not in the chain", "chain", chain, "target", target) + xlog.Debug("failover: deferring pin to a target not in the chain yet", "chain", chain, "target", target) return } if ch.pinned == target { diff --git a/core/services/failover/statesync_test.go b/core/services/failover/statesync_test.go index 7e6f7aa62..3901d3b95 100644 --- a/core/services/failover/statesync_test.go +++ b/core/services/failover/statesync_test.go @@ -176,6 +176,67 @@ var _ = Describe("Manager state sync", func() { Expect(st.Active).To(Equal("y")) }) + It("keeps a pin for a chain it does not know yet and applies it when the chain appears", func() { + b.ApplyPin("later", "y") + _, ok := b.ChainStatus("later") + Expect(ok).To(BeFalse()) + src.Put(chainCfg("later", nil, t("x"), t("y"))) + b.Sync() + st, ok := b.ChainStatus("later") + Expect(ok).To(BeTrue()) + Expect(st.Pinned).ToNot(BeNil()) + Expect(*st.Pinned).To(Equal("y")) + Expect(st.Active).To(Equal("y")) + }) + + It("re-applies a shared pin when the chain is removed and re-added", func() { + Expect(a.Pin("chain", "y")).To(Succeed()) + src.Delete("chain") + b.Sync() + _, ok := b.ChainStatus("chain") + Expect(ok).To(BeFalse()) + src.Put(chainCfg("chain", nil, t("x"), t("y"))) + b.Sync() + st, _ := b.ChainStatus("chain") + Expect(st.Pinned).ToNot(BeNil()) + Expect(*st.Pinned).To(Equal("y")) + }) + + It("applies a deferred pin once its target joins the chain", func() { + src.Put(local("z")) + b.ApplyPin("chain", "z") + st, _ := b.ChainStatus("chain") + Expect(st.Pinned).To(BeNil()) + src.Put(chainCfg("chain", nil, t("x"), t("y"), t("z"))) + b.Sync() + st, _ = b.ChainStatus("chain") + Expect(st.Pinned).ToNot(BeNil()) + Expect(*st.Pinned).To(Equal("z")) + }) + + It("hydrates a pin for an unknown chain and applies it when the chain appears", func() { + bus.pins["later"] = "y" + c := New(src, WithClock(clock), WithLeaderGate(gateFor(false))) + c.SetStateSync(bus) + src.Put(chainCfg("later", nil, t("x"), t("y"))) + c.Sync() + st, ok := c.ChainStatus("later") + Expect(ok).To(BeTrue()) + Expect(st.Pinned).ToNot(BeNil()) + Expect(*st.Pinned).To(Equal("y")) + }) + + It("standalone, a removed chain drops its pin", func() { + m := New(src, WithClock(clock)) + Expect(m.Pin("chain", "y")).To(Succeed()) + src.Delete("chain") + m.Sync() + src.Put(chainCfg("chain", nil, t("x"), t("y"))) + m.Sync() + st, _ := m.ChainStatus("chain") + Expect(st.Pinned).To(BeNil()) + }) + It("standalone manager (no sync, no gate) is always leader", func() { m := New(src, WithClock(clock)) m.Tick(ctx)