mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-18 18:02:24 +00:00
whisper: use the regular handshake in libp2p
This commit is contained in:
parent
80d093ddd4
commit
489cc699e8
2 changed files with 15 additions and 7 deletions
|
|
@ -180,12 +180,14 @@ func (server *LibP2PWhisperServer) connectToPeer(p *LibP2PPeer) error {
|
||||||
stream: s,
|
stream: s,
|
||||||
}
|
}
|
||||||
p.connectionStream = &lps
|
p.connectionStream = &lps
|
||||||
|
p.ws = p.connectionStream
|
||||||
|
|
||||||
// If we got here, it means that a connection was established.
|
|
||||||
// Save the peer.
|
|
||||||
|
|
||||||
// TODO send my known list of peers
|
// TODO send my known list of peers
|
||||||
|
|
||||||
|
// Call HandlePeer to perform the handshake
|
||||||
|
go server.whisper.HandlePeer(p, p.connectionStream)
|
||||||
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -212,7 +214,7 @@ func (server *LibP2PWhisperServer) Start() error {
|
||||||
server.Peers = append(server.Peers, peer.(*LibP2PPeer))
|
server.Peers = append(server.Peers, peer.(*LibP2PPeer))
|
||||||
}
|
}
|
||||||
|
|
||||||
go server.whisper.runMessageLoop(peer, lps)
|
go server.whisper.HandlePeer(peer, lps)
|
||||||
})
|
})
|
||||||
|
|
||||||
fmt.Println("Currently having the following peers:", server.Peers)
|
fmt.Println("Currently having the following peers:", server.Peers)
|
||||||
|
|
@ -266,7 +268,7 @@ func (server *LibP2PWhisperServer) AddPeer(addr ma.Multiaddr) *LibP2PPeer {
|
||||||
ipfsaddrpart, _ := ma.NewMultiaddr(fmt.Sprintf("/ipfs/%s", pid))
|
ipfsaddrpart, _ := ma.NewMultiaddr(fmt.Sprintf("/ipfs/%s", pid))
|
||||||
ipaddr := addr.Decapsulate(ipfsaddrpart)
|
ipaddr := addr.Decapsulate(ipfsaddrpart)
|
||||||
server.Host.Peerstore().AddAddr(peerid, ipaddr, pstore.PermanentAddrTTL)
|
server.Host.Peerstore().AddAddr(peerid, ipaddr, pstore.PermanentAddrTTL)
|
||||||
newPeer := &LibP2PPeer{id: peerid}
|
newPeer := newLibP2PPeer(server.whisper, peerid, nil).(*LibP2PPeer)
|
||||||
server.Peers = append(server.Peers, newPeer)
|
server.Peers = append(server.Peers, newPeer)
|
||||||
|
|
||||||
return newPeer
|
return newPeer
|
||||||
|
|
|
||||||
|
|
@ -127,7 +127,7 @@ func New(cfg *Config) *Whisper {
|
||||||
Name: ProtocolName,
|
Name: ProtocolName,
|
||||||
Version: uint(ProtocolVersion),
|
Version: uint(ProtocolVersion),
|
||||||
Length: NumberOfMessageCodes,
|
Length: NumberOfMessageCodes,
|
||||||
Run: whisper.HandlePeer,
|
Run: whisper.HandleDevP2PPeer,
|
||||||
NodeInfo: func() interface{} {
|
NodeInfo: func() interface{} {
|
||||||
return map[string]interface{}{
|
return map[string]interface{}{
|
||||||
"version": ProtocolVersionStr,
|
"version": ProtocolVersionStr,
|
||||||
|
|
@ -626,12 +626,18 @@ func (whisper *Whisper) Stop() error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// HandlePeer is called by the underlying P2P layer when the whisper sub-protocol
|
// HandleDevP2PPeer is called by the underlying P2P layer when the whisper sub-protocol
|
||||||
// connection is negotiated.
|
// connection is negotiated.
|
||||||
func (whisper *Whisper) HandlePeer(peer *p2p.Peer, rw p2p.MsgReadWriter) error {
|
func (whisper *Whisper) HandleDevP2PPeer(peer *p2p.Peer, rw p2p.MsgReadWriter) error {
|
||||||
// Create the new peer and start tracking it
|
// Create the new peer and start tracking it
|
||||||
whisperPeer := newPeer(whisper, peer, rw)
|
whisperPeer := newPeer(whisper, peer, rw)
|
||||||
|
|
||||||
|
return whisper.HandlePeer(whisperPeer, rw)
|
||||||
|
}
|
||||||
|
|
||||||
|
// HandlePeer sets up the connection with the peer once the underlying
|
||||||
|
// layer has established it.
|
||||||
|
func (whisper *Whisper) HandlePeer(whisperPeer Peer, rw p2p.MsgReadWriter) error {
|
||||||
whisper.peerMu.Lock()
|
whisper.peerMu.Lock()
|
||||||
whisper.peers[whisperPeer] = struct{}{}
|
whisper.peers[whisperPeer] = struct{}{}
|
||||||
whisper.peerMu.Unlock()
|
whisper.peerMu.Unlock()
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue