diff --git a/common/kademlia/kademlia.go b/common/kademlia/kademlia.go new file mode 100644 index 0000000000..784ea850ed --- /dev/null +++ b/common/kademlia/kademlia.go @@ -0,0 +1,801 @@ +package kademlia + +import ( + "encoding/json" + "fmt" + "io/ioutil" + "math/rand" + "os" + "sort" + "strings" + "sync" + "time" + + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/logger" + "github.com/ethereum/go-ethereum/logger/glog" +) + +const ( + bucketSize = 20 + maxProx = 255 + connRetryExp = 2 +) + +var ( + purgeInterval = 42 * time.Hour + initialRetryInterval = 42 * 100 * time.Millisecond +) + +type Kademlia struct { + // immutable baseparam + addr Address + + // adjustable parameters + MaxProx int + ProxBinSize int + BucketSize int + PurgeInterval time.Duration + InitialRetryInterval time.Duration + ConnRetryExp int + + nodeDB [][]*NodeRecord + nodeIndex map[Address]*NodeRecord + dbcursors []int + // state + proxLimit int + proxSize int + + // + count int + buckets []*bucket + + dblock sync.RWMutex + lock sync.RWMutex + quitC chan bool +} + +type Address common.Hash + +func (a Address) String() string { + return fmt.Sprintf("%x", a[:]) +} + +type Node interface { + Addr() Address + Url() string + LastActive() time.Time + Drop() +} + +// allow inactive peers under +type NodeRecord struct { + Addr Address `json:address` + Url string `json:url` + Active int64 `json:active` + After int64 `json:after` + after time.Time + checked time.Time + node Node +} + +func (self *NodeRecord) setActive() { + if self.node != nil { + self.Active = self.node.LastActive().Unix() + } +} + +func (self *NodeRecord) setChecked() { + self.checked = time.Now() +} + +// persisted node record database () +type kadDB struct { + Address Address `json:address` + Nodes [][]*NodeRecord `json:nodes` +} + +// public constructor +// hash is a byte slice of length equal to self.HashBytes +func New() *Kademlia { + return &Kademlia{} +} + +// accessor for KAD self address +func (self *Kademlia) Addr() Address { + return self.addr +} + +// accessor for KAD self count +// TODO: either memoize or lock +func (self *Kademlia) Count() (sum int) { + for _, b := range self.buckets { + sum += len(b.nodes) + } + return + // return self.count +} + +// accessor for KAD offline db count +func (self *Kademlia) DBCount() int { + return len(self.nodeIndex) +} + +// kademlia table + kaddb table displayed with ascii +func (self *Kademlia) String() string { + + var rows []string + // rows = append(rows, fmt.Sprintf("KΛÐΞMLIΛ basenode address: %064x\n population: %d (%d)", self.addr[:], self.Count(), self.DBCount())) + rows = append(rows, "=========================================================================") + rows = append(rows, fmt.Sprintf("%v : MaxProx: %d, ProxBinSize: %d, BucketSize: %d, proxLimit: %d, proxSize: %d", time.Now(), self.MaxProx, self.ProxBinSize, self.BucketSize, self.proxLimit, self.proxSize)) + + for i, b := range self.buckets { + + if i == self.proxLimit { + rows = append(rows, fmt.Sprintf("===================== PROX LIMIT: %d =================================", i)) + } + row := []string{fmt.Sprintf("%03d", i), fmt.Sprintf("%2d", len(b.nodes))} + var k int + c := self.dbcursors[i] + for ; k < len(b.nodes); k++ { + p := b.nodes[(c+k)%len(b.nodes)] + row = append(row, fmt.Sprintf("%s", p.Addr().String()[:8])) + if k == 3 { + break + } + } + for ; k < 3; k++ { + row = append(row, " ") + } + row = append(row, fmt.Sprintf("| %2d %2d", len(self.nodeDB[i]), self.dbcursors[i])) + + for j, p := range self.nodeDB[i] { + row = append(row, fmt.Sprintf("%08x", p.Addr[:4])) + if j == 2 { + break + } + } + rows = append(rows, strings.Join(row, " ")) + if i == self.MaxProx { + break + } + } + rows = append(rows, "=========================================================================") + return strings.Join(rows, "\n") +} + +// Start brings up a pool of entries potentially from an offline persisted source +// and sets default values for optional parameters +func (self *Kademlia) Start(addr Address) error { + self.lock.Lock() + defer self.lock.Unlock() + if self.quitC != nil { + return nil + } + self.addr = addr + if self.MaxProx == 0 { + self.MaxProx = maxProx + } + if self.BucketSize == 0 { + self.BucketSize = bucketSize + } + if self.InitialRetryInterval == 0 { + self.InitialRetryInterval = initialRetryInterval + } + if self.PurgeInterval == 0 { + self.PurgeInterval = purgeInterval + } + if self.ConnRetryExp == 0 { + self.ConnRetryExp = connRetryExp + } + // runtime parameters + if self.ProxBinSize == 0 { + self.ProxBinSize = self.BucketSize + } + + self.buckets = make([]*bucket, self.MaxProx+1) + for i, _ := range self.buckets { + self.buckets[i] = &bucket{size: self.BucketSize} // will initialise bucket{int(0),[]Node(nil),sync.Mutex} + } + + self.nodeDB = make([][]*NodeRecord, self.MaxProx+1) + self.dbcursors = make([]int, self.MaxProx+1) + self.nodeIndex = make(map[Address]*NodeRecord) + + self.quitC = make(chan bool) + glog.V(logger.Info).Infof("[KΛÐ] started") + return nil +} + +// Stop saves the routing table into a persistant form +func (self *Kademlia) Stop(path string) (err error) { + self.lock.Lock() + defer self.lock.Unlock() + if self.quitC == nil { + return + } + close(self.quitC) + self.quitC = nil + + if len(path) > 0 { + err = self.Save(path) + if err != nil { + glog.V(logger.Warn).Infof("[KΛÐ]: unable to save node records: %v", err) + } else { + glog.V(logger.Info).Infof("[KΛÐ]: node records saved to '%v'", path) + } + } + return +} + +// RemoveNode is the entrypoint where nodes are taken offline +func (self *Kademlia) RemoveNode(node Node) (err error) { + self.lock.Lock() + defer self.lock.Unlock() + var found bool + index := self.proximityBin(node.Addr()) + bucket := self.buckets[index] + for i := 0; i < len(bucket.nodes); i++ { + if node.Addr() == bucket.nodes[i].Addr() { + found = true + bucket.nodes = append(bucket.nodes[:i], bucket.nodes[(i+1):]...) + } + } + + if found { + glog.V(logger.Info).Infof("[KΛÐ]: remove node %v from table", node) + + self.count-- + if len(bucket.nodes) < bucket.size { + err = fmt.Errorf("insufficient nodes (%v) in bucket %v", len(bucket.nodes), index) + } + self.adjustProxLess(index) + + r := self.nodeIndex[node.Addr()] + r.node = nil + now := time.Now() + r.after = now + r.After = now.Unix() + r.Active = now.Unix() + + } + + return +} + +// AddNode is the entry point where new nodes are registered +func (self *Kademlia) AddNode(node Node) (err error) { + + self.lock.Lock() + defer self.lock.Unlock() + + index := self.proximityBin(node.Addr()) + + // insert in kademlia table of active nodes + bucket := self.buckets[index] + // if bucket is full insertion replaces the worst node + // TODO probably should give priority to peers with active traffic + if worst, pos := bucket.insert(node); worst != nil { + glog.V(logger.Info).Infof("[KΛÐ]: replace node %v (%d) with node %v", worst, pos, node) + // no prox adjustment needed + // do not change count + } else { + glog.V(logger.Info).Infof("[KΛÐ]: add new node %v to table", node) + self.count++ + self.adjustProxMore(index) + } + + // insert in kaddb, kademlia node record database + self.dblock.Lock() + defer self.dblock.Unlock() + record, found := self.nodeIndex[node.Addr()] + if !found { + glog.V(logger.Info).Infof("[KΛÐ]: add new record %v to node db", node) + record = &NodeRecord{ + Addr: node.Addr(), + Url: node.Url(), + } + self.nodeIndex[node.Addr()] = record + self.nodeDB[index] = append(self.nodeDB[index], record) + } + record.node = node + record.setActive() + record.setChecked() + + return + +} + +// 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 ProxBinSize +// adjust Prox (proxLimit and proxSize after an insertion of add nodes into bucket r) +func (self *Kademlia) adjustProxMore(r int) { + if r >= self.proxLimit { + exLimit := self.proxLimit + exSize := self.proxSize + self.proxSize++ + + var i int + for i = self.proxLimit; i < self.MaxProx && len(self.buckets[i].nodes) > 0 && self.proxSize-len(self.buckets[i].nodes) > self.ProxBinSize; i++ { + self.proxSize -= len(self.buckets[i].nodes) + } + self.proxLimit = i + + glog.V(logger.Detail).Infof("[KΛÐ]: Max Prox Bin: Lower Limit: %v (was %v): Bin Size: %v (was %v)", self.proxLimit, exLimit, self.proxSize, exSize) + } +} + +func (self *Kademlia) adjustProxLess(r int) { + exLimit := self.proxLimit + exSize := self.proxSize + if r >= self.proxLimit { + self.proxSize-- + } + + if r < self.proxLimit && len(self.buckets[r].nodes) == 0 { + for i := self.proxLimit - 1; i > r; i-- { + self.proxSize += len(self.buckets[i].nodes) + } + self.proxLimit = r + } else if self.proxLimit > 0 && r >= self.proxLimit-1 { + var i int + for i = self.proxLimit - 1; i > 0 && len(self.buckets[i].nodes)+self.proxSize <= self.ProxBinSize; i-- { + self.proxSize += len(self.buckets[i].nodes) + } + self.proxLimit = i + } + + if exLimit != self.proxLimit || exSize != self.proxSize { + glog.V(logger.Detail).Infof("[KΛÐ]: Max Prox Bin: Lower Limit: %v (was %v): Bin Size: %v (was %v)", self.proxLimit, exLimit, self.proxSize, exSize) + } +} + +/* +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. +*/ +func (self *Kademlia) GetNodes(target Address, max int) []Node { + return self.getNodes(target, max).nodes +} + +func (self *Kademlia) getNodes(target Address, max int) (r nodesByDistance) { + self.lock.RLock() + defer self.lock.RUnlock() + r.target = target + index := self.proximityBin(target) + start := index + var down bool + if index >= self.proxLimit { + index = self.proxLimit + start = self.MaxProx + down = true + } + var n int + limit := max + if max == 0 { + limit = 1000 + } + for { + bucket := self.buckets[start].nodes + for i := 0; i < len(bucket); i++ { + r.push(bucket[i], limit) + n++ + } + if max == 0 && start <= index && (n > 0 || start == 0) || + max > 0 && down && start <= index && (n >= limit || n == self.Count() || start == 0) { + break + } + if down { + start-- + } else { + if start == self.MaxProx { + if index == 0 { + break + } + start = index - 1 + down = true + } else { + start++ + } + } + } + glog.V(logger.Detail).Infof("[KΛÐ]: serve %d (=<%d) nodes for target lookup %v (PO%d)", n, self.MaxProx, target, index) + return +} + +// AddNodeRecords adds node records to kaddb (persisted node record db) +func (self *Kademlia) AddNodeRecords(nrs []*NodeRecord) { + self.dblock.Lock() + defer self.dblock.Unlock() + var n int + var nodes []*NodeRecord + for _, node := range nrs { + _, found := self.nodeIndex[node.Addr] + if !found && node.Addr != self.addr { + node.setChecked() + self.nodeIndex[node.Addr] = node + index := self.proximityBin(node.Addr) + dbcursor := self.dbcursors[index] + nodes = self.nodeDB[index] + newnodes := make([]*NodeRecord, len(nodes)+1) + copy(newnodes[:], nodes[:dbcursor]) + newnodes[dbcursor] = node + copy(newnodes[dbcursor+1:], nodes[dbcursor:]) + self.nodeDB[index] = newnodes + n++ + } + } + glog.V(logger.Detail).Infof("[KΛÐ]: received %d node records, added %d new", len(nrs), n) +} + +/* +GetNodeRecord return one node record with the highest priority for desired +connection. +This is used to pick candidates for live nodes that are most wanted for +a higly connected low centrality network structure for Swarm which best suits +for a Kademlia-style routing. + +The candidate is chosen using the following strategy. +We check for missing online nodes in the buckets for 1 upto Max BucketSize rounds. +On each round we proceed from the low to high proximity order buckets. +If the number of active nodes (=connected peers) is < rounds, then start looking +for a known candidate. To determine if there is a candidate to recommend the +node record database row corresponding to the bucket is checked. +If the row cursor is on position i, the ith element in the row is chosen. +If the record is scheduled not to be retried before NOW, the next element is taken. +If the record is scheduled can be retried, it is set as checked, scheduled for +checking and is returned. The time of the next check is in X (duration) such that +X = ConnRetryExp * delta where delta is the time past since the last check and +ConnRetryExp is constant obsoletion factor. (Note that when node records are added +from peer messages, they are marked as checked and placed at the cursor, ie. +given priority over older entries). Entries which were checked more than +purgeInterval ago are deleted from the kaddb row. If no candidate is found after +a full round of checking the next bucket up is considered. If no candidate is +found when we reach the maximum-proximity bucket, the next round starts. + +node record a is more favoured to b a > b iff a is a passive node (record of +offline past peer) +|proxBin(a)| < |proxBin(b)| +|| (proxBin(a) < proxBin(b) && |proxBin(a)| == |proxBin(b)|) +|| (proxBin(a) == proxBin(b) && lastChecked(a) < lastChecked(b)) + +This has double role. Starting as naive node with empty db, this implements +Kademlia bootstrapping +As a mature node, it fills short lines. All on demand. + +The second argument returned names the first missing slot found +*/ +func (self *Kademlia) GetNodeRecord() (node *NodeRecord, proxLimit int) { + // return value -1 indicates that buckets are filled in all + proxLimit = -1 + self.dblock.RLock() + defer self.dblock.RUnlock() + + for rounds := 1; rounds <= self.BucketSize; rounds++ { + ROUND: + for po, dbrow := range self.nodeDB { + if po > self.MaxProx { + break ROUND + } + bin := self.buckets[po] + bin.lock.Lock() + if len(bin.nodes) < rounds { + if proxLimit < 0 { + proxLimit = po + } + var count int + var purge []int + n := self.dbcursors[po] + + // try node records in the relavant kaddb row (of identical prox order) + // if they are ripe for checking + ROW: + for count < len(dbrow) { + node = dbrow[n] + if (node.after == time.Time{}) { + node.after = time.Unix(node.After, 0) + } + + glog.V(logger.Detail).Infof("[KΛÐ]: kaddb record %v (PO%03d:%d) not to be retried before %d %v", node.Addr, po, n, node.After, node.after) + + // time since last known connection attempt + delta := node.checked.Unix() - node.After + if delta < 4 { + node.After = 0 + } + + if node.node == nil && node.after.Before(time.Now()) { + if node.checked.Add(self.PurgeInterval).Before(time.Now()) { + // delete node + purge = append(purge, n) + glog.V(logger.Detail).Infof("[KΛÐ]: inactive node record %v (PO%03d:%d) last check: %v, next check: %v", node.Addr, po, n, node.checked, node.after) + } else { + // scheduling next check + if node.After == 0 { + node.after = time.Now().Add(self.InitialRetryInterval) + node.After = node.after.Unix() + } else { + node.After = delta*int64(self.ConnRetryExp) + node.After + node.after = time.Unix(node.After, 0) + } + + glog.V(logger.Detail).Infof("[KΛÐ]: serve node record %v (PO%03d:%d), last check: %v, next check: %v", node.Addr, po, n, node.checked, node.after) + } + break ROW + } + glog.V(logger.Detail).Infof("[KΛÐ]: kaddb record %v (PO%03d:%d) not ready. skipped. not to be retried before: %v", node.Addr, po, n, node.after) + n++ + count++ + // cycle: n = n % len(dbrow) + if n >= len(dbrow) { + n = 0 + } + } + self.dbcursors[po] = n + self.deleteNodeRecords(po, purge...) + if node != nil { + glog.V(logger.Detail).Infof("[KΛÐ]: rounds %d: prox limit: PO%03d\n%v", rounds, proxLimit, node) + node.setChecked() + bin.lock.Unlock() + return + } + } // if len < rounds + bin.lock.Unlock() + } // for po-s + glog.V(logger.Detail).Infof("[KΛÐ]: rounds %d: proxlimit: PO%03d", rounds, proxLimit) + if proxLimit == 0 || proxLimit < 0 && self.BucketSize == rounds { + return + } + } // for round + + return +} + +// deletes the noderecords of a kaddb row corresponding to the indexes +// caller must hold the dblock, +// the call is unsafe, no index checks +func (self *Kademlia) deleteNodeRecords(row int, indexes ...int) { + var prev int + var nodes []*NodeRecord + dbrow := self.nodeDB[row] + for _, next := range indexes { + // need to adjust dbcursor + if next > 0 { + if next <= self.dbcursors[row] { + self.dbcursors[row]-- + } + nodes = append(nodes, dbrow[prev:next]...) + } + prev = next + 1 + delete(self.nodeIndex, dbrow[next].Addr) + } + self.nodeDB[row] = append(nodes, dbrow[prev:]...) +} + +// in situ mutable bucket +type bucket struct { + size int + nodes []Node + lock sync.RWMutex +} + +// nodesByDistance is a list of nodes, ordered by distance to target. +type nodesByDistance struct { + nodes []Node + target Address +} + +func sortedByDistanceTo(target Address, slice []Node) bool { + var last Address + for i, node := range slice { + if i > 0 { + if target.ProxCmp(node.Addr(), last) < 0 { + return false + } + } + last = node.Addr() + } + return true +} + +// push(node, max) adds the given node to the list, keeping the total size +// below max elements. +func (h *nodesByDistance) push(node Node, max int) { + // returns the firt index ix such that func(i) returns true + ix := sort.Search(len(h.nodes), func(i int) bool { + return h.target.ProxCmp(h.nodes[i].Addr(), node.Addr()) >= 0 + }) + + if len(h.nodes) < max { + h.nodes = append(h.nodes, node) + } + if ix < len(h.nodes) { + copy(h.nodes[ix+1:], h.nodes[ix:]) + h.nodes[ix] = node + } +} + +// insert adds a peer to a bucket either by appending to existing items if +// bucket length does not exceed bucketSize, or by replacing the worst +// Node in the bucket +func (self *bucket) insert(node Node) (dropped Node, pos int) { + self.lock.Lock() + defer self.lock.Unlock() + if len(self.nodes) >= self.size { // >= allows us to add peers beyond the bucketsize limitation + dropped, pos = self.worstNode() + if dropped != nil { + self.nodes[pos] = node + glog.V(logger.Info).Infof("[KΛÐ] dropping node %v (%d)", dropped, pos) + dropped.Drop() + return + } + } + self.nodes = append(self.nodes, node) + return +} + +// worst expunges the single worst node in a row, where worst entry is the node +// that has been inactive for the longests time +func (self *bucket) worstNode() (node Node, pos int) { + var oldest time.Time + for p, n := range self.nodes { + if (oldest == time.Time{}) || !oldest.Before(n.LastActive()) { + oldest = n.LastActive() + node = n + pos = p + } + } + return +} + +/* +Taking the proximity order relative to a fix point x classifies the points in +the space (n byte long byte sequences) into bins. Items in each are at +most half as distant from x as items in the previous bin. Given a sample of +uniformly distributed items (a hash function over arbitrary sequence) the +proximity scale maps onto series of subsets with cardinalities on a negative +exponential scale. + +It also has the property that any two item belonging to the same bin are at +most half as distant from each other as they are from x. + +If we think of random sample of items in the bins as connections in a network of interconnected nodes than relative proximity can serve as the basis for local +decisions for graph traversal where the task is to find a route between two +points. Since in every hop, the finite distance halves, there is +a guaranteed constant maximum limit on the number of hops needed to reach one +node from the other. +*/ + +func (self *Kademlia) proximityBin(other Address) (ret int) { + ret = proximity(self.addr, other) + if ret > self.MaxProx { + ret = self.MaxProx + } + return +} + +/* +Proximity(x, y) returns the proximity order of the MSB distance between x and y + +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 x^y, ie., x and y bitwise xor-ed. +the binary cast is big endian: most significant bit first (=MSB). + +Proximity(x, y) is a discrete logarithmic scaling of the MSB distance. +It is defined as the reverse rank of the integer part of the base 2 +logarithm of the distance. +It is calculated by counting the number of common leading zeros in the (MSB) +binary representation of the x^y. + +(0 farthest, 255 closest, 256 self) +*/ +func proximity(one, other Address) (ret int) { + for i := 0; i < len(one); i++ { + oxo := one[i] ^ other[i] + for j := 0; j < 8; j++ { + if (uint8(oxo)>>uint8(7-j))&0x1 != 0 { + return i*8 + j + } + } + } + return len(one) * 8 +} + +// the string form of the binary representation of an address +func (a Address) Bin() string { + var bs []string + for _, b := range a[:] { + bs = append(bs, fmt.Sprintf("%08b", b)) + } + return strings.Join(bs, "") +} + +// Address.ProxCmp compares the distances a->target and b->target. +// Returns -1 if a is closer to target, 1 if b is closer to target +// and 0 if they are equal. +func (target Address) ProxCmp(a, b Address) int { + for i := range target { + da := a[i] ^ target[i] + db := b[i] ^ target[i] + if da > db { + return 1 + } else if da < db { + return -1 + } + } + return 0 +} + +// save persists kaddb on disk (written to file on path in json format. +// save is called by Kademlia.Stop() +func (self *Kademlia) Save(path string) error { + + kad := kadDB{ + Address: self.addr, + Nodes: self.nodeDB, + } + + for _, b := range kad.Nodes { + for _, node := range b { + node.setActive() + } + } + + data, err := json.MarshalIndent(&kad, "", " ") + if err != nil { + return err + } + return ioutil.WriteFile(path, data, os.ModePerm) +} + +// Load(path) loads the node record database (kaddb) from file on path. +// TODO: urls will be supported and handled with bzz-enabled dox clienta +func (self *Kademlia) Load(path string) (err error) { + var data []byte + data, err = ioutil.ReadFile(path) + if err != nil { + return + } + var kad kadDB + err = json.Unmarshal(data, &kad) + if err != nil { + return + } + 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 new file mode 100644 index 0000000000..cd9b5f3b82 --- /dev/null +++ b/common/kademlia/kademlia_test.go @@ -0,0 +1,408 @@ +package kademlia + +import ( + "fmt" + "math/rand" + "reflect" + "testing" + "testing/quick" + "time" +) + +var ( + quickrand = rand.New(rand.NewSource(time.Now().Unix())) + quickcfgGetNodes = &quick.Config{MaxCount: 5000, Rand: quickrand} + quickcfgBootStrap = &quick.Config{MaxCount: 1000, Rand: quickrand} +) + +type testNode struct { + addr Address +} + +func (n *testNode) String() string { + return fmt.Sprintf("%x", n.addr[:]) +} + +func (n *testNode) Addr() Address { + return n.addr +} + +func (n *testNode) Drop() { +} + +func (n *testNode) Url() string { + return "" +} + +func (n *testNode) LastActive() time.Time { + return time.Now() +} + +func (n *testNode) Add(a Address) (err error) { + return nil +} + +func TestAddNode(t *testing.T) { + addr, ok := gen(Address{}, quickrand).(Address) + other, ok := gen(Address{}, quickrand).(Address) + if !ok { + t.Errorf("oops") + } + kad := New() + kad.Start(addr) + err := kad.AddNode(&testNode{addr: other}) + _ = err +} + +func TestBootstrap(t *testing.T) { + t.Parallel() + 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.BucketSize = test.BucketSize + kad.Start(test.Self) + var err error + + // t.Logf("bootstapTest MaxProx: %v BucketSize: %v\n", test.MaxProx, 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{ + Addr: 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{ + Addr: 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.Addr} + n++ + } + exp := test.BucketSize * (test.MaxProx + 1) + if kad.Count() != exp { + t.Errorf("incorrect number of peers, expected %d, got %d\n%v", exp, kad.Count(), kad) + return false + } + return true + } + if err := quick.Check(test, quickcfgBootStrap); err != nil { + t.Error(err) + } + +} + +func TestGetNodes(t *testing.T) { + t.Parallel() + + test := func(test *getNodesTest) bool { + // for any node kad.le, Target and N + kad := New() + kad.MaxProx = 10 + kad.Start(test.Self) + var err error + // t.Logf("getNodesTest %v: %v\n", len(test.All), test) + for _, node := range test.All { + err = kad.AddNode(node) + if err != nil { + t.Errorf("backend not accepting node") + return false + } + } + + if len(test.All) == 0 || test.N == 0 { + return true + } + nodes := kad.GetNodes(test.Target, test.N) + + // check that the number of results is min(N, kad.len) + wantN := test.N + if tlen := kad.Count(); tlen < test.N { + wantN = tlen + } + + if len(nodes) != wantN { + t.Errorf("wrong number of nodes: got %d, want %d", len(nodes), wantN) + return false + } + + if hasDuplicates(nodes) { + t.Errorf("result contains duplicates") + return false + } + + if !sortedByDistanceTo(test.Target, nodes) { + t.Errorf("result is not sorted by distance to target") + return false + } + + // check that the result nodes have minimum distance to target. + farthestResult := nodes[len(nodes)-1].Addr() + for i, b := range kad.buckets { + for j, n := range b.nodes { + if contains(nodes, n.Addr()) { + continue // don't run the check below for nodes in result + } + if test.Target.ProxCmp(n.Addr(), farthestResult) < 0 { + _ = i * j + t.Errorf("kad.le contains node that is closer to target but it's not in result") + // t.Logf("bucket %v, item %v\n", i, j) + // t.Logf(" Target: %x", test.Target) + // t.Logf(" Farthest Result: %x", farthestResult) + // t.Logf(" ID: %x (%d)", n.Addr(), kad.proximityBin(n.Addr())) + return false + } + } + } + return true + } + if err := quick.Check(test, quickcfgGetNodes); err != nil { + t.Error(err) + } +} + +type proxTest struct { + add bool + index int + addr Address +} + +var ( + addresses []Address +) + +func TestProxAdjust(t *testing.T) { + t.Parallel() + r := rand.New(rand.NewSource(time.Now().UnixNano())) + self := gen(Address{}, r).(Address) + + kad := New() + kad.MaxProx = 10 + kad.Start(self) + var err error + for i := 0; i < 100; i++ { + a := gen(Address{}, r).(Address) + addresses = append(addresses, a) + err = kad.AddNode(&testNode{addr: a}) + if err != nil { + t.Errorf("backend not accepting node") + return + } + if !kad.proxCheck(t) { + return + } + } + + test := func(test *proxTest) bool { + node := &testNode{test.addr} + if test.add { + kad.AddNode(node) + } else { + kad.RemoveNode(node) + } + return kad.proxCheck(t) + } + if err := quick.Check(test, quickcfgGetNodes); err != nil { + t.Error(err) + } +} + +func TestSaveLoad(t *testing.T) { + r := rand.New(rand.NewSource(time.Now().UnixNano())) + addresses := gen([]Address{}, r).([]Address) + self := addresses[0] + kad := New() + kad.MaxProx = 10 + kad.Start(self) + var err error + for _, a := range addresses[1:] { + err = kad.AddNode(&testNode{addr: a}) + if err != nil { + t.Errorf("backend not accepting node") + return + } + } + nodes := kad.GetNodes(self, 100) + path := "/tmp/bzz.peers" + kad.Stop(path) + kad = New() + kad.Start(self) + kad.Load(path) + for _, b := range kad.nodeDB { + for _, node := range b { + node.node = &testNode{node.Addr} + err = kad.AddNode(node.node) + if err != nil { + t.Errorf("backend not accepting node") + return + } + } + } + loadednodes := kad.GetNodes(self, 100) + for i, node := range loadednodes { + if nodes[i].Addr() != node.Addr() { + t.Errorf("node mismatch at %d/%d", i, len(nodes)) + } + } +} + +func (self *Kademlia) proxCheck(t *testing.T) bool { + var sum, i int + var b *bucket + for i, b = range self.buckets { + l := len(b.nodes) + // if we are in the high prox multibucket + if i >= self.proxLimit { + sum += l + } else if l == 0 { + t.Errorf("bucket %d empty, yet proxLimit is %d\n%v", len(b.nodes), self.proxLimit, self) + return false + } + } + // check if merged high prox bucket does not exceed size + if sum > 0 { + // if sum > self.ProxBinSize { + // t.Errorf("bucket %d is empty, yet proxSize is %d\n%v", i, self.proxSize, self) + // return false + // } + if sum != self.proxSize { + 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 + BucketSize int + Self Address +} + +func (*bootstrapTest) Generate(rand *rand.Rand, size int) reflect.Value { + t := &bootstrapTest{ + Self: gen(Address{}, rand).(Address), + MaxProx: 10 + rand.Intn(3), + BucketSize: rand.Intn(3) + 1, + } + return reflect.ValueOf(t) +} + +type getNodesTest struct { + Self Address + Target Address + All []Node + N int +} + +func (c getNodesTest) String() string { + return fmt.Sprintf("A: %064x\nT: %064x\n(%d)\n", c.Self[:], c.Target[:], c.N) +} + +func (*getNodesTest) Generate(rand *rand.Rand, size int) reflect.Value { + t := &getNodesTest{ + Self: gen(Address{}, rand).(Address), + Target: gen(Address{}, rand).(Address), + N: rand.Intn(bucketSize), + } + for _, a := range gen([]Address{}, rand).([]Address) { + t.All = append(t.All, &testNode{addr: a}) + } + return reflect.ValueOf(t) +} + +func (*proxTest) Generate(rand *rand.Rand, size int) reflect.Value { + var add bool + if rand.Intn(1) == 0 { + add = true + } + var t *proxTest + if add { + t = &proxTest{ + addr: gen(Address{}, rand).(Address), + add: add, + } + } else { + t = &proxTest{ + index: rand.Intn(len(addresses)), + add: add, + } + } + return reflect.ValueOf(t) +} + +func hasDuplicates(slice []Node) bool { + seen := make(map[Address]bool) + for _, node := range slice { + if seen[node.Addr()] { + return true + } + seen[node.Addr()] = true + } + return false +} + +func contains(nodes []Node, addr Address) bool { + for _, n := range nodes { + if n.Addr() == addr { + return true + } + } + return false +} + +// gen wraps quick.Value so it's easier to use. +// it generates a random value of the given value's type. +func gen(typ interface{}, rand *rand.Rand) interface{} { + v, ok := quick.Value(reflect.TypeOf(typ), rand) + if !ok { + panic(fmt.Sprintf("couldn't generate random value of type %T", typ)) + } + return v.Interface() +} + +func (Address) Generate(rand *rand.Rand, size int) reflect.Value { + var id Address + // m := rand.Intn(len(id)) + for i := 0; i < len(id); i++ { + id[i] = byte(rand.Uint32()) + } + return reflect.ValueOf(id) +}