From bdb4f9dd83668edb5e5c93f9f2c5b3d50fbc2709 Mon Sep 17 00:00:00 2001 From: zelig Date: Tue, 6 Jan 2015 01:56:26 +0000 Subject: [PATCH] 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 --- cmd/mist/ui_lib.go | 2 +- eth/backend.go | 15 +- eth/peer_util.go | 23 --- javascript/javascript_runtime.go | 2 +- p2p/TODO | 2 - p2p/peer.go | 80 ++++++---- p2p/peer_selector.go | 252 +++++++++++++++++++++---------- p2p/protocol.go | 43 ++++-- p2p/protocol_test.go | 21 ++- p2p/server.go | 62 ++++++-- 10 files changed, 324 insertions(+), 178 deletions(-) delete mode 100644 eth/peer_util.go diff --git a/cmd/mist/ui_lib.go b/cmd/mist/ui_lib.go index 0aabb87d09..843c168de3 100644 --- a/cmd/mist/ui_lib.go +++ b/cmd/mist/ui_lib.go @@ -195,7 +195,7 @@ func (ui *UiLib) Connect(button qml.Object) { } 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) } } diff --git a/eth/backend.go b/eth/backend.go index 065a4f7d84..7f397eef41 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -2,7 +2,6 @@ package eth import ( "fmt" - "net" "sync" "github.com/ethereum/go-ethereum/core" @@ -21,6 +20,8 @@ const ( seedNodeAddress = "poc-7.ethdev.com:30300" ) +var seednodeId []byte = nil + type Config struct { Name string Version string @@ -248,7 +249,7 @@ func (s *Ethereum) Start(seed bool) error { // TODO: read peers here if seed { 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 } } @@ -257,14 +258,8 @@ func (s *Ethereum) Start(seed bool) error { return nil } -func (self *Ethereum) SuggestPeer(addr string) error { - netaddr, err := net.ResolveTCPAddr("tcp", addr) - if err != nil { - logger.Errorf("couldn't resolve %s:", addr, err) - return err - } - - self.net.SuggestPeer(netaddr.IP, netaddr.Port, nil) +func (self *Ethereum) SuggestPeer(addr string, pubkey []byte) error { + self.net.SuggestPeer(addr, pubkey) return nil } diff --git a/eth/peer_util.go b/eth/peer_util.go deleted file mode 100644 index 6cf80cde29..0000000000 --- a/eth/peer_util.go +++ /dev/null @@ -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 -} diff --git a/javascript/javascript_runtime.go b/javascript/javascript_runtime.go index af1405049f..d0ccc9b2b2 100644 --- a/javascript/javascript_runtime.go +++ b/javascript/javascript_runtime.go @@ -203,7 +203,7 @@ func (self *JSRE) addPeer(call otto.FunctionCall) otto.Value { if err != nil { return otto.FalseValue() } - self.ethereum.SuggestPeer(host) + self.ethereum.SuggestPeer(host, nil) return otto.TrueValue() } diff --git a/p2p/TODO b/p2p/TODO index ad8ce170be..cd5891b766 100644 --- a/p2p/TODO +++ b/p2p/TODO @@ -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 diff --git a/p2p/peer.go b/p2p/peer.go index 0d7eec9f46..de21f60dd9 100644 --- a/p2p/peer.go +++ b/p2p/peer.go @@ -45,8 +45,48 @@ func (d peerAddr) String() string { return fmt.Sprintf("%v:%d", d.IP, d.Port) } -func (d *peerAddr) RlpData() interface{} { - return []interface{}{string(d.IP), d.Port, d.Pubkey} +func (d peerAddr) RlpData() interface{} { + 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. @@ -85,11 +125,11 @@ type Peer struct { // These fields are kept so base protocol can access them. // TODO: this should be one or more interfaces - ourID ClientIdentity // client id of the Server - ourListenAddr *peerAddr // listen addr of Server, nil if not listening - newPeerAddr chan<- *peerAddr // tell server about received peers - otherPeers func() []*Peer // should return the list of all peers - pubkeyHook func(*peerAddr) error // called at end of handshake to validate pubkey + ourID ClientIdentity // client id of the Server + ourListenAddr *peerAddr // listen addr of Server, nil if not listening + addPeer func(*peerAddr) error // tell server about received peers + getPeers func(...[]byte) []*peerAddr // should return the list of all peers + pubkeyHook func(*peerAddr) error // called at end of handshake to validate pubkey } // 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 { p := newPeer(conn, server.Protocols, dialAddr) p.ourID = server.Identity - p.newPeerAddr = server.peerConnect - p.otherPeers = server.Peers + p.addPeer = server.AddPeer + p.getPeers = server.GetPeers p.pubkeyHook = server.verifyPeer p.runBaseProtocol = true @@ -460,25 +500,3 @@ func (r *eofSignal) Read(buf []byte) (int, error) { } 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 -} diff --git a/p2p/peer_selector.go b/p2p/peer_selector.go index e60b19eee8..81c644ad40 100644 --- a/p2p/peer_selector.go +++ b/p2p/peer_selector.go @@ -1,23 +1,85 @@ package p2p import ( + "encoding/json" + "fmt" + "path" + "sync" "time" + + "github.com/ethereum/go-ethereum/ethutil" ) -type PeerSelector interface { - SuggestPeer(addr *peerAddr) (ok bool) - // AddPeer(addr *peerAddr) (ok bool) - GetPeers(target []byte) []*peerAddr - Start() - Stop() +type peerInfo interface { + Addr() *peerAddr + Hash() []byte + // Pubkey() []byte + LastActive() time.Time + Disconnect() error + Connect() error +} + +type peerSelector interface { + AddPeer(peer peerInfo) error + GetPeers(target ...[]byte) []*peerAddr + Start() error + Stop() error } type BaseSelector struct { - DirPath string + DirPath string + getPeers func() []*peerAddr + peers []peerInfo } -func (self *BaseSelector) SuggestPeer(addr *peerAddr) bool { - return true +func (self *BaseSelector) AddPeer(peer peerInfo) error { + 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 ( @@ -27,36 +89,102 @@ const ( ) type Cademlia struct { - rows [hashBits]*row + hash []byte hashBits int rowLength int - maxAge time.Duration - index map[string]*peerData + // index map[string]peerInfo + rows [hashBits]*row + + depth int + + maxAge time.Duration + purgeInterval time.Duration + + lock sync.RWMutex + quitC chan bool } -type row struct { - length int - row []*peerData - lock sync.RWMutex +func newCademlia(hash []byte) *Cademlia { + return &Cademlia{ + hash: hash, + hashBits: hashBits, + rowLength: rowLength, + maxAge: maxAge * time.Second, + rows: [hashBits]*row{}, + // index: make(map[string]peerInfo), + } } -func (self *row) addresses() (addrs []*peerAddr) { - self.lock.RLock() - defer self.lock.RUnlock() - for _, p := range self.row { - addrs = append(addrs, p.addr) +func (self *Cademlia) Start() error { + self.lock.Lock() + defer self.lock.Unlock() + if self.quitC != nil { + 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 } -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() defer self.lock.Unlock() - peerData := &peerData{addr: addr} if len(self.row) >= self.length { - self.row[self.worst()] = peerData + self.row[self.worst()] = entry } else { - self.row = append(self.row, peerData) + self.row = append(self.row, entry) ok = true } return @@ -64,85 +192,47 @@ func (self *row) insert(addr *peerAddr) (ok bool) { func (self *row) worst() (index int) { var oldest time.Time - for i, p := range self.row { - if oldest == nil || p.addr.LastSeen().Before(oldest) { - oldest = p.addr.LastSeen + for i, entry := range self.row { + if (oldest == time.Time{}) || entry.peer.LastActive().Before(oldest) { + oldest = entry.peer.LastActive() index = i } } return } -func (self *row) purge(maxAge time.Time) { - var newRow []*peerData - for _, p := range self.row { - if !p.addr.LastSeen().Before(maxAge) { - newRow = append(newRow, p) +func (self *row) purge(recently time.Time) { + var newRow []*entry + for _, entry := range self.row { + if !entry.peer.LastActive().Before(recently) { + newRow = append(newRow, entry) + } else { + entry.peer.Disconnect() } } } -type peerData struct { - addr *peerAddr - hash []byte +func Hash(key []byte) []byte { + return key } -func Hash([]byte) []byte { - -} - -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) updateDepth() { } func (self *Cademlia) purgeLoop() { ticker := time.Tick(self.purgeInterval) for { select { + case <-self.quitC: + return case <-ticker: 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) { for i := 0; i < len(one); 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) { 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++ { if (xor[i]>>uint8(7-j))&0x1 != 0 { return i*8 + j } } } - return len(self.hash)*8 - 1 + return self.hashBits*8 - 1 } diff --git a/p2p/protocol.go b/p2p/protocol.go index dd8cbc4ecd..356a89413a 100644 --- a/p2p/protocol.go +++ b/p2p/protocol.go @@ -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 -func (bp *baseProtocol) loop(quit <-chan error) error { +func (bp *baseProtocol) loop(quit <-chan error, lastActiveC chan time.Time) error { ping := time.NewTimer(pingTimeout) activity := bp.peer.activity.Subscribe(time.Time{}) lastActive := time.Time{} @@ -125,6 +129,7 @@ func (bp *baseProtocol) loop(quit <-chan error) error { select { case err = <-quit: return err + case lastActiveC <- lastActive: case <-getPeersTick.C: err = bp.rw.EncodeMsg(getPeersMsg) case event := <-activity.Chan(): @@ -169,15 +174,29 @@ func (bp *baseProtocol) handle(rw MsgReadWriter) error { case pongMsg: case getPeersMsg: - peers := bp.peer.PeerList() - // this is dangerous. the spec says that we should _delay_ - // sending the response if no new information is available. - // this means that would need to send a response later when - // new peers become available. - // - // TODO: add event mechanism to notify baseProtocol for new peers - if len(peers) > 0 { - return bp.rw.EncodeMsg(peersMsg, peers...) + var target [][]byte + if err := msg.Decode(&target); err != nil { + return newPeerError(errInvalidMsg, "%v", err) + } + + peers := bp.peer.getPeers(target...) + if len(target) == 0 { + // then add ourselves to the list + 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: @@ -187,7 +206,7 @@ func (bp *baseProtocol) handle(rw MsgReadWriter) error { } for _, addr := range peers { bp.peer.Debugf("received peer suggestion: %v", addr) - bp.peer.newPeerAddr <- addr + bp.peer.addPeer(addr) } default: diff --git a/p2p/protocol_test.go b/p2p/protocol_test.go index ce25b3e1b5..38a3094957 100644 --- a/p2p/protocol_test.go +++ b/p2p/protocol_test.go @@ -31,7 +31,7 @@ func newTestPeer() (peer *Peer) { peer.pubkeyHook = func(*peerAddr) error { return nil } peer.ourID = &peerId{} peer.listenAddr = &peerAddr{} - peer.otherPeers = func() []*Peer { return nil } + peer.getPeers = func(...[]byte) []*peerAddr { return nil } return } @@ -58,25 +58,32 @@ func TestBaseProtocolPeers(t *testing.T) { for got = range addrChan { own = append(own, got) } - 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) + if len(own) < 1 { + 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() }() // run first peer peer1 := newTestPeer() peer1.ourListenAddr = ownAddr - peer1.otherPeers = func() []*Peer { - pl := make([]*Peer, len(cannedPeerList)) + peer1.getPeers = func(...[]byte) []*peerAddr { + pl := make([]*peerAddr, len(cannedPeerList)) for i, addr := range cannedPeerList { - pl[i] = &Peer{listenAddr: addr} + pl[i] = addr } return pl } go runBaseProtocol(peer1, rw1) // run second peer 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 { t.Errorf("peer2 terminated with unexpected error: %v", err) } diff --git a/p2p/server.go b/p2p/server.go index 441eec6711..dff4738ad5 100644 --- a/p2p/server.go +++ b/p2p/server.go @@ -63,7 +63,7 @@ type Server struct { NoDial bool // peer selector - PeerSelector PeerSelector + PeerSelector peerSelector // Hook for testing. This is useful because we can inhibit // the whole protocol stack. @@ -107,8 +107,34 @@ func (srv *Server) Peers() (peers []*Peer) { return } -func (srv *Server) GetPeers(target []byte) (peers []*Peer) { - // delegate to selector +// GetActivePeers returns addresses of all connected peers. +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. @@ -118,14 +144,30 @@ func (srv *Server) PeerCount() int { return srv.peerCount } -// SuggestPeer injects an address into the outbound address pool. -func (srv *Server) SuggestPeer(ip net.IP, port int, nodeID []byte) { - addr := &peerAddr{IP: ip, Port: uint64(port), Pubkey: nodeID} - select { - case srv.peerConnect <- addr: - default: // don't block - srvlog.Warnf("peer suggestion %v ignored", addr) +// SuggestPeer is a convenient method that does dns resolution +// and passes on the request to AddPeer +func (srv *Server) SuggestPeer(addr string, pubkey []byte) error { + netaddr, err := net.ResolveTCPAddr("tcp", addr) + if err != nil { + srvlog.Errorf("couldn't resolve %s:", addr, err) + 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.