diff --git a/p2p/simulations/network.go b/p2p/simulations/network.go
index 706c77dab2..752103de7d 100644
--- a/p2p/simulations/network.go
+++ b/p2p/simulations/network.go
@@ -255,7 +255,7 @@ func (self *Network) newConn(oneID, otherID discover.NodeID) (*Conn, error) {
if other == nil {
return nil, fmt.Errorf("other %v does not exist", other)
}
- distance, _ := pot.NewBytesVal(one.Addr(), nil).PO(pot.NewBytesVal(other.Addr(), nil), 0)
+ distance, _ := pot.DefaultPof(256)(one.Addr(), other.Addr(), 0)
return &Conn{
One: oneID,
Other: otherID,
diff --git a/pot/address.go b/pot/address.go
index 05f2f82167..8a396dac3b 100644
--- a/pot/address.go
+++ b/pot/address.go
@@ -13,6 +13,8 @@
//
// You should have received a copy of the GNU Lesser General Public License
// along with the go-ethereum library. If not, see .
+
+// Package pot see doc.go
package pot
import (
@@ -26,35 +28,40 @@ import (
)
var (
- zeroAddr = &common.Hash{}
- zerosHex = zeroAddr.Hex()[2:]
zerosBin = Address{}.Bin()
)
-var (
- addrlen = keylen
-)
-
+// Address is an alias for common.Hash
type Address common.Hash
+// NewAddressFromBytes constructs an Address from a byte slice
+func NewAddressFromBytes(b []byte) Address {
+ h := common.Hash{}
+ copy(h[:], b)
+ return Address(h)
+}
+
func (a Address) String() string {
return fmt.Sprintf("%x", a[:])
}
+// MarshalJSON Address serialisation
func (a *Address) MarshalJSON() (out []byte, err error) {
return []byte(`"` + a.String() + `"`), nil
}
+// UnmarshalJSON Address deserialisation
func (a *Address) UnmarshalJSON(value []byte) error {
*a = Address(common.HexToHash(string(value[1 : len(value)-1])))
return nil
}
-// the string form of the binary representation of an address (only first 8 bits)
+// Bin returns the string form of the binary representation of an address (only first 8 bits)
func (a Address) Bin() string {
return ToBin(a[:])
}
+// ToBin converts a byteslice to the string binary representation
func ToBin(a []byte) string {
var bs []string
for _, b := range a {
@@ -63,6 +70,7 @@ func ToBin(a []byte) string {
return strings.Join(bs, "")
}
+// Bytes returns the Address as a byte slice
func (a Address) Bytes() []byte {
return a[:]
}
@@ -107,23 +115,23 @@ func posProximity(one, other Address, pos int) (ret int, eq bool) {
return len(one) * 8, true
}
-// Address.ProxCmp compares the distances a->target and b->target.
+// 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 {
+func (a Address) ProxCmp(x, y Address) int {
+ for i := range a {
+ dx := x[i] ^ a[i]
+ dy := y[i] ^ a[i]
+ if dx > dy {
return 1
- } else if da < db {
+ } else if dx < dy {
return -1
}
}
return 0
}
-// randomAddressAt(address, prox) generates a random address
+// 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) {
@@ -148,71 +156,12 @@ func RandomAddressAt(self Address, prox int) (addr Address) {
return
}
-// KeyRange(a0, a1, proxLimit) returns the address inclusive address
-// range that contain addresses closer to one than other
-// func KeyRange(one, other Address, proxLimit int) (start, stop Address) {
-// prox := proximity(one, other)
-// if prox >= proxLimit {
-// prox = proxLimit
-// }
-// start = CommonBitsAddrByte(one, other, byte(0x00), prox)
-// stop = CommonBitsAddrByte(one, other, byte(0xff), prox)
-// return
-// }
-
-func CommonBitsAddrF(self, other Address, f func() byte, p int) (addr Address) {
- prox, _ := proximity(self, other)
- var pos int
- if p <= prox {
- prox = p
- }
- pos = prox / 8
- addr = self
- trans := byte(prox % 8)
- var transbytea byte
- if p > prox {
- transbytea = byte(0x7f)
- } else {
- transbytea = byte(0xff)
- }
- transbytea >>= trans
- transbyteb := transbytea ^ byte(0xff)
- addrpos := addr[pos]
- addrpos &= transbyteb
- if p > prox {
- addrpos ^= byte(0x80 >> trans)
- }
- addrpos |= transbytea & f()
- addr[pos] = addrpos
- for i := pos + 1; i < len(addr); i++ {
- addr[i] = f()
- }
-
- return
-}
-
-func CommonBitsAddr(self, other Address, prox int) (addr Address) {
- return CommonBitsAddrF(self, other, func() byte { return byte(rand.Intn(255)) }, prox)
-}
-
-func CommonBitsAddrByte(self, other Address, b byte, prox int) (addr Address) {
- return CommonBitsAddrF(self, other, func() byte { return b }, prox)
-}
-
-// randomAddressAt() generates a random address
+// RandomAddress generates a random address
func RandomAddress() Address {
return RandomAddressAt(Address{}, -1)
}
-// wraps an Address to implement the PotVal interface
-type HashAddress struct {
- Address
-}
-
-func (a *HashAddress) String() string {
- return a.Address.Bin()
-}
-
+// NewAddressFromString creates a byte slice from a string in binary representation
func NewAddressFromString(s string) []byte {
ha := [32]byte{}
@@ -227,85 +176,16 @@ func NewAddressFromString(s string) []byte {
return ha[:]
}
-func NewHashAddress(s string) *HashAddress {
- ha := NewAddressFromString(s)
- h := common.Hash{}
- copy(h[:], ha)
- return &HashAddress{Address(h)}
-}
-
-func NewHashAddressFromBytes(b []byte) *HashAddress {
- h := common.Hash{}
- copy(h[:], b)
- return &HashAddress{Address(h)}
-}
-
-// PO(addr, pos) return the proximity order of addr wrt to
-// the pinned address of the tree
-// assuming it is greater than or equal to pos
-func (self *HashAddress) PO(val PotVal, pos int) (po int, eq bool) {
- return posProximity(self.Address, val.(*HashAddress).Address, pos)
-}
-
-type BoolAddress struct {
- addr []bool
-}
-
-func NewBoolAddress(s string) *BoolAddress {
- return NewBoolAddressXOR(s, zerosBin[:len(s)])
-}
-
-func NewBoolAddressXOR(s, t string) *BoolAddress {
- if len(s) != len(t) {
- panic("lengths do not match")
- }
- addr := make([]bool, len(s))
- for i, _ := range addr {
- addr[i] = s[i] != t[i]
- }
- return &BoolAddress{addr}
-}
-
-func (self *BoolAddress) String() string {
- a := self.addr
- s := []byte(zerosBin)[:len(a)]
- for i, one := range a {
- if one {
- s[i] = byte('1')
- }
- }
- return string(s)
-}
-
-func (self *BoolAddress) PO(val PotVal, pos int) (po int, eq bool) {
- a := self.addr
- b := val.(*BoolAddress).addr
- for po = pos; po < len(b); po++ {
- if a[po] != b[po] {
- return po, false
- }
- }
- return po, true
-}
-
+// BytesAddress is an interface for elements addressable by a byte slice
type BytesAddress interface {
Address() []byte
}
-type bytesAddress struct {
- bytes []byte
- toBytes func(v AnyVal) []byte
-}
-
-func NewBytesVal(v AnyVal, f func(v AnyVal) []byte) *bytesAddress {
- if f == nil {
- f = ToBytes
+// ToBytes turns the Val into bytes
+func ToBytes(v Val) []byte {
+ if v == nil {
+ return nil
}
- b := f(v)
- return &bytesAddress{b, f}
-}
-
-func ToBytes(v AnyVal) []byte {
b, ok := v.([]byte)
if !ok {
ba, ok := v.(BytesAddress)
@@ -317,15 +197,17 @@ func ToBytes(v AnyVal) []byte {
return b
}
-func (a *bytesAddress) String() string {
- return fmt.Sprintf("%08b", a.bytes)
-}
-func (a *bytesAddress) Address() []byte {
- return a.bytes
-}
-
-func (a *bytesAddress) PO(val PotVal, i int) (int, bool) {
- return proximityOrder(a.bytes, a.toBytes(val), i)
+// DefaultPof returns a proximity order operator function
+// where all
+func DefaultPof(max int) func(one, other Val, pos int) (int, bool) {
+ return func(one, other Val, pos int) (int, bool) {
+ po, eq := proximityOrder(ToBytes(one), ToBytes(other), pos)
+ if po >= max {
+ eq = true
+ po = max
+ }
+ return po, eq
+ }
}
func proximityOrder(one, other []byte, pos int) (int, bool) {
@@ -346,3 +228,17 @@ func proximityOrder(one, other []byte, pos int) (int, bool) {
}
return len(one) * 8, true
}
+
+// Label displays the node's key in binary format
+func Label(v Val) string {
+ if v == nil {
+ return ""
+ }
+ if s, ok := v.(fmt.Stringer); ok {
+ return s.String()
+ }
+ if b, ok := v.([]byte); ok {
+ return ToBin(b)
+ }
+ panic(fmt.Sprintf("unsupported value type %T", v))
+}
diff --git a/pot/address_test.go b/pot/address_test.go
index b590201f41..b9b3f1f4e3 100644
--- a/pot/address_test.go
+++ b/pot/address_test.go
@@ -15,127 +15,128 @@
// along with the go-ethereum library. If not, see .
package pot
-import (
- "math/rand"
- "reflect"
- "testing"
-
- "github.com/ethereum/go-ethereum/common"
-)
-
-func (Address) Generate(rand *rand.Rand, size int) reflect.Value {
- var id Address
- for i := 0; i < len(id); i++ {
- id[i] = byte(uint8(rand.Intn(255)))
- }
- return reflect.ValueOf(id)
-}
-
-func TestCommonBitsAddrF(t *testing.T) {
- a := Address(common.HexToHash("0x0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
- b := Address(common.HexToHash("0x8123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
- c := Address(common.HexToHash("0x4123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
- d := Address(common.HexToHash("0x0023456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
- e := Address(common.HexToHash("0x01A3456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
- ab := CommonBitsAddrF(a, b, func() byte { return byte(0x00) }, 10)
- expab := Address(common.HexToHash("0x8000000000000000000000000000000000000000000000000000000000000000"))
-
- if ab != expab {
- t.Fatalf("%v != %v", ab, expab)
- }
- ac := CommonBitsAddrF(a, c, func() byte { return byte(0x00) }, 10)
- expac := Address(common.HexToHash("0x4000000000000000000000000000000000000000000000000000000000000000"))
-
- if ac != expac {
- t.Fatalf("%v != %v", ac, expac)
- }
- ad := CommonBitsAddrF(a, d, func() byte { return byte(0x00) }, 10)
- expad := Address(common.HexToHash("0x0000000000000000000000000000000000000000000000000000000000000000"))
-
- if ad != expad {
- t.Fatalf("%v != %v", ad, expad)
- }
- ae := CommonBitsAddrF(a, e, func() byte { return byte(0x00) }, 10)
- expae := Address(common.HexToHash("0x0180000000000000000000000000000000000000000000000000000000000000"))
-
- if ae != expae {
- t.Fatalf("%v != %v", ae, expae)
- }
- acf := CommonBitsAddrF(a, c, func() byte { return byte(0xff) }, 10)
- expacf := Address(common.HexToHash("0x7fffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff"))
-
- if acf != expacf {
- t.Fatalf("%v != %v", acf, expacf)
- }
- aeo := CommonBitsAddrF(a, e, func() byte { return byte(0x00) }, 2)
- expaeo := Address(common.HexToHash("0x0000000000000000000000000000000000000000000000000000000000000000"))
-
- if aeo != expaeo {
- t.Fatalf("%v != %v", aeo, expaeo)
- }
- aep := CommonBitsAddrF(a, e, func() byte { return byte(0xff) }, 2)
- expaep := Address(common.HexToHash("0x3fffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff"))
-
- if aep != expaep {
- t.Fatalf("%v != %v", aep, expaep)
- }
-
-}
-
-func TestRandomAddressAt(t *testing.T) {
- var a Address
- for i := 0; i < 100; i++ {
- a = RandomAddress()
- prox := rand.Intn(255)
- b := RandomAddressAt(a, prox)
- p, _ := proximity(a, b)
- if p != prox {
- t.Fatalf("incorrect address prox(%v, %v) == %v (expected %v)", a, b, p, prox)
- }
- }
-}
-
-const (
- maxTestPOs = 1000
- testPOkeylen = 9
-)
-
-func TestPOs(t *testing.T) {
- for i := 0; i < maxTestPOs; i++ {
- length := rand.Intn(256) + 1
- v0 := RandomAddress().Bin()[:length]
- v1 := RandomAddress().Bin()[:length]
- a0 := NewBoolAddress(v0)
- a1 := NewBoolAddress(v1)
- b0 := NewHashAddress(v0)
- b1 := NewHashAddress(v1)
- pos := rand.Intn(length) + 1
- apo, aeq := a0.PO(a1, pos)
- bpo, beq := b0.PO(b1, pos)
- if bpo == 256 {
- bpo = length
- }
- a0s := a0.String()
- if a0s != v0 {
- t.Fatalf("incorrect bool address. expected %v, got %v", v0, a0s)
- }
- a1s := a1.String()
- if a1s != v1 {
- t.Fatalf("incorrect bool address. expected %v, got %v", v1, a1s)
- }
- b0s := b0.String()[:length]
- if b0s != v0 {
- t.Fatalf("incorrect hash address. expected %v, got %v", v0, b0s)
- }
- b1s := b1.String()[:length]
- if b1s != v1 {
- t.Fatalf("incorrect hash address. expected %v, got %v", v1, b1s)
- }
- if apo != bpo {
- t.Fatalf("PO does not match for %v X %v (pos: %v): expected %v, got %v", v0, v1, pos, apo, bpo)
- }
- if aeq != beq {
- t.Fatalf("PO equality does not match for %v X %v (pos: %v): expected %v, got %v", v0, v1, pos, aeq, beq)
- }
- }
-}
+//
+// import (
+// "math/rand"
+// "reflect"
+// "testing"
+//
+// "github.com/ethereum/go-ethereum/common"
+// )
+// //
+// // func (Address) Generate(rand *rand.Rand, size int) reflect.Value {
+// // var id Address
+// // for i := 0; i < len(id); i++ {
+// // id[i] = byte(uint8(rand.Intn(255)))
+// // }
+// // return reflect.ValueOf(id)
+// // }
+//
+// func TestCommonBitsAddrF(t *testing.T) {
+// a := Address(common.HexToHash("0x0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
+// b := Address(common.HexToHash("0x8123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
+// c := Address(common.HexToHash("0x4123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
+// d := Address(common.HexToHash("0x0023456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
+// e := Address(common.HexToHash("0x01A3456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"))
+// ab := CommonBitsAddrF(a, b, func() byte { return byte(0x00) }, 10)
+// expab := Address(common.HexToHash("0x8000000000000000000000000000000000000000000000000000000000000000"))
+//
+// if ab != expab {
+// t.Fatalf("%v != %v", ab, expab)
+// }
+// ac := CommonBitsAddrF(a, c, func() byte { return byte(0x00) }, 10)
+// expac := Address(common.HexToHash("0x4000000000000000000000000000000000000000000000000000000000000000"))
+//
+// if ac != expac {
+// t.Fatalf("%v != %v", ac, expac)
+// }
+// ad := CommonBitsAddrF(a, d, func() byte { return byte(0x00) }, 10)
+// expad := Address(common.HexToHash("0x0000000000000000000000000000000000000000000000000000000000000000"))
+//
+// if ad != expad {
+// t.Fatalf("%v != %v", ad, expad)
+// }
+// ae := CommonBitsAddrF(a, e, func() byte { return byte(0x00) }, 10)
+// expae := Address(common.HexToHash("0x0180000000000000000000000000000000000000000000000000000000000000"))
+//
+// if ae != expae {
+// t.Fatalf("%v != %v", ae, expae)
+// }
+// acf := CommonBitsAddrF(a, c, func() byte { return byte(0xff) }, 10)
+// expacf := Address(common.HexToHash("0x7fffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff"))
+//
+// if acf != expacf {
+// t.Fatalf("%v != %v", acf, expacf)
+// }
+// aeo := CommonBitsAddrF(a, e, func() byte { return byte(0x00) }, 2)
+// expaeo := Address(common.HexToHash("0x0000000000000000000000000000000000000000000000000000000000000000"))
+//
+// if aeo != expaeo {
+// t.Fatalf("%v != %v", aeo, expaeo)
+// }
+// aep := CommonBitsAddrF(a, e, func() byte { return byte(0xff) }, 2)
+// expaep := Address(common.HexToHash("0x3fffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff"))
+//
+// if aep != expaep {
+// t.Fatalf("%v != %v", aep, expaep)
+// }
+//
+// }
+//
+// func TestRandomAddressAt(t *testing.T) {
+// var a Address
+// for i := 0; i < 100; i++ {
+// a = RandomAddress()
+// prox := rand.Intn(255)
+// b := RandomAddressAt(a, prox)
+// p, _ := proximity(a, b)
+// if p != prox {
+// t.Fatalf("incorrect address prox(%v, %v) == %v (expected %v)", a, b, p, prox)
+// }
+// }
+// }
+//
+// const (
+// maxTestPOs = 1000
+// testPOkeylen = 9
+// )
+//
+// func TestPOs(t *testing.T) {
+// for i := 0; i < maxTestPOs; i++ {
+// length := rand.Intn(256) + 1
+// v0 := RandomAddress().Bin()[:length]
+// v1 := RandomAddress().Bin()[:length]
+// a0 := NewBoolAddress(v0)
+// a1 := NewBoolAddress(v1)
+// b0 := NewHashAddress(v0)
+// b1 := NewHashAddress(v1)
+// pos := rand.Intn(length) + 1
+// apo, aeq := a0.PO(a1, pos)
+// bpo, beq := b0.PO(b1, pos)
+// if bpo == 256 {
+// bpo = length
+// }
+// a0s := a0.String()
+// if a0s != v0 {
+// t.Fatalf("incorrect bool address. expected %v, got %v", v0, a0s)
+// }
+// a1s := a1.String()
+// if a1s != v1 {
+// t.Fatalf("incorrect bool address. expected %v, got %v", v1, a1s)
+// }
+// b0s := b0.String()[:length]
+// if b0s != v0 {
+// t.Fatalf("incorrect hash address. expected %v, got %v", v0, b0s)
+// }
+// b1s := b1.String()[:length]
+// if b1s != v1 {
+// t.Fatalf("incorrect hash address. expected %v, got %v", v1, b1s)
+// }
+// if apo != bpo {
+// t.Fatalf("PO does not match for %v X %v (pos: %v): expected %v, got %v", v0, v1, pos, apo, bpo)
+// }
+// if aeq != beq {
+// t.Fatalf("PO equality does not match for %v X %v (pos: %v): expected %v, got %v", v0, v1, pos, aeq, beq)
+// }
+// }
+// }
diff --git a/pot/doc.go b/pot/doc.go
index 484590175f..5f01a91c17 100644
--- a/pot/doc.go
+++ b/pot/doc.go
@@ -1,5 +1,21 @@
+// Copyright 2017 The go-ethereum Authors
+// This file is part of the go-ethereum library.
+//
+// The go-ethereum library is free software: you can redistribute it and/or modify
+// it under the terms of the GNU Lesser General Public License as published by
+// the Free Software Foundation, either version 3 of the License, or
+// (at your option) any later version.
+//
+// The go-ethereum library is distributed in the hope that it will be useful,
+// but WITHOUT ANY WARRANTY; without even the implied warranty of
+// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+// GNU Lesser General Public License for more details.
+//
+// You should have received a copy of the GNU Lesser General Public License
+// along with the go-ethereum library. If not, see .
+
/*
-POT: proximity order tree implements a container similar to a binary tree.
+Package pot (proximity order tree) implements a container similar to a binary tree.
Value types implement the PoVal interface which provides the PO (proximity order)
comparison operator.
Each fork in the trie is itself a value. Values of the subtree contained under
diff --git a/pot/pot.go b/pot/pot.go
index 6d7b15fce9..2ebce82cec 100644
--- a/pot/pot.go
+++ b/pot/pot.go
@@ -13,6 +13,8 @@
//
// You should have received a copy of the GNU Lesser General Public License
// along with the go-ethereum library. If not, see .
+
+// Package pot see doc.go
package pot
import (
@@ -33,66 +35,69 @@ type Pot struct {
// pot is the node type (same for root, branching node and leaf)
type pot struct {
- pin PotVal
+ pin Val
bins []*pot
size int
po int
+ pof Pof
}
-// PotVal is the interface the generic container item should implement
-type PotVal interface {
- PO(PotVal, int) (po int, eq bool)
- String() string
-}
+// Val is the element type for pots
+type Val interface{}
-type AnyVal interface{}
+// Pof is the proximity order function
+type Pof func(Val, Val, int) (int, bool)
-// Pot constructor. Requires value of type PotVal to pin
-// and po to point to a span in the PotVal key
+// NewPot constructor. Requires value of type Val to pin
+// and po to point to a span in the Val key
// The pinned item counts towards the size
-func NewPot(v PotVal, po int) *Pot {
+func NewPot(v Val, po int, pof Pof) *Pot {
var size int
if v != nil {
size++
}
+ if pof == nil {
+ pof = DefaultPof(keylen)
+ }
return &Pot{
pot: &pot{
pin: v,
po: po,
size: size,
+ pof: pof,
},
}
}
-// Pin() returns the pinned element (key) of the Pot
-func (t *Pot) Pin() PotVal {
+// Pin returns the pinned element (key) of the Pot
+func (t *Pot) Pin() Val {
return t.pin
}
-// Size() returns the number of values in the Pot
+// Size returns the number of values in the Pot
func (t *Pot) Size() int {
t.lock.RLock()
defer t.lock.RUnlock()
return t.size
}
-// Add(v) inserts v into the Pot and
+// Add inserts v into the Pot and
// returns the proximity order of v and a boolean
// indicating if the item was found
// Add locks the Pot while using applicative add on its pot
-func (t *Pot) Add(val PotVal) (po int, found bool) {
+func (t *Pot) Add(val Val) (po int, found bool) {
t.lock.Lock()
defer t.lock.Unlock()
t.pot, po, found = add(t.pot, val)
return po, found
}
-// Add(t, v) returns a new Pot that contains all the elements of t
+// Add called on (t, v) returns a new Pot that contains all the elements of t
// plus the value v, using the applicative add
// the second return value is the proximity order of the inserted element
// the third is boolean indicating if the item was found
// it only readlocks the Pot while reading its pot
-func Add(t *Pot, val PotVal) (*Pot, int, bool) {
+func Add(t *Pot, val Val) (*Pot, int, bool) {
t.lock.RLock()
n := t.pot
t.lock.RUnlock()
@@ -100,25 +105,28 @@ func Add(t *Pot, val PotVal) (*Pot, int, bool) {
return &Pot{pot: r}, po, found
}
-func add(t *pot, val PotVal) (*pot, int, bool) {
+func (t *pot) clone() *pot {
+ return &pot{
+ pin: t.pin,
+ size: t.size,
+ po: t.po,
+ bins: t.bins,
+ pof: t.pof,
+ }
+}
+
+func add(t *pot, val Val) (*pot, int, bool) {
var r *pot
if t == nil || t.pin == nil {
- r = &pot{
- pin: val,
- size: t.size + 1,
- po: t.po,
- bins: t.bins,
- }
+ r = t.clone()
+ r.pin = val
+ r.size++
return r, 0, false
}
- po, found := t.pin.PO(val, t.po)
+ po, found := t.pof(t.pin, val, t.po)
if found {
- r = &pot{
- pin: val,
- size: t.size,
- po: t.po,
- bins: t.bins,
- }
+ r = t.clone()
+ r.pin = val
return r, po, true
}
@@ -147,6 +155,7 @@ func add(t *pot, val PotVal) (*pot, int, bool) {
pin: val,
size: 1,
po: po,
+ pof: t.pof,
}
}
@@ -158,28 +167,29 @@ func add(t *pot, val PotVal) (*pot, int, bool) {
size: size,
po: t.po,
bins: bins,
+ pof: t.pof,
}
return r, po, found
}
-// T.Re move(v) deletes v from the Pot and returns
+// Remove called on (v) deletes v from the Pot and returns
// the proximity order of v and a boolean value indicating
// if the value was found
// Remove locks Pot while using applicative remove on its pot
-func (t *Pot) Remove(val PotVal) (po int, found bool) {
+func (t *Pot) Remove(val Val) (po int, found bool) {
t.lock.Lock()
defer t.lock.Unlock()
t.pot, po, found = remove(t.pot, val)
return po, found
}
-// Remove(t, v) returns a new Pot that contains all the elements of t
+// Remove called on (t, v) returns a new Pot that contains all the elements of t
// minus the value v, using the applicative remove
// the second return value is the proximity order of the inserted element
// the third is boolean indicating if the item was found
// it only readlocks the Pot while reading its pot
-func Remove(t *Pot, v PotVal) (*Pot, int, bool) {
+func Remove(t *Pot, v Val) (*Pot, int, bool) {
t.lock.RLock()
n := t.pot
t.lock.RUnlock()
@@ -187,14 +197,15 @@ func Remove(t *Pot, v PotVal) (*Pot, int, bool) {
return &Pot{pot: r}, po, found
}
-func remove(t *pot, val PotVal) (r *pot, po int, found bool) {
+func remove(t *pot, val Val) (r *pot, po int, found bool) {
size := t.size
- po, found = t.pin.PO(val, t.po)
+ po, found = t.pof(t.pin, val, t.po)
if found {
size--
if size == 0 {
r = &pot{
- po: t.po,
+ po: t.po,
+ pof: t.pof,
}
return r, po, true
}
@@ -205,6 +216,7 @@ func remove(t *pot, val PotVal) (r *pot, po int, found bool) {
bins: append(t.bins[:i], last.bins...),
size: size,
po: t.po,
+ pof: t.pof,
}
return r, t.po, true
}
@@ -237,58 +249,56 @@ func remove(t *pot, val PotVal) (r *pot, po int, found bool) {
size: size,
po: t.po,
bins: bins,
+ pof: t.pof,
}
return r, po, found
}
-// Swap(k, f) looks up the item at k
+// Swap called on (k, f) looks up the item at k
// and applies the function f to the value v at k or nil if the item is not found
// if f returns nil, the element is removed
// if f returns v' <> v then v' is inserted into the Pot
// if v' == v the pot is not changed
// it panics if v'.PO(k, 0) says v and k are not equal
-func (t *Pot) Swap(val AnyVal, f func(v PotVal) PotVal) (po int, found bool, change bool) {
+func (t *Pot) Swap(val Val, f func(v Val) Val) (po int, found bool, change bool) {
t.lock.Lock()
defer t.lock.Unlock()
- ba := NewBytesVal(val, nil)
var t0 *pot
- t0, po, found, change = swap(t.pot, ba, f)
+ t0, po, found, change = swap(t.pot, val, f)
if change {
t.pot = t0
}
return po, found, change
}
-func swap(t *pot, k PotVal, f func(v PotVal) PotVal) (r *pot, po int, found bool, change bool) {
- var val PotVal
+func swap(t *pot, k Val, f func(v Val) Val) (r *pot, po int, found bool, change bool) {
+ var val Val
if t == nil || t.pin == nil {
val = f(nil)
if val == nil {
return t, t.po, false, false
}
- if _, eq := val.PO(k, t.po); !eq {
- panic("value key mismatch")
- }
- r = &pot{
- pin: val,
- size: t.size + 1,
- po: t.po,
- bins: t.bins,
- }
+ // if _, eq := t.pof(k, t.pin, t.po); !eq {
+ // panic("value key mismatch")
+ // }
+ r = t.clone()
+ r.pin = val
+ r.size++
return r, t.po, false, true
}
size := t.size
if k == nil {
panic("k is nil")
}
- po, found = k.PO(t.pin, t.po)
+ po, found = t.pof(k, t.pin, t.po)
if found {
val = f(t.pin)
if val == nil {
size--
if size == 0 {
r = &pot{
- po: t.po,
+ po: t.po,
+ pof: t.pof,
}
return r, po, true, true
}
@@ -299,18 +309,15 @@ func swap(t *pot, k PotVal, f func(v PotVal) PotVal) (r *pot, po int, found bool
bins: append(t.bins[:i], last.bins...),
size: size,
po: t.po,
+ pof: t.pof,
}
return r, t.po, true, true
// remove element
} else if val == t.pin {
return nil, po, true, false
} else { // add element
- r = &pot{
- pin: val,
- size: t.size,
- po: t.po,
- bins: t.bins,
- }
+ r = t.clone()
+ r.pin = val
return r, po, true, true
}
}
@@ -344,6 +351,7 @@ func swap(t *pot, k PotVal, f func(v PotVal) PotVal) (r *pot, po int, found bool
pin: val,
size: 1,
po: po,
+ pof: t.pof,
}
}
@@ -357,12 +365,13 @@ func swap(t *pot, k PotVal, f func(v PotVal) PotVal) (r *pot, po int, found bool
size: size,
po: t.po,
bins: bins,
+ pof: t.pof,
}
return r, po, found, true
}
-// t0.Merge(t1) changes t0 to contain all the elements of t1
+// Merge called on (t1) changes t0 to contain all the elements of t1
// it locks t0, but only readlocks t1 while taking its pot
// uses applicative union
func (t *Pot) Merge(t1 *Pot) (c int) {
@@ -388,9 +397,7 @@ func Union(t0, t1 *Pot) (*Pot, int) {
t1.lock.RUnlock()
p, c := union(n0, n1)
- return &Pot{
- pot: p,
- }, c
+ return &Pot{pot: p}, c
}
func union(t0, t1 *pot) (*pot, int) {
@@ -400,7 +407,7 @@ func union(t0, t1 *pot) (*pot, int) {
if t1 == nil || t1.size == 0 {
return t0, 0
}
- var pin PotVal
+ var pin Val
var bins []*pot
var mis []int
wg := &sync.WaitGroup{}
@@ -411,7 +418,7 @@ func union(t0, t1 *pot) (*pot, int) {
var i0, i1 int
var common int
- po, eq := pin0.PO(pin1, 0)
+ po, eq := t0.pof(pin0, pin1, 0)
for {
l0 := len(bins0)
@@ -486,6 +493,7 @@ func union(t0, t1 *pot) (*pot, int) {
bins: bins0[i:],
size: size0 + 1,
po: po,
+ pof: t0.pof,
}
bins2 := []*pot{np}
@@ -499,7 +507,7 @@ func union(t0, t1 *pot) (*pot, int) {
bins2 = append(bins2, n0.bins...)
pin0 = pin1
pin1 = n0.pin
- po, eq = pin0.PO(pin1, n0.po)
+ po, eq = t0.pof(pin0, pin1, n0.po)
}
bins0 = bins1
@@ -518,21 +526,22 @@ func union(t0, t1 *pot) (*pot, int) {
bins: bins,
size: t0.size + t1.size - common,
po: t0.po,
+ pof: t0.pof,
}
return n, common
}
-// Each(f) is a synchronous iterator over the bins of a node
+// Each called with (f) is a synchronous iterator over the bins of a node
// respecting an ordering
// proximity > pinnedness
-func (t *Pot) Each(f func(PotVal, int) bool) bool {
+func (t *Pot) Each(f func(Val, int) bool) bool {
t.lock.RLock()
n := t.pot
t.lock.RUnlock()
return n.each(f)
}
-func (t *pot) each(f func(PotVal, int) bool) bool {
+func (t *pot) each(f func(Val, int) bool) bool {
var next bool
for _, n := range t.bins {
if n == nil {
@@ -546,7 +555,7 @@ func (t *pot) each(f func(PotVal, int) bool) bool {
return f(t.pin, t.po)
}
-// EachFrom(f, start) is a synchronous iterator over the elements of a pot
+// EachFrom called with (f, start) is a synchronous iterator over the elements of a pot
// within the inclusive range starting from proximity order start
// the function argument is passed the value and the proximity order wrt the root pin
// it does NOT include the pinned item of the root
@@ -554,14 +563,14 @@ func (t *pot) each(f func(PotVal, int) bool) bool {
// proximity > pinnedness
// the iteration ends if the function return false or there are no more elements
// end of a po range can be implemented since po is passed to the function
-func (t *Pot) EachFrom(f func(PotVal, int) bool, po int) bool {
+func (t *Pot) EachFrom(f func(Val, int) bool, po int) bool {
t.lock.RLock()
n := t.pot
t.lock.RUnlock()
return n.eachFrom(f, po)
}
-func (t *pot) eachFrom(f func(PotVal, int) bool, po int) bool {
+func (t *pot) eachFrom(f func(Val, int) bool, po int) bool {
var next bool
_, lim := t.getPos(po)
for i := lim; i < len(t.bins); i++ {
@@ -578,18 +587,18 @@ func (t *pot) eachFrom(f func(PotVal, int) bool, po int) bool {
// subtree passing the proximity order and the size
// the iteration continues until the function's return value is false
// or there are no more subtries
-func (t *Pot) EachBin(val PotVal, po int, f func(int, int, func(func(val PotVal, i int) bool) bool) bool) {
+func (t *Pot) EachBin(val Val, po int, f func(int, int, func(func(val Val, i int) bool) bool) bool) {
t.lock.RLock()
n := t.pot
t.lock.RUnlock()
n.eachBin(val, po, f)
}
-func (t *pot) eachBin(val PotVal, po int, f func(int, int, func(func(val PotVal, i int) bool) bool) bool) {
+func (t *pot) eachBin(val Val, po int, f func(int, int, func(func(val Val, i int) bool) bool) bool) {
if t == nil || t.size == 0 {
return
}
- spr, _ := t.pin.PO(val, t.po)
+ spr, _ := t.pof(t.pin, val, t.po)
_, lim := t.getPos(spr)
var size int
var n *pot
@@ -604,7 +613,7 @@ func (t *pot) eachBin(val PotVal, po int, f func(int, int, func(func(val PotVal,
}
}
if lim == len(t.bins) {
- f(spr, 1, func(g func(PotVal, int) bool) bool {
+ f(spr, 1, func(g func(Val, int) bool) bool {
return g(t.pin, spr)
})
return
@@ -616,8 +625,8 @@ func (t *pot) eachBin(val PotVal, po int, f func(int, int, func(func(val PotVal,
spo++
size += n.size
}
- if !f(spr, t.size-size, func(g func(PotVal, int) bool) bool {
- return t.eachFrom(func(v PotVal, j int) bool {
+ if !f(spr, t.size-size, func(g func(Val, int) bool) bool {
+ return t.eachFrom(func(v Val, j int) bool {
return g(v, spr)
}, spo)
}) {
@@ -628,17 +637,17 @@ func (t *pot) eachBin(val PotVal, po int, f func(int, int, func(func(val PotVal,
}
}
-// syncronous iterator over neighbours of any target val
+// EachNeighbour is a syncronous iterator over neighbours of any target val
// the order of elements retrieved reflect proximity order to the target
// TODO: add maximum proxbin to start range of iteration
-func (t *Pot) EachNeighbour(val PotVal, f func(PotVal, int) bool) bool {
+func (t *Pot) EachNeighbour(val Val, f func(Val, int) bool) bool {
t.lock.RLock()
n := t.pot
t.lock.RUnlock()
return n.eachNeighbour(val, f)
}
-func (t *pot) eachNeighbour(val PotVal, f func(PotVal, int) bool) bool {
+func (t *pot) eachNeighbour(val Val, f func(Val, int) bool) bool {
if t == nil || t.size == 0 {
return false
}
@@ -647,7 +656,7 @@ func (t *pot) eachNeighbour(val PotVal, f func(PotVal, int) bool) bool {
var n *pot
ir := l
il := l
- po, eq := t.pin.PO(val, t.po)
+ po, eq := t.pof(t.pin, val, t.po)
if !eq {
n, il = t.getPos(po)
if n != nil {
@@ -667,7 +676,7 @@ func (t *pot) eachNeighbour(val PotVal, f func(PotVal, int) bool) bool {
}
for i := l - 1; i > ir; i-- {
- next = t.bins[i].each(func(v PotVal, _ int) bool {
+ next = t.bins[i].each(func(v Val, _ int) bool {
return f(v, po)
})
if !next {
@@ -677,7 +686,7 @@ func (t *pot) eachNeighbour(val PotVal, f func(PotVal, int) bool) bool {
for i := il - 1; i >= 0; i-- {
n := t.bins[i]
- next = n.each(func(v PotVal, _ int) bool {
+ next = n.each(func(v Val, _ int) bool {
return f(v, n.po)
})
if !next {
@@ -687,7 +696,7 @@ func (t *pot) eachNeighbour(val PotVal, f func(PotVal, int) bool) bool {
return true
}
-// EachNeighnbourAsync(val, max, maxPos, f, wait) is an asyncronous iterator
+// EachNeighbourAsync called on (val, max, maxPos, f, wait) is an asyncronous iterator
// over elements not closer than maxPos wrt val.
// val does not need to be match an element of the pot, but if it does, and
// maxPos is keylength than it is included in the iteration
@@ -698,7 +707,7 @@ func (t *pot) eachNeighbour(val PotVal, f func(PotVal, int) bool) bool {
// or if the entire there are no nodes not closer than maxPos that is not visited
// if wait is true, the iterator returns only if all calls to f are finished
// TODO: implement minPos for proper prox range iteration
-func (t *Pot) EachNeighbourAsync(val PotVal, max int, maxPos int, f func(PotVal, int), wait bool) {
+func (t *Pot) EachNeighbourAsync(val Val, max int, maxPos int, f func(Val, int), wait bool) {
t.lock.RLock()
n := t.pot
t.lock.RUnlock()
@@ -715,15 +724,14 @@ func (t *Pot) EachNeighbourAsync(val PotVal, max int, maxPos int, f func(PotVal,
}
}
-func (t *pot) eachNeighbourAsync(val PotVal, max int, maxPos int, f func(PotVal, int), wg *sync.WaitGroup) (extra int) {
-
+func (t *pot) eachNeighbourAsync(val Val, max int, maxPos int, f func(Val, int), wg *sync.WaitGroup) (extra int) {
l := len(t.bins)
var n *pot
il := l
ir := l
// ic := l
- po, eq := t.pin.PO(val, t.po)
+ po, eq := t.pof(t.pin, val, t.po)
// if po is too close, set the pivot branch (pom) to maxPos
pom := po
@@ -799,7 +807,7 @@ func (t *pot) eachNeighbourAsync(val PotVal, max int, maxPos int, f func(PotVal,
wg.Add(m)
}
go func(pn *pot, pm int) {
- pn.each(func(v PotVal, _ int) bool {
+ pn.each(func(v Val, _ int) bool {
if wg != nil {
defer wg.Done()
}
@@ -826,7 +834,7 @@ func (t *pot) eachNeighbourAsync(val PotVal, max int, maxPos int, f func(PotVal,
wg.Add(m)
}
go func(pn *pot, pm int) {
- pn.each(func(v PotVal, _ int) bool {
+ pn.each(func(v Val, _ int) bool {
if wg != nil {
defer wg.Done()
}
@@ -840,7 +848,7 @@ func (t *pot) eachNeighbourAsync(val PotVal, max int, maxPos int, f func(PotVal,
return max + extra
}
-// getPos(n) returns the forking node at PO n and its index if it exists
+// getPos called on (n) returns the forking node at PO n and its index if it exists
// otherwise nil
// caller is suppoed to hold the lock
func (t *pot) getPos(po int) (n *pot, i int) {
@@ -856,7 +864,7 @@ func (t *pot) getPos(po int) (n *pot, i int) {
return nil, len(t.bins)
}
-// need(m, max, extra) uses max m out of extra, and then max
+// need called on (m, max, extra) uses max m out of extra, and then max
// if needed, returns the adjusted counts
func need(m, max, extra int) (int, int, int) {
if m <= extra {
diff --git a/pot/pot_test.go b/pot/pot_test.go
index 263dd7942d..1d4dd76966 100644
--- a/pot/pot_test.go
+++ b/pot/pot_test.go
@@ -37,55 +37,45 @@ func init() {
log.Root().SetHandler(log.LvlFilterHandler(log.LvlError, log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
}
-type testBVAddr struct {
- *HashAddress
- i int
-}
-
-func NewTestBVAddr(s string, i int) *testBVAddr {
- return &testBVAddr{NewHashAddress(s), i}
-}
-
-func (a *testBVAddr) String() string {
- return a.HashAddress.String()[:keylen]
-}
-
-func (self *testBVAddr) PO(val PotVal, po int) (int, bool) {
- return self.HashAddress.PO(val.(*testBVAddr).HashAddress, po)
-}
-
type testAddr struct {
- *BoolAddress
+ a []byte
i int
}
-func NewTestAddr(s string, i int) *testAddr {
- return &testAddr{NewBoolAddress(s), i}
+//
+// func newTestAddr(s string, i int) *testAddr {
+// b := NewAddressFromString(s)
+// h := &common.Hash{}
+// copy((*h)[:], b)
+// a := Address(*h)
+// return &testAddr{&a, i}
+// }
+
+func newTestAddr(s string, i int) *testAddr {
+ return &testAddr{NewAddressFromString(s), i}
}
-func (self *testAddr) PO(val PotVal, po int) (int, bool) {
- return self.BoolAddress.PO(val.(*testAddr).BoolAddress, po)
+func (a *testAddr) Address() []byte {
+ return a.a
}
-func str(v PotVal) string {
- if v == nil {
- return ""
- }
- return v.(*testAddr).String()
+func (a *testAddr) String() string {
+ return Label(a.a)
+ // return a.Address.String()[:keylen]
}
func randomTestAddr(n int, i int) *testAddr {
v := RandomAddress().Bin()[:n]
- return NewTestAddr(v, i)
+ return newTestAddr(v, i)
}
-func randomTestBVAddr(n int, i int) *testBVAddr {
+func randomtestAddr(n int, i int) *testAddr {
v := RandomAddress().Bin()[:n]
- return NewTestBVAddr(v, i)
+ return newTestAddr(v, i)
}
func indexes(t *Pot) (i []int, po []int) {
- t.Each(func(v PotVal, p int) bool {
+ t.Each(func(v Val, p int) bool {
a := v.(*testAddr)
i = append(i, a.i)
po = append(po, p)
@@ -96,16 +86,16 @@ func indexes(t *Pot) (i []int, po []int) {
func testAdd(t *Pot, n int, values ...string) {
for i, val := range values {
- t.Add(NewTestAddr(val, i+n))
+ t.Add(newTestAddr(val, i+n))
}
}
// func RandomBoolAddress()
func TestPotAdd(t *testing.T) {
- n := NewPot(NewTestAddr("001111", 0), 0)
+ n := NewPot(newTestAddr("00111100", 0), 0, DefaultPof(8))
// Pin set correctly
- exp := "001111"
- got := str(n.Pin())[:6]
+ exp := "00111100"
+ got := Label(n.Pin())[:8]
if got != exp {
t.Fatalf("incorrect pinned value. Expected %v, got %v", exp, got)
}
@@ -116,7 +106,7 @@ func TestPotAdd(t *testing.T) {
t.Fatalf("incorrect number of elements in Pot. Expected %v, got %v", expi, goti)
}
- testAdd(n, 1, "011111", "001111", "011111", "000111")
+ testAdd(n, 1, "01111100", "00111100", "01111100", "00011100")
// check size
goti = n.Size()
expi = 3
@@ -138,15 +128,15 @@ func TestPotAdd(t *testing.T) {
// func RandomBoolAddress()
func TestPotRemove(t *testing.T) {
- n := NewPot(NewTestAddr("001111", 0), 0)
- n.Remove(NewTestAddr("001111", 0))
- exp := ""
- got := str(n.Pin())
+ n := NewPot(newTestAddr("00111100", 0), 0, DefaultPof(8))
+ n.Remove(newTestAddr("00111100", 0))
+ exp := ""
+ got := Label(n.Pin())
if got != exp {
t.Fatalf("incorrect pinned value. Expected %v, got %v", exp, got)
}
- testAdd(n, 1, "000000", "011111", "001111", "000111")
- n.Remove(NewTestAddr("001111", 0))
+ testAdd(n, 1, "00000000", "01111100", "00111100", "00011100")
+ n.Remove(newTestAddr("00111100", 0))
goti := n.Size()
expi := 3
if goti != expi {
@@ -164,7 +154,7 @@ func TestPotRemove(t *testing.T) {
t.Fatalf("incorrect po-s in iteration over Pot. Expected %v, got %v", exp, got)
}
// remove again
- n.Remove(NewTestAddr("001111", 0))
+ n.Remove(newTestAddr("00111100", 0))
inds, po = indexes(n)
got = fmt.Sprintf("%v", inds)
exp = "[2 4]"
@@ -175,31 +165,32 @@ func TestPotRemove(t *testing.T) {
}
func TestPotSwap(t *testing.T) {
+ // t.Skip("")
max := maxEachNeighbour
- n := NewPot(nil, 0)
- var m []*testBVAddr
+ n := NewPot(nil, 0, nil)
+ var m []*testAddr
for j := 0; j < 2*max; {
- v := randomTestBVAddr(keylen, j)
+ v := randomtestAddr(keylen, j)
_, found := n.Add(v)
if !found {
m = append(m, v)
j++
}
}
- k := make(map[string]*testBVAddr)
+ k := make(map[string]*testAddr)
for j := 0; j < max; {
- v := randomTestBVAddr(keylen, 1)
- _, found := k[v.String()]
+ v := randomtestAddr(keylen, 1)
+ _, found := k[Label(v)]
if !found {
- k[v.String()] = v
+ k[Label(v)] = v
j++
}
}
for _, v := range k {
m = append(m, v)
}
- f := func(v PotVal) PotVal {
- tv := v.(*testBVAddr)
+ f := func(v Val) Val {
+ tv := v.(*testAddr)
if tv.i < max {
return nil
}
@@ -207,7 +198,7 @@ func TestPotSwap(t *testing.T) {
return v
}
for _, val := range m {
- n.Swap(val, func(v PotVal) PotVal {
+ n.Swap(val, func(v Val) Val {
if v == nil {
return val
}
@@ -215,9 +206,9 @@ func TestPotSwap(t *testing.T) {
})
}
sum := 0
- n.Each(func(v PotVal, i int) bool {
+ n.Each(func(v Val, i int) bool {
sum++
- tv := v.(*testBVAddr)
+ tv := v.(*testAddr)
if tv.i > 1 {
t.Fatalf("item value incorrect, expected 0, got %v", tv.i)
}
@@ -231,10 +222,10 @@ func TestPotSwap(t *testing.T) {
}
}
-func checkPo(val PotVal) func(PotVal, int) error {
- return func(v PotVal, po int) error {
+func checkPo(val Val, pof Pof) func(Val, int) error {
+ return func(v Val, po int) error {
// check the po
- exp, _ := val.PO(v, 0)
+ exp, _ := pof(val, v, 0)
if po != exp {
return fmt.Errorf("incorrect prox order for item %v in neighbour iteration for %v. Expected %v, got %v", v, val, exp, po)
}
@@ -242,9 +233,9 @@ func checkPo(val PotVal) func(PotVal, int) error {
}
}
-func checkOrder(val PotVal) func(PotVal, int) error {
+func checkOrder(val Val) func(Val, int) error {
var po int = keylen
- return func(v PotVal, p int) error {
+ return func(v Val, p int) error {
if po < p {
return fmt.Errorf("incorrect order for item %v in neighbour iteration for %v. PO %v > %v (previous max)", v, val, p, po)
}
@@ -253,26 +244,26 @@ func checkOrder(val PotVal) func(PotVal, int) error {
}
}
-func checkValues(m map[string]bool, val PotVal) func(PotVal, int) error {
- return func(v PotVal, po int) error {
- duplicate, ok := m[v.String()]
+func checkValues(m map[string]bool, val Val) func(Val, int) error {
+ return func(v Val, po int) error {
+ duplicate, ok := m[Label(v)]
if !ok {
return fmt.Errorf("alien value %v", v)
}
if duplicate {
return fmt.Errorf("duplicate value returned: %v", v)
}
- m[v.String()] = true
+ m[Label(v)] = true
return nil
}
}
var errNoCount = errors.New("not count")
-func testPotEachNeighbour(n *Pot, val PotVal, expCount int, fs ...func(PotVal, int) error) error {
+func testPotEachNeighbour(n *Pot, val Val, expCount int, fs ...func(Val, int) error) error {
var err error
var count int
- n.EachNeighbour(val, func(v PotVal, po int) bool {
+ n.EachNeighbour(val, func(v Val, po int) bool {
for _, f := range fs {
err = f(v, po)
if err != nil {
@@ -297,15 +288,14 @@ const (
)
func TestPotMergeOne(t *testing.T) {
- pot1 := NewPot(nil, 0)
- pot1.Add(NewTestAddr("10", 0))
- pot1.Add(NewTestAddr("00", 0))
- pot2 := NewPot(nil, 0)
- pot2.Add(NewTestAddr("01", 0))
- log.Debug(fmt.Sprintf("\n%v\n%v", pot2, pot1))
+ pot1 := NewPot(nil, 0, DefaultPof(2))
+ pot1.Add(newTestAddr("10", 0))
+ pot1.Add(newTestAddr("00", 0))
+ pot2 := NewPot(nil, 0, DefaultPof(2))
+ pot2.Add(newTestAddr("01", 0))
pot1.Merge(pot2)
count := 0
- pot1.Each(func(val PotVal, i int) bool {
+ pot1.Each(func(val Val, i int) bool {
count++
return true
})
@@ -315,25 +305,25 @@ func TestPotMergeOne(t *testing.T) {
}
func TestPotMergeCommon(t *testing.T) {
- vs := make([]*testBVAddr, mergeTestCount)
+ vs := make([]*testAddr, mergeTestCount)
+ // vs := make([]*testAddr, mergeTestCount)
for i := 0; i < maxEachNeighbourTests; i++ {
for i := 0; i < len(vs); i++ {
- vs[i] = randomTestBVAddr(keylen, i)
+ vs[i] = randomtestAddr(keylen, i)
}
max0 := rand.Intn(mergeTestChoose) + 1
max1 := rand.Intn(mergeTestChoose) + 1
- n0 := NewPot(nil, 0)
- n1 := NewPot(nil, 0)
+ n0 := NewPot(nil, 0, nil)
+ n1 := NewPot(nil, 0, nil)
log.Trace(fmt.Sprintf("round %v: %v - %v", i, max0, max1))
m := make(map[string]bool)
for j := 0; j < max0; {
r := rand.Intn(max0)
v := vs[r]
- // v := randomTestBVAddr(keylen, j)
_, found := n0.Add(v)
if !found {
- m[v.String()] = false
+ m[Label(v)] = false
j++
}
}
@@ -346,10 +336,10 @@ func TestPotMergeCommon(t *testing.T) {
if !found {
j++
}
- _, found = m[v.String()]
+ _, found = m[Label(v)]
if !found {
expAdded++
- m[v.String()] = false
+ m[Label(v)] = false
}
}
if i < 6 {
@@ -372,8 +362,8 @@ func TestPotMergeCommon(t *testing.T) {
if !checkDuplicates(n.pot) {
t.Fatalf("%v: merged pot contains duplicates: \n%v", i, n)
}
- for k, _ := range m {
- _, found := n.Add(NewTestBVAddr(k, 0))
+ for k := range m {
+ _, found := n.Add(newTestAddr(k, 0))
if !found {
t.Fatalf("%v: merged pot (size:%v, added: %v) missing element %v\n%v", i, size, added, k, n)
}
@@ -385,32 +375,32 @@ func TestPotMergeScale(t *testing.T) {
for i := 0; i < maxEachNeighbourTests; i++ {
max0 := rand.Intn(maxEachNeighbour) + 1
max1 := rand.Intn(maxEachNeighbour) + 1
- n0 := NewPot(nil, 0)
- n1 := NewPot(nil, 0)
+ n0 := NewPot(nil, 0, nil)
+ n1 := NewPot(nil, 0, nil)
log.Trace(fmt.Sprintf("round %v: %v - %v", i, max0, max1))
m := make(map[string]bool)
for j := 0; j < max0; {
- v := randomTestBVAddr(keylen, j)
- // v := randomTestBVAddr(keylen, j)
+ v := randomtestAddr(keylen, j)
+ // v := randomtestAddr(keylen, j)
_, found := n0.Add(v)
if !found {
- m[v.String()] = false
+ m[Label(v)] = false
j++
}
}
expAdded := 0
for j := 0; j < max1; {
- v := randomTestBVAddr(keylen, j)
- // v := randomTestBVAddr(keylen, j)
+ v := randomtestAddr(keylen, j)
+ // v := randomtestAddr(keylen, j)
_, found := n1.Add(v)
if !found {
j++
}
- _, found = m[v.String()]
+ _, found = m[Label(v)]
if !found {
expAdded++
- m[v.String()] = false
+ m[Label(v)] = false
}
}
if i < 6 {
@@ -433,8 +423,8 @@ func TestPotMergeScale(t *testing.T) {
if !checkDuplicates(n.pot) {
t.Fatalf("%v: merged pot contains duplicates: \n%v", i, n)
}
- for k, _ := range m {
- _, found := n.Add(NewTestBVAddr(k, 0))
+ for k := range m {
+ _, found := n.Add(newTestAddr(k, 0))
if !found {
t.Fatalf("%v: merged pot (size:%v, added: %v) missing element %v\n%v", i, size, added, k, n)
}
@@ -443,7 +433,7 @@ func TestPotMergeScale(t *testing.T) {
}
func checkDuplicates(t *pot) bool {
- var po int = -1
+ po := -1
for _, c := range t.bins {
if c == nil {
return false
@@ -460,13 +450,13 @@ func TestPotEachNeighbourSync(t *testing.T) {
for i := 0; i < maxEachNeighbourTests; i++ {
max := rand.Intn(maxEachNeighbour/2) + maxEachNeighbour/2
pin := randomTestAddr(keylen, 0)
- n := NewPot(pin, 0)
+ n := NewPot(pin, 0, nil)
m := make(map[string]bool)
- m[pin.String()] = false
+ m[Label(pin)] = false
for j := 1; j <= max; j++ {
v := randomTestAddr(keylen, j)
n.Add(v)
- m[v.String()] = false
+ m[Label(v)] = false
}
size := n.Size()
@@ -476,14 +466,14 @@ func TestPotEachNeighbourSync(t *testing.T) {
count := rand.Intn(size/2) + size/2
val := randomTestAddr(keylen, max+1)
log.Trace(fmt.Sprintf("%v: pin: %v, size: %v, val: %v, count: %v", i, n.Pin(), size, val, count))
- err := testPotEachNeighbour(n, val, count, checkPo(val), checkOrder(val), checkValues(m, val))
+ err := testPotEachNeighbour(n, val, count, checkPo(val, n.pof), checkOrder(val), checkValues(m, val))
if err != nil {
t.Fatal(err)
}
minPoFound := keylen
maxPoNotFound := 0
for k, found := range m {
- po, _ := val.PO(NewTestAddr(k, 0), 0)
+ po, _ := n.pof(val, newTestAddr(k, 0), 0)
if found {
if po < minPoFound {
minPoFound = po
@@ -503,8 +493,8 @@ func TestPotEachNeighbourSync(t *testing.T) {
func TestPotEachNeighbourAsync(t *testing.T) {
for i := 0; i < maxEachNeighbourTests; i++ {
max := rand.Intn(maxEachNeighbour/2) + maxEachNeighbour/2
- n := NewPot(randomTestAddr(keylen, 0), 0)
- var size int = 1
+ n := NewPot(randomTestAddr(keylen, 0), 0, nil)
+ size := 1
for j := 1; j <= max; j++ {
v := randomTestAddr(keylen, j)
_, found := n.Add(v)
@@ -526,32 +516,24 @@ func TestPotEachNeighbourAsync(t *testing.T) {
maxPos := rand.Intn(keylen)
log.Trace(fmt.Sprintf("%v: pin: %v, size: %v, val: %v, count: %v, maxPos: %v", i, n.Pin(), size, val, count, maxPos))
msize := 0
- remember := func(v PotVal, po int) error {
- // mu.Lock()
- // defer mu.Unlock()
+ remember := func(v Val, po int) error {
if po > maxPos {
- // log.Trace(fmt.Sprintf("NOT ADD %v", v))
return errNoCount
}
- // log.Trace(fmt.Sprintf("ADD %v, %v", v, msize))
- m[v.String()] = true
+ m[Label(v)] = true
msize++
return nil
}
if i == 0 {
continue
}
- err := testPotEachNeighbour(n, val, count, remember)
- if err != nil {
- log.Error(err.Error())
- }
+ testPotEachNeighbour(n, val, count, remember)
d := 0
- forget := func(v PotVal, po int) {
+ forget := func(v Val, po int) {
mu.Lock()
defer mu.Unlock()
d++
- // log.Trace(fmt.Sprintf("DEL %v", v))
- delete(m, v.String())
+ delete(m, Label(v))
}
n.EachNeighbourAsync(val, count, maxPos, forget, true)
@@ -567,7 +549,7 @@ func TestPotEachNeighbourAsync(t *testing.T) {
func benchmarkEachNeighbourSync(t *testing.B, max, count int, d time.Duration) {
t.ReportAllocs()
pin := randomTestAddr(keylen, 0)
- n := NewPot(pin, 0)
+ n := NewPot(pin, 0, nil)
for j := 1; j <= max; {
v := randomTestAddr(keylen, j)
_, found := n.Add(v)
@@ -579,7 +561,7 @@ func benchmarkEachNeighbourSync(t *testing.B, max, count int, d time.Duration) {
for i := 0; i < t.N; i++ {
val := randomTestAddr(keylen, max+1)
m := 0
- n.EachNeighbour(val, func(v PotVal, po int) bool {
+ n.EachNeighbour(val, func(v Val, po int) bool {
time.Sleep(d)
m++
if m == count {
@@ -597,7 +579,7 @@ func benchmarkEachNeighbourSync(t *testing.B, max, count int, d time.Duration) {
func benchmarkEachNeighbourAsync(t *testing.B, max, count int, d time.Duration) {
t.ReportAllocs()
pin := randomTestAddr(keylen, 0)
- n := NewPot(pin, 0)
+ n := NewPot(pin, 0, nil)
for j := 1; j <= max; {
v := randomTestAddr(keylen, j)
_, found := n.Add(v)
@@ -608,7 +590,7 @@ func benchmarkEachNeighbourAsync(t *testing.B, max, count int, d time.Duration)
t.ResetTimer()
for i := 0; i < t.N; i++ {
val := randomTestAddr(keylen, max+1)
- n.EachNeighbourAsync(val, count, keylen, func(v PotVal, po int) {
+ n.EachNeighbourAsync(val, count, keylen, func(v Val, po int) {
time.Sleep(d)
}, true)
}
diff --git a/swarm/network/hive.go b/swarm/network/hive.go
index e402930569..ab492201c5 100644
--- a/swarm/network/hive.go
+++ b/swarm/network/hive.go
@@ -89,9 +89,10 @@ type Hive struct {
tick <-chan time.Time
}
-// Hive constructor embeds both arguments
+// NewHive constructs a new hive
// HiveParams: config parameters
// Overlay: Topology Driver Interface
+// StateStore: to save peers across sessions
func NewHive(params *HiveParams, overlay Overlay, store StateStore) *Hive {
return &Hive{
HiveParams: params,
@@ -105,105 +106,98 @@ func NewHive(params *HiveParams, overlay Overlay, store StateStore) *Hive {
// these are called on the p2p.Server which runs on the node
// af() returns an arbitrary ticker channel
// rw is a read writer for json configs
-func (self *Hive) Start(server *p2p.Server) error {
- if self.store != nil {
- if err := self.loadPeers(); err != nil {
+func (h *Hive) Start(server *p2p.Server) error {
+ if h.store != nil {
+ if err := h.loadPeers(); err != nil {
return err
}
}
- self.more = make(chan bool, 1)
- self.quit = make(chan bool)
- log.Debug("hive started")
+ h.more = make(chan bool, 1)
+ h.quit = make(chan bool)
// this loop is doing bootstrapping and maintains a healthy table
- go self.keepAlive()
+ go h.keepAlive()
go func() {
// each iteration, ask kademlia about most preferred peer
- for more := range self.more {
+ for more := range h.more {
if !more {
// receiving false closes the loop while allowing parallel routines
// to attempt to write to more (remove Peer when shutting down)
return
}
- log.Debug("hive delegate to overlay driver: suggest addr to connect to")
+ log.Trace(fmt.Sprintf("%x: hive delegate to overlay driver: suggest addr to connect to", h.BaseAddr()[:4]))
// log.Trace("hive delegate to overlay driver: suggest addr to connect to")
- addr, order, want := self.SuggestPeer()
- if self.Discovery {
+ addr, order, want := h.SuggestPeer()
+ if h.Discovery {
if addr != nil {
- log.Info(fmt.Sprintf("========> connect to bee %v", addr))
+ log.Trace(fmt.Sprintf("%x ========> connect to bee %x", h.BaseAddr()[:4], addr.Address()[:4]))
under, err := discover.ParseNode(string(addr.(Addr).Under()))
if err == nil {
server.AddPeer(under)
} else {
- log.Error(fmt.Sprintf("===X====> connect to bee %v failed: invalid node URL: %v", addr, err))
+ log.Error(fmt.Sprintf("%x ===X====> connect to bee %x failed: invalid node URL: %v", h.BaseAddr()[:4], addr.Address()[:4], err))
}
} else {
- log.Trace("cannot suggest peers")
+ log.Trace(fmt.Sprintf("%x cannot suggest peers", h.BaseAddr()[:4]))
}
if want {
- log.Debug(fmt.Sprintf("========> request peers nearest %v", addr))
- RequestOrder(self.Overlay, uint8(order), self.PeersBroadcastSetSize, self.MaxPeersPerRequest)
+ log.Trace(fmt.Sprintf("%x ========> request peers for PO%0d", h.BaseAddr()[:4], order))
+ RequestOrder(h.Overlay, uint8(order), h.PeersBroadcastSetSize, h.MaxPeersPerRequest)
}
}
- log.Info(fmt.Sprintf("%v", self))
- select {
- case <-self.quit:
- return
- default:
- }
+ log.Trace(fmt.Sprintf("%v", h))
}
}()
return nil
}
// Stop terminates the updateloop and saves the peers
-func (self *Hive) Stop() {
- if self.store != nil {
- self.savePeers()
+func (h *Hive) Stop() {
+ if h.store != nil {
+ h.savePeers()
}
// closing toggle channel quits the updateloop
- close(self.quit)
+ close(h.quit)
}
-func (self *Hive) Run(p *bzzPeer) error {
- dp := NewDiscovery(p, self)
+// Run protocol run function
+func (h *Hive) Run(p *bzzPeer) error {
+ dp := NewDiscovery(p, h)
log.Debug(fmt.Sprintf("to add new bee %v", p))
- self.On(dp)
- self.wake()
- defer self.wake()
- defer self.Off(dp)
+ h.On(dp)
+ h.wake()
+ defer h.wake()
+ defer h.Off(dp)
return p.Run(dp.HandleMsg)
}
// NodeInfo function is used by the p2p.server RPC interface to display
// protocol specific node information
-func (self *Hive) NodeInfo() interface{} {
- return self.String()
+func (h *Hive) NodeInfo() interface{} {
+ return h.String()
}
// PeerInfo function is used by the p2p.server RPC interface to display
// protocol specific information any connected peer referred to by their NodeID
-func (self *Hive) PeerInfo(id discover.NodeID) interface{} {
- self.lock.Lock()
- defer self.lock.Unlock()
+func (h *Hive) PeerInfo(id discover.NodeID) interface{} {
+ h.lock.Lock()
+ defer h.lock.Unlock()
addr := NewAddrFromNodeID(id)
return interface{}(addr)
}
-func (self *Hive) Register(peers chan OverlayAddr) error {
- defer self.wake()
- return self.Overlay.Register(peers)
+func (h *Hive) Register(peers chan OverlayAddr) error {
+ defer h.wake()
+ return h.Overlay.Register(peers)
}
// wake triggers
-func (self *Hive) wake() {
+func (h *Hive) wake() {
select {
- case self.more <- true:
- log.Trace("hive woken up")
- case <-self.quit:
+ case h.more <- true:
+ case <-h.quit:
default:
- log.Trace("hive already awake")
}
}
@@ -233,30 +227,30 @@ func (t *timeTicker) Ch() <-chan time.Time {
// keepAlive is a forever loop
// in its awake state it periodically triggers connection attempts
-// by writing to self.more until Kademlia Table is saturated
-// wake state is toggled by writing to self.toggle
+// by writing to h.more until Kademlia Table is saturated
+// wake state is toggled by writing to h.toggle
// it goes to sleep mode if table is saturated
// it restarts if the table becomes non-full again due to disconnections
-func (self *Hive) keepAlive() {
- if self.tick == nil {
- ticker := time.NewTicker(self.KeepAliveInterval)
+func (h *Hive) keepAlive() {
+ if h.tick == nil {
+ ticker := time.NewTicker(h.KeepAliveInterval)
defer ticker.Stop()
- self.tick = ticker.C
+ h.tick = ticker.C
}
for {
select {
- case <-self.tick:
- log.Debug("wake up: make hive alive")
- self.wake()
- case <-self.quit:
+ case <-h.tick:
+ h.wake()
+ case <-h.quit:
+ h.more <- false
return
}
}
}
// loadPeers, savePeer implement persistence callback/
-func (self *Hive) loadPeers() error {
- data, err := self.store.Load("peers")
+func (h *Hive) loadPeers() error {
+ data, err := h.store.Load("peers")
if err != nil {
return err
}
@@ -275,13 +269,13 @@ func (self *Hive) loadPeers() error {
c <- a
}
}()
- return self.Overlay.Register(c)
+ return h.Overlay.Register(c)
}
// savePeers, savePeer implement persistence callback/
-func (self *Hive) savePeers() error {
+func (h *Hive) savePeers() error {
var peers []*bzzAddr
- self.Overlay.EachAddr(nil, 256, func(pa OverlayAddr, i int) bool {
+ h.Overlay.EachAddr(nil, 256, func(pa OverlayAddr, i int) bool {
if pa == nil {
log.Warn(fmt.Sprintf("empty addr: %v", i))
return true
@@ -293,7 +287,7 @@ func (self *Hive) savePeers() error {
if err != nil {
return fmt.Errorf("could not encode peers: %v", err)
}
- if err := self.store.Save("peers", data); err != nil {
+ if err := h.store.Save("peers", data); err != nil {
return fmt.Errorf("could not save peers: %v", err)
}
return nil
diff --git a/swarm/network/kademlia.go b/swarm/network/kademlia.go
index 69e2ef6f8e..1d7663c124 100644
--- a/swarm/network/kademlia.go
+++ b/swarm/network/kademlia.go
@@ -93,11 +93,12 @@ func NewKademlia(addr []byte, params *KadParams) *Kademlia {
return &Kademlia{
base: addr,
KadParams: params,
- addrs: pot.NewPot(nil, 0),
- conns: pot.NewPot(nil, 0),
+ addrs: pot.NewPot(nil, 0, nil),
+ conns: pot.NewPot(nil, 0, nil),
}
}
+// Notifier interface type for peer allowing / requesting peer and depth notifications
type Notifier interface {
NotifyPeer(OverlayAddr, uint8) error
NotifyDepth(uint8) error
@@ -116,16 +117,14 @@ type OverlayConn interface {
Off() OverlayAddr // call to return a persitent OverlayAddr
}
+// OverlayAddr represents a kademlia peer record
type OverlayAddr interface {
OverlayPeer
Update(OverlayAddr) OverlayAddr // returns the updated version of the original
}
// entry represents a Kademlia table entry (an extension of OverlayPeer)
-// implements the pot.PotVal interface via BytesAddress, so entry can be
-// used directly as a pot element
type entry struct {
- pot.PotVal
OverlayPeer
seenAt time.Time
retries int
@@ -134,51 +133,56 @@ type entry struct {
// newEntry creates a kademlia peer from an OverlayPeer interface
func newEntry(p OverlayPeer) *entry {
return &entry{
- PotVal: pot.NewBytesVal(p, nil),
OverlayPeer: p,
seenAt: time.Now(),
}
}
-func (self *entry) addr() OverlayAddr {
- a, _ := self.OverlayPeer.(OverlayAddr)
+// Bin is the binary (bitvector) serialisation of the entry address
+func (e *entry) Bin() string {
+ return pot.ToBin(e.addr().Address())
+}
+
+// Label is a short tag for the entry for debug
+func Label(e *entry) string {
+ return fmt.Sprintf("%s (%d)", e.Bin()[:8], e.retries)
+}
+
+// Hex is the hexadecimal serialisation of the entry address
+func (e *entry) Hex() string {
+ return fmt.Sprintf("%x", e.addr().Address())
+}
+
+// String is the short tag for the entry
+func (e *entry) String() string {
+ return fmt.Sprintf("%s (%d)", e.Hex()[:4], e.retries)
+}
+
+// addr returns the kad peer record (OverlayAddr) corresponding to the entry
+func (e *entry) addr() OverlayAddr {
+ a, _ := e.OverlayPeer.(OverlayAddr)
return a
}
-func (self *entry) conn() OverlayConn {
- c, _ := self.OverlayPeer.(OverlayConn)
+// conn returns the connected peer (OverlayPeer) corresponding to the entry
+func (e *entry) conn() OverlayConn {
+ c, _ := e.OverlayPeer.(OverlayConn)
return c
}
-func (self *entry) String() string {
- return fmt.Sprintf("%x", self.OverlayPeer.Address())
-}
-
// Register enters each OverlayAddr as kademlia peer record into the
// database of known peer addresses
-func (self *Kademlia) Register(peers chan OverlayAddr) error {
- np := pot.NewPot(nil, 0)
+func (k *Kademlia) Register(peers chan OverlayAddr) error {
+ np := pot.NewPot(nil, 0, nil)
for p := range peers {
- // error if self received, peer should know better
- if bytes.Equal(p.Address(), self.base) {
- return fmt.Errorf("add peers: %x is self", self.base)
+ // error if k received, peer should know better
+ if bytes.Equal(p.Address(), k.base) {
+ return fmt.Errorf("add peers: %x is k", k.base)
}
np, _, _ = pot.Add(np, newEntry(p))
}
- com := self.addrs.Merge(np)
- log.Debug(fmt.Sprintf("merged %v peers, %v known, total: %v", np.Size(), com, self.addrs.Size()))
- // log.Trace(fmt.Sprintf("merged %v peers, %v known", np.Size(), com))
-
- // TODO: remove this check
- m := make(map[string]bool)
- self.addrs.Each(func(val pot.PotVal, i int) bool {
- _, found := m[val.String()]
- if found {
- panic("duplicate found")
- }
- m[val.String()] = true
- return true
- })
+ com := k.addrs.Merge(np)
+ log.Trace(fmt.Sprintf("%x merged %v peers, %v known, total: %v", k.BaseAddr()[:4], np.Size(), com, k.addrs.Size()))
return nil
}
@@ -186,32 +190,31 @@ func (self *Kademlia) Register(peers chan OverlayAddr) error {
// lowest bincount below depth
// naturally if there is an empty row it returns a peer for that
//
-func (self *Kademlia) SuggestPeer() (a OverlayAddr, o int, want bool) {
- minsize := self.MinBinSize
- depth := self.Depth()
- empty := self.FirstEmptyBin()
+func (k *Kademlia) SuggestPeer() (a OverlayAddr, o int, want bool) {
+ minsize := k.MinBinSize
+ depth := k.Depth()
+ // empty := k.FirstEmptyBin()
// if there is a callable neighbour within the current proxBin, connect
// this makes sure nearest neighbour set is fully connected
- log.Debug(fmt.Sprintf("candidate nearest neighbour checking above PO %v", depth))
- // log.Trace(fmt.Sprintf("candidate nearest neighbour checking above PO %v", depth))
var ppo int
- ba := pot.NewBytesVal(self.base, nil)
- self.addrs.EachNeighbour(ba, func(val pot.PotVal, po int) bool {
- a = self.callable(val)
- log.Trace(fmt.Sprintf("candidate nearest neighbour at %x: %v (%v). a == nil is %v", val.(*entry).Address(), a, po, a == nil))
+ k.addrs.EachNeighbour(k.base, func(val pot.Val, po int) bool {
+ if po < depth {
+ return false
+ }
+ a = k.callable(val)
+ log.Trace(fmt.Sprintf("%x candidate nearest neighbour at %v: %v (%v)", k.BaseAddr()[:4], val.(*entry), a, po))
ppo = po
- return a == nil && po >= depth
+ return a == nil
})
if a != nil {
- log.Debug(fmt.Sprintf("candidate nearest neighbour found: %v (%v)", a, ppo))
+ log.Trace(fmt.Sprintf("%x candidate nearest neighbour found: %v (%v)", k.BaseAddr()[:4], a, ppo))
return a, 0, false
}
- log.Debug(fmt.Sprintf("no candidate nearest neighbours to connect to (Depth: %v, minProxSize: %v) %#v", depth, self.MinProxBinSize, a))
+ log.Trace(fmt.Sprintf("%x no candidate nearest neighbours to connect to (Depth: %v, minProxSize: %v) %#v", k.BaseAddr()[:4], depth, k.MinProxBinSize, a))
var bpo []int
prev := -1
- self.conns.EachBin(pot.NewBytesVal(self.base, nil), 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool {
- log.Trace(fmt.Sprintf("check PO%02d: ", po))
+ k.conns.EachBin(k.base, 0, func(po, size int, f func(func(val pot.Val, i int) bool) bool) bool {
prev++
for ; prev < po; prev++ {
bpo = append(bpo, prev)
@@ -224,34 +227,31 @@ func (self *Kademlia) SuggestPeer() (a OverlayAddr, o int, want bool) {
return size > 0 && po < depth
})
// all buckets are full
- // minsize == self.MinBinSize
+ // minsize == k.MinBinSize
if len(bpo) == 0 {
- log.Trace(fmt.Sprintf("all bins saturated"))
+ log.Debug(fmt.Sprintf("%x: all bins saturated", k.BaseAddr()[:4]))
return nil, 0, false
}
// as long as we got candidate peers to connect to
// dont ask for new peers (want = false)
// try to select a candidate peer
- // for i := len(bpo) - 1; i >= 0; i-- {
// find the first callable peer
i := 0
nxt := bpo[0]
- self.addrs.EachBin(ba, nxt, func(po, size int, f func(func(pot.PotVal, int) bool) bool) bool {
+ k.addrs.EachBin(k.base, nxt, func(po, size int, f func(func(pot.Val, int) bool) bool) bool {
// for each bin we find callable candidate peers
- if po == nxt {
- if i == len(bpo)-1 {
- return false
- }
- i++
- nxt = bpo[i]
+ if i >= depth {
+ return false
}
- log.Trace(fmt.Sprintf("check PO%02d: ", po))
- f(func(val pot.PotVal, j int) bool {
- a = self.callable(val)
- if po == empty {
- log.Debug(fmt.Sprintf("candidate nearest neighbour found: %v (%v)", a, ppo))
+ if po == nxt {
+ i++
+ if i < len(bpo) {
+ nxt = bpo[i]
}
- return a == nil && po <= depth
+ }
+ f(func(val pot.Val, j int) bool {
+ a = k.callable(val)
+ return a == nil && i < len(bpo) && po <= depth
})
return false
})
@@ -259,22 +259,17 @@ func (self *Kademlia) SuggestPeer() (a OverlayAddr, o int, want bool) {
if a != nil {
return a, 0, false
}
- return a, bpo[i], true
- // cannot find a candidate, ask for more for this proximity bin specifically
- // o = bpo[i]
- // want = true
- // // }
- // return a, o, want
+ return a, nxt, true
}
// On inserts the peer as a kademlia peer into the live peers
-func (self *Kademlia) On(p OverlayConn) {
+func (k *Kademlia) On(p OverlayConn) {
e := newEntry(p)
- self.conns.Swap(p, func(v pot.PotVal) pot.PotVal {
+ k.conns.Swap(p, func(v pot.Val) pot.Val {
// if not found live
if v == nil {
// insert new online peer into addrs
- self.addrs.Swap(p, func(v pot.PotVal) pot.PotVal {
+ k.addrs.Swap(p, func(v pot.Val) pot.Val {
return e
})
// insert new online peer into conns
@@ -289,36 +284,34 @@ func (self *Kademlia) On(p OverlayConn) {
return
}
- depth := uint8(self.Depth())
- if depth != self.depth {
- self.depth = depth
+ depth := uint8(k.Depth())
+ if depth != k.depth {
+ k.depth = depth
} else {
depth = 0
}
go np.NotifyDepth(depth)
- f := func(val pot.PotVal, po int) {
+ f := func(val pot.Val, po int) {
dp := val.(*entry).OverlayPeer.(Notifier)
dp.NotifyPeer(p.Off(), uint8(po))
- // log.Trace(fmt.Sprintf("peer %v notified of %v (%v)", dp, p, po))
- log.Debug(fmt.Sprintf("peer %v notified of %v (%v)", dp, p, po))
+ log.Trace(fmt.Sprintf("peer %v notified of %v (%v)", dp, p, po))
if depth > 0 {
dp.NotifyDepth(depth)
- log.Debug(fmt.Sprintf("peer %v notified of new depth %v", dp, depth))
- // log.Trace(fmt.Sprintf("peer %v notified of new depth %v", dp, depth))
+ log.Trace(fmt.Sprintf("peer %v notified of new depth %v", dp, depth))
}
}
- self.conns.EachNeighbourAsync(e, 1024, 255, f, false)
+ k.conns.EachNeighbourAsync(e, 1024, 255, f, false)
}
// Off removes a peer from among live peers
-func (self *Kademlia) Off(p OverlayConn) {
- self.addrs.Swap(p, func(v pot.PotVal) pot.PotVal {
+func (k *Kademlia) Off(p OverlayConn) {
+ k.addrs.Swap(p, func(v pot.Val) pot.Val {
// v cannot be nil, must check otherwise we overwrite entry
if v == nil {
panic(fmt.Sprintf("connected peer not found %v", p))
}
- self.conns.Swap(p, func(v pot.PotVal) pot.PotVal {
+ k.conns.Swap(p, func(_ pot.Val) pot.Val {
// v cannot be nil, but no need to check
return nil
})
@@ -329,13 +322,12 @@ func (self *Kademlia) Off(p OverlayConn) {
// EachConn is an iterator with args (base, po, f) applies f to each live peer
// that has proximity order po or less as measured from the base
// if base is nil, kademlia base address is used
-func (self *Kademlia) EachConn(base []byte, o int, f func(OverlayConn, int, bool) bool) {
+func (k *Kademlia) EachConn(base []byte, o int, f func(OverlayConn, int, bool) bool) {
if len(base) == 0 {
- base = self.base
+ base = k.base
}
- p := pot.NewBytesVal(base, nil)
- depth := self.Depth()
- self.conns.EachNeighbour(p, func(val pot.PotVal, po int) bool {
+ depth := k.Depth()
+ k.conns.EachNeighbour(base, func(val pot.Val, po int) bool {
if po > o {
return true
}
@@ -343,15 +335,14 @@ func (self *Kademlia) EachConn(base []byte, o int, f func(OverlayConn, int, bool
})
}
-// EachAddr(base, po, f) is an iterator applying f to each known peer
+// EachAddr called with (base, po, f) is an iterator applying f to each known peer
// that has proximity order po or less as measured from the base
// if base is nil, kademlia base address is used
-func (self *Kademlia) EachAddr(base []byte, o int, f func(OverlayAddr, int) bool) {
+func (k *Kademlia) EachAddr(base []byte, o int, f func(OverlayAddr, int) bool) {
if len(base) == 0 {
- base = self.base
+ base = k.base
}
- p := pot.NewBytesVal(base, nil)
- self.addrs.EachNeighbour(p, func(val pot.PotVal, po int) bool {
+ k.addrs.EachNeighbour(base, func(val pot.Val, po int) bool {
if po > o {
return true
}
@@ -362,39 +353,39 @@ func (self *Kademlia) EachAddr(base []byte, o int, f func(OverlayAddr, int) bool
// Depth returns the proximity order that defines the distance of
// the nearest neighbour set with cardinality >= MinProxBinSize
// if there is altogether less than MinProxBinSize peers it returns 0
-func (self *Kademlia) Depth() (depth int) {
- if self.conns.Size() < self.MinProxBinSize {
+func (k *Kademlia) Depth() (depth int) {
+ if k.conns.Size() < k.MinProxBinSize {
return 0
}
var size int
- f := func(v pot.PotVal, i int) bool {
+ f := func(v pot.Val, i int) bool {
size++
depth = i
- return size < self.MinProxBinSize
+ return size < k.MinProxBinSize
}
- self.conns.EachNeighbour(pot.NewBytesVal(self.base, nil), f)
+ k.conns.EachNeighbour(k.base, f)
return depth
}
-func (self *Kademlia) callable(val pot.PotVal) OverlayAddr {
+// calleble when called with val,
+func (k *Kademlia) callable(val pot.Val) OverlayAddr {
e := val.(*entry)
// not callable if peer is live or exceeded maxRetries
- if e.conn() != nil || e.retries > self.MaxRetries {
+ if e.conn() != nil || e.retries > k.MaxRetries {
log.Trace(fmt.Sprintf("peer %v (%T) not callable", e, e.OverlayPeer))
return nil
}
// calculate the allowed number of retries based on time lapsed since last seen
timeAgo := time.Since(e.seenAt)
var retries int
- for delta := int(timeAgo) / self.RetryInterval; delta > 0; delta /= self.RetryExponent {
- log.Trace(fmt.Sprintf("delta: %v", delta))
+ for delta := int(timeAgo) / k.RetryInterval; delta > 0; delta /= k.RetryExponent {
retries++
}
// this is never called concurrently, so safe to increment
// peer can be retried again
if retries < e.retries {
- log.Trace(fmt.Sprintf("long time since last try (at %v) needed before retry %v, wait only warrants %v", timeAgo, e.retries, retries))
+ log.Trace(fmt.Sprintf("%v long time since last try (at %v) needed before retry %v, wait only warrants %v", e, timeAgo, e.retries, retries))
return nil
}
e.retries++
@@ -404,64 +395,61 @@ func (self *Kademlia) callable(val pot.PotVal) OverlayAddr {
}
// BaseAddr return the kademlia base addres
-func (self *Kademlia) BaseAddr() []byte {
- return self.base
+func (k *Kademlia) BaseAddr() []byte {
+ return k.base
}
// String returns kademlia table + kaddb table displayed with ascii
-func (self *Kademlia) String() string {
-
+func (k *Kademlia) String() string {
+ wsrow := " "
var rows []string
rows = append(rows, "=========================================================================")
- rows = append(rows, fmt.Sprintf("%v KΛÐΞMLIΛ hive: queen's address: %x", time.Now().UTC().Format(time.UnixDate), self.BaseAddr()[:3]))
- rows = append(rows, fmt.Sprintf("population: %d (%d), MinProxBinSize: %d, MinBinSize: %d, MaxBinSize: %d", self.conns.Size(), self.addrs.Size(), self.MinProxBinSize, self.MinBinSize, self.MaxBinSize))
+ rows = append(rows, fmt.Sprintf("%v KΛÐΞMLIΛ hive: queen's address: %x", time.Now().UTC().Format(time.UnixDate), k.BaseAddr()[:3]))
+ rows = append(rows, fmt.Sprintf("population: %d (%d), MinProxBinSize: %d, MinBinSize: %d, MaxBinSize: %d", k.conns.Size(), k.addrs.Size(), k.MinProxBinSize, k.MinBinSize, k.MaxBinSize))
- liverows := make([]string, self.MaxProxDisplay)
- peersrows := make([]string, self.MaxProxDisplay)
+ liverows := make([]string, k.MaxProxDisplay)
+ peersrows := make([]string, k.MaxProxDisplay)
var depth int
prev := -1
var depthSet bool
- rest := self.conns.Size()
- self.conns.EachBin(pot.NewBytesVal(self.base, nil), 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool {
+ rest := k.conns.Size()
+ k.conns.EachBin(k.base, 0, func(po, size int, f func(func(val pot.Val, i int) bool) bool) bool {
var rowlen int
- if po >= self.MaxProxDisplay {
- po = self.MaxProxDisplay - 1
+ if po >= k.MaxProxDisplay {
+ po = k.MaxProxDisplay - 1
}
row := []string{fmt.Sprintf("%2d", size)}
rest -= size
- f(func(val pot.PotVal, vpo int) bool {
+ f(func(val pot.Val, vpo int) bool {
e := val.(*entry)
- label := e.String()[:6]
- row = append(row, fmt.Sprintf("%s (%d)", label, e.retries))
- row = append(row, val.(*entry).String()[:6])
+ row = append(row, e.String())
rowlen++
return rowlen < 4
})
- if !depthSet && (po > prev+1 || rest < self.MinProxBinSize) {
+ if !depthSet && (po > prev+1 || rest < k.MinProxBinSize) {
depthSet = true
depth = prev + 1
}
- for ; rowlen <= 5; rowlen++ {
- row = append(row, " ")
- }
- liverows[po] = strings.Join(row, " ")
+ row = append(row, wsrow)
+ r := strings.Join(row, " ")
+ liverows[po] = r[:35]
prev = po
return true
})
- self.addrs.EachBin(pot.NewBytesVal(self.base, nil), 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool {
+ k.addrs.EachBin(k.base, 0, func(po, size int, f func(func(val pot.Val, i int) bool) bool) bool {
var rowlen int
- if po >= self.MaxProxDisplay {
- po = self.MaxProxDisplay - 1
+ if po >= k.MaxProxDisplay {
+ po = k.MaxProxDisplay - 1
}
if size < 0 {
panic("wtf")
}
row := []string{fmt.Sprintf("%2d", size)}
// we are displaying live peers too
- f(func(val pot.PotVal, vpo int) bool {
- row = append(row, val.(*entry).String()[:6])
+ f(func(val pot.Val, vpo int) bool {
+ row = append(row, val.(*entry).String())
rowlen++
return rowlen < 4
})
@@ -469,7 +457,7 @@ func (self *Kademlia) String() string {
return true
})
- for i := 0; i < self.MaxProxDisplay; i++ {
+ for i := 0; i < k.MaxProxDisplay; i++ {
if i == depth {
rows = append(rows, fmt.Sprintf("============ DEPTH: %d ==========================================", i))
}
@@ -494,15 +482,15 @@ func (self *Kademlia) String() string {
// the MaxBinSize parameter it drops the oldest n peers such that
// the bin is reduced to MinBinSize peers thus leaving slots to newly
// connecting peers
-func (self *Kademlia) Prune(c <-chan time.Time) {
+func (k *Kademlia) Prune(c <-chan time.Time) {
go func() {
for range c {
total := 0
- self.conns.EachBin(pot.NewBytesVal(self.base, nil), 0, func(po, size int, f func(func(pot.PotVal, int) bool) bool) bool {
- extra := size - self.MinBinSize
- if size > self.MaxBinSize {
+ k.conns.EachBin(k.base, 0, func(po, size int, f func(func(pot.Val, int) bool) bool) bool {
+ extra := size - k.MinBinSize
+ if size > k.MaxBinSize {
n := 0
- f(func(v pot.PotVal, po int) bool {
+ f(func(v pot.Val, po int) bool {
v.(*entry).conn().Drop(fmt.Errorf("bucket full"))
n++
return n < extra
@@ -511,25 +499,28 @@ func (self *Kademlia) Prune(c <-chan time.Time) {
}
return true
})
- log.Debug(fmt.Sprintf("pruned %v peers", total))
+ log.Trace(fmt.Sprintf("pruned %v peers", total))
}
}()
}
+// NewPeerPot just creates a new pot record OverlayAddr
func NewPeerPot(kadMinProxSize int, ids ...discover.NodeID) map[discover.NodeID][][]byte {
// create a table of all nodes for health check
- np := pot.NewPot(nil, 0)
+ np := pot.NewPot(nil, 0, nil)
for _, id := range ids {
o := ToOverlayAddr(id.Bytes())
- np, _, _ = pot.Add(np, pot.NewBytesVal(o, nil))
+ np, _, _ = pot.Add(np, o)
}
nnmap := make(map[discover.NodeID][][]byte)
for _, id := range ids {
pl := 0
var nns [][]byte
- np.EachNeighbour(pot.NewBytesVal(id.Bytes(), nil), func(val pot.PotVal, po int) bool {
- a := val.(pot.BytesAddress).Address()
+ np.EachNeighbour(id.Bytes(), func(val pot.Val, po int) bool {
+ // a := val.(pot.BytesAddress).Address()
+ // nns = append(nns, a)
+ a := val.([]byte)
nns = append(nns, a)
if len(nns) >= kadMinProxSize {
pl = po
@@ -541,9 +532,10 @@ func NewPeerPot(kadMinProxSize int, ids ...discover.NodeID) map[discover.NodeID]
return nnmap
}
-func (self *Kademlia) FirstEmptyBin() (i int) {
+// FirstEmptyBin returns the farthest proximity order (int) that has no peer records
+func (k *Kademlia) FirstEmptyBin() (i int) {
i = -1
- self.conns.EachBin(pot.NewBytesVal(self.base, nil), 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool {
+ k.conns.EachBin(k.base, 0, func(po, size int, f func(func(val pot.Val, i int) bool) bool) bool {
if po > i+1 {
i = po
return false
@@ -554,22 +546,22 @@ func (self *Kademlia) FirstEmptyBin() (i int) {
return i
}
-func (self *Kademlia) Full() bool {
- return self.FirstEmptyBin() >= self.Depth()
+// Full returns a bool if the kademlia table is healthy and complete
+func (k *Kademlia) Full() bool {
+ return k.FirstEmptyBin() >= k.Depth()
}
// Healthy reports the health state of the kademlia connectivity
-//
-func (self *Kademlia) Healthy(peers [][]byte) bool {
- return self.gotNearestNeighbours(peers) && self.Full()
+func (k *Kademlia) Healthy(peers [][]byte) bool {
+ return k.gotNearestNeighbours(peers) && k.Full()
}
-func (self *Kademlia) gotNearestNeighbours(peers [][]byte) (got bool) {
+func (k *Kademlia) gotNearestNeighbours(peers [][]byte) (got bool) {
pm := make(map[string]bool)
for _, p := range peers {
pm[string(p)] = true
}
- self.EachConn(nil, 255, func(p OverlayConn, po int, nn bool) bool {
+ k.EachConn(nil, 255, func(p OverlayConn, po int, nn bool) bool {
if !nn {
return false
}
diff --git a/swarm/network/kademlia_test.go b/swarm/network/kademlia_test.go
index 1d1ebfba17..a59c6370ab 100644
--- a/swarm/network/kademlia_test.go
+++ b/swarm/network/kademlia_test.go
@@ -28,12 +28,13 @@ import (
func init() {
h := log.LvlFilterHandler(log.LvlTrace, log.StreamHandler(os.Stderr, log.TerminalFormat(true)))
+ // h := log.LvlFilterHandler(log.LvlInfo, log.StreamHandler(os.Stderr, log.TerminalFormat(true)))
// h := log.CallerFileHandler(log.StreamHandler(os.Stderr, log.TerminalFormat(true)))
log.Root().SetHandler(h)
}
func testKadPeerAddr(s string) *bzzAddr {
- a := pot.NewHashAddress(s).Bytes()
+ a := pot.NewAddressFromString(s)
return &bzzAddr{OAddr: a, UAddr: a}
}
@@ -54,7 +55,7 @@ type testPeerNotification struct {
po uint8
}
-type testProxNotification struct {
+type testDepthNotification struct {
rec string
po uint8
}
@@ -99,7 +100,7 @@ func newTestKademlia(b string) *testKademlia {
params := NewKadParams()
params.MinBinSize = 1
params.MinProxBinSize = 2
- base := pot.NewHashAddress(b).Bytes()
+ base := pot.NewAddressFromString(b)
return &testKademlia{
NewKademlia(base, params),
false,
@@ -119,8 +120,7 @@ func (k *testKademlia) newTestKadPeer(s string) Peer {
func (k *testKademlia) On(ons ...string) *testKademlia {
for _, s := range ons {
- p := k.newTestKadPeer(s)
- k.Kademlia.On(p)
+ k.Kademlia.On(k.newTestKadPeer(s).(OverlayConn))
}
return k
}
@@ -158,6 +158,7 @@ func testSuggestPeer(t *testing.T, k *testKademlia, expAddr string, expPo int, e
if want != expWant {
return fmt.Errorf("expected SuggestPeer to want peers: %v", expWant)
}
+ // t.Logf("%v", k)
return nil
}
@@ -165,50 +166,50 @@ func binStr(a OverlayPeer) string {
if a == nil {
return ""
}
- return pot.ToBin(a.Address())[:6]
+ return pot.ToBin(a.Address())[:8]
}
func TestSuggestPeerFindPeers(t *testing.T) {
// 2 row gap, unsaturated proxbin, no callables -> want PO 0
- k := newTestKademlia("000000").On("001000")
- // k.MinProxBinSize = 2
- // k.MinBinSize = 2
+ k := newTestKademlia("00000000").On("00100000")
err := testSuggestPeer(t, k, "", 0, true)
if err != nil {
t.Fatal(err.Error())
}
// 2 row gap, saturated proxbin, no callables -> want PO 0
- k.On("000100")
+ k.On("00010000")
err = testSuggestPeer(t, k, "", 0, true)
if err != nil {
t.Fatal(err.Error())
}
// 1 row gap (1 less), saturated proxbin, no callables -> want PO 1
- k.On("100000")
+ k.On("10000000")
err = testSuggestPeer(t, k, "", 1, true)
if err != nil {
t.Fatal(err.Error())
}
// no gap (1 less), saturated proxbin, no callables -> do not want more
- k.On("010000", "001001")
+ k.On("01000000", "00100001")
err = testSuggestPeer(t, k, "", 0, false)
if err != nil {
t.Fatal(err.Error())
}
// oversaturated proxbin, > do not want more
- k.On("001001")
+ k.On("00100001")
err = testSuggestPeer(t, k, "", 0, false)
if err != nil {
t.Fatal(err.Error())
}
// reintroduce gap, disconnected peer callable
- k.Off("010000")
- err = testSuggestPeer(t, k, "010000", 0, false)
+ log.Info(k.String())
+ k.Off("01000000")
+ log.Info(k.String())
+ err = testSuggestPeer(t, k, "01000000", 0, false)
if err != nil {
t.Fatal(err.Error())
}
@@ -221,31 +222,33 @@ func TestSuggestPeerFindPeers(t *testing.T) {
}
// on and off again, peer callable again
- k.On("010000")
- k.Off("010000")
- err = testSuggestPeer(t, k, "010000", 0, false)
+ k.On("01000000")
+ k.Off("01000000")
+ err = testSuggestPeer(t, k, "01000000", 0, false)
if err != nil {
t.Fatal(err.Error())
}
- k.On("010000")
+ k.On("01000000")
// new closer peer appears, it is immediately wanted
- k.Register("000101")
- err = testSuggestPeer(t, k, "000101", 0, false)
+ k.Register("00010001")
+ err = testSuggestPeer(t, k, "00010001", 0, false)
if err != nil {
t.Fatal(err.Error())
}
// PO1 disconnects
- k.On("000101")
- k.Off("010000")
+ k.On("00010001")
+ log.Info(k.String())
+ k.Off("01000000")
+ log.Info(k.String())
// second time, gap filling
- err = testSuggestPeer(t, k, "010000", 0, false)
+ err = testSuggestPeer(t, k, "01000000", 0, false)
if err != nil {
t.Fatal(err.Error())
}
- k.On("010000")
+ k.On("01000000")
err = testSuggestPeer(t, k, "", 0, false)
if err != nil {
t.Fatal(err.Error())
@@ -257,53 +260,53 @@ func TestSuggestPeerFindPeers(t *testing.T) {
t.Fatal(err.Error())
}
- k.Register("010001")
+ k.Register("01000001")
err = testSuggestPeer(t, k, "", 0, true)
if err != nil {
t.Fatal(err.Error())
}
- k.On("100001")
+ k.On("10000001")
log.Trace("Kad:\n%v", k.String())
- err = testSuggestPeer(t, k, "010001", 0, false)
+ err = testSuggestPeer(t, k, "01000001", 0, false)
if err != nil {
t.Fatal(err.Error())
}
- k.On("100001")
- k.On("010001")
+ k.On("10000001")
+ k.On("01000001")
err = testSuggestPeer(t, k, "", 0, false)
if err != nil {
t.Fatal(err.Error())
}
k.MinBinSize = 3
- k.Register("100010")
- err = testSuggestPeer(t, k, "100010", 0, false)
+ k.Register("10000010")
+ err = testSuggestPeer(t, k, "10000010", 0, false)
if err != nil {
t.Fatal(err.Error())
}
- k.On("100010")
+ k.On("10000010")
err = testSuggestPeer(t, k, "", 1, true)
if err != nil {
t.Fatal(err.Error())
}
- k.On("010010")
+ k.On("01000010")
err = testSuggestPeer(t, k, "", 2, true)
if err != nil {
t.Fatal(err.Error())
}
- k.On("001010")
+ k.On("00100010")
err = testSuggestPeer(t, k, "", 3, true)
if err != nil {
log.Trace("Kad:\n%v", k.String())
t.Fatal(err.Error())
}
- k.On("000110")
+ k.On("00010010")
err = testSuggestPeer(t, k, "", 0, false)
if err != nil {
log.Trace("Kad:\n%v", k.String())
@@ -314,14 +317,14 @@ func TestSuggestPeerFindPeers(t *testing.T) {
func TestSuggestPeerRetries(t *testing.T) {
// 2 row gap, unsaturated proxbin, no callables -> want PO 0
- k := newTestKademlia("000000")
+ k := newTestKademlia("00000000")
cycle := 50 * time.Millisecond
k.RetryInterval = int(cycle)
k.MaxRetries = 3
k.RetryExponent = 3
- k.Register("010000")
- k.On("000001", "000010")
- err := testSuggestPeer(t, k, "010000", 0, false)
+ k.Register("01000000")
+ k.On("00000001", "00000010")
+ err := testSuggestPeer(t, k, "01000000", 0, false)
if err != nil {
t.Fatal(err.Error())
}
@@ -333,7 +336,7 @@ func TestSuggestPeerRetries(t *testing.T) {
// cycle *= time.Duration(k.RetryExponent)
time.Sleep(cycle)
- err = testSuggestPeer(t, k, "010000", 0, false)
+ err = testSuggestPeer(t, k, "01000000", 0, false)
if err != nil {
t.Fatal(err.Error())
}
@@ -345,7 +348,7 @@ func TestSuggestPeerRetries(t *testing.T) {
cycle *= time.Duration(k.RetryExponent)
time.Sleep(cycle)
- err = testSuggestPeer(t, k, "010000", 0, false)
+ err = testSuggestPeer(t, k, "01000000", 0, false)
if err != nil {
t.Fatal(err.Error())
}
@@ -357,7 +360,7 @@ func TestSuggestPeerRetries(t *testing.T) {
cycle *= time.Duration(k.RetryExponent)
time.Sleep(cycle)
- err = testSuggestPeer(t, k, "010000", 0, false)
+ err = testSuggestPeer(t, k, "01000000", 0, false)
if err != nil {
t.Fatal(err.Error())
}
@@ -378,10 +381,10 @@ func TestSuggestPeerRetries(t *testing.T) {
}
func TestPruning(t *testing.T) {
- k := newTestKademlia("000000")
- k.On("100000", "110000", "101000", "100100", "100010")
- k.On("010000", "011000", "010100", "010010", "010001")
- k.On("001000", "001100", "001010", "001001")
+ k := newTestKademlia("00000000")
+ k.On("10000000", "11000000", "10100000", "10010000", "10000010")
+ k.On("01000000", "01100000", "01000100", "01000010", "01000001")
+ k.On("00100000", "00110000", "00100010", "00100001")
k.MaxBinSize = 4
k.MinBinSize = 3
prune := make(chan time.Time)
@@ -411,10 +414,10 @@ func TestPruning(t *testing.T) {
// TODO: this is now based on just taking the first 2 peers
// in order of connecting
expDropped := []string{
- "101000",
- "110000",
- "010100",
- "011000",
+ "10100000",
+ "11000000",
+ "01000100",
+ "01100000",
}
for _, addr := range expDropped {
err := dropped[addr]
@@ -428,7 +431,7 @@ func TestPruning(t *testing.T) {
}
func TestKademliaHiveString(t *testing.T) {
- k := newTestKademlia("000000").On("010000", "001000").Register("100000", "100001")
+ k := newTestKademlia("00000000").On("01000000", "00100000").Register("10000000", "10000001")
h := k.String()
expH := "\n=========================================================================\nMon Feb 27 12:10:28 UTC 2017 KΛÐΞMLIΛ hive: queen's address: 000000\npopulation: 2 (4), MinProxBinSize: 2, MinBinSize: 1, MaxBinSize: 4\n============ PROX LIMIT: 0 ==========================================\n000 0 | 2 840000 800000\n001 1 400000 | 1 400000\n002 1 200000 | 1 200000\n003 0 | 0\n004 0 | 0\n005 0 | 0\n006 0 | 0\n007 0 | 0\n========================================================================="
if expH[100:] != h[100:] {
@@ -436,7 +439,7 @@ func TestKademliaHiveString(t *testing.T) {
}
}
-func (self *testKademlia) checkNotifications(npeers []*testPeerNotification, nprox []*testProxNotification) error {
+func (self *testKademlia) checkNotifications(npeers []*testPeerNotification, nprox []*testDepthNotification) error {
for _, pn := range npeers {
key := pn.rec + pn.addr
po, found := self.notifications[key]
@@ -460,50 +463,50 @@ func (self *testKademlia) checkNotifications(npeers []*testPeerNotification, npr
}
func TestNotifications(t *testing.T) {
- k := newTestKademlia("000000")
+ k := newTestKademlia("00000000")
k.Discovery = true
k.MinProxBinSize = 3
- k.On("010000", "001000")
+ k.On("01000000", "00100000")
time.Sleep(1000 * time.Millisecond)
err := k.checkNotifications(
[]*testPeerNotification{
- &testPeerNotification{"010000", "001000", 1},
+ &testPeerNotification{"01000000", "00100000", 1},
},
- []*testProxNotification{
- &testProxNotification{"001000", 0},
- &testProxNotification{"010000", 0},
+ []*testDepthNotification{
+ &testDepthNotification{"00100000", 0},
+ &testDepthNotification{"01000000", 0},
},
)
if err != nil {
t.Fatal(err.Error())
}
- k = k.On("100000")
+ k = k.On("10000000")
time.Sleep(100 * time.Millisecond)
k.checkNotifications(
[]*testPeerNotification{
- &testPeerNotification{"010000", "100000", 0},
- &testPeerNotification{"001000", "100000", 0},
+ &testPeerNotification{"01000000", "10000000", 0},
+ &testPeerNotification{"00100000", "10000000", 0},
},
- []*testProxNotification{
- &testProxNotification{"100000", 0},
+ []*testDepthNotification{
+ &testDepthNotification{"10000000", 0},
},
)
- k = k.On("010001")
+ k = k.On("01000001")
time.Sleep(100 * time.Millisecond)
k.checkNotifications(
[]*testPeerNotification{
- &testPeerNotification{"010000", "010001", 5},
- &testPeerNotification{"001000", "010001", 1},
- &testPeerNotification{"100000", "010001", 0},
+ &testPeerNotification{"01000000", "01000001", 5},
+ &testPeerNotification{"00100000", "01000001", 1},
+ &testPeerNotification{"10000000", "01000001", 0},
},
- []*testProxNotification{
- &testProxNotification{"100000", 0},
- &testProxNotification{"010000", 0},
- &testProxNotification{"010001", 0},
- &testProxNotification{"001000", 0},
+ []*testDepthNotification{
+ &testDepthNotification{"10000000", 0},
+ &testDepthNotification{"01000000", 0},
+ &testDepthNotification{"01000001", 0},
+ &testDepthNotification{"00100000", 0},
},
)
}
diff --git a/swarm/network/protocol.go b/swarm/network/protocol.go
index 0aea049426..8e9bb4f595 100644
--- a/swarm/network/protocol.go
+++ b/swarm/network/protocol.go
@@ -333,15 +333,15 @@ func RandomAddr() *bzzAddr {
pubkey := crypto.FromECDSAPub(&key.PublicKey)
var id discover.NodeID
copy(id[:], pubkey[1:])
- return &bzzAddr{
- OAddr: crypto.Keccak256(pubkey[1:]),
- UAddr: id[:],
- }
+ return NewAddrFromNodeID(id)
}
// NewNodeIDFromAddr transforms the underlay address to an adapters.NodeID
func NewNodeIDFromAddr(addr Addr) discover.NodeID {
- return discover.MustBytesID(addr.Under())
+ log.Info(fmt.Sprintf("uaddr=%s", string(addr.Under())))
+ node := discover.MustParseNode(string(addr.Under()))
+ // return discover.MustBytesID(addr.Under())
+ return node.ID
}
// NewAddrFromNodeID constucts a bzzAddr from a discover.NodeID
diff --git a/swarm/network/simulations/discovery/discovery_test.go b/swarm/network/simulations/discovery/discovery_test.go
index fb038ed2a7..7bfb2aef72 100644
--- a/swarm/network/simulations/discovery/discovery_test.go
+++ b/swarm/network/simulations/discovery/discovery_test.go
@@ -57,7 +57,7 @@ func TestDiscoverySimulationSimAdapter(t *testing.T) {
}
func testDiscoverySimulation(t *testing.T, adapter adapters.NodeAdapter) {
- // create 10 node network
+ // create network
nodeCount := 10
net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{
ID: "0",
diff --git a/swarm/network/simulations/overlay.go b/swarm/network/simulations/overlay.go
index ede71057c6..6570d30ac0 100644
--- a/swarm/network/simulations/overlay.go
+++ b/swarm/network/simulations/overlay.go
@@ -209,8 +209,9 @@ func main() {
config := &simulations.ServerConfig{
NewAdapter: func() adapters.NodeAdapter { return adapters.NewSimAdapter(services) },
- DefaultMockerID: "bootNet",
- Mockers: mockers,
+ DefaultMockerID: "randomNodes",
+ // DefaultMockerID: "bootNet",
+ Mockers: mockers,
}
log.Info("starting simulation server on 0.0.0.0:8888...")
diff --git a/swarm/pss/pss.go b/swarm/pss/pss.go
index 58e6df735c..9ac35d9f47 100644
--- a/swarm/pss/pss.go
+++ b/swarm/pss/pss.go
@@ -133,8 +133,8 @@ func (self *Pss) Protocols() []p2p.Protocol {
func (self *Pss) Run(p *p2p.Peer, rw p2p.MsgReadWriter) error {
pp := protocols.NewPeer(p, rw, pssSpec)
id := p.ID()
- h := pot.NewHashAddressFromBytes(network.ToOverlayAddr(id[:]))
- self.fwdPool[h.Address] = pp
+ a := pot.NewAddressFromBytes(network.ToOverlayAddr(id[:]))
+ self.fwdPool[a] = pp
return pp.Run(self.handlePssMsg)
}
@@ -315,9 +315,8 @@ func (self *Pss) Forward(msg *PssMsg) error {
// send with kademlia
// find the closest peer to the recipient and attempt to send
self.Overlay.EachConn(msg.To, 256, func(op network.OverlayConn, po int, isproxbin bool) bool {
- //p, ok := op.(senderPeer)
- h := pot.NewHashAddressFromBytes(op.Address())
- pp := self.fwdPool[h.Address]
+ a := pot.NewAddressFromBytes(op.Address())
+ pp := self.fwdPool[a]
addr := self.Overlay.BaseAddr()
sendMsg := fmt.Sprintf("%x: msg to %x via %x", common.ByteLabel(addr), common.ByteLabel(msg.To), common.ByteLabel(op.Address()))
if self.checkFwdCache(op.Address(), digest) {
@@ -459,7 +458,7 @@ func RegisterPssProtocol(ps *Pss, topic *Topic, spec *protocols.Spec, targetprot
}
func (self *PssProtocol) handle(msg []byte, p *p2p.Peer, senderAddr []byte) error {
- hashoaddr := pot.NewHashAddressFromBytes(senderAddr).Address
+ hashoaddr := pot.NewAddressFromBytes(senderAddr)
if !self.isActive(hashoaddr, *self.topic) {
rw := &PssReadWriter{
Pss: self.Pss,
diff --git a/swarm/pss/pss_test.go b/swarm/pss/pss_test.go
index 875771b463..3ee8af3116 100644
--- a/swarm/pss/pss_test.go
+++ b/swarm/pss/pss_test.go
@@ -199,8 +199,8 @@ func TestSimpleLinear(t *testing.T) {
Peer: pp,
addr: network.ToOverlayAddr(id[:]),
}
- h := pot.NewHashAddressFromBytes(bp.addr)
- ps.fwdPool[h.Address] = pp
+ a := pot.NewAddressFromBytes(bp.addr)
+ ps.fwdPool[a] = pp
ps.Overlay.On(bp)
defer ps.Overlay.Off(bp)
log.Debug(fmt.Sprintf("%v", ps.Overlay))