From ec6ca0af10f39d8e1ced02023fa254e52f69aeb8 Mon Sep 17 00:00:00 2001 From: heymingwei Date: Wed, 14 May 2025 08:37:11 +0000 Subject: [PATCH] fix(blobstore): waiting apply snapshot finish when save new raft log Signed-off-by: heymingwei --- blobstore/common/raftserver/server.go | 21 +++++++++++++++++++-- 1 file changed, 19 insertions(+), 2 deletions(-) diff --git a/blobstore/common/raftserver/server.go b/blobstore/common/raftserver/server.go index 2edad153f..86089f3b0 100644 --- a/blobstore/common/raftserver/server.go +++ b/blobstore/common/raftserver/server.go @@ -51,8 +51,9 @@ type RaftServer interface { // to raft storage concurrently; the application must read // raftDone before assuming the raft messages are stable. type apply struct { - entries []pb.Entry - snapshot pb.Snapshot + entries []pb.Entry + snapshot pb.Snapshot + snapFinishCh chan struct{} } type raftServer struct { @@ -329,6 +330,11 @@ func (s *raftServer) raftApply() { s.applyEntries(entries) s.applySnapshotFinish(snap) s.applyWait.Trigger(s.store.Applied()) + + // only take effect when has snapshot + if ap.snapFinishCh != nil { + close(ap.snapFinishCh) + } case snapMsg := <-s.snapMsgc: go s.processSnapshotMessage(snapMsg) case snap := <-s.snapshotC: @@ -515,6 +521,10 @@ func (s *raftServer) raftStart() { snapshot: rd.Snapshot, } + if !raft.IsEmptySnap(ap.snapshot) { + ap.snapFinishCh = make(chan struct{}) + } + select { case s.applyc <- ap: case <-s.stopc: @@ -522,6 +532,13 @@ func (s *raftServer) raftStart() { } s.tr.Send(s.processMessages(rd.Messages)) + // waiting apply snapshot is finished + if ap.snapFinishCh != nil { + log.Debugf("waiting apply snapshot finish...") + <-ap.snapFinishCh + log.Debugf("apply snapshot finish...") + } + err := s.store.Save(rd.HardState, rd.Entries) if err != nil { log.Panicf("save raft entries error: %v", err)