diff --git a/p2p/peer.go b/p2p/peer.go index e48c71914b..f00d7cbc4b 100644 --- a/p2p/peer.go +++ b/p2p/peer.go @@ -11,7 +11,6 @@ import ( "sync" "time" - "github.com/ethereum/go-ethereum/ethutil" "github.com/ethereum/go-ethereum/event" "github.com/ethereum/go-ethereum/logger" ) @@ -462,9 +461,10 @@ func (r *eofSignal) Read(buf []byte) (int, error) { return n, err } -func (peer *Peer) PeerList() []ethutil.RlpEncodable { +func (peer *Peer) PeerList() []interface{} { peers := peer.otherPeers() - ds := make([]ethutil.RlpEncodable, 0, len(peers)) + fmt.Printf("address length: %v\n", len(peers)) + ds := make([]interface{}, 0, len(peers)) for _, p := range peers { p.infolock.Lock() addr := p.listenAddr @@ -478,7 +478,8 @@ func (peer *Peer) PeerList() []ethutil.RlpEncodable { ds = append(ds, addr) } ourAddr := peer.ourListenAddr - if ourAddr != nil && !ourAddr.IP.IsLoopback() && !ourAddr.IP.IsUnspecified() { + if ourAddr != nil && !ourAddr.IP.IsUnspecified() { + // if ourAddr != nil && !ourAddr.IP.IsLoopback() && !ourAddr.IP.IsUnspecified() { ds = append(ds, ourAddr) } fmt.Printf("address length: %v\n", len(ds)) diff --git a/p2p/protocol.go b/p2p/protocol.go index f0e5480897..381f09dfc5 100644 --- a/p2p/protocol.go +++ b/p2p/protocol.go @@ -88,20 +88,25 @@ type baseProtocol struct { func runBaseProtocol(peer *Peer, rw MsgReadWriter) error { bp := &baseProtocol{rw, peer} - if err := bp.doHandshake(rw); err != nil { + errc := make(chan error, 1) + go func() { errc <- rw.WriteMsg(bp.handshakeMsg()) }() + if err := bp.readHandshake(); err != nil { + return err + } + // handle write error + if err := <-errc; err != nil { return err } // run main loop - quit := make(chan error, 1) go func() { for { if err := bp.handle(rw); err != nil { - quit <- err + errc <- err break } } }() - return bp.loop(quit) + return bp.loop(errc) } var pingTimeout = 2 * time.Second @@ -175,7 +180,7 @@ func (bp *baseProtocol) handle(rw MsgReadWriter) error { // // TODO: add event mechanism to notify baseProtocol for new peers if len(peers) > 0 { - return bp.rw.EncodeMsg(peersMsg, peers) + return bp.rw.EncodeMsg(peersMsg, peers...) } case peersMsg: @@ -194,14 +199,9 @@ func (bp *baseProtocol) handle(rw MsgReadWriter) error { return nil } -func (bp *baseProtocol) doHandshake(rw MsgReadWriter) error { - // send our handshake - if err := rw.WriteMsg(bp.handshakeMsg()); err != nil { - return err - } - +func (bp *baseProtocol) readHandshake() error { // read and handle remote handshake - msg, err := rw.ReadMsg() + msg, err := bp.rw.ReadMsg() if err != nil { return err } @@ -211,12 +211,10 @@ func (bp *baseProtocol) doHandshake(rw MsgReadWriter) error { if msg.Size > baseProtocolMaxMsgSize { return newPeerError(errMisc, "message too big") } - var hs handshake if err := msg.Decode(&hs); err != nil { return err } - // validate handshake info if hs.Version != baseProtocolVersion { return newPeerError(errP2PVersionMismatch, "Require protocol %d, received %d\n", @@ -239,9 +237,7 @@ func (bp *baseProtocol) doHandshake(rw MsgReadWriter) error { if err := bp.peer.pubkeyHook(pa); err != nil { return newPeerError(errPubkeyForbidden, "%v", err) } - // TODO: remove Caps with empty name - var addr *peerAddr if hs.ListenPort != 0 { addr = newPeerAddr(bp.peer.conn.RemoteAddr(), hs.NodeID) diff --git a/p2p/protocol_test.go b/p2p/protocol_test.go index 0844fe7fdf..5a0793cc95 100644 --- a/p2p/protocol_test.go +++ b/p2p/protocol_test.go @@ -2,6 +2,8 @@ package p2p import ( "fmt" + "net" + "reflect" "testing" "github.com/ethereum/go-ethereum/crypto" @@ -24,7 +26,7 @@ func (self *peerId) Pubkey() (pubkey []byte) { return } -func testPeerFree() (peer *Peer) { +func newTestPeer() (peer *Peer) { peer = NewPeer(&peerId{}, []Cap{}) peer.pubkeyHook = func(*peerAddr) error { return nil } peer.ourID = &peerId{} @@ -38,6 +40,7 @@ func TestBaseProtocolPeers(t *testing.T) { {IP: net.ParseIP("1.2.3.4"), Port: 2222, Pubkey: []byte{}}, {IP: net.ParseIP("5.6.7.8"), Port: 3333, Pubkey: []byte{}}, } + var ownAddr *peerAddr = &peerAddr{IP: net.ParseIP("1.3.5.7"), Port: 1111, Pubkey: []byte{}} rw1, rw2 := MsgPipe() // run matcher, close pipe when addresses have arrived addrChan := make(chan *peerAddr, len(cannedPeerList)) @@ -51,16 +54,18 @@ func TestBaseProtocolPeers(t *testing.T) { } close(addrChan) var own []*peerAddr - for _, got = range addrChan { + var got *peerAddr + for got = range addrChan { own = append(own, got) } - if len(own) != 1 || !reflect.DeepEqual(own[0], ourAddr) { - t.Errorf("mismatch: peers own address is incorrectly or not given, got %v, want %#v", ownAddr, own) + if len(own) != 1 || !reflect.DeepEqual(ownAddr, own[0]) { + t.Errorf("mismatch: peers own address is incorrectly or not given, got %v, want %#v", ownAddr) } rw2.Close() }() // run first peer - peer1 := testPeer() + peer1 := newTestPeer() + peer1.ourListenAddr = ownAddr peer1.otherPeers = func() []*Peer { pl := make([]*Peer, len(cannedPeerList)) for i, addr := range cannedPeerList { @@ -70,7 +75,7 @@ func TestBaseProtocolPeers(t *testing.T) { } go runBaseProtocol(peer1, rw1) // run second peer - peer2 := testPeer() + peer2 := newTestPeer() peer2.newPeerAddr = addrChan // feed peer suggestions into matcher if err := runBaseProtocol(peer2, rw2); err != ErrPipeClosed { t.Errorf("peer2 terminated with unexpected error: %v", err)