Skip to content
Merged
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
1 change: 0 additions & 1 deletion instance.go
Original file line number Diff line number Diff line change
Expand Up @@ -539,7 +539,6 @@ func (i *Instance) createEpochConfig(validators common.Nodes) (*epochConfig, err
// TODO: For simplicity, we use the same value for all timeouts. If needed we can expand the config.
MaxProposalWait: i.Config.ParameterConfig.MaxNetworkDelay * 2, // 1 proposal + 1 vote
MaxRebroadcastWait: i.Config.ParameterConfig.MaxNetworkDelay * 2,
FinalizeRebroadcastTimeout: i.Config.ParameterConfig.MaxNetworkDelay * 2,
MaxRoundWindow: i.Config.ParameterConfig.MaxRoundWindow,
ID: i.Config.ID,
RandomSource: source, // Seed the random source from crypto/rand
Expand Down
2 changes: 1 addition & 1 deletion simplex/empty_block_builder_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ import (
)

// blocksUntilDone mimics the production ShouldBuildEmptyBlock contract
// (Epoch.haveUnFinalizedButNotarizedSuffix): it blocks until the given context is
// (Epoch.shouldBuildEmptyBlock): it blocks until the given context is
// cancelled and then reports whether the cancellation was because the empty-block
// timeout elapsed.
func blocksUntilDone(ctx context.Context) bool {
Expand Down
113 changes: 55 additions & 58 deletions simplex/epoch.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import (
var (
ErrAlreadyStarted = errors.New("epoch already started")
errNotarizationBlockMismatch = errors.New("notarization block header mismatches stored round block header")
errAlreadyTimedOutOnRound = errors.New("already timed out on this round")
)

const (
Expand Down Expand Up @@ -70,7 +71,6 @@ type EpochConfig struct {
MaxRoundWindow uint64
MaxReplicationResponseSize int
MaxRebroadcastWait time.Duration
FinalizeRebroadcastTimeout time.Duration
QCDeserializer common.QCDeserializer
Logger common.Logger
ID common.NodeID
Expand Down Expand Up @@ -110,7 +110,6 @@ type Epoch struct {
validatorsToPKs map[string][]byte
rounds map[uint64]*Round
emptyVotes map[uint64]*EmptyVoteSet
oldestNotFinalizedNotarization NotarizationTime
futureMessages messagesFromNode
round uint64 // The current round we notarize
monitor *Monitor
Expand Down Expand Up @@ -138,7 +137,6 @@ func (e *Epoch) AdvanceTime(t time.Time) {
e.monitor.AdvanceTime(t)
e.replicationState.AdvanceTime(t)
e.timeoutHandler.Tick(t)
e.oldestNotFinalizedNotarization.CheckForNotFinalizedNotarizedBlocks(t)
}

// HandleMessage notifies the engine about a reception of a message.
Expand Down Expand Up @@ -218,7 +216,6 @@ func (e *Epoch) init() error {
if err := e.maybeAssignDefaultConfig(); err != nil {
return err
}
e.initOldestNotFinalizedNotarization()
e.oneTimeVerifier = NewOneTimeVerifier(e.Logger)
scheduler := common.NewScheduler(e.Logger, DefaultProcessingBlocks)
e.blockVerificationScheduler = common.NewBlockVerificationScheduler(e.Logger, DefaultProcessingBlocks, scheduler)
Expand Down Expand Up @@ -247,7 +244,7 @@ func (e *Epoch) init() error {
e.futureMessages[string(node)] = make(map[uint64]*messagesForRound)
}
e.blockBuilder = &EmptyBlockBuilder{
ShouldBuildEmptyBlock: e.haveUnFinalizedButNotarizedSuffix,
ShouldBuildEmptyBlock: e.shouldBuildEmptyBlock,
Timeout: e.MaxProposalWait,
BB: e.BlockBuilder,
}
Expand All @@ -261,32 +258,7 @@ func (e *Epoch) init() error {
return e.setMetadataFromStorage()
}

func (e *Epoch) initOldestNotFinalizedNotarization() {
rebroadcastFinalizationVotes := func() {
e.lock.Lock()
defer e.lock.Unlock()

if err := e.rebroadcastPastFinalizeVotes(); err != nil {
e.Logger.Error("Could not rebroadcast past finalization votes", zap.Error(err))
}
}
e.oldestNotFinalizedNotarization = NewNotarizationTime(
e.FinalizeRebroadcastTimeout,
e.haveNotFinalizedNotarizedRound,
rebroadcastFinalizationVotes, e.getRound)
}

func (e *Epoch) getRound() uint64 {
e.lock.Lock()
defer e.lock.Unlock()

return e.round
}

func (e *Epoch) maybeAssignDefaultConfig() error {
if e.FinalizeRebroadcastTimeout == 0 {
e.FinalizeRebroadcastTimeout = DefaultFinalizeVoteRebroadcastTimeout
}
if e.MaxProposalWait == 0 {
e.MaxProposalWait = DefaultMaxProposalWaitTime
}
Expand Down Expand Up @@ -329,15 +301,20 @@ func (e *Epoch) Start() error {
return nil
}

func (e *Epoch) haveUnFinalizedButNotarizedSuffix(ctx context.Context) bool {
func (e *Epoch) shouldBuildEmptyBlock(ctx context.Context) bool {
<-ctx.Done()

if errors.Is(context.Cause(ctx), common.ErrShouldBuildEmptyBlock) {
e.lock.Lock()
defer e.lock.Unlock()

r := e.getHighestRound()
return r != nil && r.finalization == nil
for r, round := range e.rounds {
didNotTimeOutOnRound := !e.haveWeAlreadyTimedOutOnThisRound(r)
notarizedButNotFinalized := round.notarization != nil && round.finalization == nil
if didNotTimeOutOnRound && notarizedButNotFinalized {
return true
}
}
}

return false
Expand Down Expand Up @@ -1442,12 +1419,17 @@ func (e *Epoch) rebroadcastPastFinalizeVotes() error {
continue
}

if e.haveWeAlreadyTimedOutOnThisRound(r) {
e.Logger.Debug("Round already timed out when rebroadcasting finalize votes", zap.Uint64("round", r))
continue
}

var finalizeVoteMessage *common.Message
// Try to re-use finalization we created if possible, else create it.
if vote, exists := round.finalizeVotes[string(e.ID)]; exists {
finalizeVoteMessage = &common.Message{FinalizeVote: vote}
} else {
_, msg, err := e.constructFinalizeVoteMessage(round.notarization.Vote.BlockHeader)
_, msg, err := e.maybeConstructFinalizeVoteMessage(round.notarization.Vote.BlockHeader)
if err != nil {
return err
}
Expand Down Expand Up @@ -2419,6 +2401,10 @@ func (e *Epoch) createNotarizedBlockVerificationTask(block common.Block, notariz
return md.Digest
}

if err := e.finalizeVoteForReplicatedNotarization(md.Round); err != nil {
e.Logger.Error("Failed to finalize vote for replicated notarization", zap.Uint64("round", md.Round), zap.Error(err))
}

err = e.processReplicationState()
if err != nil {
e.haltedError = err
Expand Down Expand Up @@ -3113,6 +3099,34 @@ func (e *Epoch) increaseRound() {
e.round++
}

// finalizeVoteForReplicatedNotarization casts our finalize vote for a round whose notarization we learned
// of through replication instead of by collecting votes.
func (e *Epoch) finalizeVoteForReplicatedNotarization(r uint64) error {
round, exists := e.rounds[r]
if !exists || round.notarization == nil || round.finalization != nil {
return nil
}

if e.haveWeAlreadyTimedOutOnThisRound(r) {
e.Logger.Debug("Not finalize voting for a replicated notarization of a round we timed out on", zap.Uint64("round", r))
return nil
}

md := round.notarization.Vote.BlockHeader
finalizeVote, finalizeVoteMsg, err := e.maybeConstructFinalizeVoteMessage(md)
if err != nil {
return err
}
e.broadcast(finalizeVoteMsg)

e.Logger.Debug("Broadcasting finalize vote for a replicated notarization",
zap.Uint64("round", md.Round),
zap.Uint64("seq", md.Seq),
zap.Stringer("digest", md.Digest))

return e.handleFinalizeVoteMessage(&finalizeVote, e.ID)
}

func (e *Epoch) doNotarized(r uint64) error {
if e.haveWeAlreadyTimedOutOnThisRound(r) {
e.Logger.Info("We have already timed out on this round, will not finalize it", zap.Uint64("round", r))
Expand All @@ -3124,7 +3138,7 @@ func (e *Epoch) doNotarized(r uint64) error {

md := block.BlockHeader()

finalizeVote, finalizeVoteMsg, err := e.constructFinalizeVoteMessage(md)
finalizeVote, finalizeVoteMsg, err := e.maybeConstructFinalizeVoteMessage(md)
if err != nil {
return err
}
Expand All @@ -3141,7 +3155,12 @@ func (e *Epoch) doNotarized(r uint64) error {
return errors.Join(err1, err2)
}

func (e *Epoch) constructFinalizeVoteMessage(md common.BlockHeader) (common.FinalizeVote, *common.Message, error) {
func (e *Epoch) maybeConstructFinalizeVoteMessage(md common.BlockHeader) (common.FinalizeVote, *common.Message, error) {
if e.haveWeAlreadyTimedOutOnThisRound(md.Round) {
e.Logger.Error("Will not cast a finalize vote for a round we timed out on)", zap.Uint64("round", md.Round))
return common.FinalizeVote{}, nil, errAlreadyTimedOutOnRound
}

f := common.ToBeSignedFinalization{BlockHeader: md}
signature, err := f.Sign(e.Signer)
if err != nil {
Expand Down Expand Up @@ -3464,28 +3483,6 @@ func (e *Epoch) locateQuorumRecordByRound(targetRound uint64) *common.VerifiedQu
return qr
}

func (e *Epoch) haveNotFinalizedNotarizedRound() (uint64, bool) {
e.lock.Lock()
defer e.lock.Unlock()

var minRoundNum uint64
var found bool
for _, round := range e.rounds {
if round.finalization != nil || round.notarization == nil {
continue
}

if !found {
minRoundNum = round.num
found = true
} else if round.num < minRoundNum {
minRoundNum = round.num
}
}

return minRoundNum, found
}

func (e *Epoch) handleBlockDigestRequest(req *common.BlockDigestRequest, from common.NodeID) error {
e.Logger.Debug("Received block digest request", zap.Stringer("from", from), zap.Uint64("seq", req.Seq))
block, notarizationOrFinalization, ok := e.locateBlock(req.Seq, req.Digest[:])
Expand Down
76 changes: 73 additions & 3 deletions simplex/epoch_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2440,13 +2440,17 @@ func TestNotarizedNotFinalizedTipCausesEmptyBlockProposal(t *testing.T) {
name: "round 0 finalized, round 1 only notarized, round 2 only notarized",
finalizedRounds: []bool{true, false, false},
},
{
name: "round 0 finalized, round 1 only notarized, round 2 finalized",
finalizedRounds: []bool{true, false, true},
},
} {
t.Run(testCase.name, func(t *testing.T) {
nodes := []NodeID{{1}, {2}, {3}, {4}}
// Pick the node's ID such that it will be the leader in the next round.
nodeID := nodes[len(testCase.finalizedRounds)]

bb := testutil.NewTestBlockBuilder()
bb := testutil.NewTestControlledBlockBuilder(t)
recordingComm := &recordingComm{
Communication: testutil.NewNoopComm(nodes),
BroadcastMessages: make(chan *Message, 100),
Expand All @@ -2466,7 +2470,7 @@ func TestNotarizedNotFinalizedTipCausesEmptyBlockProposal(t *testing.T) {

for r, finalized := range testCase.finalizedRounds {
if finalized {
notarizeAndFinalizeRound(t, e, bb)
notarizeAndFinalizeRound(t, e, &bb.TestBlockBuilder)
continue
}
block := notarizeRoundNotFinalized(t, e, nodes, uint64(r))
Expand Down Expand Up @@ -2532,10 +2536,13 @@ func TestNotarizedNotFinalizedTipStuckLeaderCausesEmptyNotarization(t *testing.T
name: "round 0 finalized, round 1 only notarized, round 2 only notarized",
finalizedRounds: []bool{true, false, false},
},
{
name: "round 0 finalized, round 1 only notarized, round 2 finalized",
finalizedRounds: []bool{true, false, true},
},
} {
t.Run(testCase.name, func(t *testing.T) {
nodes := []NodeID{{1}, {2}, {3}, {4}}

bb := testutil.NewTestBlockBuilder()
conf, wal, _ := testutil.DefaultTestNodeEpochConfig(t, nodes[0], testutil.NewNoopComm(nodes), bb)
conf.MaxProposalWait = 50 * time.Millisecond
Expand Down Expand Up @@ -3414,3 +3421,66 @@ func TestEpochLeaderRecordsTimeoutDespiteStaleBlacklistOnEpochChange(t *testing.
require.True(t, expected.Equals(&actual),
"the timeout of nodes[2] was not recorded: expected %s, got %s", expected.String(), actual.String())
}

// TestRebroadcastDoesNotFinalizeVoteOnTimedOutRound asserts that a node never casts a finalize vote for a
// round it has cast an empty vote for. Doing both lets an empty notarization and a finalization form for
// the same round, since the two quorums may then overlap only in nodes that voted both ways.
func TestRebroadcastDoesNotFinalizeVoteOnTimedOutRound(t *testing.T) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

nice test 👍

bb := testutil.NewTestBlockBuilder()
nodes := []NodeID{{1}, {2}, {3}, {4}}
comm := &recordingComm{
Communication: testutil.NewNoopComm(nodes),
BroadcastMessages: make(chan *Message, 1000),
}
conf, wal, _ := testutil.DefaultTestNodeEpochConfig(t, nodes[0], comm, bb)
conf.ReplicationEnabled = true

e, err := NewEpoch(conf)
require.NoError(t, err)
t.Cleanup(e.Stop)
require.NoError(t, e.Start())

// Round 0: we lead, and the round is notarized and finalized normally.
notarizeAndFinalizeRound(t, e, bb)
require.Equal(t, uint64(1), e.Metadata().Round)

// Round 1: the leader never proposes, so we time out and cast an empty vote.
const timedOutRound = uint64(1)
bb.BlockShouldBeBuilt <- struct{}{}
now := conf.StartTime
testutil.WaitForBlockProposerTimeout(t, e, &now, timedOutRound)
require.True(t, wal.ContainsEmptyVote(timedOutRound))

// The other three nodes notarized a block for round 1 regardless, and we learn of it through
// replication. We store it so the round can still be finalized, but we must not vote to finalize it.
block := testutil.NewTestBlock(e.Metadata(), emptyBlacklist)
sigAggr := e.SignatureAggregatorCreator(e.Comm.Validators())
notarization, err := testutil.NewNotarization(e.Logger, sigAggr, block, nodes[1:])
require.NoError(t, err)
require.NoError(t, e.HandleMessage(&Message{ReplicationResponse: &ReplicationResponse{
Data: []QuorumRound{{Block: block, Notarization: &notarization}},
}}, nodes[1]))
wal.AssertNotarization(timedOutRound)
testutil.WaitToEnterRound(t, e, timedOutRound+1)

// Everything broadcast so far, including the empty vote, is not what this test is about.
for len(comm.BroadcastMessages) > 0 {
<-comm.BroadcastMessages
}

// Round 1 is now notarized but not finalized and the round no longer advances, so NotarizationTime
// rebroadcasts finalize votes. Drive its clock through several timeouts and make sure none of them is
// for the round we timed out on.
step := DefaultFinalizeVoteRebroadcastTimeout / 3
for range 15 {
now = now.Add(step)
e.AdvanceTime(now)
for len(comm.BroadcastMessages) > 0 {
msg := <-comm.BroadcastMessages
if msg.FinalizeVote != nil {
require.NotEqual(t, timedOutRound, msg.FinalizeVote.Finalization.Round,
"node cast a finalize vote for round %d after casting an empty vote for it", timedOutRound)
}
}
}
}
4 changes: 2 additions & 2 deletions simplex/replication_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1501,7 +1501,7 @@ func TestReplicationChain(t *testing.T) {

for {
numBlocks := n.Storage.NumBlocks()
if numBlocks == numNotarizations-missedNotarizations+1 {
if numBlocks >= numNotarizations-missedNotarizations+1 {
break
}
net.AdvanceTime(simplex.DefaultReplicationRequestTimeout)
Expand All @@ -1516,7 +1516,7 @@ func TestReplicationChain(t *testing.T) {
for {
numBlocks := blockFinalize3.Storage.NumBlocks()

if numBlocks == numNotarizations-missedNotarizations+1 {
if numBlocks >= numNotarizations-missedNotarizations+1 {
break
}
net.AdvanceTime(simplex.DefaultReplicationRequestTimeout)
Expand Down
Loading
Loading