-
Notifications
You must be signed in to change notification settings - Fork 178
/
finalization_distributor.go
104 lines (79 loc) · 3.49 KB
/
finalization_distributor.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
package pubsub
import (
"github.com/onflow/flow-go/consensus/hotstuff"
"sync"
"github.com/onflow/flow-go/consensus/hotstuff/model"
"github.com/onflow/flow-go/model/flow"
)
type OnBlockFinalizedConsumer = func(finalizedBlockID flow.Identifier)
type OnBlockIncorporatedConsumer = func(incorporatedBlockID flow.Identifier)
// FinalizationDistributor subscribes for finalization events from hotstuff and distributes it to subscribers
type FinalizationDistributor struct {
blockFinalizedConsumers []OnBlockFinalizedConsumer
blockIncorporatedConsumers []OnBlockIncorporatedConsumer
hotStuffFinalizationConsumers []hotstuff.FinalizationConsumer
lock sync.RWMutex
}
func NewFinalizationDistributor() *FinalizationDistributor {
return &FinalizationDistributor{
blockFinalizedConsumers: make([]OnBlockFinalizedConsumer, 0),
blockIncorporatedConsumers: make([]OnBlockIncorporatedConsumer, 0),
lock: sync.RWMutex{},
}
}
func (p *FinalizationDistributor) AddOnBlockFinalizedConsumer(consumer OnBlockFinalizedConsumer) {
p.lock.Lock()
defer p.lock.Unlock()
p.blockFinalizedConsumers = append(p.blockFinalizedConsumers, consumer)
}
func (p *FinalizationDistributor) AddOnBlockIncorporatedConsumer(consumer OnBlockIncorporatedConsumer) {
p.lock.Lock()
defer p.lock.Unlock()
p.blockIncorporatedConsumers = append(p.blockIncorporatedConsumers, consumer)
}
func (p *FinalizationDistributor) AddConsumer(consumer hotstuff.FinalizationConsumer) {
p.lock.Lock()
defer p.lock.Unlock()
p.hotStuffFinalizationConsumers = append(p.hotStuffFinalizationConsumers, consumer)
}
func (p *FinalizationDistributor) OnEventProcessed() {}
func (p *FinalizationDistributor) OnBlockIncorporated(block *model.Block) {
p.lock.RLock()
defer p.lock.RUnlock()
for _, consumer := range p.blockIncorporatedConsumers {
consumer(block.BlockID)
}
for _, consumer := range p.hotStuffFinalizationConsumers {
consumer.OnBlockIncorporated(block)
}
}
func (p *FinalizationDistributor) OnFinalizedBlock(block *model.Block) {
p.lock.RLock()
defer p.lock.RUnlock()
for _, consumer := range p.blockFinalizedConsumers {
consumer(block.BlockID)
}
for _, consumer := range p.hotStuffFinalizationConsumers {
consumer.OnFinalizedBlock(block)
}
}
func (p *FinalizationDistributor) OnDoubleProposeDetected(block1, block2 *model.Block) {
p.lock.RLock()
defer p.lock.RUnlock()
for _, consumer := range p.hotStuffFinalizationConsumers {
consumer.OnDoubleProposeDetected(block1, block2)
}
}
func (p *FinalizationDistributor) OnReceiveVote(uint64, *model.Vote) {}
func (p *FinalizationDistributor) OnReceiveProposal(uint64, *model.Proposal) {}
func (p *FinalizationDistributor) OnEnteringView(uint64, flow.Identifier) {}
func (p *FinalizationDistributor) OnQcTriggeredViewChange(*flow.QuorumCertificate, uint64) {}
func (p *FinalizationDistributor) OnProposingBlock(*model.Proposal) {}
func (p *FinalizationDistributor) OnVoting(*model.Vote) {}
func (p *FinalizationDistributor) OnQcConstructedFromVotes(*flow.QuorumCertificate) {}
func (p *FinalizationDistributor) OnStartingTimeout(*model.TimerInfo) {}
func (p *FinalizationDistributor) OnReachedTimeout(*model.TimerInfo) {}
func (p *FinalizationDistributor) OnQcIncorporated(*flow.QuorumCertificate) {}
func (p *FinalizationDistributor) OnForkChoiceGenerated(uint64, *flow.QuorumCertificate) {}
func (p *FinalizationDistributor) OnDoubleVotingDetected(*model.Vote, *model.Vote) {}
func (p *FinalizationDistributor) OnInvalidVoteDetected(*model.Vote) {}