feat(failover): share state through a sync hook and gate probes on a leader

Assisted-by: Claude:claude-opus-5-5
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
This commit is contained in:
Ettore Di Giacinto committed 2026-09-27 07:42:20 +00:00
1 parent 6738478c8a
commit 65c22c53bc
4 files changed
+512 -11

No files matched your search

+65 -9
View File
@@ -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
+45 -2
View File
@@ -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
+218
View File
@@ -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)
}
}
+184
View File
@@ -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())
})
})