fix getPeerMsg/peerMsg RLP encode/decode, logs. tests pass

This commit is contained in:
zelig 2015-01-03 14:57:31 +00:00
parent 4249bfb7dc
commit 9380644db0
3 changed files with 28 additions and 26 deletions

View file

@ -11,7 +11,6 @@ import (
"sync" "sync"
"time" "time"
"github.com/ethereum/go-ethereum/ethutil"
"github.com/ethereum/go-ethereum/event" "github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/logger" "github.com/ethereum/go-ethereum/logger"
) )
@ -462,9 +461,10 @@ func (r *eofSignal) Read(buf []byte) (int, error) {
return n, err return n, err
} }
func (peer *Peer) PeerList() []ethutil.RlpEncodable { func (peer *Peer) PeerList() []interface{} {
peers := peer.otherPeers() 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 { for _, p := range peers {
p.infolock.Lock() p.infolock.Lock()
addr := p.listenAddr addr := p.listenAddr
@ -478,7 +478,8 @@ func (peer *Peer) PeerList() []ethutil.RlpEncodable {
ds = append(ds, addr) ds = append(ds, addr)
} }
ourAddr := peer.ourListenAddr 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) ds = append(ds, ourAddr)
} }
fmt.Printf("address length: %v\n", len(ds)) fmt.Printf("address length: %v\n", len(ds))

View file

@ -88,20 +88,25 @@ type baseProtocol struct {
func runBaseProtocol(peer *Peer, rw MsgReadWriter) error { func runBaseProtocol(peer *Peer, rw MsgReadWriter) error {
bp := &baseProtocol{rw, peer} 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 return err
} }
// run main loop // run main loop
quit := make(chan error, 1)
go func() { go func() {
for { for {
if err := bp.handle(rw); err != nil { if err := bp.handle(rw); err != nil {
quit <- err errc <- err
break break
} }
} }
}() }()
return bp.loop(quit) return bp.loop(errc)
} }
var pingTimeout = 2 * time.Second 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 // TODO: add event mechanism to notify baseProtocol for new peers
if len(peers) > 0 { if len(peers) > 0 {
return bp.rw.EncodeMsg(peersMsg, peers) return bp.rw.EncodeMsg(peersMsg, peers...)
} }
case peersMsg: case peersMsg:
@ -194,14 +199,9 @@ func (bp *baseProtocol) handle(rw MsgReadWriter) error {
return nil return nil
} }
func (bp *baseProtocol) doHandshake(rw MsgReadWriter) error { func (bp *baseProtocol) readHandshake() error {
// send our handshake
if err := rw.WriteMsg(bp.handshakeMsg()); err != nil {
return err
}
// read and handle remote handshake // read and handle remote handshake
msg, err := rw.ReadMsg() msg, err := bp.rw.ReadMsg()
if err != nil { if err != nil {
return err return err
} }
@ -211,12 +211,10 @@ func (bp *baseProtocol) doHandshake(rw MsgReadWriter) error {
if msg.Size > baseProtocolMaxMsgSize { if msg.Size > baseProtocolMaxMsgSize {
return newPeerError(errMisc, "message too big") return newPeerError(errMisc, "message too big")
} }
var hs handshake var hs handshake
if err := msg.Decode(&hs); err != nil { if err := msg.Decode(&hs); err != nil {
return err return err
} }
// validate handshake info // validate handshake info
if hs.Version != baseProtocolVersion { if hs.Version != baseProtocolVersion {
return newPeerError(errP2PVersionMismatch, "Require protocol %d, received %d\n", 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 { if err := bp.peer.pubkeyHook(pa); err != nil {
return newPeerError(errPubkeyForbidden, "%v", err) return newPeerError(errPubkeyForbidden, "%v", err)
} }
// TODO: remove Caps with empty name // TODO: remove Caps with empty name
var addr *peerAddr var addr *peerAddr
if hs.ListenPort != 0 { if hs.ListenPort != 0 {
addr = newPeerAddr(bp.peer.conn.RemoteAddr(), hs.NodeID) addr = newPeerAddr(bp.peer.conn.RemoteAddr(), hs.NodeID)

View file

@ -2,6 +2,8 @@ package p2p
import ( import (
"fmt" "fmt"
"net"
"reflect"
"testing" "testing"
"github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/crypto"
@ -24,7 +26,7 @@ func (self *peerId) Pubkey() (pubkey []byte) {
return return
} }
func testPeerFree() (peer *Peer) { func newTestPeer() (peer *Peer) {
peer = NewPeer(&peerId{}, []Cap{}) peer = NewPeer(&peerId{}, []Cap{})
peer.pubkeyHook = func(*peerAddr) error { return nil } peer.pubkeyHook = func(*peerAddr) error { return nil }
peer.ourID = &peerId{} 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("1.2.3.4"), Port: 2222, Pubkey: []byte{}},
{IP: net.ParseIP("5.6.7.8"), Port: 3333, 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() rw1, rw2 := MsgPipe()
// run matcher, close pipe when addresses have arrived // run matcher, close pipe when addresses have arrived
addrChan := make(chan *peerAddr, len(cannedPeerList)) addrChan := make(chan *peerAddr, len(cannedPeerList))
@ -51,16 +54,18 @@ func TestBaseProtocolPeers(t *testing.T) {
} }
close(addrChan) close(addrChan)
var own []*peerAddr var own []*peerAddr
for _, got = range addrChan { var got *peerAddr
for got = range addrChan {
own = append(own, got) own = append(own, got)
} }
if len(own) != 1 || !reflect.DeepEqual(own[0], ourAddr) { 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, own) t.Errorf("mismatch: peers own address is incorrectly or not given, got %v, want %#v", ownAddr)
} }
rw2.Close() rw2.Close()
}() }()
// run first peer // run first peer
peer1 := testPeer() peer1 := newTestPeer()
peer1.ourListenAddr = ownAddr
peer1.otherPeers = func() []*Peer { peer1.otherPeers = func() []*Peer {
pl := make([]*Peer, len(cannedPeerList)) pl := make([]*Peer, len(cannedPeerList))
for i, addr := range cannedPeerList { for i, addr := range cannedPeerList {
@ -70,7 +75,7 @@ func TestBaseProtocolPeers(t *testing.T) {
} }
go runBaseProtocol(peer1, rw1) go runBaseProtocol(peer1, rw1)
// run second peer // run second peer
peer2 := testPeer() peer2 := newTestPeer()
peer2.newPeerAddr = addrChan // feed peer suggestions into matcher peer2.newPeerAddr = addrChan // feed peer suggestions into matcher
if err := runBaseProtocol(peer2, rw2); err != ErrPipeClosed { if err := runBaseProtocol(peer2, rw2); err != ErrPipeClosed {
t.Errorf("peer2 terminated with unexpected error: %v", err) t.Errorf("peer2 terminated with unexpected error: %v", err)