From 03d4401f4769430b0c2c1a62003237a51e3b6a89 Mon Sep 17 00:00:00 2001 From: Guillaume Ballet Date: Sat, 24 Mar 2018 17:29:27 +0100 Subject: [PATCH] whisper: Open stream to the client + send registration msg --- cmd/wnode/main.go | 7 ++- whisper/whisperv6/doc.go | 1 + whisper/whisperv6/libp2p_glue.go | 76 ++++++++++++++++++++++++++++---- whisper/whisperv6/whisper.go | 1 + 4 files changed, 75 insertions(+), 10 deletions(-) diff --git a/cmd/wnode/main.go b/cmd/wnode/main.go index fff337024a..660b2901e2 100644 --- a/cmd/wnode/main.go +++ b/cmd/wnode/main.go @@ -208,7 +208,7 @@ func initialize() { if *useLibP2P { var libp2pbootstrap *string for libp2pbootstrap == nil { - libp2pbootstrap = scanLineA("Please enter the bootstrap node's addres: ") + libp2pbootstrap = scanLineA("Please enter the bootstrap node's address: ") } libp2pbootaddr, err := multiaddr.NewMultiaddr(*libp2pbootstrap) if err != nil { @@ -293,7 +293,10 @@ func initialize() { } if *useLibP2P { - server, err = whisper.NewLibP2PWhisperServer() + var port uint + var a string + fmt.Scanf(*argIP, "%s:%d", &a, &port) + server, err = whisper.NewLibP2PWhisperServer(port, shh) if err != nil { utils.Fatalf("Error starting the libp2p client: %v", err) } diff --git a/whisper/whisperv6/doc.go b/whisper/whisperv6/doc.go index 8cde3a719d..09b646873f 100644 --- a/whisper/whisperv6/doc.go +++ b/whisper/whisperv6/doc.go @@ -48,6 +48,7 @@ const ( messagesCode = 1 // normal whisper message powRequirementCode = 2 // PoW requirement bloomFilterExCode = 3 // bloom filter exchange + lp2pPeerCode = 125 // libp2p-based peer exchange protocol p2pRequestCode = 126 // peer-to-peer message, used by Dapp protocol p2pMessageCode = 127 // peer-to-peer message (to be consumed by the peer, but not forwarded any further) NumberOfMessageCodes = 128 diff --git a/whisper/whisperv6/libp2p_glue.go b/whisper/whisperv6/libp2p_glue.go index 41dcb21eb4..eeb44219ad 100644 --- a/whisper/whisperv6/libp2p_glue.go +++ b/whisper/whisperv6/libp2p_glue.go @@ -31,6 +31,7 @@ import ( peer "github.com/libp2p/go-libp2p-peer" ma "github.com/multiformats/go-multiaddr" set "gopkg.in/fatih/set.v0" + pstore "github.com/libp2p/go-libp2p-peerstore" ) // LibP2PStream is a wrapper used to implement the MsgReadWriter @@ -120,6 +121,8 @@ type LibP2PPeer struct { *PeerBase id peer.ID + + connectionStream *LibP2PStream } func newLibP2PPeer(w *Whisper, pid peer.ID, rw p2p.MsgReadWriter) Peer { @@ -135,6 +138,7 @@ func newLibP2PPeer(w *Whisper, pid peer.ID, rw p2p.MsgReadWriter) Peer { fullNode: true, }, pid, + nil, } } @@ -156,6 +160,8 @@ type LibP2PWhisperServer struct { Host host.Host Peers []*LibP2PPeer + + whisper *Whisper } // Start starts the server @@ -172,23 +178,66 @@ func (server *LibP2PWhisperServer) Start() error { } } - whisper := server.Peers[0].host lps := &LibP2PStream{stream} // Unknown peer if peer == nil { - peer = newLibP2PPeer(whisper, pid, lps) + peer = newLibP2PPeer(server.whisper, pid, lps) // TODO check critical section server.Peers = append(server.Peers, peer.(*LibP2PPeer)) } - whisper.runMessageLoop(peer, lps) + server.whisper.runMessageLoop(peer, lps) }) + + fmt.Println("Currently having the following peers:", server.Peers) + + // Open a stream to every peer currently known + for _, p := range server.Peers { + fmt.Println("opening stream to peer: ", p) + // Do not dial to self - actually, this should raise + // an exception rather than be ignored + // if pid == server.Host.ID() { + // continue + // } + + // Create a stream with the peer + s, err := server.Host.NewStream(context.Background(), p.id, WhisperProtocolString) + if err != nil { + panic(err) + } + + // Save the stream + lps := LibP2PStream{ + stream: s, + } + p.connectionStream = &lps + + // Send a first message to notify the remote peer we want + // to connect + connectMsg := p2p.Msg{ + Code: lp2pPeerCode, + Size: 0, + Payload: bytes.NewReader([]byte{}), + } + + err = lps.WriteMsg(connectMsg) + if err != nil { + // XXX + panic(err) + } + } + return nil } // Stop stops the server func (server *LibP2PWhisperServer) Stop() { + for _, p := range server.Peers { + // TODO send disconnect message + p.connectionStream.stream.Close() + } + server.Host.Close() } @@ -205,6 +254,7 @@ func (server *LibP2PWhisperServer) Enode() string { return fullAddr.String() } +// AddPeer is a helper function to add peers to the server func (server *LibP2PWhisperServer) AddPeer(addr ma.Multiaddr) *LibP2PPeer { fmt.Println("Adding peer: ", addr) pid, err := addr.ValueForProtocol(ma.P_IPFS) @@ -212,15 +262,25 @@ func (server *LibP2PWhisperServer) AddPeer(addr ma.Multiaddr) *LibP2PPeer { // XXX panic(err) } - return &LibP2PPeer{id: peer.ID(pid)} + peerid, err := peer.IDB58Decode(pid) + if err != nil { + panic(err) + } + ipfsaddrpart, _ := ma.NewMultiaddr(fmt.Sprintf("/ipfs/%s", pid)) + ipaddr := addr.Decapsulate(ipfsaddrpart) + server.Host.Peerstore().AddAddr(peerid, ipaddr, pstore.PermanentAddrTTL) + newPeer := &LibP2PPeer{id: peer.ID(pid)} + server.Peers = append(server.Peers, newPeer) + + return newPeer } // NewLibP2PWhisperServer creates a new WhisperServer with // a libp2p backend. -func NewLibP2PWhisperServer() (WhisperServer, error) { +func NewLibP2PWhisperServer(port uint, whisper *Whisper) (WhisperServer, error) { priv, _, err := crypto.GenerateKeyPair(crypto.Ed25519, 384) opts := []libp2p.Option{ - libp2p.ListenAddrStrings(fmt.Sprintf("/ip4/0.0.0.0/tcp/%d", WhisperPort)), + libp2p.ListenAddrStrings(fmt.Sprintf("/ip4/0.0.0.0/tcp/%d", port)), libp2p.Identity(priv), } @@ -229,6 +289,6 @@ func NewLibP2PWhisperServer() (WhisperServer, error) { return nil, fmt.Errorf("Error setting up the libp2p network: %s", err) } - server := &LibP2PWhisperServer{h, []*LibP2PPeer{}} + server := &LibP2PWhisperServer{h, []*LibP2PPeer{}, whisper} return server, nil -} +} \ No newline at end of file diff --git a/whisper/whisperv6/whisper.go b/whisper/whisperv6/whisper.go index 2bfd32b96a..e4132639d6 100644 --- a/whisper/whisperv6/whisper.go +++ b/whisper/whisperv6/whisper.go @@ -718,6 +718,7 @@ func (whisper *Whisper) runMessageLoop(p Peer, rw p2p.MsgReadWriter) error { return errors.New("invalid bloom filter exchange message") } p.setBloomFilter(bloom) + case lp2pPeerCode: case p2pMessageCode: // peer-to-peer message, sent directly to peer bypassing PoW checks, etc. // this message is not supposed to be forwarded to other peers, and