Skip to content

[Bug] Cross-node subscriptions lost when a raft snapshot is installed - #173

Merged
wind-c merged 1 commit into
wind-c:mainfrom
goingforstudying-ctrl:fix/raft-snapshot-restore-routing
Sep 3, 2026
Merged

wind-c merged 1 commit into
wind-c:mainfrom
goingforstudying-ctrl:fix/raft-snapshot-restore-routing

Conversation

@goingforstudying-ctrl

Copy link
Copy Markdown
Contributor

Was reading through the cluster raft code and I think the GetAll concurrency fix in d5c33e1 accidentally broke snapshot restore.

KV.GetAll now returns a copy of the routing table (correctly, it fixed the concurrent map access during Persist), but both restore paths still decode into its return value:

gob.NewDecoder(ir).Decode(f.GetAll())

GetAll hands back a pointer to a throwaway copy, so the snapshot decodes into thin air and the live routing table never changes. Both backends have it: hashicorp/fsm.go Restore and etcd/kvstore.go recoverFromSnapshot.

The visible effect: a node that installs a snapshot comes back with an empty routing table. That happens when a follower falls far enough behind that the leader sends InstallSnapshot, when a partitioned node rejoins, or on every graceful restart, since Peer.Stop() takes a snapshot right before shutdown and raft replays it on the way back up. After that, Lookup returns nothing for those filters, pickNodes finds no remote nodes, and publishes stop being forwarded to subscribers on other nodes until they happen to re-subscribe. For long-lived bridged clients that may be never. Nothing logs an error, messages just silently don't route.

Raft's docs say Restore must discard previous state and replace it with the snapshot, so I added KV.Restore(io.Reader) which decodes into a fresh map and swaps it in under lock, and pointed both backends at it. Replace rather than merge matters here: a node restoring an older snapshot shouldn't keep filters the snapshot doesn't know about.

Repro is a plain round trip through the real Persist/Restore path:

src := NewFsm(nil)
src.Add("topic/a", "node1")
snap, _ := src.Snapshot()
snap.Persist(sink) // capture the bytes
dst := NewFsm(nil)
dst.Restore(io.NopCloser(bytes.NewReader(sink.Bytes())))
dst.Lookup("topic/a") // nil before this fix, [node1] after

Added round-trip tests for the KV and both backends, including corrupt-snapshot cases (a failed decode leaves the existing state untouched). go test ./cluster/... all passes, including the existing peer tests that spin up real raft nodes.

One thing I left alone: notifyReplay replays every restored filter including ones the local node itself owns. That was the behavior before and changing it felt out of scope, but flagging it in case it matters.

KV.GetAll returns a copy since d5c33e1, so decoding a snapshot into
it never reached the live routing table. Both restore paths
(hashicorp Fsm.Restore, etcd recoverFromSnapshot) decoded into the
throwaway copy, leaving the node with an empty routing table after
installing a snapshot. Remote subscriptions stopped receiving
forwarded publishes until they re-subscribed.

Add KV.Restore which decodes into a fresh map and swaps it in under
lock, replacing prior state as raft snapshot semantics require.
@wind-c
wind-c merged commit 40fdcac into wind-c:main Sep 3, 2026
1 check passed
@wind-c

wind-c commented Sep 3, 2026

Copy link
Copy Markdown
Owner

Thanks @goingforstudying-ctrl

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants