From dbc22690e5e16617342bb4679fc4bc72acb215d2 Mon Sep 17 00:00:00 2001 From: zelig Date: Sun, 21 Dec 2014 19:09:47 +0000 Subject: [PATCH] initial commit of p2p peer selection/cademlia --- p2p/TODO | 3 + p2p/peer_selector.go | 163 +++++++++++++++++++++++++++++++++++++++++++ p2p/server.go | 21 +++++- 3 files changed, 184 insertions(+), 3 deletions(-) create mode 100644 p2p/TODO create mode 100644 p2p/peer_selector.go diff --git a/p2p/TODO b/p2p/TODO new file mode 100644 index 0000000000..ad8ce170be --- /dev/null +++ b/p2p/TODO @@ -0,0 +1,3 @@ +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_selector.go b/p2p/peer_selector.go new file mode 100644 index 0000000000..e60b19eee8 --- /dev/null +++ b/p2p/peer_selector.go @@ -0,0 +1,163 @@ +package p2p + +import ( + "time" +) + +type PeerSelector interface { + SuggestPeer(addr *peerAddr) (ok bool) + // AddPeer(addr *peerAddr) (ok bool) + GetPeers(target []byte) []*peerAddr + Start() + Stop() +} + +type BaseSelector struct { + DirPath string +} + +func (self *BaseSelector) SuggestPeer(addr *peerAddr) bool { + return true +} + +const ( + hashBits = 160 + rowLength = 10 + maxAge = 1 +) + +type Cademlia struct { + rows [hashBits]*row + hashBits int + rowLength int + maxAge time.Duration + index map[string]*peerData +} + +type row struct { + length int + row []*peerData + lock sync.RWMutex +} + +func (self *row) addresses() (addrs []*peerAddr) { + self.lock.RLock() + defer self.lock.RUnlock() + for _, p := range self.row { + addrs = append(addrs, p.addr) + } + return +} + +func (self *row) insert(addr *peerAddr) (ok bool) { + self.lock.Lock() + defer self.lock.Unlock() + peerData := &peerData{addr: addr} + if len(self.row) >= self.length { + self.row[self.worst()] = peerData + } else { + self.row = append(self.row, peerData) + ok = true + } + return +} + +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 + 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) + } + } +} + +type peerData struct { + addr *peerAddr + hash []byte +} + +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) purgeLoop() { + ticker := time.Tick(self.purgeInterval) + for { + select { + case <-ticker: + for _, r := range self.rows { + r.purge(time.Since(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] + } + return +} + +func (self *Cademlia) commonPrefixLength(other []byte) (ret int) { + xor := Xor(self.hash, other) + for i := 0; i < len(self.hash); 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 +} diff --git a/p2p/server.go b/p2p/server.go index 3267812343..e91a83061c 100644 --- a/p2p/server.go +++ b/p2p/server.go @@ -62,6 +62,9 @@ type Server struct { // If NoDial is true, the server will not dial any peers. NoDial bool + // peer selector + PeerSelector PeerSelector + // Hook for testing. This is useful because we can inhibit // the whole protocol stack. newPeerFunc peerFunc @@ -104,6 +107,10 @@ func (srv *Server) Peers() (peers []*Peer) { return } +func (srv *Server) GetPeers(target []byte) (peers []*Peer) { + // delegate to selector +} + // PeerCount returns the number of connected peers. func (srv *Server) PeerCount() int { srv.lock.RLock() @@ -113,9 +120,13 @@ func (srv *Server) PeerCount() int { // SuggestPeer injects an address into the outbound address pool. func (srv *Server) SuggestPeer(ip net.IP, port int, nodeID []byte) { - select { - case srv.peerConnect <- &peerAddr{ip, uint64(port), nodeID}: - default: // don't block + addr := &peerAddr{ip, uint64(port), nodeID, time.Now()} + ok := srv.PeerSelector.SuggestPeer(addr) + if ok { + select { + case srv.peerConnect <- addr: + default: // don't block + } } } @@ -185,6 +196,10 @@ func (srv *Server) Start() (err error) { srvlog.Warnln("I will be kind-of useless, neither dialing nor listening.") } + if srv.PeerSelector == nil { + srv.PeerSelector = &BaseSelector{} + } + // make all slots available for i := range srv.peers { srv.peerSlots <- i