fix(failover): roll back a pin that could not be shared

Pin and Unpin apply the change locally first for read-your-writes. When
the shared write then failed, the local pin stayed, so this frontend
served a target the others did not. Restore the previous pin on error,
unless a newer change arrived meanwhile.

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 ed9951cf7c
commit 4f801dfc99
2 files changed
+58 -4

No files matched your search

+26 -4
View File
@@ -587,17 +587,34 @@ func (m *Manager) Pin(chain, target string) error {
m.unlockAndFlush()
return fmt.Errorf("%w: %q", ErrTargetNotInChain, target)
}
prev := m.pins[chain]
ch.pinned = target
m.pins[chain] = target
m.recomputeLocked(ch, ReasonManual)
s := m.sync
m.unlockAndFlush()
if s != nil {
return s.SetPin(chain, target)
if s == nil {
return nil
}
if err := s.SetPin(chain, target); err != nil {
m.rollbackPin(chain, target, prev)
return err
}
return nil
}
// rollbackPin restores prev after a pin change that the other frontends
// never saw: serving it here alone would split the cluster. A newer change
// made meanwhile is kept.
func (m *Manager) rollbackPin(chain, applied, prev string) {
m.mu.Lock()
current := m.pins[chain]
m.mu.Unlock()
if current == applied {
m.ApplyPin(chain, prev)
}
}
func (m *Manager) Unpin(chain string) error {
m.mu.Lock()
ch := m.chainLocked(chain)
@@ -605,13 +622,18 @@ func (m *Manager) Unpin(chain string) error {
m.unlockAndFlush()
return fmt.Errorf("%w: %q", ErrChainNotFound, chain)
}
prev := m.pins[chain]
ch.pinned = ""
delete(m.pins, chain)
m.recomputeLocked(ch, ReasonManual)
s := m.sync
m.unlockAndFlush()
if s != nil {
return s.ClearPin(chain)
if s == nil {
return nil
}
if err := s.ClearPin(chain); err != nil {
m.rollbackPin(chain, "", prev)
return err
}
return nil
}
+32
View File
@@ -52,6 +52,12 @@ func (l *loopSync) Pins() map[string]string {
return out
}
// failingPinSync is a StateSync whose pin writes fail (the DB is down).
type failingPinSync struct{ *loopSync }
func (failingPinSync) SetPin(string, string) error { return errBoom }
func (failingPinSync) ClearPin(string) error { return errBoom }
var _ = Describe("Manager state sync", func() {
var (
clock *fakeClock
@@ -172,6 +178,32 @@ var _ = Describe("Manager state sync", func() {
Expect(st.Pinned).ToNot(BeNil())
})
It("rolls a pin back when the shared write fails", func() {
m := New(src, WithClock(clock))
m.SetStateSync(failingPinSync{bus})
Expect(m.Pin("chain", "y")).To(MatchError(errBoom))
st, _ := m.ChainStatus("chain")
Expect(st.Pinned).To(BeNil())
})
It("restores the previous pin when a re-pin or an unpin fails to share", func() {
m := New(src, WithClock(clock))
m.SetStateSync(bus)
Expect(m.Pin("chain", "y")).To(Succeed())
m.SetStateSync(failingPinSync{bus})
Expect(m.Pin("chain", "x")).To(MatchError(errBoom))
st, _ := m.ChainStatus("chain")
Expect(st.Pinned).ToNot(BeNil())
Expect(*st.Pinned).To(Equal("y"))
Expect(m.Unpin("chain")).To(MatchError(errBoom))
st, _ = m.ChainStatus("chain")
Expect(st.Pinned).ToNot(BeNil())
Expect(*st.Pinned).To(Equal("y"))
Expect(st.Active).To(Equal("y"))
})
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.