diff --git a/internal/bft/completion_signal_test.go b/internal/bft/completion_signal_test.go new file mode 100644 index 00000000..031f0a8e --- /dev/null +++ b/internal/bft/completion_signal_test.go @@ -0,0 +1,200 @@ +// 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) { + view := &View{abortChan: make(chan struct{})} + controller := newDecidingController(t, view) + + requireReturns(t, func() { + controller.decide(newDecision(t, view)) + }) + assert.Equal(t, uint64(1), controller.getCurrentDecisionsInView()) +} + +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(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. + attempt := &inFlightAttempt{ + decideCh: make(chan struct{}), + syncCh: make(chan struct{}), + viewRef: inFlightView, + } + + requireReturns(t, func() { + 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. + attempt := &inFlightAttempt{ + decideCh: make(chan struct{}), + syncCh: make(chan struct{}), + } + + 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()) { + 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") +} + +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/controller.go b/internal/bft/controller.go index 8a192518..8adb15ab 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 @@ -172,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 @@ -542,9 +533,18 @@ 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 + } + 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() @@ -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() @@ -797,7 +810,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 +894,18 @@ 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: // In case we are in the middle of shutting down, @@ -895,9 +914,9 @@ 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 + case <-view.AbortChan(): // If we stopped the view, abort delivery } } @@ -918,6 +937,8 @@ type decision struct { proposal types.Proposal signatures []types.Signature requests []types.RequestInfo + 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 diff --git a/internal/bft/viewchanger.go b/internal/bft/viewchanger.go index 9a98a8b2..d28eea89 100644 --- a/internal/bft/viewchanger.go +++ b/internal/bft/viewchanger.go @@ -48,6 +48,27 @@ 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 { + 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 @@ -76,11 +97,9 @@ type ViewChanger struct { Pruner Pruner // for the in flight proposal view - ViewSequences *atomic.Value - inFlightDecideChan chan struct{} - inFlightSyncChan chan struct{} - inFlightView *View - inFlightViewLock sync.RWMutex + ViewSequences *atomic.Value + inFlightView *View + inFlightViewLock sync.RWMutex Ticker <-chan time.Time lastTick time.Time @@ -148,9 +167,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 +1235,14 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success inFlightViewLatestSeq := proposalMD.LatestSequence v.inFlightViewLock.Lock() + attempt := &inFlightAttempt{ + 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 +1252,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 +1267,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)) @@ -1277,7 +1302,14 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success <-v.Ticker inFlightView.Start() - defer inFlightView.Abort() + defer func() { + inFlightView.Abort() + v.inFlightViewLock.Lock() + if v.inFlightView == inFlightView { + v.inFlightView = nil + } + v.inFlightViewLock.Unlock() + }() v.inFlightViewLock.Unlock() @@ -1286,10 +1318,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 +1337,10 @@ func (v *ViewChanger) commitInFlightProposal(proposal *protos.Proposal) (success } } -// Decide delivers to the application and informs the view changer after delivery -func (v *ViewChanger) Decide(proposal types.Proposal, signatures []types.Signature, requests []types.RequestInfo) { - v.inFlightView.stop() +// 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) { + 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,11 +1355,10 @@ func (v *ViewChanger) Decide(proposal types.Proposal, signatures []types.Signatu } v.Pruner.MaybePruneRevokedRequests() + // The waiter may have already left (timed out or stopped), so never block on it. select { - case v.inFlightDecideChan <- struct{}{}: - return - case <-v.stopChan: - return + case attempt.decideCh <- struct{}{}: + default: } } @@ -1335,12 +1367,16 @@ 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 -func (v *ViewChanger) Sync() { +// 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() - v.inFlightSyncChan <- struct{}{} + // The waiter may have already left (timed out or stopped), so never block on it. + select { + case attempt.syncCh <- struct{}{}: + default: + } } // HandleViewMessage passes a message to the in flight proposal view if applicable