|
|
@ -35,11 +35,7 @@ func SendMessage(peer Peer, msg []byte) { |
|
|
|
send(peer.Ip, peer.Port, content) |
|
|
|
send(peer.Ip, peer.Port, content) |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
// BroadcastMessage sends the message to a list of peers
|
|
|
|
func BoadcastToPeers(peers []Peer, content []byte) { |
|
|
|
func BroadcastMessage(peers []Peer, msg []byte) { |
|
|
|
|
|
|
|
// Construct broadcast p2p message
|
|
|
|
|
|
|
|
content := ConstructP2pMessage(byte(17), msg) |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
var wg sync.WaitGroup |
|
|
|
var wg sync.WaitGroup |
|
|
|
wg.Add(len(peers)) |
|
|
|
wg.Add(len(peers)) |
|
|
|
|
|
|
|
|
|
|
@ -56,6 +52,14 @@ func BroadcastMessage(peers []Peer, msg []byte) { |
|
|
|
log.Info("Broadcasting Down", "time spent", time.Now().Sub(start).Seconds()) |
|
|
|
log.Info("Broadcasting Down", "time spent", time.Now().Sub(start).Seconds()) |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// BroadcastMessage sends the message to a list of peers
|
|
|
|
|
|
|
|
func BroadcastMessage(peers []Peer, msg []byte) { |
|
|
|
|
|
|
|
// Construct broadcast p2p message
|
|
|
|
|
|
|
|
content := ConstructP2pMessage(byte(17), msg) |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
BoadcastToPeers(peers, content) |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
func SelectMyPeers(peers []Peer, min int, max int) []Peer { |
|
|
|
func SelectMyPeers(peers []Peer, min int, max int) []Peer { |
|
|
|
res := []Peer{} |
|
|
|
res := []Peer{} |
|
|
|
for _, peer := range peers { |
|
|
|
for _, peer := range peers { |
|
|
@ -72,37 +76,16 @@ func BroadcastMessageFromLeader(peers []Peer, msg []byte) { |
|
|
|
content := ConstructP2pMessage(byte(17), msg) |
|
|
|
content := ConstructP2pMessage(byte(17), msg) |
|
|
|
|
|
|
|
|
|
|
|
peers = SelectMyPeers(peers, 0, MAX_BROADCAST-1) |
|
|
|
peers = SelectMyPeers(peers, 0, MAX_BROADCAST-1) |
|
|
|
var wg sync.WaitGroup |
|
|
|
BoadcastToPeers(peers, content) |
|
|
|
wg.Add(len(peers)) |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
for _, peer := range peers { |
|
|
|
|
|
|
|
peerCopy := peer |
|
|
|
|
|
|
|
go func() { |
|
|
|
|
|
|
|
defer wg.Done() |
|
|
|
|
|
|
|
send(peerCopy.Ip, peerCopy.Port, content) |
|
|
|
|
|
|
|
}() |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
wg.Wait() |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
// BroadcastMessage sends the message to a list of peers from a leader.
|
|
|
|
// BroadcastMessage sends the message to a list of peers from a leader.
|
|
|
|
func BroadcastMessageFromValidator(Self Peer, peers []Peer, msg []byte) { |
|
|
|
func BroadcastMessageFromValidator(selfPeer Peer, peers []Peer, msg []byte) { |
|
|
|
|
|
|
|
|
|
|
|
// Construct broadcast p2p message
|
|
|
|
// Construct broadcast p2p message
|
|
|
|
content := ConstructP2pMessage(byte(17), msg) |
|
|
|
content := ConstructP2pMessage(byte(17), msg) |
|
|
|
|
|
|
|
|
|
|
|
peers = SelectMyPeers(peers, 0, MAX_BROADCAST-1) |
|
|
|
peers = SelectMyPeers(peers, (selfPeer.ValidatorID+1)*MAX_BROADCAST, (selfPeer.ValidatorID+2)*MAX_BROADCAST-1) |
|
|
|
var wg sync.WaitGroup |
|
|
|
BoadcastToPeers(peers, content) |
|
|
|
wg.Add(len(peers)) |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
for _, peer := range peers { |
|
|
|
|
|
|
|
peerCopy := peer |
|
|
|
|
|
|
|
go func() { |
|
|
|
|
|
|
|
defer wg.Done() |
|
|
|
|
|
|
|
send(peerCopy.Ip, peerCopy.Port, content) |
|
|
|
|
|
|
|
}() |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
wg.Wait() |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
// ConstructP2pMessage constructs the p2p message as [messageType, contentSize, content]
|
|
|
|
// ConstructP2pMessage constructs the p2p message as [messageType, contentSize, content]
|
|
|
|