diff --git a/core/services/failover/distsync/distsync.go b/core/services/failover/distsync/distsync.go index 5dba41527..825b495c6 100644 --- a/core/services/failover/distsync/distsync.go +++ b/core/services/failover/distsync/distsync.go @@ -59,8 +59,7 @@ func New(ctx context.Context, nats messaging.MessagingClient, pins syncstate.Sto // 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. + // periodic ReconcilePins carries the repair into the chains. var reconcile time.Duration if pinStore != nil { reconcile = pinReconcileInterval diff --git a/core/services/failover/distsync/distsync_test.go b/core/services/failover/distsync/distsync_test.go index 6a4a6cb93..4f4f37985 100644 --- a/core/services/failover/distsync/distsync_test.go +++ b/core/services/failover/distsync/distsync_test.go @@ -62,6 +62,14 @@ func (s *fakeSource) GetAllModelsConfigs() []config.ModelConfig { type memPinStore struct { mu sync.Mutex data map[string]distsync.PinRecord + // down, when set, makes every call fail, simulating a database outage. + down error +} + +func (s *memPinStore) setDown(err error) { + s.mu.Lock() + defer s.mu.Unlock() + s.down = err } func newMemPinStore() *memPinStore { return &memPinStore{data: map[string]distsync.PinRecord{}} } @@ -69,6 +77,9 @@ func newMemPinStore() *memPinStore { return &memPinStore{data: map[string]distsy func (s *memPinStore) List(context.Context) ([]distsync.PinRecord, error) { s.mu.Lock() defer s.mu.Unlock() + if s.down != nil { + return nil, s.down + } out := make([]distsync.PinRecord, 0, len(s.data)) for _, v := range s.data { out = append(out, v) @@ -79,6 +90,9 @@ func (s *memPinStore) List(context.Context) ([]distsync.PinRecord, error) { func (s *memPinStore) Upsert(_ context.Context, v distsync.PinRecord) error { s.mu.Lock() defer s.mu.Unlock() + if s.down != nil { + return s.down + } s.data[v.Chain] = v return nil } @@ -86,6 +100,9 @@ func (s *memPinStore) Upsert(_ context.Context, v distsync.PinRecord) error { func (s *memPinStore) Delete(_ context.Context, k string) error { s.mu.Lock() defer s.mu.Unlock() + if s.down != nil { + return s.down + } delete(s.data, k) return nil } @@ -178,6 +195,39 @@ var _ = Describe("distsync", func() { Expect(st.Pinned).To(BeNil()) }) + It("does not re-apply a pin whose write failed during a database outage", func() { + a, _ := newManager("a") + b, _ := newManager("b") + a.Tick(ctx) + b.Tick(ctx) + + pinStore.setDown(errBoom) + Expect(a.Pin("chain", "y")).To(MatchError(errBoom)) + st, _ := a.ChainStatus("chain") + Expect(st.Pinned).To(BeNil(), "a failed pin write must be rolled back") + + // The periodic pin re-sync runs while the database is still down: it + // must not resurrect the pin the rollback just undid. + a.ReconcilePins() + st, _ = a.ChainStatus("chain") + Expect(st.Pinned).To(BeNil(), "the failed pin must not come back on this frontend") + st, _ = b.ChainStatus("chain") + Expect(st.Pinned).To(BeNil(), "the failed pin must never reach a peer") + + // Same for an unpin that fails: the pin stays in force everywhere. + pinStore.setDown(nil) + Expect(a.Pin("chain", "y")).To(Succeed()) + pinStore.setDown(errBoom) + Expect(a.Unpin("chain")).To(MatchError(errBoom)) + a.ReconcilePins() + st, _ = a.ChainStatus("chain") + Expect(st.Pinned).ToNot(BeNil()) + Expect(*st.Pinned).To(Equal("y")) + st, _ = b.ChainStatus("chain") + Expect(st.Pinned).ToNot(BeNil()) + Expect(*st.Pinned).To(Equal("y")) + }) + It("a trip on B makes A's plan skip the target", func() { a, _ := newManager("a") b, _ := newManager("b") diff --git a/core/services/syncstate/syncstate.go b/core/services/syncstate/syncstate.go index 809177d40..5aa69470f 100644 --- a/core/services/syncstate/syncstate.go +++ b/core/services/syncstate/syncstate.go @@ -54,7 +54,7 @@ type delta[K comparable, V any] struct { } // SyncedMap is a cross-replica in-memory map. A local write (Set/Delete) updates -// memory, the optional durable Store, then broadcasts a delta to peers. A peer's +// the optional durable Store, then memory, then broadcasts a delta to peers. A peer's // delta updates memory only and fires OnApply - it never re-broadcasts and never // writes the Store. That structural split is the echo-loop guard (same pattern as // galleryop.mergeStatus / OpCache.applyStart): receiving your own broadcast just @@ -141,34 +141,40 @@ func (m *SyncedMap[K, V]) Close() error { return nil } -// Set updates the value locally, writes through the Store, then broadcasts. -// Per the data-flow contract the Store write happens under the lock so memory and -// durable state move together; the broadcast is best-effort after unlocking. +// Set writes through the Store, then updates the value locally, then +// broadcasts. The Store write comes first and happens under the lock so memory +// and durable state move together: when it fails, Set returns the error with +// memory and peers untouched. Keeping an unpersisted value in memory would let +// this replica serve it (and a caller that re-reads the map re-apply it) while +// the Store and every other replica disagree, until the next re-hydrate. +// The broadcast is best-effort after unlocking. func (m *SyncedMap[K, V]) Set(ctx context.Context, v V) error { k := m.cfg.Key(v) m.mu.Lock() - m.data[k] = v if m.cfg.Store != nil { if err := m.cfg.Store.Upsert(ctx, v); err != nil { m.mu.Unlock() return err } } + m.data[k] = v m.mu.Unlock() m.publish(opSet, k, v) return nil } -// Delete removes the key locally, deletes it from the Store, then broadcasts. +// Delete deletes the key from the Store, then removes it locally, then +// broadcasts. A failed Store delete leaves memory and peers untouched, for the +// same reason as Set. func (m *SyncedMap[K, V]) Delete(ctx context.Context, k K) error { m.mu.Lock() - delete(m.data, k) if m.cfg.Store != nil { if err := m.cfg.Store.Delete(ctx, k); err != nil { m.mu.Unlock() return err } } + delete(m.data, k) m.mu.Unlock() var zero V m.publish(opDelete, k, zero) diff --git a/core/services/syncstate/syncstate_test.go b/core/services/syncstate/syncstate_test.go index 1e31db41b..e8daf6014 100644 --- a/core/services/syncstate/syncstate_test.go +++ b/core/services/syncstate/syncstate_test.go @@ -2,6 +2,7 @@ package syncstate_test import ( "context" + "errors" "sync" . "github.com/onsi/ginkgo/v2" @@ -34,6 +35,8 @@ type fakeStore struct { upsertCalls int deleteCalls int listCalls int + // fail, when set, makes Upsert and Delete return it without writing. + fail error } func newFakeStore(seed ...*job) *fakeStore { @@ -59,6 +62,9 @@ func (s *fakeStore) Upsert(_ context.Context, j *job) error { s.mu.Lock() defer s.mu.Unlock() s.upsertCalls++ + if s.fail != nil { + return s.fail + } s.data[j.ID] = j return nil } @@ -67,6 +73,9 @@ func (s *fakeStore) Delete(_ context.Context, k string) error { s.mu.Lock() defer s.mu.Unlock() s.deleteCalls++ + if s.fail != nil { + return s.fail + } delete(s.data, k) return nil } @@ -207,6 +216,41 @@ var _ = Describe("SyncedMap", func() { }) }) + Describe("failed Store write", func() { + It("leaves memory and peers untouched when Set or Delete cannot persist", func() { + bus := testutil.NewFakeBus() + storeA := newFakeStore(&job{ID: "kept", Status: "running"}) + a := syncstate.New(syncstate.Config[string, *job]{Name: stateName, Key: jobKey, Nats: bus, Store: storeA}) + b := syncstate.New(syncstate.Config[string, *job]{Name: stateName, Key: jobKey, Nats: bus}) + Expect(a.Start(ctx)).To(Succeed()) + Expect(b.Start(ctx)).To(Succeed()) + defer func() { + Expect(a.Close()).To(Succeed()) + Expect(b.Close()).To(Succeed()) + }() + + errDown := errors.New("database down") + storeA.mu.Lock() + storeA.fail = errDown + storeA.mu.Unlock() + + Expect(a.Set(ctx, &job{ID: "new", Status: "running"})).To(MatchError(errDown)) + _, ok := a.Get("new") + Expect(ok).To(BeFalse(), "a value that was not persisted must not be served") + _, ok = b.Get("new") + Expect(ok).To(BeFalse(), "a value that was not persisted must not be broadcast") + + Expect(a.Set(ctx, &job{ID: "kept", Status: "done"})).To(MatchError(errDown)) + got, ok := a.Get("kept") + Expect(ok).To(BeTrue()) + Expect(got.Status).To(Equal("running"), "a failed overwrite must keep the persisted value") + + Expect(a.Delete(ctx, "kept")).To(MatchError(errDown)) + _, ok = a.Get("kept") + Expect(ok).To(BeTrue(), "a failed delete must keep the persisted value") + }) + }) + Describe("OnApply hook", func() { It("fires with the correct op and key on an applied delta", func() { bus := testutil.NewFakeBus()