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
60 changes: 37 additions & 23 deletions adapters.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,15 @@ type Communication struct {
Broadcaster
}

func newCommunication(sender Sender, broadcaster Broadcaster, validators common.Nodes) *Communication {
c := &Communication{
Sender: sender,
Broadcaster: broadcaster,
}
c.SetValidators(validators)
return c
}

func (c *Communication) SetValidators(nodes common.Nodes) {
c.nodes.Store(nodes)
}
Expand All @@ -32,48 +41,53 @@ func (c *Communication) Validators() common.Nodes {
return nodes
}

// EpochAwareStorage is a wrapper around Storage that is aware of epoch changes.
// Upon an epoch change, it will ignore blocks from previous epochs
// and will call the onEpochChange callback when a new epoch is detected.
type EpochAwareStorage struct {
// InstanceStorage is a wrapper around Storage that skips indexing Telocks
// and delegates post-index handling to a caller-provided onIndex hook.
type InstanceStorage struct {
// CachedStorage is used to ensure that we prune the cache on Index.
*CachedStorage

msm *metadata.StateMachine
onEpochChange func(seq uint64, validators common.Nodes) error
epoch uint64
msm *metadata.StateMachine

onIndex func(block *ParsedBlock) error
}

func NewInstanceStorage(storage *CachedStorage, msm *metadata.StateMachine, onIndex func(block *ParsedBlock) error) *InstanceStorage {
return &InstanceStorage{
CachedStorage: storage,
msm: msm,
onIndex: onIndex,
}
}

func (e *EpochAwareStorage) Retrieve(seq uint64) (common.VerifiedBlock, common.Finalization, error) {
block, finalization, err := e.GetBlock(seq)
func (s *InstanceStorage) Retrieve(seq uint64) (common.VerifiedBlock, common.Finalization, error) {
block, finalization, err := s.GetBlock(seq)
if err != nil {
return nil, common.Finalization{}, err
}
parsedBlock := &ParsedBlock{
msm: e.msm,
msm: s.msm,
StateMachineBlock: block,
}
return parsedBlock, *finalization, nil
}

func (e *EpochAwareStorage) Index(ctx context.Context, block common.VerifiedBlock, certificate common.Finalization) error {
if block.BlockHeader().Epoch < e.epoch {
// This is a Telock from a previous epoch, so we ignore it and do not index it.
func (s *InstanceStorage) Index(ctx context.Context, block common.VerifiedBlock, certificate common.Finalization) error {
pb, ok := block.(*ParsedBlock)
if !ok {
return fmt.Errorf("expected ParsedBlock, got %T", block)
}

// A Telock only extends time until the epoch transition finalizes, so we never index it.
if pb.Type() == metadata.BlockTypeTelock {
return nil
}

if err := e.CachedStorage.Index(ctx, block, certificate); err != nil {
if err := s.CachedStorage.Index(ctx, block, certificate); err != nil {
return err
}
// This is a sealing block, and it is not the zero block
if block.SealingBlockInfo() != nil && block.SealingBlockInfo().PrevSealingBlockHash != [32]byte{} {
if err := e.onEpochChange(block.BlockHeader().Seq, block.SealingBlockInfo().ValidatorSet); err != nil {
return err
}
// We are now in a new epoch, so we update the epoch number to prevent indexing Telocks from the previous epoch.
e.epoch = block.BlockHeader().Seq
}
return nil

return s.onIndex(pb)
}

// cachedBlock is a wrapper around ParsedBlock that caches the block in the CachedStorage upon verification.
Expand Down
Loading
Loading