peer selection

- SuggestPeer lookup convenience method from backend to server
- now takes (host, pubkey) arguments
- earlier suggest peer -> AddPeer which now calls the peer selector
- delete peer_util (save, restore functionality belongs to peer selector)
- peerRecord implements peerInfo (a record entry object for peer info)
- introduce lastActive timestamping via the protocol pingpong loop
- instead of connect channel, peer simply takes 2 functions (addPeer, getPeers) this is inline with blockPool etc
- protocol.PeerList now many phases (server gives active peers including own address, selector proximity peers, and protocol filters for recipient peer and does encoding)
- fix and simplify peer_selector
This commit is contained in:
zelig 2015-01-06 01:56:26 +00:00
parent 1fb47697e1
commit bdb4f9dd83
10 changed files with 324 additions and 178 deletions

View file

@ -195,7 +195,7 @@ func (ui *UiLib) Connect(button qml.Object) {
} }
func (ui *UiLib) ConnectToPeer(addr string) { func (ui *UiLib) ConnectToPeer(addr string) {
if err := ui.eth.SuggestPeer(addr); err != nil { if err := ui.eth.SuggestPeer(addr, []byte{}); err != nil {
guilogger.Infoln(err) guilogger.Infoln(err)
} }
} }

View file

@ -2,7 +2,6 @@ package eth
import ( import (
"fmt" "fmt"
"net"
"sync" "sync"
"github.com/ethereum/go-ethereum/core" "github.com/ethereum/go-ethereum/core"
@ -21,6 +20,8 @@ const (
seedNodeAddress = "poc-7.ethdev.com:30300" seedNodeAddress = "poc-7.ethdev.com:30300"
) )
var seednodeId []byte = nil
type Config struct { type Config struct {
Name string Name string
Version string Version string
@ -248,7 +249,7 @@ func (s *Ethereum) Start(seed bool) error {
// TODO: read peers here // TODO: read peers here
if seed { if seed {
logger.Infof("Connect to seed node %v", seedNodeAddress) logger.Infof("Connect to seed node %v", seedNodeAddress)
if err := s.SuggestPeer(seedNodeAddress); err != nil { if err := s.SuggestPeer(seedNodeAddress, seednodeId); err != nil {
return err return err
} }
} }
@ -257,14 +258,8 @@ func (s *Ethereum) Start(seed bool) error {
return nil return nil
} }
func (self *Ethereum) SuggestPeer(addr string) error { func (self *Ethereum) SuggestPeer(addr string, pubkey []byte) error {
netaddr, err := net.ResolveTCPAddr("tcp", addr) self.net.SuggestPeer(addr, pubkey)
if err != nil {
logger.Errorf("couldn't resolve %s:", addr, err)
return err
}
self.net.SuggestPeer(netaddr.IP, netaddr.Port, nil)
return nil return nil
} }

View file

@ -1,23 +0,0 @@
package eth
import (
"encoding/json"
"github.com/ethereum/go-ethereum/ethutil"
)
func WritePeers(path string, addresses []string) {
if len(addresses) > 0 {
data, _ := json.MarshalIndent(addresses, "", " ")
ethutil.WriteFile(path, data)
}
}
func ReadPeers(path string) (ips []string, err error) {
var data string
data, err = ethutil.ReadAllFile(path)
if err != nil {
json.Unmarshal([]byte(data), &ips)
}
return
}

View file

@ -203,7 +203,7 @@ func (self *JSRE) addPeer(call otto.FunctionCall) otto.Value {
if err != nil { if err != nil {
return otto.FalseValue() return otto.FalseValue()
} }
self.ethereum.SuggestPeer(host) self.ethereum.SuggestPeer(host, nil)
return otto.TrueValue() return otto.TrueValue()
} }

View file

@ -1,3 +1 @@
Define peer address' time field and lastSeen function.
Overwrite time field when connect and disconnect.
Change protocol getPeers message response conditional on target. peerList if noTarget else server.getPeers Change protocol getPeers message response conditional on target. peerList if noTarget else server.getPeers

View file

@ -45,8 +45,48 @@ func (d peerAddr) String() string {
return fmt.Sprintf("%v:%d", d.IP, d.Port) return fmt.Sprintf("%v:%d", d.IP, d.Port)
} }
func (d *peerAddr) RlpData() interface{} { func (d peerAddr) RlpData() interface{} {
return []interface{}{string(d.IP), d.Port, d.Pubkey} return []interface{}{d.IP, d.Port, d.Pubkey}
}
type peerRecord struct {
addr *peerAddr
hash []byte
lastActive time.Time
lastActiveC chan time.Time
peer *Peer
}
func (self *peerRecord) Addr() *peerAddr {
return self.addr
}
func (self *peerRecord) Hash() []byte {
if self.hash == nil {
self.hash = Hash(self.addr.Pubkey)
}
return self.hash
}
func (self *peerRecord) LastActive() (lastActive time.Time) {
var ok bool
select {
case lastActive, ok = <-self.lastActiveC:
if ok {
self.lastActive = lastActive
}
default:
lastActive = self.lastActive
}
return
}
func (self *peerRecord) Connect() error {
return nil
}
func (self *peerRecord) Disconnect() error {
return nil
} }
// Peer represents a remote peer. // Peer represents a remote peer.
@ -85,11 +125,11 @@ type Peer struct {
// These fields are kept so base protocol can access them. // These fields are kept so base protocol can access them.
// TODO: this should be one or more interfaces // TODO: this should be one or more interfaces
ourID ClientIdentity // client id of the Server ourID ClientIdentity // client id of the Server
ourListenAddr *peerAddr // listen addr of Server, nil if not listening ourListenAddr *peerAddr // listen addr of Server, nil if not listening
newPeerAddr chan<- *peerAddr // tell server about received peers addPeer func(*peerAddr) error // tell server about received peers
otherPeers func() []*Peer // should return the list of all peers getPeers func(...[]byte) []*peerAddr // should return the list of all peers
pubkeyHook func(*peerAddr) error // called at end of handshake to validate pubkey pubkeyHook func(*peerAddr) error // called at end of handshake to validate pubkey
} }
// NewPeer returns a peer for testing purposes. // NewPeer returns a peer for testing purposes.
@ -104,8 +144,8 @@ func NewPeer(id ClientIdentity, caps []Cap) *Peer {
func newServerPeer(server *Server, conn net.Conn, dialAddr *peerAddr) *Peer { func newServerPeer(server *Server, conn net.Conn, dialAddr *peerAddr) *Peer {
p := newPeer(conn, server.Protocols, dialAddr) p := newPeer(conn, server.Protocols, dialAddr)
p.ourID = server.Identity p.ourID = server.Identity
p.newPeerAddr = server.peerConnect p.addPeer = server.AddPeer
p.otherPeers = server.Peers p.getPeers = server.GetPeers
p.pubkeyHook = server.verifyPeer p.pubkeyHook = server.verifyPeer
p.runBaseProtocol = true p.runBaseProtocol = true
@ -460,25 +500,3 @@ func (r *eofSignal) Read(buf []byte) (int, error) {
} }
return n, err return n, err
} }
func (peer *Peer) PeerList() []interface{} {
peers := peer.otherPeers()
ds := make([]interface{}, 0, len(peers))
for _, p := range peers {
p.infolock.Lock()
addr := p.listenAddr
p.infolock.Unlock()
// filter out this peer and peers that are not listening or
// have not completed the handshake.
// TODO: track previously sent peers and exclude them as well.
if p == peer || addr == nil {
continue
}
ds = append(ds, addr)
}
ourAddr := peer.ourListenAddr
if ourAddr != nil && !ourAddr.IP.IsLoopback() && !ourAddr.IP.IsUnspecified() {
ds = append(ds, ourAddr)
}
return ds
}

View file

@ -1,23 +1,85 @@
package p2p package p2p
import ( import (
"encoding/json"
"fmt"
"path"
"sync"
"time" "time"
"github.com/ethereum/go-ethereum/ethutil"
) )
type PeerSelector interface { type peerInfo interface {
SuggestPeer(addr *peerAddr) (ok bool) Addr() *peerAddr
// AddPeer(addr *peerAddr) (ok bool) Hash() []byte
GetPeers(target []byte) []*peerAddr // Pubkey() []byte
Start() LastActive() time.Time
Stop() Disconnect() error
Connect() error
}
type peerSelector interface {
AddPeer(peer peerInfo) error
GetPeers(target ...[]byte) []*peerAddr
Start() error
Stop() error
} }
type BaseSelector struct { type BaseSelector struct {
DirPath string DirPath string
getPeers func() []*peerAddr
peers []peerInfo
} }
func (self *BaseSelector) SuggestPeer(addr *peerAddr) bool { func (self *BaseSelector) AddPeer(peer peerInfo) error {
return true return nil
}
func (self *BaseSelector) GetPeers(target ...[]byte) []*peerAddr {
return self.getPeers()
}
func (self *BaseSelector) Start() error {
if len(self.DirPath) > 0 {
path := path.Join(self.DirPath, "peers.json")
peers, err := ReadPeers(path)
if err != nil {
return err
}
self.peers = peers
}
return nil
}
func (self *BaseSelector) Stop() error {
if len(self.DirPath) > 0 {
path := path.Join(self.DirPath, "peers.json")
if err := WritePeers(path, self.peers); err != nil {
return err
}
}
return nil
}
func WritePeers(path string, addresses []peerInfo) error {
if len(addresses) > 0 {
data, err := json.MarshalIndent(addresses, "", " ")
if err == nil {
ethutil.WriteFile(path, data)
}
return err
}
return nil
}
func ReadPeers(path string) (peers []peerInfo, err error) {
var data string
data, err = ethutil.ReadAllFile(path)
if err == nil {
json.Unmarshal([]byte(data), &peers)
}
return
} }
const ( const (
@ -27,36 +89,102 @@ const (
) )
type Cademlia struct { type Cademlia struct {
rows [hashBits]*row hash []byte
hashBits int hashBits int
rowLength int rowLength int
maxAge time.Duration // index map[string]peerInfo
index map[string]*peerData rows [hashBits]*row
depth int
maxAge time.Duration
purgeInterval time.Duration
lock sync.RWMutex
quitC chan bool
} }
type row struct { func newCademlia(hash []byte) *Cademlia {
length int return &Cademlia{
row []*peerData hash: hash,
lock sync.RWMutex hashBits: hashBits,
rowLength: rowLength,
maxAge: maxAge * time.Second,
rows: [hashBits]*row{},
// index: make(map[string]peerInfo),
}
} }
func (self *row) addresses() (addrs []*peerAddr) { func (self *Cademlia) Start() error {
self.lock.RLock() self.lock.Lock()
defer self.lock.RUnlock() defer self.lock.Unlock()
for _, p := range self.row { if self.quitC != nil {
addrs = append(addrs, p.addr) return nil
}
self.quitC = make(chan bool)
go self.purgeLoop()
return nil
}
func (self *Cademlia) Stop() {
self.lock.Lock()
defer self.lock.Unlock()
if self.quitC == nil {
return
}
close(self.quitC)
self.quitC = nil
}
func (self *Cademlia) AddPeer(peer peerInfo) (err error) {
index := self.commonPrefixLength(peer.Hash())
row := self.rows[index]
needed := row.insert(&entry{peer: peer})
if needed {
if index >= self.depth {
self.updateDepth()
}
} else {
err = fmt.Errorf("no worse peer found")
} }
return return
} }
func (self *row) insert(addr *peerAddr) (ok bool) { func (self *Cademlia) GetPeers(target []byte) (peers []*peerAddr) {
index := self.commonPrefixLength(target)
var entries []*entry
if index >= self.depth {
for i := self.depth; i < self.hashBits; i++ {
entries = append(entries, self.rows[i].row...)
}
} else {
entries = self.rows[index].row
}
for _, entry := range entries {
peers = append(peers, entry.peer.Addr())
}
return
}
type entry struct {
peer peerInfo
// metadata
}
type row struct {
length int
row []*entry
lock sync.RWMutex
}
func (self *row) insert(entry *entry) (ok bool) {
self.lock.Lock() self.lock.Lock()
defer self.lock.Unlock() defer self.lock.Unlock()
peerData := &peerData{addr: addr}
if len(self.row) >= self.length { if len(self.row) >= self.length {
self.row[self.worst()] = peerData self.row[self.worst()] = entry
} else { } else {
self.row = append(self.row, peerData) self.row = append(self.row, entry)
ok = true ok = true
} }
return return
@ -64,85 +192,47 @@ func (self *row) insert(addr *peerAddr) (ok bool) {
func (self *row) worst() (index int) { func (self *row) worst() (index int) {
var oldest time.Time var oldest time.Time
for i, p := range self.row { for i, entry := range self.row {
if oldest == nil || p.addr.LastSeen().Before(oldest) { if (oldest == time.Time{}) || entry.peer.LastActive().Before(oldest) {
oldest = p.addr.LastSeen oldest = entry.peer.LastActive()
index = i index = i
} }
} }
return return
} }
func (self *row) purge(maxAge time.Time) { func (self *row) purge(recently time.Time) {
var newRow []*peerData var newRow []*entry
for _, p := range self.row { for _, entry := range self.row {
if !p.addr.LastSeen().Before(maxAge) { if !entry.peer.LastActive().Before(recently) {
newRow = append(newRow, p) newRow = append(newRow, entry)
} else {
entry.peer.Disconnect()
} }
} }
} }
type peerData struct { func Hash(key []byte) []byte {
addr *peerAddr return key
hash []byte
} }
func Hash([]byte) []byte { func (self *Cademlia) updateDepth() {
}
func (self *Cademlia) prefixLength(other []byte) {
}
func newCademlia() *Cademlia {
return &Cademlia{
hashBits: hashBits,
rowLength: rowLength,
maxAge: maxAge * time.Second,
rows: make([hashBits]*row),
index: make(map[string]*peerData),
}
}
func (self *Cademlia) Start() {
go self.purgeLoop()
} }
func (self *Cademlia) purgeLoop() { func (self *Cademlia) purgeLoop() {
ticker := time.Tick(self.purgeInterval) ticker := time.Tick(self.purgeInterval)
for { for {
select { select {
case <-self.quitC:
return
case <-ticker: case <-ticker:
for _, r := range self.rows { for _, r := range self.rows {
r.purge(time.Since(self.maxAge)) r.purge(time.Now().Add(-self.maxAge))
} }
} }
} }
} }
func (self *Cademlia) SuggestPeer(addr *peerAddr) bool {
index := self.commonPrefixLength(Hash(addr.Pubkey))
row := self.rows[index]
longer := row.insert(addr)
if index >= self.depth && longer {
self.updateDepth()
}
return longer
}
func (self *Cademlia) GetPeers(target []byte) (peers []*peerAddr) {
index := self.prefixLength(target)
if index >= self.depth {
for i := self.depth; i < hashBits; i++ {
peers = append(peers, self.rows[i].addresses())
}
} else {
peers = self.rows[index].addresses()
}
return
}
func Xor(one, other []byte) (xor []byte) { func Xor(one, other []byte) (xor []byte) {
for i := 0; i < len(one); i++ { for i := 0; i < len(one); i++ {
xor[i] = one[i] ^ other[i] xor[i] = one[i] ^ other[i]
@ -152,12 +242,12 @@ func Xor(one, other []byte) (xor []byte) {
func (self *Cademlia) commonPrefixLength(other []byte) (ret int) { func (self *Cademlia) commonPrefixLength(other []byte) (ret int) {
xor := Xor(self.hash, other) xor := Xor(self.hash, other)
for i := 0; i < len(self.hash); i++ { for i := 0; i < self.hashBits; i++ {
for j := 0; j < 8; j++ { for j := 0; j < 8; j++ {
if (xor[i]>>uint8(7-j))&0x1 != 0 { if (xor[i]>>uint8(7-j))&0x1 != 0 {
return i*8 + j return i*8 + j
} }
} }
} }
return len(self.hash)*8 - 1 return self.hashBits*8 - 1
} }

View file

@ -105,12 +105,16 @@ func runBaseProtocol(peer *Peer, rw MsgReadWriter) error {
} }
} }
}() }()
return bp.loop(errc) var lastActiveC chan time.Time
if bp.peer.listenAddr != nil {
lastActiveC = bp.peer.listenAddr.lastActiveC
}
return bp.loop(errc, lastActiveC)
} }
var pingTimeout = 2 * time.Second var pingTimeout = 2 * time.Second
func (bp *baseProtocol) loop(quit <-chan error) error { func (bp *baseProtocol) loop(quit <-chan error, lastActiveC chan time.Time) error {
ping := time.NewTimer(pingTimeout) ping := time.NewTimer(pingTimeout)
activity := bp.peer.activity.Subscribe(time.Time{}) activity := bp.peer.activity.Subscribe(time.Time{})
lastActive := time.Time{} lastActive := time.Time{}
@ -125,6 +129,7 @@ func (bp *baseProtocol) loop(quit <-chan error) error {
select { select {
case err = <-quit: case err = <-quit:
return err return err
case lastActiveC <- lastActive:
case <-getPeersTick.C: case <-getPeersTick.C:
err = bp.rw.EncodeMsg(getPeersMsg) err = bp.rw.EncodeMsg(getPeersMsg)
case event := <-activity.Chan(): case event := <-activity.Chan():
@ -169,15 +174,29 @@ func (bp *baseProtocol) handle(rw MsgReadWriter) error {
case pongMsg: case pongMsg:
case getPeersMsg: case getPeersMsg:
peers := bp.peer.PeerList() var target [][]byte
// this is dangerous. the spec says that we should _delay_ if err := msg.Decode(&target); err != nil {
// sending the response if no new information is available. return newPeerError(errInvalidMsg, "%v", err)
// this means that would need to send a response later when }
// new peers become available.
// peers := bp.peer.getPeers(target...)
// TODO: add event mechanism to notify baseProtocol for new peers if len(target) == 0 {
if len(peers) > 0 { // then add ourselves to the list
return bp.rw.EncodeMsg(peersMsg, peers...) ourAddr := bp.peer.ourListenAddr
if ourAddr != nil && !ourAddr.IP.IsLoopback() && !ourAddr.IP.IsUnspecified() {
peers = append(peers, ourAddr)
}
}
ds := make([]interface{}, 0, len(peers))
// encode and filter out requesting peer
for _, addr := range peers {
if addr != bp.peer.listenAddr {
ds = append(ds, addr)
}
}
if len(ds) > 0 {
return bp.rw.EncodeMsg(peersMsg, ds...)
} }
case peersMsg: case peersMsg:
@ -187,7 +206,7 @@ func (bp *baseProtocol) handle(rw MsgReadWriter) error {
} }
for _, addr := range peers { for _, addr := range peers {
bp.peer.Debugf("received peer suggestion: %v", addr) bp.peer.Debugf("received peer suggestion: %v", addr)
bp.peer.newPeerAddr <- addr bp.peer.addPeer(addr)
} }
default: default:

View file

@ -31,7 +31,7 @@ func newTestPeer() (peer *Peer) {
peer.pubkeyHook = func(*peerAddr) error { return nil } peer.pubkeyHook = func(*peerAddr) error { return nil }
peer.ourID = &peerId{} peer.ourID = &peerId{}
peer.listenAddr = &peerAddr{} peer.listenAddr = &peerAddr{}
peer.otherPeers = func() []*Peer { return nil } peer.getPeers = func(...[]byte) []*peerAddr { return nil }
return return
} }
@ -58,25 +58,32 @@ func TestBaseProtocolPeers(t *testing.T) {
for got = range addrChan { for got = range addrChan {
own = append(own, got) own = append(own, got)
} }
if len(own) != 1 || !reflect.DeepEqual(ownAddr, own[0]) { if len(own) < 1 {
t.Errorf("mismatch: peers own address is incorrectly or not given, got %v, want %#v", ownAddr) t.Errorf("mismatch: peers own address not given")
} else {
if !reflect.DeepEqual(ownAddr, own[0]) {
t.Errorf("mismatch: peers own address is incorrectly or not given, got %v, want %#v", own[0], ownAddr)
}
} }
rw2.Close() rw2.Close()
}() }()
// run first peer // run first peer
peer1 := newTestPeer() peer1 := newTestPeer()
peer1.ourListenAddr = ownAddr peer1.ourListenAddr = ownAddr
peer1.otherPeers = func() []*Peer { peer1.getPeers = func(...[]byte) []*peerAddr {
pl := make([]*Peer, len(cannedPeerList)) pl := make([]*peerAddr, len(cannedPeerList))
for i, addr := range cannedPeerList { for i, addr := range cannedPeerList {
pl[i] = &Peer{listenAddr: addr} pl[i] = addr
} }
return pl return pl
} }
go runBaseProtocol(peer1, rw1) go runBaseProtocol(peer1, rw1)
// run second peer // run second peer
peer2 := newTestPeer() peer2 := newTestPeer()
peer2.newPeerAddr = addrChan // feed peer suggestions into matcher peer2.addPeer = func(addr *peerAddr) error {
addrChan <- addr // feed peer suggestions into matcher
return nil
}
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)
} }

View file

@ -63,7 +63,7 @@ type Server struct {
NoDial bool NoDial bool
// peer selector // peer selector
PeerSelector PeerSelector PeerSelector peerSelector
// Hook for testing. This is useful because we can inhibit // Hook for testing. This is useful because we can inhibit
// the whole protocol stack. // the whole protocol stack.
@ -107,8 +107,34 @@ func (srv *Server) Peers() (peers []*Peer) {
return return
} }
func (srv *Server) GetPeers(target []byte) (peers []*Peer) { // GetActivePeers returns addresses of all connected peers.
// delegate to selector func (srv *Server) GetActivePeers() (peers []*peerAddr) {
for _, peer := range srv.Peers() {
if peer != nil {
peer.infolock.Lock()
addr := peer.listenAddr
peer.infolock.Unlock()
// filter out this peer and peers that are not listening or
// have not completed the handshake.
// TODO: track previously sent peers and exclude them as well.
if addr == nil {
continue
}
peers = append(peers, addr)
}
}
return peers
}
// GetPeers returns addresses near target if given or supported by the client
// or falls back to all actively connected peers
func (srv *Server) GetPeers(target ...[]byte) (peers []*peerAddr) {
if len(target) == 1 { // delegate to selector
srv.PeerSelector.GetPeers(target[0])
return peers
} else {
return srv.GetActivePeers()
}
} }
// PeerCount returns the number of connected peers. // PeerCount returns the number of connected peers.
@ -118,14 +144,30 @@ func (srv *Server) PeerCount() int {
return srv.peerCount return srv.peerCount
} }
// SuggestPeer injects an address into the outbound address pool. // SuggestPeer is a convenient method that does dns resolution
func (srv *Server) SuggestPeer(ip net.IP, port int, nodeID []byte) { // and passes on the request to AddPeer
addr := &peerAddr{IP: ip, Port: uint64(port), Pubkey: nodeID} func (srv *Server) SuggestPeer(addr string, pubkey []byte) error {
select { netaddr, err := net.ResolveTCPAddr("tcp", addr)
case srv.peerConnect <- addr: if err != nil {
default: // don't block srvlog.Errorf("couldn't resolve %s:", addr, err)
srvlog.Warnf("peer suggestion %v ignored", addr) return err
} }
peerAddr := &peerAddr{IP: netaddr.IP, Port: uint64(netaddr.Port), Pubkey: pubkey}
return srv.AddPeer(peerAddr)
}
// AddPeer takes a peerAddr address as argument.
// If not found among connected peers turns to the peerSelector
// to decide if it is a worthwhile connection
func (srv *Server) AddPeer(addr *peerAddr) (err error) {
// need to look up nodeID first
peer := &peerRecord{addr: addr}
if err = srv.PeerSelector.AddPeer(peer); err == nil {
srvlog.Infof("peer %v accepted by peer selection", peer)
} else {
srvlog.Infof("peer %v rejected by peer selection", peer)
}
return
} }
// Broadcast sends an RLP-encoded message to all connected peers. // Broadcast sends an RLP-encoded message to all connected peers.