feat(failover): sync pins, target health and chain state over NATS

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 5f98dfcf95
commit 7b88674fc5
4 files changed
+492

No files matched your search

+158
View File
@@ -0,0 +1,158 @@
package distsync
import (
"context"
"errors"
"fmt"
"reflect"
"time"
"github.com/mudler/LocalAI/core/services/failover"
"github.com/mudler/LocalAI/core/services/messaging"
"github.com/mudler/LocalAI/core/services/syncstate"
"github.com/mudler/xlog"
)
// compile-time assertions.
var (
_ syncstate.Store[string, PinRecord] = (*PinStore)(nil)
_ failover.StateSync = (*Sync)(nil)
)
// 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.
// - failover.targets and failover.chains are ephemeral live-health state,
// NATS-only with no Store and no Reconcile: there is nothing durable to
// hydrate from, and a Reconcile tick would re-hydrate them empty and wipe
// live state clean off a running frontend. A late joiner instead catches
// up from the leader's periodic Republish (see failover.Manager.Republish).
type Sync struct {
pins *syncstate.SyncedMap[string, PinRecord]
targets *syncstate.SyncedMap[string, failover.TargetSnapshot]
chains *syncstate.SyncedMap[string, failover.ChainSnapshot]
}
// New builds and starts the three maps, then attaches the result to m via
// SetStateSync so any already-durable pins hydrate onto m immediately.
func New(ctx context.Context, nats messaging.MessagingClient, pins syncstate.Store[string, PinRecord], m *failover.Manager) (*Sync, error) {
s := &Sync{}
// pins is already typed as the Store interface (the brief fixes this
// signature), so a caller holding a nil *PinStore (e.g. standalone mode,
// no DB configured) and passing it straight through boxes it into a
// non-nil interface wrapping a nil pointer - "pins != nil" alone would
// not catch that, and the SyncedMap would then try to hydrate/write
// through a nil *gorm.DB. isNilStore catches both that and a literal nil
// argument, the same defense finetune/service.go gets for free by taking
// a concrete *distributed.FineTuneStore and nil-checking before boxing it.
var pinStore syncstate.Store[string, PinRecord]
if !isNilStore(pins) {
pinStore = pins
}
s.pins = syncstate.New(syncstate.Config[string, PinRecord]{
Name: "failover.pins",
Key: func(p PinRecord) string { return p.Chain },
Nats: nats,
Store: pinStore,
OnApply: func(op string, chain string, v PinRecord) {
if op == "delete" {
m.ApplyPin(chain, "")
return
}
m.ApplyPin(chain, v.Target)
},
})
if err := s.pins.Start(ctx); err != nil {
return nil, fmt.Errorf("distsync: starting pins map: %w", err)
}
s.targets = syncstate.New(syncstate.Config[string, failover.TargetSnapshot]{
Name: "failover.targets",
Key: func(t failover.TargetSnapshot) string { return t.Target },
Nats: nats,
OnApply: func(_ string, _ string, v failover.TargetSnapshot) {
m.ApplyTarget(v)
},
})
if err := s.targets.Start(ctx); err != nil {
_ = s.pins.Close()
return nil, fmt.Errorf("distsync: starting targets map: %w", err)
}
s.chains = syncstate.New(syncstate.Config[string, failover.ChainSnapshot]{
Name: "failover.chains",
Key: func(c failover.ChainSnapshot) string { return c.Chain },
Nats: nats,
OnApply: func(_ string, _ string, v failover.ChainSnapshot) {
m.ApplyChain(v)
},
})
if err := s.chains.Start(ctx); err != nil {
_ = s.targets.Close()
_ = s.pins.Close()
return nil, fmt.Errorf("distsync: starting chains map: %w", err)
}
m.SetStateSync(s)
return s, nil
}
// isNilStore reports whether pins is nil - either a literal nil argument, or
// the classic Go footgun of a non-nil interface value wrapping a nil
// pointer (e.g. a nil *PinStore passed in directly). Both must disable the
// durable Store the same way, so syncstate hydrates from nothing rather than
// panicking on a nil *gorm.DB the first time it dereferences it.
func isNilStore(pins syncstate.Store[string, PinRecord]) bool {
if pins == nil {
return true
}
v := reflect.ValueOf(pins)
return v.Kind() == reflect.Ptr && v.IsNil()
}
// Close releases all three maps' subscriptions and background workers.
func (s *Sync) Close() error {
return errors.Join(s.chains.Close(), s.targets.Close(), s.pins.Close())
}
// PublishTarget shares a target health transition. StateSync's methods
// return no error to the manager, and a publish must never block it (the
// manager calls this outside its lock precisely so a synchronous NATS echo
// is safe) - so a failure here is logged and dropped; a missed publish
// self-heals on the leader's next Republish.
func (s *Sync) PublishTarget(t failover.TargetSnapshot) {
if err := s.targets.Set(context.Background(), t); err != nil {
xlog.Warn("distsync: publishing target state failed", "target", t.Target, "error", err)
}
}
// PublishChain shares the leader's decision for a chain.
func (s *Sync) PublishChain(c failover.ChainSnapshot) {
if err := s.chains.Set(context.Background(), c); err != nil {
xlog.Warn("distsync: publishing chain state failed", "chain", c.Chain, "error", err)
}
}
// SetPin durably persists and broadcasts a pin.
func (s *Sync) SetPin(chain, target string) error {
return s.pins.Set(context.Background(), PinRecord{Chain: chain, Target: target, UpdatedAt: time.Now()})
}
// ClearPin durably removes and broadcasts a pin's removal.
func (s *Sync) ClearPin(chain string) error {
return s.pins.Delete(context.Background(), chain)
}
// Pins returns every known pin (chain -> target), for Manager.SetStateSync's
// hydrate-on-attach.
func (s *Sync) Pins() map[string]string {
out := make(map[string]string)
for chain, rec := range s.pins.Snapshot() {
out[chain] = rec.Target
}
return out
}
@@ -0,0 +1,13 @@
package distsync_test
import (
"testing"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
)
func TestDistsync(t *testing.T) {
RegisterFailHandler(Fail)
RunSpecs(t, "Distsync test suite")
}
@@ -0,0 +1,258 @@
package distsync_test
import (
"context"
"errors"
"sync"
"time"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
"github.com/mudler/LocalAI/core/config"
"github.com/mudler/LocalAI/core/services/failover"
"github.com/mudler/LocalAI/core/services/failover/distsync"
"github.com/mudler/LocalAI/core/services/testutil"
)
// fakeSource is a minimal failover.ConfigSource with one chain "chain" of two
// targets: "x" (remote, primary) and "y" (local warm, fallback). It is
// shared across the managers in a spec, mirroring how frontends in a real
// deployment read the same config loader.
type fakeSource struct {
mu sync.Mutex
cfgs map[string]config.ModelConfig
}
func newChainSource() *fakeSource {
return &fakeSource{cfgs: map[string]config.ModelConfig{
"x": {Name: "x", Backend: "cloud-proxy"},
"y": {Name: "y", Backend: "llama-cpp"},
"chain": {Name: "chain", Failover: &config.FailoverConfig{Targets: []config.FailoverTarget{
{Model: "x"},
{Model: "y", Warm: true},
}}},
}}
}
func (s *fakeSource) GetModelConfig(n string) (config.ModelConfig, bool) {
s.mu.Lock()
defer s.mu.Unlock()
c, ok := s.cfgs[n]
return c, ok
}
func (s *fakeSource) GetAllModelsConfigs() []config.ModelConfig {
s.mu.Lock()
defer s.mu.Unlock()
out := make([]config.ModelConfig, 0, len(s.cfgs))
for _, c := range s.cfgs {
out = append(out, c)
}
return out
}
// memPinStore is an in-memory syncstate.Store[string, distsync.PinRecord]
// shared by several distsync.Sync instances the way a real DB would be, so a
// spec can build a "late joiner" that hydrates from what earlier instances
// already wrote through.
type memPinStore struct {
mu sync.Mutex
data map[string]distsync.PinRecord
}
func newMemPinStore() *memPinStore { return &memPinStore{data: map[string]distsync.PinRecord{}} }
func (s *memPinStore) List(context.Context) ([]distsync.PinRecord, error) {
s.mu.Lock()
defer s.mu.Unlock()
out := make([]distsync.PinRecord, 0, len(s.data))
for _, v := range s.data {
out = append(out, v)
}
return out, nil
}
func (s *memPinStore) Upsert(_ context.Context, v distsync.PinRecord) error {
s.mu.Lock()
defer s.mu.Unlock()
s.data[v.Chain] = v
return nil
}
func (s *memPinStore) Delete(_ context.Context, k string) error {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.data, k)
return nil
}
// leaderGateFor grants leadership to exactly one named manager at a time
// (whatever *leader currently holds), mirroring an advisory-lock leader loop
// where only one frontend probes and decides chains.
func leaderGateFor(name string, leader *string) failover.LeaderGate {
return func(_ context.Context, fn func()) bool {
if *leader != name {
return false
}
fn()
return true
}
}
var errBoom = errors.New("boom")
var _ = Describe("distsync", func() {
var (
ctx context.Context
bus *testutil.FakeBus
src *fakeSource
pinStore *memPinStore
leader string
)
BeforeEach(func() {
ctx = context.Background()
bus = testutil.NewFakeBus()
src = newChainSource()
pinStore = newMemPinStore()
leader = "a"
})
// newManager wires a fresh failover.Manager to a fresh distsync.Sync on
// the shared bus and pin store, named so leaderGateFor can grant or deny
// it leadership.
newManager := func(name string) (*failover.Manager, *distsync.Sync) {
m := failover.New(src, failover.WithLeaderGate(leaderGateFor(name, &leader)))
s, err := distsync.New(ctx, bus, pinStore, m)
Expect(err).ToNot(HaveOccurred())
return m, s
}
It("a pin on A is visible on B and survives a new instance C built from the same store", func() {
a, _ := newManager("a")
b, _ := newManager("b")
a.Tick(ctx)
b.Tick(ctx)
Expect(a.Pin("chain", "y")).To(Succeed())
stB, ok := b.ChainStatus("chain")
Expect(ok).To(BeTrue())
Expect(stB.Pinned).ToNot(BeNil())
Expect(*stB.Pinned).To(Equal("y"))
// C is built after the pin was already written through to the shared
// store, and before ever ticking: SetStateSync (inside distsync.New)
// hydrates C's pins from s.Pins(), so the chain is created pinned the
// first time anything asks for it.
c, _ := newManager("c")
stC, ok := c.ChainStatus("chain")
Expect(ok).To(BeTrue())
Expect(stC.Pinned).ToNot(BeNil())
Expect(*stC.Pinned).To(Equal("y"))
})
It("a trip on B makes A's plan skip the target", func() {
a, _ := newManager("a")
b, _ := newManager("b")
a.Tick(ctx)
b.Tick(ctx)
b.ReportFailure("x", errBoom)
att, err := a.Plan("chain")
Expect(err).ToNot(HaveOccurred())
Expect(att.Target()).To(Equal("y"))
})
It("the leader's chain switch reaches the follower", func() {
a, _ := newManager("a") // leader
b, _ := newManager("b") // follower
a.Tick(ctx)
b.Tick(ctx)
a.ReportFailure("x", errBoom)
st, ok := b.ChainStatus("chain")
Expect(ok).To(BeTrue())
Expect(st.Active).To(Equal("y"))
})
It("late joiner converges on heartbeat", func() {
a, _ := newManager("a") // leader
a.Tick(ctx)
a.ReportFailure("x", errBoom)
// C joins after the trip: failover.targets/chains are NATS-only with
// no Store, so C's hydrate on Start sees nothing and it starts out
// believing every target is healthy.
c, _ := newManager("c") // follower
c.Tick(ctx)
before, ok := c.ChainStatus("chain")
Expect(ok).To(BeTrue())
Expect(before.Active).To(Equal("x"))
a.Republish()
after, ok := c.ChainStatus("chain")
Expect(ok).To(BeTrue())
Expect(after.Active).To(Equal("y"))
})
It("guards against a typed-nil PinStore passed as the Store interface", func() {
var nilStore *distsync.PinStore // deliberately typed, deliberately nil
m := failover.New(src, failover.WithLeaderGate(leaderGateFor("a", &leader)))
_, err := distsync.New(ctx, bus, nilStore, m)
Expect(err).ToNot(HaveOccurred())
m.Tick(ctx)
Expect(m.Pin("chain", "y")).To(Succeed())
st, ok := m.ChainStatus("chain")
Expect(ok).To(BeTrue())
Expect(st.Pinned).ToNot(BeNil())
Expect(*st.Pinned).To(Equal("y"))
})
It("Close stops all three maps without erroring", func() {
_, s := newManager("a")
Expect(s.Close()).To(Succeed())
})
It("PinStore round-trips through a real sqlite-backed gorm DB", func() {
db, err := gorm.Open(sqlite.Open("file::memory:?cache=shared"), &gorm.Config{})
Expect(err).ToNot(HaveOccurred())
store, err := distsync.NewPinStore(db)
Expect(err).ToNot(HaveOccurred())
recs, err := store.List(ctx)
Expect(err).ToNot(HaveOccurred())
Expect(recs).To(BeEmpty())
Expect(store.Upsert(ctx, distsync.PinRecord{Chain: "chain", Target: "y", UpdatedAt: time.Now()})).To(Succeed())
recs, err = store.List(ctx)
Expect(err).ToNot(HaveOccurred())
Expect(recs).To(HaveLen(1))
Expect(recs[0].Chain).To(Equal("chain"))
Expect(recs[0].Target).To(Equal("y"))
// Upsert again on the same key updates rather than duplicating.
Expect(store.Upsert(ctx, distsync.PinRecord{Chain: "chain", Target: "x", UpdatedAt: time.Now()})).To(Succeed())
recs, err = store.List(ctx)
Expect(err).ToNot(HaveOccurred())
Expect(recs).To(HaveLen(1))
Expect(recs[0].Target).To(Equal("x"))
Expect(store.Delete(ctx, "chain")).To(Succeed())
recs, err = store.List(ctx)
Expect(err).ToNot(HaveOccurred())
Expect(recs).To(BeEmpty())
})
})
@@ -0,0 +1,63 @@
// Package distsync wires failover.StateSync to syncstate.SyncedMap so target
// health, chain decisions and pins are shared across frontends over NATS,
// with pins durable in a small gorm-backed table.
package distsync
import (
"context"
"fmt"
"time"
"github.com/mudler/LocalAI/core/services/advisorylock"
"gorm.io/gorm"
)
// PinRecord is the durable form of a chain pin. It doubles as the value type
// of the "failover.pins" SyncedMap, so a hydrate (Store.List) needs no
// conversion and the wire delta carries the exact row.
type PinRecord struct {
Chain string `gorm:"primaryKey" json:"chain"`
Target string `json:"target"`
UpdatedAt time.Time `json:"updated_at"`
}
// TableName pins the table name independent of the Go type name.
func (PinRecord) TableName() string { return "failover_pins" }
// PinStore is gorm-backed durable storage for chain pins, implementing
// syncstate.Store[string, PinRecord] (asserted in distsync.go).
type PinStore struct {
db *gorm.DB
}
// NewPinStore migrates the failover_pins table under the schema-migrate
// advisory lock - the same guard jobs.NewJobStore uses - so several
// frontends starting at once do not race on the migration, and returns a
// ready-to-use store.
func NewPinStore(db *gorm.DB) (*PinStore, error) {
if err := advisorylock.WithLockCtx(context.Background(), db, advisorylock.KeySchemaMigrate, func() error {
return db.AutoMigrate(&PinRecord{})
}); err != nil {
return nil, fmt.Errorf("distsync: migrating pin table: %w", err)
}
return &PinStore{db: db}, nil
}
// List returns every durable pin, for hydrate on Start.
func (s *PinStore) List(ctx context.Context) ([]PinRecord, error) {
var out []PinRecord
if err := s.db.WithContext(ctx).Find(&out).Error; err != nil {
return nil, err
}
return out, nil
}
// Upsert writes a pin through. Save inserts or updates by primary key.
func (s *PinStore) Upsert(ctx context.Context, v PinRecord) error {
return s.db.WithContext(ctx).Save(&v).Error
}
// Delete removes a pin by chain name.
func (s *PinStore) Delete(ctx context.Context, k string) error {
return s.db.WithContext(ctx).Delete(&PinRecord{Chain: k}).Error
}