From 7b88674fc58c67cb796729210af10062d539c6d3 Mon Sep 17 00:00:00 2001 From: Ettore Di Giacinto Date: Sat, 26 Sep 2026 22:11:54 +0000 Subject: [PATCH] feat(failover): sync pins, target health and chain state over NATS Assisted-by: Claude:claude-opus-5-5 Signed-off-by: Ettore Di Giacinto --- core/services/failover/distsync/distsync.go | 158 +++++++++++ .../failover/distsync/distsync_suite_test.go | 13 + .../failover/distsync/distsync_test.go | 258 ++++++++++++++++++ core/services/failover/distsync/pinstore.go | 63 +++++ 4 files changed, 492 insertions(+) create mode 100644 core/services/failover/distsync/distsync.go create mode 100644 core/services/failover/distsync/distsync_suite_test.go create mode 100644 core/services/failover/distsync/distsync_test.go create mode 100644 core/services/failover/distsync/pinstore.go diff --git a/core/services/failover/distsync/distsync.go b/core/services/failover/distsync/distsync.go new file mode 100644 index 000000000..349aef14f --- /dev/null +++ b/core/services/failover/distsync/distsync.go @@ -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 +} diff --git a/core/services/failover/distsync/distsync_suite_test.go b/core/services/failover/distsync/distsync_suite_test.go new file mode 100644 index 000000000..37a6d7092 --- /dev/null +++ b/core/services/failover/distsync/distsync_suite_test.go @@ -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") +} diff --git a/core/services/failover/distsync/distsync_test.go b/core/services/failover/distsync/distsync_test.go new file mode 100644 index 000000000..57a205f38 --- /dev/null +++ b/core/services/failover/distsync/distsync_test.go @@ -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()) + }) +}) diff --git a/core/services/failover/distsync/pinstore.go b/core/services/failover/distsync/pinstore.go new file mode 100644 index 000000000..eaa682431 --- /dev/null +++ b/core/services/failover/distsync/pinstore.go @@ -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 +}