diff --git a/.agents/distributed-state.md b/.agents/distributed-state.md index 5936636ea..00173edfe 100644 --- a/.agents/distributed-state.md +++ b/.agents/distributed-state.md @@ -50,6 +50,15 @@ up. If your map has no durable backing, a leader (or another privileged writer) must republish its live state after a reconnect instead of relying on `Reconcile` to recover it. +**Gotcha:** a hydrate (on `Start`, after a NATS reconnect, and on every +`Reconcile` tick) replaces the map's contents **without firing `OnApply`**. +If `OnApply` feeds derived state (the failover manager's pins, a cache, a +running process), that state stays stale after a reconnect or a repaired +missed delta. Re-sync the derived state from the map's `Snapshot()`: on a +periodic tick, and/or from an `OnReconnect` callback registered after the +map's `Start` (callbacks run in registration order, so the map has already +re-hydrated). `failover.Manager.ReconcilePins` is the example. + ### 2. Single-runner Use `advisorylock.RunLeaderLoop` or `advisorylock.TryWithLockCtx` @@ -117,7 +126,9 @@ When your PR adds or changes state that lives longer than a single request: (`advisorylock`) / stateless / documented per-instance - [ ] If shared: `Store` added if the state must survive a cluster restart, or a `Loader` for standalone rehydration; if neither, a leader - republishes after reconnect instead of relying on bare `Reconcile` + republishes after reconnect instead of relying on bare `Reconcile`; + state derived through `OnApply` re-syncs from `Snapshot()` after a + hydrate (hydrate fires no `OnApply`) - [ ] If single-runner: new lock key added to `keys.go`; `HeldLock` chosen over `RunLeaderLoop`/`TryWithLockCtx` if leadership must be sticky across ticks diff --git a/core/services/failover/distsync/distsync.go b/core/services/failover/distsync/distsync.go index f7bb3dfce..5dba41527 100644 --- a/core/services/failover/distsync/distsync.go +++ b/core/services/failover/distsync/distsync.go @@ -19,11 +19,15 @@ var ( _ failover.StateSync = (*Sync)(nil) ) +// pinReconcileInterval is how often the pins map re-reads the DB. +const pinReconcileInterval = 30 * time.Second + // 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. +// Start, not just from whatever peers happen to broadcast afterwards, +// and re-reads it every pinReconcileInterval to repair a missed delta. // - failover.targets and failover.chains are ephemeral live-health state, // NATS-only with no Store and no Reconcile: with neither set, a Reconcile // tick's hydrate is a no-op (nothing durable to pull from), so it could @@ -53,11 +57,20 @@ func New(ctx context.Context, nats messaging.MessagingClient, pins syncstate.Sto pinStore = pins } + // 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. + var reconcile time.Duration + if pinStore != nil { + reconcile = pinReconcileInterval + } s.pins = syncstate.New(syncstate.Config[string, PinRecord]{ - Name: "failover.pins", - Key: func(p PinRecord) string { return p.Chain }, - Nats: nats, - Store: pinStore, + Name: "failover.pins", + Key: func(p PinRecord) string { return p.Chain }, + Nats: nats, + Store: pinStore, + Reconcile: reconcile, OnApply: func(op string, chain string, v PinRecord) { if op == "delete" { m.ApplyPin(chain, "") @@ -98,6 +111,12 @@ func New(ctx context.Context, nats messaging.MessagingClient, pins syncstate.Sto } m.SetStateSync(s) + // The pins map re-hydrates from the DB on reconnect without OnApply, so + // hand the manager the result. Registered after the map's own callback, + // which runs first. + if r, ok := nats.(interface{ OnReconnect(func()) }); ok { + r.OnReconnect(m.ReconcilePins) + } return s, nil } diff --git a/core/services/failover/distsync/distsync_test.go b/core/services/failover/distsync/distsync_test.go index 57a205f38..6a4a6cb93 100644 --- a/core/services/failover/distsync/distsync_test.go +++ b/core/services/failover/distsync/distsync_test.go @@ -156,6 +156,28 @@ var _ = Describe("distsync", func() { Expect(*stC.Pinned).To(Equal("y")) }) + It("re-applies pins that changed while B missed the deltas, after a reconnect", func() { + a, _ := newManager("a") + b, _ := newManager("b") + a.Tick(ctx) + b.Tick(ctx) + + // Written to the DB with no broadcast: B's map only learns it from + // the reconnect re-hydrate, which fires no OnApply. + Expect(pinStore.Upsert(ctx, distsync.PinRecord{Chain: "chain", Target: "y", UpdatedAt: time.Now()})).To(Succeed()) + bus.TriggerReconnect() + st, _ := b.ChainStatus("chain") + Expect(st.Pinned).ToNot(BeNil()) + Expect(*st.Pinned).To(Equal("y")) + + Expect(pinStore.Delete(ctx, "chain")).To(Succeed()) + bus.TriggerReconnect() + st, _ = b.ChainStatus("chain") + Expect(st.Pinned).To(BeNil()) + st, _ = a.ChainStatus("chain") + Expect(st.Pinned).To(BeNil()) + }) + It("a trip on B makes A's plan skip the target", func() { a, _ := newManager("a") b, _ := newManager("b") diff --git a/core/services/failover/schedule.go b/core/services/failover/schedule.go index c910cb9b4..c851db835 100644 --- a/core/services/failover/schedule.go +++ b/core/services/failover/schedule.go @@ -23,10 +23,14 @@ func (m *Manager) Run(ctx context.Context) { case <-ticker.C: m.Tick(ctx) // Publishes are fire-and-forget, so a frontend that missed one - // (restart, dropped message) converges within ten seconds. + // (restart, dropped message) converges within ten seconds. Pins + // are not republished: every frontend re-reads the shared set. m.ticks++ - if m.ticks%10 == 0 && m.IsLeader() { - m.Republish() + if m.ticks%10 == 0 { + m.ReconcilePins() + if m.IsLeader() { + m.Republish() + } } } } diff --git a/core/services/failover/statesync.go b/core/services/failover/statesync.go index a14378a04..49515d394 100644 --- a/core/services/failover/statesync.go +++ b/core/services/failover/statesync.go @@ -69,7 +69,35 @@ func (m *Manager) SetStateSync(s StateSync) { if s == nil { return } - for chain, target := range s.Pins() { + m.ReconcilePins() +} + +// ReconcilePins makes this frontend's pins match the shared pin set. Pins +// normally arrive as deltas, but a re-hydrate of the shared set (after a NATS +// reconnect, or a missed delta repaired from the DB) changes it without +// delivering them; a frontend left with a stale pin would serve it while the +// others do not. The scheduler runs this periodically on every frontend. +func (m *Manager) ReconcilePins() { + m.mu.Lock() + s := m.sync + m.mu.Unlock() + if s == nil { + return + } + // Read outside the lock: the store may need I/O. + want := s.Pins() + m.mu.Lock() + var stale []string + for chain := range m.pins { + if _, ok := want[chain]; !ok { + stale = append(stale, chain) + } + } + m.mu.Unlock() + for _, chain := range stale { + m.ApplyPin(chain, "") + } + for chain, target := range want { m.ApplyPin(chain, target) } } diff --git a/core/services/failover/statesync_test.go b/core/services/failover/statesync_test.go index d379bb674..cbe5814a0 100644 --- a/core/services/failover/statesync_test.go +++ b/core/services/failover/statesync_test.go @@ -140,6 +140,38 @@ var _ = Describe("Manager state sync", func() { Expect(st.Pinned).ToNot(BeNil()) }) + It("ReconcilePins converges a frontend that missed a pin and an unpin", func() { + // The shared pin set changes without B hearing the delta, as after a + // NATS reconnect whose re-hydrate fires no OnApply. + bus.mu.Lock() + bus.pins["chain"] = "y" + bus.mu.Unlock() + b.ReconcilePins() + st, _ := b.ChainStatus("chain") + Expect(st.Pinned).ToNot(BeNil()) + Expect(*st.Pinned).To(Equal("y")) + Expect(st.Active).To(Equal("y")) + + bus.mu.Lock() + delete(bus.pins, "chain") + bus.mu.Unlock() + b.ReconcilePins() + st, _ = b.ChainStatus("chain") + Expect(st.Pinned).To(BeNil()) + }) + + It("ReconcilePins leaves a pin for a chain this frontend does not know yet", func() { + bus.mu.Lock() + bus.pins["later"] = "y" + bus.mu.Unlock() + b.ReconcilePins() + src.Put(chainCfg("later", nil, t("x"), t("y"))) + b.Sync() + st, ok := b.ChainStatus("later") + Expect(ok).To(BeTrue()) + Expect(st.Pinned).ToNot(BeNil()) + }) + 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.