fix(failover): re-sync pins on every frontend after a missed delta

A NATS reconnect re-hydrates the pins map from the DB without OnApply, so
a frontend that missed an unpin kept serving the old pin and flip-flopped
with the leader's republish. Every frontend now reconciles the manager's
pins with the shared set every ten ticks and after a reconnect, and the
pins map re-reads the DB every 30 s to repair a delta dropped without a
reconnect.

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:21 +00:00
1 parent fe8519d7b5
commit ed9951cf7c
6 files changed
+126 -10

No files matched your search

+12 -1
View File
@@ -50,6 +50,15 @@ up. If your map has no durable backing, a leader (or another privileged
writer) must republish its live state after a reconnect instead of relying on
`Reconcile` to recover it.
**Gotcha:** a hydrate (on `Start`, after a NATS reconnect, and on every
`Reconcile` tick) replaces the map's contents **without firing `OnApply`**.
If `OnApply` feeds derived state (the failover manager's pins, a cache, a
running process), that state stays stale after a reconnect or a repaired
missed delta. Re-sync the derived state from the map's `Snapshot()`: on a
periodic tick, and/or from an `OnReconnect` callback registered after the
map's `Start` (callbacks run in registration order, so the map has already
re-hydrated). `failover.Manager.ReconcilePins` is the example.
### 2. Single-runner
Use `advisorylock.RunLeaderLoop` or `advisorylock.TryWithLockCtx`
@@ -117,7 +126,9 @@ When your PR adds or changes state that lives longer than a single request:
(`advisorylock`) / stateless / documented per-instance
- [ ] If shared: `Store` added if the state must survive a cluster restart,
or a `Loader` for standalone rehydration; if neither, a leader
republishes after reconnect instead of relying on bare `Reconcile`
republishes after reconnect instead of relying on bare `Reconcile`;
state derived through `OnApply` re-syncs from `Snapshot()` after a
hydrate (hydrate fires no `OnApply`)
- [ ] If single-runner: new lock key added to `keys.go`; `HeldLock` chosen
over `RunLeaderLoop`/`TryWithLockCtx` if leadership must be sticky
across ticks
+24 -5
View File
@@ -19,11 +19,15 @@ var (
_ failover.StateSync = (*Sync)(nil)
)
// pinReconcileInterval is how often the pins map re-reads the DB.
const pinReconcileInterval = 30 * time.Second
// Sync is a failover.StateSync backed by three syncstate.SyncedMaps:
//
// - failover.pins keeps the durable source of truth: Store-backed (when a
// PinStore is given) so a late joiner hydrates every pin from the DB on
// Start, not just from whatever peers happen to broadcast afterwards.
// Start, not just from whatever peers happen to broadcast afterwards,
// and re-reads it every pinReconcileInterval to repair a missed delta.
// - failover.targets and failover.chains are ephemeral live-health state,
// NATS-only with no Store and no Reconcile: with neither set, a Reconcile
// tick's hydrate is a no-op (nothing durable to pull from), so it could
@@ -53,11 +57,20 @@ func New(ctx context.Context, nats messaging.MessagingClient, pins syncstate.Sto
pinStore = pins
}
// A delta dropped without a reconnect would leave this map stale until
// the next reconnect; re-reading the DB repairs it, and the manager's
// periodic ReconcilePins carries the repair into the chains. It also
// drops a pin that a failed Store write left in memory.
var reconcile time.Duration
if pinStore != nil {
reconcile = pinReconcileInterval
}
s.pins = syncstate.New(syncstate.Config[string, PinRecord]{
Name: "failover.pins",
Key: func(p PinRecord) string { return p.Chain },
Nats: nats,
Store: pinStore,
Name: "failover.pins",
Key: func(p PinRecord) string { return p.Chain },
Nats: nats,
Store: pinStore,
Reconcile: reconcile,
OnApply: func(op string, chain string, v PinRecord) {
if op == "delete" {
m.ApplyPin(chain, "")
@@ -98,6 +111,12 @@ func New(ctx context.Context, nats messaging.MessagingClient, pins syncstate.Sto
}
m.SetStateSync(s)
// The pins map re-hydrates from the DB on reconnect without OnApply, so
// hand the manager the result. Registered after the map's own callback,
// which runs first.
if r, ok := nats.(interface{ OnReconnect(func()) }); ok {
r.OnReconnect(m.ReconcilePins)
}
return s, nil
}
@@ -156,6 +156,28 @@ var _ = Describe("distsync", func() {
Expect(*stC.Pinned).To(Equal("y"))
})
It("re-applies pins that changed while B missed the deltas, after a reconnect", func() {
a, _ := newManager("a")
b, _ := newManager("b")
a.Tick(ctx)
b.Tick(ctx)
// Written to the DB with no broadcast: B's map only learns it from
// the reconnect re-hydrate, which fires no OnApply.
Expect(pinStore.Upsert(ctx, distsync.PinRecord{Chain: "chain", Target: "y", UpdatedAt: time.Now()})).To(Succeed())
bus.TriggerReconnect()
st, _ := b.ChainStatus("chain")
Expect(st.Pinned).ToNot(BeNil())
Expect(*st.Pinned).To(Equal("y"))
Expect(pinStore.Delete(ctx, "chain")).To(Succeed())
bus.TriggerReconnect()
st, _ = b.ChainStatus("chain")
Expect(st.Pinned).To(BeNil())
st, _ = a.ChainStatus("chain")
Expect(st.Pinned).To(BeNil())
})
It("a trip on B makes A's plan skip the target", func() {
a, _ := newManager("a")
b, _ := newManager("b")
+7 -3
View File
@@ -23,10 +23,14 @@ func (m *Manager) Run(ctx context.Context) {
case <-ticker.C:
m.Tick(ctx)
// Publishes are fire-and-forget, so a frontend that missed one
// (restart, dropped message) converges within ten seconds.
// (restart, dropped message) converges within ten seconds. Pins
// are not republished: every frontend re-reads the shared set.
m.ticks++
if m.ticks%10 == 0 && m.IsLeader() {
m.Republish()
if m.ticks%10 == 0 {
m.ReconcilePins()
if m.IsLeader() {
m.Republish()
}
}
}
}
+29 -1
View File
@@ -69,7 +69,35 @@ func (m *Manager) SetStateSync(s StateSync) {
if s == nil {
return
}
for chain, target := range s.Pins() {
m.ReconcilePins()
}
// ReconcilePins makes this frontend's pins match the shared pin set. Pins
// normally arrive as deltas, but a re-hydrate of the shared set (after a NATS
// reconnect, or a missed delta repaired from the DB) changes it without
// delivering them; a frontend left with a stale pin would serve it while the
// others do not. The scheduler runs this periodically on every frontend.
func (m *Manager) ReconcilePins() {
m.mu.Lock()
s := m.sync
m.mu.Unlock()
if s == nil {
return
}
// Read outside the lock: the store may need I/O.
want := s.Pins()
m.mu.Lock()
var stale []string
for chain := range m.pins {
if _, ok := want[chain]; !ok {
stale = append(stale, chain)
}
}
m.mu.Unlock()
for _, chain := range stale {
m.ApplyPin(chain, "")
}
for chain, target := range want {
m.ApplyPin(chain, target)
}
}
+32
View File
@@ -140,6 +140,38 @@ var _ = Describe("Manager state sync", func() {
Expect(st.Pinned).ToNot(BeNil())
})
It("ReconcilePins converges a frontend that missed a pin and an unpin", func() {
// The shared pin set changes without B hearing the delta, as after a
// NATS reconnect whose re-hydrate fires no OnApply.
bus.mu.Lock()
bus.pins["chain"] = "y"
bus.mu.Unlock()
b.ReconcilePins()
st, _ := b.ChainStatus("chain")
Expect(st.Pinned).ToNot(BeNil())
Expect(*st.Pinned).To(Equal("y"))
Expect(st.Active).To(Equal("y"))
bus.mu.Lock()
delete(bus.pins, "chain")
bus.mu.Unlock()
b.ReconcilePins()
st, _ = b.ChainStatus("chain")
Expect(st.Pinned).To(BeNil())
})
It("ReconcilePins leaves a pin for a chain this frontend does not know yet", func() {
bus.mu.Lock()
bus.pins["later"] = "y"
bus.mu.Unlock()
b.ReconcilePins()
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())
})
It("SetLeaderGate gates a manager built without one", func() {
// Production builds the manager before distributed init, so the gate
// arrives through the setter rather than the option.