mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-22 12:46:44 +00:00
fix getPeerMsg/peerMsg RLP encode/decode, logs. tests pass
This commit is contained in:
parent
43877af5b5
commit
c2184512bb
3 changed files with 28 additions and 26 deletions
|
|
@ -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))
|
||||||
|
|
|
||||||
|
|
@ -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)
|
||||||
|
|
|
||||||
|
|
@ -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)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue