From d1533be6c64ee3ba8722593f1b6db26c019229dc Mon Sep 17 00:00:00 2001 From: Kevin <61072789+kevindharmawan@users.noreply.github.com> Date: Fri, 3 Jul 2026 16:07:26 -0400 Subject: [PATCH 1/4] Avoid blocking controller delivery completion Signed-off-by: Kevin <61072789+kevindharmawan@users.noreply.github.com> --- internal/bft/completion_signal_test.go | 98 ++++++++++++++++++++++++++ internal/bft/controller.go | 14 ++-- 2 files changed, 106 insertions(+), 6 deletions(-) create mode 100644 internal/bft/completion_signal_test.go diff --git a/internal/bft/completion_signal_test.go b/internal/bft/completion_signal_test.go new file mode 100644 index 00000000..74c11956 --- /dev/null +++ b/internal/bft/completion_signal_test.go @@ -0,0 +1,98 @@ +// Copyright IBM Corp. All Rights Reserved. +// +// SPDX-License-Identifier: Apache-2.0 +// + +package bft + +import ( + "testing" + "time" + + "github.com/hyperledger-labs/SmartBFT/pkg/types" + protos "github.com/hyperledger-labs/SmartBFT/smartbftprotos" + "github.com/stretchr/testify/assert" + "go.uber.org/zap" + "google.golang.org/protobuf/proto" +) + +func TestControllerDecideDoesNotBlockIfDeliveryWaiterLeft(t *testing.T) { + metadata, err := proto.Marshal(&protos.ViewMetadata{}) + assert.NoError(t, err) + + checkpoint := &types.Checkpoint{} + checkpoint.Set(types.Proposal{Metadata: metadata}, nil) + + controller := &Controller{ + ID: 2, + N: 4, + NodesList: []uint64{1, 2, 3, 4}, + Logger: zap.NewNop().Sugar(), + Deliver: applicationFunc(func(types.Proposal, []types.Signature) types.Reconfig { return types.Reconfig{} }), + Verifier: noopVerifier{}, + Checkpoint: checkpoint, + stopChan: make(chan struct{}), + } + + requireReturns(t, func() { + controller.decide(decision{ + proposal: types.Proposal{Metadata: metadata}, + delivered: make(chan struct{}), + }) + }) +} + +func requireReturns(t *testing.T, f func()) { + t.Helper() + + done := make(chan struct{}) + go func() { + defer close(done) + f() + }() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("function blocked") + } +} + +// These tests live in package bft to reach unexported completion paths, so they +// cannot use internal/bft/mocks without creating an import cycle. + +type applicationFunc func(types.Proposal, []types.Signature) types.Reconfig + +func (f applicationFunc) Deliver(proposal types.Proposal, signatures []types.Signature) types.Reconfig { + return f(proposal, signatures) +} + +type noopVerifier struct{} + +func (noopVerifier) VerificationSequence() uint64 { + return 0 +} + +func (noopVerifier) VerifyProposal(types.Proposal) ([]types.RequestInfo, error) { + panic("unexpected VerifyProposal call") +} + +func (noopVerifier) VerifyRequest([]byte) (types.RequestInfo, error) { + panic("unexpected VerifyRequest call") +} + +func (noopVerifier) VerifyConsenterSig(types.Signature, types.Proposal) ([]byte, error) { + panic("unexpected VerifyConsenterSig call") +} + +func (noopVerifier) VerifySignature(types.Signature) error { + panic("unexpected VerifySignature call") +} + +func (noopVerifier) RequestsFromProposal(types.Proposal) []types.RequestInfo { + panic("unexpected RequestsFromProposal call") +} + +func (noopVerifier) AuxiliaryData([]byte) []byte { + panic("unexpected AuxiliaryData call") +} diff --git a/internal/bft/controller.go b/internal/bft/controller.go index 8a192518..6b5ac798 100644 --- a/internal/bft/controller.go +++ b/internal/bft/controller.go @@ -132,7 +132,6 @@ type Controller struct { syncChan chan struct{} decisionChan chan decision - deliverChan chan struct{} leaderToken chan struct{} verificationSequence atomic.Uint64 @@ -542,9 +541,10 @@ func (c *Controller) decide(d decision) { } c.Logger.Debugf("Node %d delivered proposal", c.ID) c.removeDeliveredFromPool(d) - select { - case c.deliverChan <- struct{}{}: - case <-c.stopChan: + if d.delivered != nil { + close(d.delivered) + } + if c.stopped() { return } c.incrementCurrentDecisionsInView() @@ -797,7 +797,6 @@ func (c *Controller) Start(startViewNumber uint64, startProposalSequence uint64, c.stopChan = make(chan struct{}) c.leaderToken = make(chan struct{}, 1) c.decisionChan = make(chan decision, 1) - c.deliverChan = make(chan struct{}) c.viewChange = make(chan viewInfo, 1) c.abortViewChan = make(chan uint64, 1) @@ -882,11 +881,13 @@ func (c *Controller) stopped() bool { // Decide delivers the decision to the application func (c *Controller) Decide(proposal types.Proposal, signatures []types.Signature, requests []types.RequestInfo) { + delivered := make(chan struct{}) select { case c.decisionChan <- decision{ proposal: proposal, requests: requests, signatures: signatures, + delivered: delivered, }: case <-c.stopChan: // In case we are in the middle of shutting down, @@ -895,7 +896,7 @@ func (c *Controller) Decide(proposal types.Proposal, signatures []types.Signatur } select { - case <-c.deliverChan: // wait for the delivery of the decision to the application + case <-delivered: // wait for the delivery of the decision to the application case <-c.stopChan: // If we stopped the controller, abort delivery case <-c.currentViewAbortChan(): // If we stopped the view, abort delivery } @@ -918,6 +919,7 @@ type decision struct { proposal types.Proposal signatures []types.Signature requests []types.RequestInfo + delivered chan struct{} } // BroadcastConsensus broadcasts the message and informs the heartbeat monitor if necessary From 6bd292473b6bbc7b1915dfeb06ddceecc25bb4bb Mon Sep 17 00:00:00 2001 From: Kevin <61072789+kevindharmawan@users.noreply.github.com> Date: Fri, 3 Jul 2026 16:14:40 -0400 Subject: [PATCH 2/4] Avoid blocking in-flight view-change completion Signed-off-by: Kevin <61072789+kevindharmawan@users.noreply.github.com> --- internal/bft/completion_signal_test.go | 58 ++++++++++++++ internal/bft/viewchanger.go | 102 +++++++++++++++++++++---- 2 files changed, 145 insertions(+), 15 deletions(-) diff --git a/internal/bft/completion_signal_test.go b/internal/bft/completion_signal_test.go index 74c11956..8e93acbb 100644 --- a/internal/bft/completion_signal_test.go +++ b/internal/bft/completion_signal_test.go @@ -42,6 +42,44 @@ func TestControllerDecideDoesNotBlockIfDeliveryWaiterLeft(t *testing.T) { }) } +func TestViewChangerDecideDoesNotBlockIfInFlightWaiterLeft(t *testing.T) { + inFlightView := &View{abortChan: make(chan struct{})} + + viewChanger := &ViewChanger{ + Logger: zap.NewNop().Sugar(), + Application: applicationFunc(func(types.Proposal, []types.Signature) types.Reconfig { return types.Reconfig{} }), + RequestsTimer: noopRequestsTimer{}, + Pruner: noopPruner{}, + // Unbuffered and unread: the attempt's waiter has already left. + inFlightAttempt: &inFlightAttempt{ + id: 1, + decideCh: make(chan struct{}), + syncCh: make(chan struct{}), + viewRef: inFlightView, + }, + inFlightView: inFlightView, + } + + requireReturns(t, func() { + viewChanger.Decide(types.Proposal{}, nil, nil) + }) +} + +func TestViewChangerSyncDoesNotBlockIfInFlightWaiterLeft(t *testing.T) { + viewChanger := &ViewChanger{ + Logger: zap.NewNop().Sugar(), + Synchronizer: synchronizerFunc(func() {}), + // Unbuffered and unread: the attempt's waiter has already left. + inFlightAttempt: &inFlightAttempt{ + id: 1, + decideCh: make(chan struct{}), + syncCh: make(chan struct{}), + }, + } + + requireReturns(t, viewChanger.Sync) +} + func requireReturns(t *testing.T, f func()) { t.Helper() @@ -96,3 +134,23 @@ func (noopVerifier) RequestsFromProposal(types.Proposal) []types.RequestInfo { func (noopVerifier) AuxiliaryData([]byte) []byte { panic("unexpected AuxiliaryData call") } + +type noopRequestsTimer struct{} + +func (noopRequestsTimer) StopTimers() {} + +func (noopRequestsTimer) RestartTimers() {} + +func (noopRequestsTimer) RemoveRequest(types.RequestInfo) error { + return nil +} + +type noopPruner struct{} + +func (noopPruner) MaybePruneRevokedRequests() {} + +type synchronizerFunc func() + +func (f synchronizerFunc) Sync() { + f() +} diff --git a/internal/bft/viewchanger.go b/internal/bft/viewchanger.go index 9a98a8b2..45b80f86 100644 --- a/internal/bft/viewchanger.go +++ b/internal/bft/viewchanger.go @@ -48,6 +48,28 @@ type change struct { stopView bool } +type inFlightAttempt struct { + id uint64 + view uint64 + sequence uint64 + decideCh chan struct{} + syncCh chan struct{} + viewRef *View +} + +type inFlightAttemptCallbacks struct { + vc *ViewChanger + attempt *inFlightAttempt +} + +func (c *inFlightAttemptCallbacks) Decide(proposal types.Proposal, signatures []types.Signature, requests []types.RequestInfo) { + c.vc.decideInFlight(c.attempt, proposal, signatures, requests) +} + +func (c *inFlightAttemptCallbacks) Sync() { + c.vc.syncInFlight(c.attempt) +} + // ViewChanger is responsible for running the view change protocol. type ViewChanger struct { // Configuration @@ -77,8 +99,8 @@ type ViewChanger struct { // for the in flight proposal view ViewSequences *atomic.Value - inFlightDecideChan chan struct{} - inFlightSyncChan chan struct{} + inFlightAttemptSeq uint64 + inFlightAttempt *inFlightAttempt inFlightView *View inFlightViewLock sync.RWMutex @@ -148,9 +170,6 @@ func (v *ViewChanger) Start(startViewNumber uint64) { v.backOffFactor = 1 - v.inFlightDecideChan = make(chan struct{}) - v.inFlightSyncChan = make(chan struct{}) - go func() { defer v.vcDone.Done() v.ControllerStartedWG.Wait() @@ -1219,6 +1238,18 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success inFlightViewLatestSeq := proposalMD.LatestSequence v.inFlightViewLock.Lock() + v.inFlightAttemptSeq++ + attempt := &inFlightAttempt{ + id: v.inFlightAttemptSeq, + view: inFlightViewNum, + sequence: inFlightViewLatestSeq, + decideCh: make(chan struct{}, 1), + syncCh: make(chan struct{}, 1), + } + callbacks := &inFlightAttemptCallbacks{ + vc: v, + attempt: attempt, + } inFlightView := &View{ RetrieveCheckpoint: v.Checkpoint.Get, DecisionsPerLeader: v.DecisionsPerLeader, @@ -1228,9 +1259,9 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success Number: inFlightViewNum, LeaderID: v.SelfID, // so that no byzantine leader will cause a complain Quorum: v.quorum, - Decider: v, + Decider: callbacks, FailureDetector: v, - Sync: v, + Sync: callbacks, Logger: v.Logger, Comm: v.Comm, Verifier: v.Verifier, @@ -1243,6 +1274,7 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success MetricsBlacklist: v.MetricsBlacklist, MetricsView: v.MetricsView, } + attempt.viewRef = inFlightView inFlightView.MetricsView.ViewNumber.Set(float64(inFlightView.Number)) inFlightView.MetricsView.LeaderID.Set(float64(inFlightView.LeaderID)) inFlightView.MetricsView.ProposalSequence.Set(float64(inFlightView.ProposalSequence)) @@ -1250,6 +1282,7 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success inFlightView.MetricsView.Phase.Set(float64(inFlightView.Phase)) v.inFlightView = inFlightView + v.inFlightAttempt = attempt v.inFlightView.inFlightProposal = &types.Proposal{ VerificationSequence: int64(proposal.VerificationSequence), Metadata: proposal.Metadata, @@ -1277,7 +1310,17 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success <-v.Ticker inFlightView.Start() - defer inFlightView.Abort() + defer func() { + inFlightView.Abort() + v.inFlightViewLock.Lock() + if v.inFlightAttempt == attempt { + v.inFlightAttempt = nil + } + if v.inFlightView == inFlightView { + v.inFlightView = nil + } + v.inFlightViewLock.Unlock() + }() v.inFlightViewLock.Unlock() @@ -1286,10 +1329,10 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success // wait for view to finish or time out for { select { - case <-v.inFlightDecideChan: + case <-attempt.decideCh: v.Logger.Infof("In-flight view %d with latest sequence %d has committed a decision", inFlightViewNum, inFlightViewLatestSeq) return true - case <-v.inFlightSyncChan: + case <-attempt.syncCh: v.Logger.Infof("In-flight view %d with latest sequence %d has asked to sync", inFlightViewNum, inFlightViewLatestSeq) return false case now := <-v.Ticker: @@ -1305,9 +1348,22 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success } } -// Decide delivers to the application and informs the view changer after delivery +func (v *ViewChanger) currentInFlightAttempt() *inFlightAttempt { + v.inFlightViewLock.RLock() + defer v.inFlightViewLock.RUnlock() + return v.inFlightAttempt +} + +// Decide delivers to the application and informs the view changer after delivery. +// It is kept for compatibility; in-flight views use per-attempt callbacks. func (v *ViewChanger) Decide(proposal types.Proposal, signatures []types.Signature, requests []types.RequestInfo) { - v.inFlightView.stop() + v.decideInFlight(v.currentInFlightAttempt(), proposal, signatures, requests) +} + +func (v *ViewChanger) decideInFlight(attempt *inFlightAttempt, proposal types.Proposal, signatures []types.Signature, requests []types.RequestInfo) { + if attempt != nil && attempt.viewRef != nil { + attempt.viewRef.stop() + } v.Logger.Debugf("Delivering to app from Decide the last decision proposal") reconfig := v.Application.Deliver(proposal, signatures) if reconfig.InLatestDecision { @@ -1322,8 +1378,13 @@ func (v *ViewChanger) Decide(proposal types.Proposal, signatures []types.Signatu } v.Pruner.MaybePruneRevokedRequests() + if attempt == nil { + return + } select { - case v.inFlightDecideChan <- struct{}{}: + case attempt.decideCh <- struct{}{}: + return + default: return case <-v.stopChan: return @@ -1335,12 +1396,23 @@ func (v *ViewChanger) Complain(viewNum uint64, stopView bool) { v.Logger.Panicf("Node %d has complained while in the view for the in flight proposal", v.SelfID) } -// Sync calls the synchronizer and informs the view changer of the sync +// Sync calls the synchronizer and informs the view changer of the sync. +// It is kept for compatibility; in-flight views use per-attempt callbacks. func (v *ViewChanger) Sync() { + v.syncInFlight(v.currentInFlightAttempt()) +} + +func (v *ViewChanger) syncInFlight(attempt *inFlightAttempt) { // the in flight proposal view asked to sync v.Logger.Debugf("Node %d is calling sync because the in flight proposal view has asked to sync", v.SelfID) v.Synchronizer.Sync() - v.inFlightSyncChan <- struct{}{} + if attempt == nil { + return + } + select { + case attempt.syncCh <- struct{}{}: + default: + } } // HandleViewMessage passes a message to the in flight proposal view if applicable From 72dd0a2f5fbdf9d4116cf4509aae2a7be1de1a57 Mon Sep 17 00:00:00 2001 From: Kevin <61072789+kevindharmawan@users.noreply.github.com> Date: Thu, 27 Aug 2026 15:56:23 -0400 Subject: [PATCH 3/4] Remove unused in-flight attempt state and dead callbacks Signed-off-by: Kevin <61072789+kevindharmawan@users.noreply.github.com> --- internal/bft/viewchanger.go | 58 +++++++------------------------------ 1 file changed, 11 insertions(+), 47 deletions(-) diff --git a/internal/bft/viewchanger.go b/internal/bft/viewchanger.go index 45b80f86..d28eea89 100644 --- a/internal/bft/viewchanger.go +++ b/internal/bft/viewchanger.go @@ -48,10 +48,9 @@ type change struct { stopView bool } +// inFlightAttempt is one run of the in-flight proposal view, so that the view's +// completion is signaled to the waiter of that run only. type inFlightAttempt struct { - id uint64 - view uint64 - sequence uint64 decideCh chan struct{} syncCh chan struct{} viewRef *View @@ -98,11 +97,9 @@ type ViewChanger struct { Pruner Pruner // for the in flight proposal view - ViewSequences *atomic.Value - inFlightAttemptSeq uint64 - inFlightAttempt *inFlightAttempt - inFlightView *View - inFlightViewLock sync.RWMutex + ViewSequences *atomic.Value + inFlightView *View + inFlightViewLock sync.RWMutex Ticker <-chan time.Time lastTick time.Time @@ -1238,11 +1235,7 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success inFlightViewLatestSeq := proposalMD.LatestSequence v.inFlightViewLock.Lock() - v.inFlightAttemptSeq++ attempt := &inFlightAttempt{ - id: v.inFlightAttemptSeq, - view: inFlightViewNum, - sequence: inFlightViewLatestSeq, decideCh: make(chan struct{}, 1), syncCh: make(chan struct{}, 1), } @@ -1282,7 +1275,6 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success inFlightView.MetricsView.Phase.Set(float64(inFlightView.Phase)) v.inFlightView = inFlightView - v.inFlightAttempt = attempt v.inFlightView.inFlightProposal = &types.Proposal{ VerificationSequence: int64(proposal.VerificationSequence), Metadata: proposal.Metadata, @@ -1313,9 +1305,6 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success defer func() { inFlightView.Abort() v.inFlightViewLock.Lock() - if v.inFlightAttempt == attempt { - v.inFlightAttempt = nil - } if v.inFlightView == inFlightView { v.inFlightView = nil } @@ -1348,22 +1337,10 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success } } -func (v *ViewChanger) currentInFlightAttempt() *inFlightAttempt { - v.inFlightViewLock.RLock() - defer v.inFlightViewLock.RUnlock() - return v.inFlightAttempt -} - -// Decide delivers to the application and informs the view changer after delivery. -// It is kept for compatibility; in-flight views use per-attempt callbacks. -func (v *ViewChanger) Decide(proposal types.Proposal, signatures []types.Signature, requests []types.RequestInfo) { - v.decideInFlight(v.currentInFlightAttempt(), proposal, signatures, requests) -} - +// decideInFlight delivers the decision of the in-flight proposal view to the application +// and informs the waiter of the given attempt after delivery. func (v *ViewChanger) decideInFlight(attempt *inFlightAttempt, proposal types.Proposal, signatures []types.Signature, requests []types.RequestInfo) { - if attempt != nil && attempt.viewRef != nil { - attempt.viewRef.stop() - } + attempt.viewRef.stop() v.Logger.Debugf("Delivering to app from Decide the last decision proposal") reconfig := v.Application.Deliver(proposal, signatures) if reconfig.InLatestDecision { @@ -1378,16 +1355,10 @@ func (v *ViewChanger) decideInFlight(attempt *inFlightAttempt, proposal types.Pr } v.Pruner.MaybePruneRevokedRequests() - if attempt == nil { - return - } + // The waiter may have already left (timed out or stopped), so never block on it. select { case attempt.decideCh <- struct{}{}: - return default: - return - case <-v.stopChan: - return } } @@ -1396,19 +1367,12 @@ func (v *ViewChanger) Complain(viewNum uint64, stopView bool) { v.Logger.Panicf("Node %d has complained while in the view for the in flight proposal", v.SelfID) } -// Sync calls the synchronizer and informs the view changer of the sync. -// It is kept for compatibility; in-flight views use per-attempt callbacks. -func (v *ViewChanger) Sync() { - v.syncInFlight(v.currentInFlightAttempt()) -} - +// syncInFlight calls the synchronizer and informs the waiter of the given attempt of the sync. func (v *ViewChanger) syncInFlight(attempt *inFlightAttempt) { // the in flight proposal view asked to sync v.Logger.Debugf("Node %d is calling sync because the in flight proposal view has asked to sync", v.SelfID) v.Synchronizer.Sync() - if attempt == nil { - return - } + // The waiter may have already left (timed out or stopped), so never block on it. select { case attempt.syncCh <- struct{}{}: default: From a167920679bf9565e8ca180f6d66e1623df3f1ba Mon Sep 17 00:00:00 2001 From: Kevin <61072789+kevindharmawan@users.noreply.github.com> Date: Thu, 27 Aug 2026 15:56:41 -0400 Subject: [PATCH 4/4] Skip view bookkeeping for decisions of a stopped view Signed-off-by: Kevin <61072789+kevindharmawan@users.noreply.github.com> --- internal/bft/completion_signal_test.go | 114 +++++++++++++++++-------- internal/bft/controller.go | 39 ++++++--- 2 files changed, 108 insertions(+), 45 deletions(-) diff --git a/internal/bft/completion_signal_test.go b/internal/bft/completion_signal_test.go index 8e93acbb..031f0a8e 100644 --- a/internal/bft/completion_signal_test.go +++ b/internal/bft/completion_signal_test.go @@ -17,67 +17,111 @@ import ( ) func TestControllerDecideDoesNotBlockIfDeliveryWaiterLeft(t *testing.T) { - metadata, err := proto.Marshal(&protos.ViewMetadata{}) - assert.NoError(t, err) + view := &View{abortChan: make(chan struct{})} + controller := newDecidingController(t, view) - checkpoint := &types.Checkpoint{} - checkpoint.Set(types.Proposal{Metadata: metadata}, nil) + requireReturns(t, func() { + controller.decide(newDecision(t, view)) + }) + assert.Equal(t, uint64(1), controller.getCurrentDecisionsInView()) +} - controller := &Controller{ - ID: 2, - N: 4, - NodesList: []uint64{1, 2, 3, 4}, - Logger: zap.NewNop().Sugar(), - Deliver: applicationFunc(func(types.Proposal, []types.Signature) types.Reconfig { return types.Reconfig{} }), - Verifier: noopVerifier{}, - Checkpoint: checkpoint, - stopChan: make(chan struct{}), - } +func TestControllerDecideSkipsViewBookkeepingIfDecidingViewIsGone(t *testing.T) { + t.Run("view aborted", func(t *testing.T) { + view := &View{abortChan: make(chan struct{})} + controller := newDecidingController(t, view) + view.stop() - requireReturns(t, func() { - controller.decide(decision{ - proposal: types.Proposal{Metadata: metadata}, - delivered: make(chan struct{}), + requireReturns(t, func() { + controller.decide(newDecision(t, view)) }) + assert.Equal(t, uint64(0), controller.getCurrentDecisionsInView()) + assert.Empty(t, controller.leaderToken) + }) + + t.Run("view replaced", func(t *testing.T) { + view := &View{abortChan: make(chan struct{})} + view.stop() + controller := newDecidingController(t, &View{abortChan: make(chan struct{})}) + + requireReturns(t, func() { + controller.decide(newDecision(t, view)) + }) + assert.Equal(t, uint64(0), controller.getCurrentDecisionsInView()) + assert.Empty(t, controller.leaderToken) }) } func TestViewChangerDecideDoesNotBlockIfInFlightWaiterLeft(t *testing.T) { inFlightView := &View{abortChan: make(chan struct{})} - viewChanger := &ViewChanger{ Logger: zap.NewNop().Sugar(), Application: applicationFunc(func(types.Proposal, []types.Signature) types.Reconfig { return types.Reconfig{} }), RequestsTimer: noopRequestsTimer{}, Pruner: noopPruner{}, - // Unbuffered and unread: the attempt's waiter has already left. - inFlightAttempt: &inFlightAttempt{ - id: 1, - decideCh: make(chan struct{}), - syncCh: make(chan struct{}), - viewRef: inFlightView, - }, - inFlightView: inFlightView, + } + // Unbuffered and unread: the attempt's waiter has already left. + attempt := &inFlightAttempt{ + decideCh: make(chan struct{}), + syncCh: make(chan struct{}), + viewRef: inFlightView, } requireReturns(t, func() { - viewChanger.Decide(types.Proposal{}, nil, nil) + viewChanger.decideInFlight(attempt, types.Proposal{}, nil, nil) }) + assert.True(t, inFlightView.Stopped()) } func TestViewChangerSyncDoesNotBlockIfInFlightWaiterLeft(t *testing.T) { viewChanger := &ViewChanger{ Logger: zap.NewNop().Sugar(), Synchronizer: synchronizerFunc(func() {}), - // Unbuffered and unread: the attempt's waiter has already left. - inFlightAttempt: &inFlightAttempt{ - id: 1, - decideCh: make(chan struct{}), - syncCh: make(chan struct{}), - }, + } + // Unbuffered and unread: the attempt's waiter has already left. + attempt := &inFlightAttempt{ + decideCh: make(chan struct{}), + syncCh: make(chan struct{}), } - requireReturns(t, viewChanger.Sync) + requireReturns(t, func() { + viewChanger.syncInFlight(attempt) + }) +} + +// newDecidingController returns a controller whose current view is the given view, +// and whose own ID is the leader of the current view, so that a delivered decision +// would normally acquire the leader token. +func newDecidingController(t *testing.T, currView Proposer) *Controller { + metadata, err := proto.Marshal(&protos.ViewMetadata{}) + assert.NoError(t, err) + + checkpoint := &types.Checkpoint{} + checkpoint.Set(types.Proposal{Metadata: metadata}, nil) + + return &Controller{ + ID: 1, + N: 4, + NodesList: []uint64{1, 2, 3, 4}, + Logger: zap.NewNop().Sugar(), + Deliver: applicationFunc(func(types.Proposal, []types.Signature) types.Reconfig { return types.Reconfig{} }), + Verifier: noopVerifier{}, + Checkpoint: checkpoint, + currView: currView, + stopChan: make(chan struct{}), + leaderToken: make(chan struct{}, 1), + } +} + +func newDecision(t *testing.T, view Proposer) decision { + metadata, err := proto.Marshal(&protos.ViewMetadata{}) + assert.NoError(t, err) + + return decision{ + proposal: types.Proposal{Metadata: metadata}, + view: view, + delivered: make(chan struct{}), + } } func requireReturns(t *testing.T, f func()) { diff --git a/internal/bft/controller.go b/internal/bft/controller.go index 6b5ac798..8adb15ab 100644 --- a/internal/bft/controller.go +++ b/internal/bft/controller.go @@ -171,14 +171,6 @@ func (c *Controller) currentViewStopped() bool { return view.Stopped() } -func (c *Controller) currentViewAbortChan() <-chan struct{} { - c.currViewLock.RLock() - view := c.currView - c.currViewLock.RUnlock() - - return view.AbortChan() -} - func (c *Controller) currentViewLeader() uint64 { c.currViewLock.RLock() view := c.currView @@ -547,6 +539,14 @@ func (c *Controller) decide(d decision) { if c.stopped() { return } + if !c.isRunningCurrentView(d.view) { + // The view that decided was aborted or replaced while the decision was waiting + // to be delivered, so the decision must not be accounted to the current view, + // and the current view must not be rotated or lead because of it. + c.Logger.Debugf("Node %d delivered a decision from a view that is no longer running, skipping view bookkeeping", c.ID) + c.MaybePruneRevokedRequests() + return + } c.incrementCurrentDecisionsInView() md := &protos.ViewMetadata{} @@ -566,6 +566,19 @@ func (c *Controller) decide(d decision) { } } +// isRunningCurrentView returns whether the given view is the current view and has not been stopped. +func (c *Controller) isRunningCurrentView(view Proposer) bool { + if view == nil { + c.Logger.Panicf("Decision without a deciding view") + } + + c.currViewLock.RLock() + currView := c.currView + c.currViewLock.RUnlock() + + return view == currView && !view.Stopped() +} + func (c *Controller) checkIfRotate(blacklist []uint64) bool { view := c.getCurrentViewNumber() decisionsInView := c.getCurrentDecisionsInView() @@ -881,12 +894,17 @@ func (c *Controller) stopped() bool { // Decide delivers the decision to the application func (c *Controller) Decide(proposal types.Proposal, signatures []types.Signature, requests []types.RequestInfo) { + c.currViewLock.RLock() + view := c.currView + c.currViewLock.RUnlock() + delivered := make(chan struct{}) select { case c.decisionChan <- decision{ proposal: proposal, requests: requests, signatures: signatures, + view: view, delivered: delivered, }: case <-c.stopChan: @@ -898,7 +916,7 @@ func (c *Controller) Decide(proposal types.Proposal, signatures []types.Signatur select { case <-delivered: // wait for the delivery of the decision to the application case <-c.stopChan: // If we stopped the controller, abort delivery - case <-c.currentViewAbortChan(): // If we stopped the view, abort delivery + case <-view.AbortChan(): // If we stopped the view, abort delivery } } @@ -919,7 +937,8 @@ type decision struct { proposal types.Proposal signatures []types.Signature requests []types.RequestInfo - delivered chan struct{} + view Proposer // the view that decided + delivered chan struct{} // closed once the decision was delivered to the application } // BroadcastConsensus broadcasts the message and informs the heartbeat monitor if necessary