Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
200 changes: 200 additions & 0 deletions internal/bft/completion_signal_test.go
Original file line number Diff line number Diff line change
@@ -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()
}
51 changes: 36 additions & 15 deletions internal/bft/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -132,7 +132,6 @@ type Controller struct {

syncChan chan struct{}
decisionChan chan decision
deliverChan chan struct{}
leaderToken chan struct{}
verificationSequence atomic.Uint64

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Behavior change worth confirming: previously decide blocked on the unbuffered deliverChan send when the waiter had left (view aborted), so the post-delivery bookkeeping below (incrementCurrentDecisionsInView, checkIfRotate/changeView, MaybePruneRevokedRequests, acquireLeaderToken) never ran in the aborted case. It now always runs once delivery happened. This unblocks the run loop (the real fix) and is defensible because the decision was in fact delivered, but it means a rotation changeView and leader-token acquisition can now fire for a view that is being aborted, racing the in-progress view change. changeView guards against going backwards (latestView > newViewNumber), which mitigates it, but please confirm the extra increment/rotate cannot momentarily double-count or start a spurious rotation view before the pending abort/view-change is processed.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in a167920.decide now records the view that produced the decision, and decide checks after delivery that this view is still current and not stopped. Otherwise it only prunes revoked requests and returns, skipping the increment, checkIfRotate/changeView, and acquireLeaderToken. The check runs on the controller run loop, the same goroutine that handles abortView and changeView, so a view change already processed is always observed first. For a dead view the result matches the pre-PR behavior, where the blocked send prevented the bookkeeping, minus the hang.

TestControllerDecideSkipsViewBookkeepingIfDecidingViewIsGone covers the aborted and replaced cases.

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()
Expand All @@ -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()
Expand Down Expand Up @@ -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)

Expand Down Expand Up @@ -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,
Expand All @@ -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
}
}

Expand All @@ -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
Expand Down
Loading
Loading