(peerIdx uint32, data []byte)
| 401 | } |
| 402 | |
| 403 | func (self *Server) sendToPeer(peerIdx uint32, data []byte) error { |
| 404 | peer := self.peerPool.getPeer(peerIdx) |
| 405 | if peer == nil { |
| 406 | return fmt.Errorf("send peer failed: failed to get peer %d", peerIdx) |
| 407 | } |
| 408 | msg := &p2pmsg.ConsensusPayload{ |
| 409 | Data: data, |
| 410 | Owner: self.account.PublicKey, |
| 411 | } |
| 412 | |
| 413 | sink := common.NewZeroCopySink(nil) |
| 414 | msg.SerializationUnsigned(sink) |
| 415 | msg.Signature, _ = signature.Sign(self.account, sink.Bytes()) |
| 416 | |
| 417 | cons := msgpack.NewConsensus(msg) |
| 418 | p2pid, present := self.peerPool.getP2pId(peerIdx) |
| 419 | if present { |
| 420 | self.p2p.Transmit(p2pid, cons) |
| 421 | } else { |
| 422 | log.Errorf("sendToPeer transmit failed index:%d", peerIdx) |
| 423 | } |
| 424 | return nil |
| 425 | } |
| 426 | |
| 427 | func (self *Server) broadcast(msg ConsensusMsg) { |
| 428 | self.msgSendC <- &SendMsgEvent{ |
no test coverage detected