diff --git a/core/services/failover/manager.go b/core/services/failover/manager.go index 847584d62..1412d7edf 100644 --- a/core/services/failover/manager.go +++ b/core/services/failover/manager.go @@ -60,6 +60,21 @@ type Manager struct { hasChains atomic.Bool // probes counts running probes; only tests wait on it. probes sync.WaitGroup + + // sync shares state with other frontends; nil when standalone. + sync StateSync + // gate grants probing and chain decisions to one frontend; nil means + // this manager is always the leader. + gate LeaderGate + leader bool + // applying is set while a peer's target state is applied, so the + // transition is not published back to the peers. + applying bool + // pending holds publishes queued under the lock; unlockAndFlush runs + // them after unlocking because the sync layer can call back into Apply*. + pending []func() + // ticks counts Run's ticks for the periodic republish; only Run uses it. + ticks int } type targetState struct { @@ -73,6 +88,9 @@ type targetState struct { lastProbe time.Time lastActivity time.Time lastError string + // since and reason describe the last state change, for snapshots. + since time.Time + reason Reason // probing is set while a probe runs, so the scheduler does not start a // second one for the same target. probing bool @@ -90,6 +108,9 @@ type chainState struct { activeSince time.Time pinned string state ChainState + // adopted is set once a follower received the leader's decision for this + // chain; from then on it stops choosing the active target itself. + adopted bool } func New(src ConfigSource, opts ...Option) *Manager { @@ -103,6 +124,7 @@ func New(src ConfigSource, opts ...Option) *Manager { for _, o := range opts { o(m) } + m.leader = m.gate == nil return m } @@ -112,7 +134,7 @@ func (m *Manager) Sync() { m.mu.Lock() m.syncLocked() warm, deliver := m.takeWarmLocked() - m.mu.Unlock() + m.unlockAndFlush() if deliver && m.onWarm != nil { m.onWarm(warm) } @@ -264,7 +286,7 @@ func (m *Manager) targetStates() map[string]TargetState { // the scheduler calls this on every tick. func (m *Manager) Reevaluate() { m.mu.Lock() - defer m.mu.Unlock() + defer m.unlockAndFlush() for _, ch := range m.chains { m.recomputeLocked(ch, "") } @@ -277,6 +299,7 @@ func (m *Manager) setTargetLocked(ts *targetState, to TargetState, reason Reason from := ts.state now := m.clock.Now() ts.state = to + ts.since, ts.reason = now, reason switch to { case StateDown: ts.downSince = now @@ -287,11 +310,19 @@ func (m *Manager) setTargetLocked(ts *targetState, to TargetState, reason Reason ts.failures = nil } m.emitLocked(Event{Type: EventTargetState, Target: ts.name, From: string(from), To: string(to), Reason: reason, Error: errMsg, At: now}) + m.queuePublishTargetLocked(ts, from) } // recomputeLocked picks the active target. override replaces the reason of a // resulting switch (pin and unpin are always "manual"). func (m *Manager) recomputeLocked(ch *chainState, override Reason) { + if m.sync != nil && !m.leader && ch.adopted && ch.pinned == "" { + // The leader decides; deciding here too would let frontends serve + // different targets. A pin is exempt: it fixes the active target + // the same way on every frontend, and applying it at once gives the + // caller read-your-writes. + return + } now := m.clock.Now() prev := ch.active next := prev @@ -333,6 +364,7 @@ func (m *Manager) recomputeLocked(ch *chainState, override Reason) { default: state = ChainFallback } + changed := next != prev || state != ch.state switch { case next != prev: ch.active = next @@ -348,6 +380,17 @@ func (m *Manager) recomputeLocked(ch *chainState, override Reason) { m.emitLocked(Event{Type: EventChainSwitched, Chain: ch.name, From: ch.targets[prev], To: ch.targets[next], State: string(state), Reason: ReasonRecovery, At: now}) } ch.state = state + if changed { + pub := reason + if next == prev { + // Same reasons as the events above for a state-only change. + pub = ReasonRecovery + if state == ChainDegraded { + pub = ReasonDegraded + } + } + m.queuePublishChainLocked(ch, pub) + } } func (m *Manager) recomputeForLocked(target string) { @@ -373,7 +416,7 @@ type Attempt struct { // order; a pinned chain only the pinned target. func (m *Manager) Plan(chain string) (*Attempt, error) { m.mu.Lock() - defer m.mu.Unlock() + defer m.unlockAndFlush() ch := m.chainLocked(chain) if ch == nil { return nil, fmt.Errorf("%w: %q", ErrChainNotFound, chain) @@ -447,7 +490,7 @@ func (a *Attempt) Succeed() { a.m.ReportSuccess(a.Target()) } func (m *Manager) ReportFailure(target string, err error) { m.mu.Lock() - defer m.mu.Unlock() + defer m.unlockAndFlush() ts := m.targetLocked(target) if ts == nil { return @@ -483,7 +526,7 @@ func (m *Manager) recordFailureLocked(ts *targetState, msg string) { func (m *Manager) ReportSuccess(target string) { m.mu.Lock() - defer m.mu.Unlock() + defer m.unlockAndFlush() ts := m.targetLocked(target) if ts == nil { return @@ -515,37 +558,50 @@ func (m *Manager) recordPassLocked(ts *targetState) { } } +// Pin takes effect here at once (read-your-writes), then is shared with the +// other frontends. func (m *Manager) Pin(chain, target string) error { m.mu.Lock() - defer m.mu.Unlock() ch := m.chainLocked(chain) if ch == nil { + m.unlockAndFlush() return fmt.Errorf("%w: %q", ErrChainNotFound, chain) } if !slices.Contains(ch.targets, target) { + m.unlockAndFlush() return fmt.Errorf("%w: %q", ErrTargetNotInChain, target) } ch.pinned = target m.recomputeLocked(ch, ReasonManual) + s := m.sync + m.unlockAndFlush() + if s != nil { + return s.SetPin(chain, target) + } return nil } func (m *Manager) Unpin(chain string) error { m.mu.Lock() - defer m.mu.Unlock() ch := m.chainLocked(chain) if ch == nil { + m.unlockAndFlush() return fmt.Errorf("%w: %q", ErrChainNotFound, chain) } ch.pinned = "" m.recomputeLocked(ch, ReasonManual) + s := m.sync + m.unlockAndFlush() + if s != nil { + return s.ClearPin(chain) + } return nil } // Status returns every chain, sorted by name. func (m *Manager) Status() []ChainStatus { m.mu.Lock() - defer m.mu.Unlock() + defer m.unlockAndFlush() m.syncLocked() names := make([]string, 0, len(m.chains)) for name := range m.chains { @@ -561,7 +617,7 @@ func (m *Manager) Status() []ChainStatus { func (m *Manager) ChainStatus(name string) (ChainStatus, bool) { m.mu.Lock() - defer m.mu.Unlock() + defer m.unlockAndFlush() ch := m.chainLocked(name) if ch == nil { return ChainStatus{}, false diff --git a/core/services/failover/schedule.go b/core/services/failover/schedule.go index 08f138d00..739c48b42 100644 --- a/core/services/failover/schedule.go +++ b/core/services/failover/schedule.go @@ -22,6 +22,12 @@ func (m *Manager) Run(ctx context.Context) { return case <-ticker.C: m.Tick(ctx) + // Publishes are fire-and-forget, so a frontend that missed one + // (restart, dropped message) converges within ten seconds. + m.ticks++ + if m.ticks%10 == 0 && m.IsLeader() { + m.Republish() + } } } } @@ -34,8 +40,45 @@ func (m *Manager) Run(ctx context.Context) { // hangs until its timeout) must not delay probing and fail-back of every // other chain. Each probe applies its own result, and a target whose probe is // still running is skipped until it ends. +// +// With a leader gate, only the leader probes and decides chains; followers +// still recompute, which only moves chains that have not yet adopted a +// leader decision. func (m *Manager) Tick(ctx context.Context) { m.Sync() + if m.gate == nil { + m.lead(ctx) + return + } + if m.gate(ctx, func() { m.lead(ctx) }) { + return + } + m.mu.Lock() + m.leader = false + m.mu.Unlock() + m.Reevaluate() +} + +// lead is the leader's share of a tick. +func (m *Manager) lead(ctx context.Context) { + m.mu.Lock() + became := !m.leader + m.leader = true + if became { + // The previous leader owned the warm-set callback's effects; + // deliver the set again so this frontend takes them over. + m.warmPending = true + } + warm, deliver := m.takeWarmLocked() + m.unlockAndFlush() + if deliver && m.onWarm != nil { + m.onWarm(warm) + } + if became { + // Followers hold the old leader's view; send ours at once instead of + // letting them wait for the periodic republish. + m.Republish() + } for _, j := range m.dueProbes() { m.probes.Add(1) go func(j probeJob) { @@ -61,7 +104,7 @@ type probeJob struct { func (m *Manager) dueProbes() []probeJob { m.mu.Lock() - defer m.mu.Unlock() + defer m.unlockAndFlush() now := m.clock.Now() var jobs []probeJob for _, ts := range m.targets { @@ -139,7 +182,7 @@ func (m *Manager) endProbe(target string) { func (m *Manager) applyProbe(j probeJob, err error) { m.mu.Lock() - defer m.mu.Unlock() + defer m.unlockAndFlush() ts := m.targets[j.target] if ts == nil { return diff --git a/core/services/failover/statesync.go b/core/services/failover/statesync.go new file mode 100644 index 000000000..d63a01b96 --- /dev/null +++ b/core/services/failover/statesync.go @@ -0,0 +1,218 @@ +package failover + +import ( + "context" + "slices" + "time" + + "github.com/mudler/xlog" +) + +// TargetSnapshot is one target's health as shared between frontends. +type TargetSnapshot struct { + Target string `json:"target"` + State TargetState `json:"state"` + Reason Reason `json:"reason"` + Error string `json:"error,omitempty"` + ConsecutiveOK int `json:"consecutive_ok"` + Since time.Time `json:"since"` +} + +// ChainSnapshot is one chain's active target as decided by the leader. +type ChainSnapshot struct { + Chain string `json:"chain"` + Active string `json:"active"` + ActiveSince time.Time `json:"active_since"` + State ChainState `json:"state"` + Reason Reason `json:"reason"` +} + +// StateSync shares failover state between frontends. Implementations may +// deliver a publish back to the publisher synchronously (NATS echoes), so the +// manager never calls it while holding its lock. +type StateSync interface { + PublishTarget(TargetSnapshot) + PublishChain(ChainSnapshot) + SetPin(chain, target string) error + ClearPin(chain string) error + Pins() map[string]string +} + +// LeaderGate runs fn only on the one frontend that holds leadership and +// reports whether it did. Probing and chain decisions happen on the leader +// only, so N frontends do not probe every target N times or disagree on the +// active target. +type LeaderGate func(ctx context.Context, fn func()) bool + +// WithLeaderGate makes the manager probe and decide chains only while the gate +// grants leadership. Without it the manager is always the leader. +func WithLeaderGate(g LeaderGate) Option { return func(m *Manager) { m.gate = g } } + +// SetStateSync attaches the sync layer and hydrates pins from it. Pins are +// read outside the lock because the store may need I/O. +func (m *Manager) SetStateSync(s StateSync) { + m.mu.Lock() + m.sync = s + m.mu.Unlock() + if s == nil { + return + } + for chain, target := range s.Pins() { + m.ApplyPin(chain, target) + } +} + +// IsLeader reports whether this manager probed and decided chains at its last +// tick. A standalone manager (no gate) is always the leader. +func (m *Manager) IsLeader() bool { + m.mu.Lock() + defer m.mu.Unlock() + return m.leader +} + +// unlockAndFlush releases the lock and then runs the publishes queued while +// it was held: the sync layer may call straight back into Apply*, which takes +// the lock again. +func (m *Manager) unlockAndFlush() { + pending := m.pending + m.pending = nil + m.mu.Unlock() + for _, f := range pending { + f() + } +} + +func (m *Manager) targetSnapshotLocked(ts *targetState) TargetSnapshot { + since := ts.since + if ts.state == StateDown { + since = ts.downSince + } + return TargetSnapshot{ + Target: ts.name, State: ts.state, Reason: ts.reason, Error: ts.lastError, + ConsecutiveOK: ts.consecutiveOK, Since: since, + } +} + +func (m *Manager) chainSnapshotLocked(ch *chainState, reason Reason) ChainSnapshot { + return ChainSnapshot{ + Chain: ch.name, Active: ch.targets[ch.active], ActiveSince: ch.activeSince, + State: ch.state, Reason: reason, + } +} + +// queuePublishTargetLocked shares a local target transition. Missing is a fact +// about this frontend's config, not about the target, so it is never shared. +func (m *Manager) queuePublishTargetLocked(ts *targetState, from TargetState) { + if m.sync == nil || m.applying || ts.state == StateMissing || from == StateMissing { + return + } + s, snap := m.sync, m.targetSnapshotLocked(ts) + m.pending = append(m.pending, func() { s.PublishTarget(snap) }) +} + +func (m *Manager) queuePublishChainLocked(ch *chainState, reason Reason) { + if m.sync == nil || !m.leader { + return + } + s, snap := m.sync, m.chainSnapshotLocked(ch, 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. +func (m *Manager) ApplyTarget(s TargetSnapshot) { + m.mu.Lock() + defer m.unlockAndFlush() + ts := m.targetLocked(s.Target) + if ts == nil || ts.state == StateMissing || s.State == StateMissing { + return + } + if ts.state == s.State { + return + } + m.applying = true + m.setTargetLocked(ts, s.State, s.Reason, s.Error) + m.applying = false + // Keep the publisher's clock, so dwell timers (cold recovery, fail-back) + // run from when the target actually changed, not from when we heard. + if !s.Since.IsZero() { + ts.since = s.Since + if s.State == StateDown { + ts.downSince = s.Since + } + } + ts.consecutiveOK = s.ConsecutiveOK + if s.Error != "" { + ts.lastError = s.Error + } + m.recomputeForLocked(ts.name) +} + +// ApplyChain adopts the leader's decision for a chain. The leader ignores it: +// it is the source of these decisions, and a late publish from a previous +// leader must not undo its own. +func (m *Manager) ApplyChain(s ChainSnapshot) { + m.mu.Lock() + defer m.unlockAndFlush() + if m.leader { + return + } + ch := m.chainLocked(s.Chain) + if ch == nil { + return + } + next := slices.Index(ch.targets, s.Active) + if next < 0 { + // The frontends disagree on the chain's targets while a config + // change propagates; keep deciding locally until they agree. + xlog.Debug("failover: ignoring chain state for an unknown target", "chain", s.Chain, "target", s.Active) + return + } + prev := ch.active + changed := next != prev || s.State != ch.state + ch.active, ch.activeSince, ch.state, ch.adopted = next, s.ActiveSince, s.State, true + if changed { + m.emitLocked(Event{Type: EventChainSwitched, Chain: ch.name, From: ch.targets[prev], To: ch.targets[next], State: string(s.State), Reason: s.Reason, At: m.clock.Now()}) + } +} + +// ApplyPin sets (target != "") or clears a pin set on any frontend. +func (m *Manager) ApplyPin(chain, target string) { + m.mu.Lock() + defer m.unlockAndFlush() + 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) + return + } + if ch.pinned == target { + return + } + ch.pinned = target + m.recomputeLocked(ch, ReasonManual) +} + +// Republish sends every target and chain state, so frontends that missed a +// publish (joined late, dropped a message) converge. Only the leader's view +// is authoritative, so followers do nothing. +func (m *Manager) Republish() { + m.mu.Lock() + defer m.unlockAndFlush() + if m.sync == nil || !m.leader { + return + } + s := m.sync + for _, ts := range m.targets { + if ts.state == StateMissing { + continue + } + snap := m.targetSnapshotLocked(ts) + m.pending = append(m.pending, func() { s.PublishTarget(snap) }) + } + for _, ch := range m.chains { + m.queuePublishChainLocked(ch, ReasonInitial) + } +} diff --git a/core/services/failover/statesync_test.go b/core/services/failover/statesync_test.go new file mode 100644 index 000000000..7e6f7aa62 --- /dev/null +++ b/core/services/failover/statesync_test.go @@ -0,0 +1,184 @@ +package failover + +import ( + "context" + "sync" + "time" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" +) + +// loopSync is an in-process StateSync that delivers every publish to all +// managers synchronously, including the publisher (like NATS echo). +type loopSync struct { + mu sync.Mutex + peers []*Manager + pins map[string]string +} + +func (l *loopSync) add(m *Manager) { l.mu.Lock(); l.peers = append(l.peers, m); l.mu.Unlock() } +func (l *loopSync) each(f func(*Manager)) { + l.mu.Lock() + ps := append([]*Manager(nil), l.peers...) + l.mu.Unlock() + for _, p := range ps { + f(p) + } +} +func (l *loopSync) PublishTarget(s TargetSnapshot) { l.each(func(m *Manager) { m.ApplyTarget(s) }) } +func (l *loopSync) PublishChain(s ChainSnapshot) { l.each(func(m *Manager) { m.ApplyChain(s) }) } +func (l *loopSync) SetPin(c, t string) error { + l.mu.Lock() + l.pins[c] = t + l.mu.Unlock() + l.each(func(m *Manager) { m.ApplyPin(c, t) }) + return nil +} +func (l *loopSync) ClearPin(c string) error { + l.mu.Lock() + delete(l.pins, c) + l.mu.Unlock() + l.each(func(m *Manager) { m.ApplyPin(c, "") }) + return nil +} +func (l *loopSync) Pins() map[string]string { + l.mu.Lock() + defer l.mu.Unlock() + out := map[string]string{} + for k, v := range l.pins { + out[k] = v + } + return out +} + +var _ = Describe("Manager state sync", func() { + var ( + clock *fakeClock + src *fakeSource + bus *loopSync + a, b *Manager + leaderIsA bool + gateFor func(isA bool) LeaderGate + ctx = context.Background() + ) + + BeforeEach(func() { + clock = newFakeClock() + src = newFakeSource(remote("x"), local("y"), chainCfg("chain", nil, t("x"), t("y"))) + bus = &loopSync{pins: map[string]string{}} + leaderIsA = true + gateFor = func(isA bool) LeaderGate { + return func(_ context.Context, fn func()) bool { + if isA != leaderIsA { + return false + } + fn() + return true + } + } + a = New(src, WithClock(clock), WithLeaderGate(gateFor(true))) + b = New(src, WithClock(clock), WithLeaderGate(gateFor(false))) + bus.add(a) + bus.add(b) + a.SetStateSync(bus) + b.SetStateSync(bus) + a.Tick(ctx) + b.Tick(ctx) + }) + + It("echo of own publish is a no-op and emits one event", func() { + events, cancel := a.Subscribe(16) + defer cancel() + a.ReportFailure("x", errBoom) + n := 0 + for _, e := range drain(events) { + if e.Type == EventTargetState && e.Target == "x" { + n++ + } + } + Expect(n).To(Equal(1)) + }) + + It("a trip on one frontend is skipped by the other's plan", func() { + b.ReportFailure("x", errBoom) + att, err := a.Plan("chain") + Expect(err).ToNot(HaveOccurred()) + Expect(att.Target()).To(Equal("y")) + }) + + It("followers adopt the leader's chain state and emit the switch", func() { + events, cancel := b.Subscribe(16) + defer cancel() + a.ReportFailure("x", errBoom) // leader recomputes and publishes chain state + st, _ := b.ChainStatus("chain") + Expect(st.Active).To(Equal("y")) + var sw []Event + for _, e := range drain(events) { + if e.Type == EventChainSwitched { + sw = append(sw, e) + } + } + Expect(sw).ToNot(BeEmpty()) + }) + + It("a pin on one frontend applies on all", func() { + Expect(b.Pin("chain", "y")).To(Succeed()) + st, _ := a.ChainStatus("chain") + Expect(st.Pinned).ToNot(BeNil()) + Expect(*st.Pinned).To(Equal("y")) + Expect(a.Unpin("chain")).To(Succeed()) + st, _ = b.ChainStatus("chain") + Expect(st.Pinned).To(BeNil()) + }) + + It("hydrates pins when the sync is attached", func() { + bus.pins["chain"] = "y" + c := New(src, WithClock(clock), WithLeaderGate(gateFor(false))) + c.SetStateSync(bus) + st, _ := c.ChainStatus("chain") + Expect(st.Pinned).ToNot(BeNil()) + }) + + It("only the leader probes", func() { + pa, pb := &fakeProber{fail: map[string]error{}}, &fakeProber{fail: map[string]error{}} + a = New(src, WithClock(clock), WithProber(pa), WithLeaderGate(gateFor(true))) + b = New(src, WithClock(clock), WithProber(pb), WithLeaderGate(gateFor(false))) + a.SetStateSync(bus) + b.SetStateSync(bus) + a.Tick(ctx) + b.Tick(ctx) + Eventually(func() int { return len(pa.take()) }).Should(BeNumerically(">", 0)) + Consistently(func() int { return len(pb.take()) }, 200*time.Millisecond).Should(Equal(0)) + Expect(a.IsLeader()).To(BeTrue()) + Expect(b.IsLeader()).To(BeFalse()) + }) + + It("new leader keeps activeSince across a leadership move", func() { + a.ReportFailure("x", errBoom) + before, _ := b.ChainStatus("chain") + leaderIsA = false + clock.Advance(5 * time.Second) + a.Tick(ctx) + b.Tick(ctx) + after, _ := b.ChainStatus("chain") + Expect(after.ActiveSince).To(Equal(before.ActiveSince)) + Expect(b.IsLeader()).To(BeTrue()) + }) + + It("republish sends every target and chain", func() { + c := New(src, WithClock(clock), WithLeaderGate(gateFor(false))) + bus.add(c) + c.SetStateSync(bus) + a.ReportFailure("x", errBoom) + a.Republish() + st, _ := c.ChainStatus("chain") + Expect(st.Active).To(Equal("y")) + }) + + It("standalone manager (no sync, no gate) is always leader", func() { + m := New(src, WithClock(clock)) + m.Tick(ctx) + Expect(m.IsLeader()).To(BeTrue()) + }) +})