fix(syncstate): do not keep a value in memory when its Store write fails

SyncedMap.Set and Delete changed memory before they wrote the Store.
When the write failed they returned the error, but the unpersisted
value stayed in memory. Callers that re-read the map then applied it
again. The failover manager rolled back a failed pin, but its periodic
pin re-sync read the pin back from the map and re-applied it on that
frontend until the database came back.

Write the Store first and change memory only after it succeeds. A
failed Set or Delete now leaves memory and peers as they were, so the
map always matches what the Store holds. This also stops a failed
create of a fine-tune, quantization or agent task from leaving an
orphan entry that the API had reported as failed.

Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
Assisted-by: Claude:claude-opus-5-5 [Claude Code]
This commit is contained in:
Ettore Di Giacinto committed 2026-09-27 19:57:13 +00:00
1 parent 9bbcde4b1b
commit 7804e4591c
4 files changed
+108 -9

No files matched your search

+1 -2
View File
@@ -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
@@ -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")
+13 -7
View File
@@ -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)
+44
View File
@@ -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()