diff --git a/bzz/hive.go b/bzz/hive.go index d1408e6675..9c49b68185 100644 --- a/bzz/hive.go +++ b/bzz/hive.go @@ -3,6 +3,7 @@ package bzz import ( // "fmt" + "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common/kademlia" ) @@ -24,8 +25,10 @@ type peer struct { // to keep the nodetable uptodate type hive struct { + addr kademlia.Address kad *kademlia.Kademlia path string + ping chan bool } func newHive(hivepath string) *hive { @@ -35,32 +38,55 @@ func newHive(hivepath string) *hive { } } -func (self *hive) start(address kademlia.Address) (err error) { +func (self *hive) start(address kademlia.Address, connectPeer func(string) error) (err error) { + self.ping = make(chan bool) + self.addr = address self.kad.Start(address) err = self.kad.Load(self.path) if err != nil { dpaLogger.Warnf("Warning: error reading kademlia node db (skipping): %v", err) err = nil } - // go func() { - // for { - // select { - // case <-timer: - // case <-subscr: - // } - // maxpeers := 4 - // self.getPeerEntries(maxpeers) - // } - // }() + go func() { + for _ = range self.ping { + node, full := self.kad.GetNodeRecord() + if node != nil { + if len(node.Url) > 0 { + connectPeer(node.Url) + } else if !full { + // a random peer is taken + peers := self.kad.GetNodes(kademlia.RandomAddress(), 1) + if len(peers) > 0 { + req := &retrieveRequestMsgData{ + Key: Key(common.Hash(kademlia.RandomAddressAt(self.addr, 0)).Bytes()), + } + peers[0].(peer).retrieve(req) + } + } + } + } + }() return } +func (self *hive) stop() error { + close(self.ping) + return self.kad.Stop(self.path) +} + func (self *hive) addPeer(p peer) { self.kad.AddNode(p) + // self lookup + req := &retrieveRequestMsgData{ + Key: Key(common.Hash(self.addr).Bytes()), + } + p.retrieve(req) + self.ping <- true } func (self *hive) removePeer(p peer) { self.kad.RemoveNode(p) + self.ping <- false } // Retrieve a list of live peers that are closer to target than us @@ -91,14 +117,3 @@ func (self *hive) addPeerEntries(req *peersMsgData) { } self.kad.AddNodeRecords(nrs) } - -// called to ask periodically for preferences -// Kademlia ideally maintains a queue of prioritized nodes -func (self *hive) getPeerEntries(max int) (resp *peersMsgData, err error) { - nrs, err := self.kad.GetNodeRecords(max) - for _, n := range nrs { - _ = n - // resp // build response from kademlia noderecords - } - return -} diff --git a/bzz/netstore.go b/bzz/netstore.go index dcb38e14c1..6fddd75691 100644 --- a/bzz/netstore.go +++ b/bzz/netstore.go @@ -61,15 +61,19 @@ func NewNetStore(path, hivepath string) (netstore *NetStore, err error) { return } -func (self *NetStore) Start(node *discover.Node) (err error) { +func (self *NetStore) Start(node *discover.Node, connectPeer func(string) error) (err error) { self.self = node - err = self.hive.start(kademlia.Address(node.Sha())) + err = self.hive.start(kademlia.Address(node.Sha()), connectPeer) if err != nil { return } return } +func (self *NetStore) Stop() (err error) { + return self.hive.stop() +} + func (self *NetStore) Put(entry *Chunk) { chunk, err := self.localStore.Get(entry.Key) dpaLogger.Debugf("NetStore.Put: localStore.Get returned with %v.", err) diff --git a/common/kademlia/kademlia.go b/common/kademlia/kademlia.go index c845f25851..338e89f20a 100644 --- a/common/kademlia/kademlia.go +++ b/common/kademlia/kademlia.go @@ -1,15 +1,14 @@ package kademlia import ( - "fmt" - "sort" - // "math" "encoding/json" + "fmt" "io/ioutil" + "math/rand" "os" + "sort" "strings" "sync" - "sync/atomic" "time" "github.com/ethereum/go-ethereum/common" @@ -19,8 +18,9 @@ import ( var kadlogger = logger.NewLogger("KΛÐ") const ( - bucketSize = 20 - maxProx = 255 + minBucketSize = 1 + bucketSize = 20 + maxProx = 255 ) type Kademlia struct { @@ -28,12 +28,12 @@ type Kademlia struct { addr Address // adjustable parameters - MaxProx int - MaxProxBinSize int - BucketSize int - currentMaxBucketSize int - nodeDB [][]*NodeRecord - nodeIndex map[Address]*NodeRecord + MaxProx int + ProxBinSize int + BucketSize int + MinBucketSize int + nodeDB [][]*NodeRecord + nodeIndex map[Address]*NodeRecord GetNode func(int) @@ -108,9 +108,12 @@ func (self *Kademlia) Start(addr Address) error { if self.BucketSize == 0 { self.BucketSize = bucketSize } + if self.MinBucketSize == 0 { + self.MinBucketSize = minBucketSize + } // runtime parameters - if self.MaxProxBinSize == 0 { - self.MaxProxBinSize = self.BucketSize + if self.ProxBinSize == 0 { + self.ProxBinSize = self.BucketSize } self.buckets = make([]*bucket, self.MaxProx+1) @@ -159,7 +162,7 @@ func (self *Kademlia) RemoveNode(node Node) (err error) { if len(bucket.nodes) < bucket.size { err = fmt.Errorf("insufficient nodes (%v) in bucket %v", len(bucket.nodes), index) } - if len(bucket.nodes) == 0 { + if len(bucket.nodes) == 0 || index >= self.proxLimit { self.adjustProx(index, -1) } // async callback to notify user that bucket needs filling @@ -214,22 +217,24 @@ func (self *Kademlia) AddNode(node Node) (err error) { // adjust Prox (proxLimit and proxSize after an insertion of add nodes into bucket r) func (self *Kademlia) adjustProx(r int, add int) { + var i int switch { - case add > 0 && r == self.proxLimit: - self.proxLimit += add - for ; self.proxLimit < self.MaxProx && len(self.buckets[self.proxLimit].nodes) > 0; self.proxLimit++ { - self.proxSize -= len(self.buckets[self.proxLimit].nodes) - } - case add > 0 && r > self.proxLimit && self.proxSize+add > self.MaxProxBinSize: - self.proxLimit++ - self.proxSize -= len(self.buckets[r].nodes) - add - case add > 0 && r > self.proxLimit: + case add > 0 && r >= self.proxLimit: self.proxSize += add + for i = self.proxLimit; i < self.MaxProx && len(self.buckets[i].nodes) > 0 && self.proxSize > self.ProxBinSize; i++ { + self.proxSize -= len(self.buckets[i].nodes) + } + self.proxLimit = i case add < 0 && r < self.proxLimit && len(self.buckets[r].nodes) == 0: - for i := self.proxLimit - 1; i > r; i-- { + for i = self.proxLimit - 1; i > r; i-- { self.proxSize += len(self.buckets[i].nodes) } self.proxLimit = r + case add < 0 && self.proxLimit > 0 && r >= self.proxLimit-1: + for i = self.proxLimit - 1; len(self.buckets[i].nodes)+self.proxSize <= self.ProxBinSize; i-- { + self.proxSize += len(self.buckets[i].nodes) + } + self.proxLimit = i } } @@ -238,7 +243,7 @@ GetNodes(target) returns the list of nodes belonging to the same proximity bin as the target. The most proximate bin will be the union of the bins between proxLimit and MaxProx. proxLimit is dynamically adjusted so that 1) there is no empty buckets in bin < proxLimit and 2) the sum of all items are the maximum -possible but lower than MaxProxBinSize +possible but lower than ProxBinSize */ func (self *Kademlia) GetNodes(target Address, max int) []Node { return self.getNodes(target, max).nodes @@ -298,49 +303,58 @@ func (self *Kademlia) AddNodeRecords(nrs []*NodeRecord) { _, found := self.nodeIndex[node.Address] if !found { self.nodeIndex[node.Address] = node - index := self.proximityBin(node.Address) + index := proximity(self.addr, node.Address) self.nodeDB[index] = append(self.nodeDB[index], node) } } } /* -GetNodeRecords gives back an at most max length slice of node records -in order of decreasing priority for desired connection +GetNodeRecord gives back a node record with the highest priority for desired +connection Used to pick candidates for live nodes to satisfy Kademlia network for Swarm -Does a round robin on buckets starting from 0 to proxLimit then back -on each round i we inspect if live-nodes fill the bucket -if len(nodes) + i < currentMaxBucketSize, then take ith element in corresponding +if len(nodes) < MinBucketSize, then take ith element in corresponding db row ordered by reputation (active time?) +node record a is more favoured to b a > b iff +|proxBin(a)| < |proxBin(b)| +|| proxBin(a) < proxBin(b) && |proxBin(a)| < MinBucketSize +|| lastActive(a) < lastActive(b) This has double role. Starting as naive node with empty db, this implements Kademlia bootstrapping As a mature node, it manages quickly fill in blanks or short lines All on demand */ -func (self *Kademlia) GetNodeRecords(max int) (nrs []*NodeRecord, err error) { - var round int - for max > 0 { - for i, b := range self.buckets { - if len(b.nodes)+round < self.currentMaxBucketSize { - if nr := self.getNodeRecord(i, round); nr != nil { - nrs = append(nrs) +func (self *Kademlia) GetNodeRecord() (*NodeRecord, bool) { + full := true + for i, b := range self.nodeDB { + if i >= self.MaxProx { + break + } + if len(self.buckets[i].nodes) < self.MinBucketSize { + full = false + for _, node := range b { + if node.node == nil { + return node, full } } } - round++ - max-- } - return -} - -func (self *Kademlia) getNodeRecord(row, col int) (nr *NodeRecord) { - if row >= 0 && row < len(self.nodeDB) && - col >= 0 && col < len(self.nodeDB[row]) { - nr = self.nodeDB[row][col] + for i, b := range self.nodeDB { + if i > self.MaxProx { + break + } + if len(self.buckets[i].nodes) < self.BucketSize { + full = false + for _, node := range b { + if node.node == nil { + return node, full + } + } + } } - return + return nil, full } // in situ mutable bucket @@ -401,21 +415,19 @@ func (self *bucket) insert(node Node) (err error) { self.lock.Lock() defer self.lock.Unlock() if len(self.nodes) >= self.size { // >= allows us to add peers beyond the bucketsize limitation - worst := self.worstNode() - self.nodes[worst] = node - } else { - self.nodes = append(self.nodes, node) + self.worstNode().Drop() // assumes self.size > 0 } + self.nodes = append(self.nodes, node) return } // worst expunges the single worst entry in a row, where worst entry is with a peer that has not been active the longests -func (self *bucket) worstNode() (index int) { +func (self *bucket) worstNode() (node Node) { var oldest time.Time - for i, node := range self.nodes { + for _, n := range self.nodes { if (oldest == time.Time{}) || node.LastActive().Before(oldest) { - oldest = node.LastActive() - index = i + oldest = n.LastActive() + node = n } } return @@ -452,6 +464,9 @@ The distance metric MSB(x, y) of two equal length byte sequences x an y is the value of the binary integer cast of the xor-ed byte sequence (most significant bit first). proximity(x, y) counts the common zeros in the front of this distance measure. +which is equivalent to the reverse rank of the integer part of the base 2 +logarithm of the distance +called proximity belt (0 farthest, 255 closest, 256 self) */ func proximity(one, other Address) (ret int) { for i := 0; i < len(one); i++ { @@ -485,16 +500,6 @@ func (self *Kademlia) DB() [][]*NodeRecord { return self.nodeDB } -func (n *NodeRecord) bumpActive() { - stamp := time.Now().Unix() - atomic.StoreInt64(&n.Active, stamp) -} - -func (n *NodeRecord) LastActive() time.Time { - stamp := atomic.LoadInt64(&n.Active) - return time.Unix(stamp, 0) -} - // save persists all peers encountered func (self *Kademlia) Save(path string) error { @@ -528,9 +533,38 @@ func (self *Kademlia) Load(path string) (err error) { if err != nil { return } - self.nodeDB = kad.Nodes if self.addr != kad.Address { return fmt.Errorf("invalid kad db: address mismatch, expected %v, got %v", self.addr, kad.Address) } + self.nodeDB = kad.Nodes return } + +// randomAddressAt(address, prox) generates a random address +// at proximity order prox relative to address +// if prox is negative a random address is generated +func RandomAddressAt(self Address, prox int) (addr Address) { + addr = self + var pos int + if prox >= 0 { + pos = prox / 8 + trans := prox % 8 + transbytea := byte(0) + for j := 0; j <= trans; j++ { + transbytea |= 1 << uint8(7-j) + } + flipbyte := byte(1 << uint8(7-trans)) + transbyteb := transbytea ^ byte(255) + randbyte := byte(rand.Intn(255)) + addr[pos] = ((addr[pos] & transbytea) ^ flipbyte) | randbyte&transbyteb + } + for i := pos + 1; i < len(addr); i++ { + addr[i] = byte(rand.Intn(255)) + } + return +} + +// randomAddressAt() generates a random address +func RandomAddress() Address { + return RandomAddressAt(Address{}, -1) +} diff --git a/common/kademlia/kademlia_test.go b/common/kademlia/kademlia_test.go index c0e78f351d..7796d8ca7d 100644 --- a/common/kademlia/kademlia_test.go +++ b/common/kademlia/kademlia_test.go @@ -15,8 +15,9 @@ import ( ) var ( - quickrand = rand.New(rand.NewSource(time.Now().Unix())) - quickcfg = &quick.Config{MaxCount: 5000, Rand: quickrand} + quickrand = rand.New(rand.NewSource(time.Now().Unix())) + quickcfgGetNodes = &quick.Config{MaxCount: 5000, Rand: quickrand} + quickcfgBootStrap = &quick.Config{MaxCount: 1000, Rand: quickrand} ) var once sync.Once @@ -67,6 +68,82 @@ func TestAddNode(t *testing.T) { _ = err } +func TestBootstrap(t *testing.T) { + t.Parallel() + LogInit(logger.DebugLevel) + r := rand.New(rand.NewSource(time.Now().UnixNano())) + + test := func(test *bootstrapTest) bool { + // for any node kad.le, Target and N + kad := New() + kad.MaxProx = test.MaxProx + kad.MinBucketSize = test.MinBucketSize + kad.BucketSize = test.BucketSize + kad.Start(test.Self) + var err error + + t.Logf("bootstapTest MaxProx: %v MinBucketSize: %v BucketSize: %v\n", test.MaxProx, test.MinBucketSize, test.BucketSize) + + addr := gen(Address{}, r).(Address) + prox := proximity(addr, test.Self) + + for p := 0; p <= prox; p++ { + var nrs []*NodeRecord + for i := 0; i < test.BucketSize; i++ { + nrs = append(nrs, &NodeRecord{ + Address: RandomAddressAt(test.Self, p), + }) + } + kad.AddNodeRecords(nrs) + } + + node := &testNode{addr} + + n := 0 + for n < 100 { + err = kad.AddNode(node) + if err != nil { + t.Errorf("backend not accepting node") + return false + } + var nrs []*NodeRecord + prox := proximity(node.addr, test.Self) + for i := 0; i < test.BucketSize; i++ { + nrs = append(nrs, &NodeRecord{ + Address: RandomAddressAt(test.Self, prox+1), + }) + } + kad.AddNodeRecords(nrs) + + var lens []int + for i := 0; i <= test.MaxProx; i++ { + lens = append(lens, len(kad.buckets[i].nodes)) + } + + record, _ := kad.GetNodeRecord() + if record == nil { + t.Logf("after round %d, no more node records needed", n) + break + } + node = &testNode{record.Address} + n++ + } + exp := test.BucketSize * (test.MaxProx + 1) + if kad.Count() != exp { + t.Errorf("incorrect number of peers, expceted %d, got %d", exp, kad.Count()) + } + return true + if n < 1 { + t.Errorf("incorrect number of rounds, expceted %d, got %d", 0, n) + } + return true + } + if err := quick.Check(test, quickcfgBootStrap); err != nil { + t.Error(err) + } + +} + func TestGetNodes(t *testing.T) { t.Parallel() LogInit(logger.DebugLevel) @@ -131,7 +208,7 @@ func TestGetNodes(t *testing.T) { } return true } - if err := quick.Check(test, quickcfg); err != nil { + if err := quick.Check(test, quickcfgGetNodes); err != nil { t.Error(err) } } @@ -178,7 +255,7 @@ func TestProxAdjust(t *testing.T) { } return kad.proxCheck(t) } - if err := quick.Check(test, quickcfg); err != nil { + if err := quick.Check(test, quickcfgGetNodes); err != nil { t.Error(err) } } @@ -285,7 +362,7 @@ func (self *Kademlia) proxCheck(t *testing.T) bool { } // check if merged high prox bucket does not exceed size if sum > 0 { - if sum > self.MaxProxBinSize { + if sum > self.ProxBinSize { t.Errorf("bucket %d is empty, yet proxSize is %d", i, self.proxSize) return false } @@ -293,10 +370,32 @@ func (self *Kademlia) proxCheck(t *testing.T) bool { t.Errorf("proxSize incorrect, expected %v, got %v", sum, self.proxSize) return false } + if self.proxLimit > 0 && sum+len(self.buckets[self.proxLimit-1].nodes) < self.ProxBinSize { + t.Errorf("proxBinSize incorrect, expected %v got %v", sum, self.proxSize) + return false + } } return true } +type bootstrapTest struct { + MaxProx int + MinBucketSize int + BucketSize int + Self Address +} + +func (*bootstrapTest) Generate(rand *rand.Rand, size int) reflect.Value { + m := rand.Intn(2) + 1 + t := &bootstrapTest{ + Self: gen(Address{}, rand).(Address), + MaxProx: 10 + rand.Intn(3), + MinBucketSize: m, + BucketSize: rand.Intn(3) + m, + } + return reflect.ValueOf(t) +} + type getNodesTest struct { Self Address Target Address diff --git a/eth/backend.go b/eth/backend.go index 8e8bdb1314..050456da66 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -468,7 +468,7 @@ func (s *Ethereum) Start() error { if s.DPA != nil { s.DPA.Start() - s.netStore.Start(s.net.Self()) + s.netStore.Start(s.net.Self(), s.AddPeer) go bzz.StartHttpServer(s.DPA) } @@ -543,6 +543,13 @@ func (s *Ethereum) Stop() { s.whisper.Stop() } + if s.DPA != nil { + s.DPA.Stop() + } + if s.netStore != nil { + s.netStore.Stop() + } + glog.V(logger.Info).Infoln("Server stopped") close(s.shutdownChan) }