|
|
@ -64,30 +64,68 @@ func (node *Node) BroadcastCXReceipts(newBlock *types.Block, lastCommits []byte) |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// BroadcastCXReceiptsWithShardID broadcasts cross shard receipts to given ToShardID
|
|
|
|
|
|
|
|
func (node *Node) BroadcastCXReceiptsWithShardID(newBlock *types.Block, lastCommits []byte, toShardID uint32) { |
|
|
|
|
|
|
|
//#### Read payload data from committed msg
|
|
|
|
|
|
|
|
if len(lastCommits) <= 96 { |
|
|
|
|
|
|
|
utils.Logger().Debug().Int("lastCommitsLen", len(lastCommits)).Msg("[BroadcastCXReceipts] lastCommits Not Enough Length") |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
commitSig := make([]byte, 96) |
|
|
|
|
|
|
|
commitBitmap := make([]byte, len(lastCommits)-96) |
|
|
|
|
|
|
|
offset := 0 |
|
|
|
|
|
|
|
copy(commitSig[:], lastCommits[offset:offset+96]) |
|
|
|
|
|
|
|
offset += 96 |
|
|
|
|
|
|
|
copy(commitBitmap[:], lastCommits[offset:]) |
|
|
|
|
|
|
|
//#### END Read payload data from committed msg
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
epoch := newBlock.Header().Epoch() |
|
|
|
|
|
|
|
shardingConfig := core.ShardingSchedule.InstanceForEpoch(epoch) |
|
|
|
|
|
|
|
shardNum := int(shardingConfig.NumShards()) |
|
|
|
|
|
|
|
myShardID := node.Consensus.ShardID |
|
|
|
|
|
|
|
utils.Logger().Info().Int("shardNum", shardNum).Uint32("myShardID", myShardID).Uint64("blockNum", newBlock.NumberU64()).Msg("[BroadcastCXReceipts]") |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
cxReceipts, err := node.Blockchain().ReadCXReceipts(toShardID, newBlock.NumberU64(), newBlock.Hash(), false) |
|
|
|
|
|
|
|
if err != nil || len(cxReceipts) == 0 { |
|
|
|
|
|
|
|
utils.Logger().Warn().Err(err).Uint32("ToShardID", toShardID).Int("numCXReceipts", len(cxReceipts)).Msg("[BroadcastCXReceipts] No ReadCXReceipts found") |
|
|
|
|
|
|
|
return |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
merkleProof, err := node.Blockchain().CXMerkleProof(toShardID, newBlock) |
|
|
|
|
|
|
|
if err != nil { |
|
|
|
|
|
|
|
utils.Logger().Warn().Uint32("ToShardID", toShardID).Msg("[BroadcastCXReceipts] Unable to get merkleProof") |
|
|
|
|
|
|
|
return |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
utils.Logger().Info().Uint32("ToShardID", toShardID).Msg("[BroadcastCXReceipts] ReadCXReceipts and MerkleProof Found") |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
groupID := p2p.ShardID(toShardID) |
|
|
|
|
|
|
|
go node.host.SendMessageToGroups([]p2p.GroupID{p2p.NewGroupIDByShardID(groupID)}, host.ConstructP2pMessage(byte(0), proto_node.ConstructCXReceiptsProof(cxReceipts, merkleProof, newBlock.Header(), commitSig, commitBitmap))) |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
// BroadcastMissingCXReceipts broadcasts missing cross shard receipts per request
|
|
|
|
// BroadcastMissingCXReceipts broadcasts missing cross shard receipts per request
|
|
|
|
func (node *Node) BroadcastMissingCXReceipts() { |
|
|
|
func (node *Node) BroadcastMissingCXReceipts() { |
|
|
|
sendNextTime := []common.Hash{} |
|
|
|
sendNextTime := []core.CxEntry{} |
|
|
|
it := node.CxPool.Pool().Iterator() |
|
|
|
it := node.CxPool.Pool().Iterator() |
|
|
|
for blockHash := range it.C { |
|
|
|
for entry := range it.C { |
|
|
|
blk := node.Blockchain().GetBlockByHash(blockHash.(common.Hash)) |
|
|
|
cxEntry := entry.(core.CxEntry) |
|
|
|
|
|
|
|
toShardID := cxEntry.ToShardID |
|
|
|
|
|
|
|
blk := node.Blockchain().GetBlockByHash(cxEntry.BlockHash) |
|
|
|
if blk == nil { |
|
|
|
if blk == nil { |
|
|
|
continue |
|
|
|
continue |
|
|
|
} |
|
|
|
} |
|
|
|
blockNum := blk.NumberU64() |
|
|
|
blockNum := blk.NumberU64() |
|
|
|
nextBlk := node.Blockchain().GetBlockByNumber(blockNum + 1) |
|
|
|
nextHeader := node.Blockchain().GetHeaderByNumber(blockNum + 1) |
|
|
|
if nextBlk == nil { |
|
|
|
if nextHeader == nil { |
|
|
|
sendNextTime = append(sendNextTime, blockHash.(common.Hash)) |
|
|
|
sendNextTime = append(sendNextTime, cxEntry) |
|
|
|
continue |
|
|
|
continue |
|
|
|
} |
|
|
|
} |
|
|
|
sig := nextBlk.Header().LastCommitSignature() |
|
|
|
sig := nextHeader.LastCommitSignature() |
|
|
|
bitmap := nextBlk.Header().LastCommitBitmap() |
|
|
|
bitmap := nextHeader.LastCommitBitmap() |
|
|
|
lastCommits := append(sig[:], bitmap...) |
|
|
|
lastCommits := append(sig[:], bitmap...) |
|
|
|
go node.BroadcastCXReceipts(blk, lastCommits) |
|
|
|
go node.BroadcastCXReceiptsWithShardID(blk, lastCommits, toShardID) |
|
|
|
} |
|
|
|
} |
|
|
|
node.CxPool.Clear() |
|
|
|
node.CxPool.Clear() |
|
|
|
// this should not happen or maybe happen for impatient user
|
|
|
|
// this should not happen or maybe happen for impatient user
|
|
|
|
for _, blockHash := range sendNextTime { |
|
|
|
for _, entry := range sendNextTime { |
|
|
|
node.CxPool.Add(blockHash) |
|
|
|
node.CxPool.Add(entry) |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|