mirror of
https://github.com/libp2p/go-libp2p-peerstore.git
synced 2026-08-22 15:33:27 +08:00
Abstract Peerstore tests, fix bugs
This commit is contained in:
@@ -250,6 +250,10 @@ func (mgr *AddrManager) ClearAddrs(p peer.ID) {
|
||||
}
|
||||
|
||||
func (mgr *AddrManager) AddrStream(ctx context.Context, p peer.ID) <-chan ma.Multiaddr {
|
||||
mgr.addrmu.Lock()
|
||||
defer mgr.addrmu.Unlock()
|
||||
mgr.init()
|
||||
|
||||
baseaddrslice := mgr.addrs[p]
|
||||
initial := make([]ma.Multiaddr, 0, len(baseaddrslice))
|
||||
for _, a := range baseaddrslice {
|
||||
@@ -266,6 +270,7 @@ type AddrSubManager struct {
|
||||
|
||||
func NewAddrSubManager() *AddrSubManager {
|
||||
return &AddrSubManager{
|
||||
mu: sync.RWMutex{},
|
||||
subs: make(map[peer.ID][]*addrSub),
|
||||
}
|
||||
}
|
||||
@@ -292,8 +297,8 @@ func (mgr *AddrSubManager) removeSub(p peer.ID, s *addrSub) {
|
||||
}
|
||||
|
||||
func (mgr *AddrSubManager) BroadcastAddr(p peer.ID, addr ma.Multiaddr) {
|
||||
mgr.mu.RLock()
|
||||
defer mgr.mu.RUnlock()
|
||||
mgr.mu.Lock()
|
||||
defer mgr.mu.Unlock()
|
||||
|
||||
subs, ok := mgr.subs[p]
|
||||
if !ok {
|
||||
@@ -313,7 +318,11 @@ func (mgr *AddrSubManager) AddrStream(ctx context.Context, p peer.ID, initial []
|
||||
|
||||
out := make(chan ma.Multiaddr)
|
||||
|
||||
mgr.subs[p] = append(mgr.subs[p], sub)
|
||||
if _, ok := mgr.subs[p]; ok {
|
||||
mgr.subs[p] = append(mgr.subs[p], sub)
|
||||
} else {
|
||||
mgr.subs[p] = []*addrSub{sub}
|
||||
}
|
||||
|
||||
sort.Sort(addr.AddrList(initial))
|
||||
|
||||
|
||||
@@ -36,11 +36,11 @@ func peerAddressKey(p *peer.ID, addr *ma.Multiaddr) (ds.Key, error) {
|
||||
if err != nil {
|
||||
return ds.Key{}, nil
|
||||
}
|
||||
return ds.NewKey(p.Pretty()).ChildString(hash.B58String()), nil
|
||||
return ds.NewKey(peer.IDB58Encode(*p)).ChildString(hash.B58String()), nil
|
||||
}
|
||||
|
||||
func peerIDFromKey(key ds.Key) (peer.ID, error) {
|
||||
idstring := key.Parent().Type()
|
||||
idstring := key.Parent().Name()
|
||||
return peer.IDB58Decode(idstring)
|
||||
}
|
||||
|
||||
@@ -78,7 +78,7 @@ func (mgr *DatastoreAddrManager) SetAddrs(p peer.ID, addrs []ma.Multiaddr, ttl t
|
||||
mgr.ds.Delete(key)
|
||||
continue
|
||||
}
|
||||
if has, err := mgr.ds.Has(key); err != nil && has == false {
|
||||
if has, err := mgr.ds.Has(key); err != nil || !has {
|
||||
mgr.subsManager.BroadcastAddr(p, addr)
|
||||
}
|
||||
if err := mgr.ds.Put(key, addr.Bytes()); err != nil {
|
||||
|
||||
@@ -9,6 +9,7 @@ import (
|
||||
|
||||
"context"
|
||||
|
||||
"github.com/ipfs/go-datastore"
|
||||
"github.com/ipfs/go-ds-badger"
|
||||
"github.com/libp2p/go-libp2p-peer"
|
||||
ma "github.com/multiformats/go-multiaddr"
|
||||
@@ -54,6 +55,22 @@ func testHas(t *testing.T, exp, act []ma.Multiaddr) {
|
||||
}
|
||||
}
|
||||
|
||||
func setupBadgerDatastore(t *testing.T) (datastore.Datastore, func()) {
|
||||
dataPath, err := ioutil.TempDir(os.TempDir(), "badger")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ds, err := badger.NewDatastore(dataPath, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
closer := func() {
|
||||
ds.Close()
|
||||
os.RemoveAll(dataPath)
|
||||
}
|
||||
return ds, closer
|
||||
}
|
||||
|
||||
func setupBadgerAddrManager(t *testing.T) (*BadgerAddrManager, func()) {
|
||||
dataPath, err := ioutil.TempDir(os.TempDir(), "badger")
|
||||
if err != nil {
|
||||
@@ -71,19 +88,11 @@ func setupBadgerAddrManager(t *testing.T) (*BadgerAddrManager, func()) {
|
||||
}
|
||||
|
||||
func setupDatastoreAddrManager(t *testing.T) (*DatastoreAddrManager, func()) {
|
||||
dataPath, err := ioutil.TempDir(os.TempDir(), "badger")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ds, err := badger.NewDatastore(dataPath, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ds, closeDB := setupBadgerDatastore(t)
|
||||
mgr := NewDatastoreAddrManager(context.Background(), ds, 100*time.Microsecond)
|
||||
closer := func() {
|
||||
mgr.Stop()
|
||||
ds.Close()
|
||||
os.RemoveAll(dataPath)
|
||||
closeDB()
|
||||
}
|
||||
return mgr, closer
|
||||
}
|
||||
|
||||
@@ -10,7 +10,8 @@ import (
|
||||
|
||||
"os"
|
||||
|
||||
peer "github.com/libp2p/go-libp2p-peer"
|
||||
"github.com/libp2p/go-libp2p-crypto"
|
||||
"github.com/libp2p/go-libp2p-peer"
|
||||
ma "github.com/multiformats/go-multiaddr"
|
||||
)
|
||||
|
||||
@@ -27,13 +28,34 @@ func getAddrs(t *testing.T, n int) []ma.Multiaddr {
|
||||
return addrs
|
||||
}
|
||||
|
||||
func TestAddrStream(t *testing.T) {
|
||||
func runTestWithPeerstores(t *testing.T, testFunc func(*testing.T, Peerstore)) {
|
||||
t.Helper()
|
||||
t.Log("NewPeerstore")
|
||||
ps1 := NewPeerstore()
|
||||
testFunc(t, ps1)
|
||||
|
||||
t.Log("NewPeerstoreDatastore")
|
||||
ps2, closer2 := setupDatastorePeerstore(t)
|
||||
defer closer2()
|
||||
testFunc(t, ps2)
|
||||
}
|
||||
|
||||
func setupDatastorePeerstore(t *testing.T) (Peerstore, func()) {
|
||||
ds, closeDB := setupBadgerDatastore(t)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
ps := NewPeerstoreDatastore(ctx, ds)
|
||||
closer := func() {
|
||||
cancel()
|
||||
closeDB()
|
||||
}
|
||||
return ps, closer
|
||||
}
|
||||
|
||||
func testAddrStream(t *testing.T, ps Peerstore) {
|
||||
addrs := getAddrs(t, 100)
|
||||
|
||||
pid := peer.ID("testpeer")
|
||||
|
||||
ps := NewPeerstore()
|
||||
|
||||
ps.AddAddrs(pid, addrs[:10], time.Hour)
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
@@ -105,12 +127,14 @@ func TestAddrStream(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestGetStreamBeforePeerAdded(t *testing.T) {
|
||||
func TestAddrStream(t *testing.T) {
|
||||
runTestWithPeerstores(t, testAddrStream)
|
||||
}
|
||||
|
||||
func testGetStreamBeforePeerAdded(t *testing.T, ps Peerstore) {
|
||||
addrs := getAddrs(t, 10)
|
||||
pid := peer.ID("testpeer")
|
||||
|
||||
ps := NewPeerstore()
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
@@ -156,12 +180,14 @@ func TestGetStreamBeforePeerAdded(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAddrStreamDuplicates(t *testing.T) {
|
||||
func TestGetStreamBeforePeerAdded(t *testing.T) {
|
||||
runTestWithPeerstores(t, testGetStreamBeforePeerAdded)
|
||||
}
|
||||
|
||||
func testAddrStreamDuplicates(t *testing.T, ps Peerstore) {
|
||||
addrs := getAddrs(t, 10)
|
||||
pid := peer.ID("testpeer")
|
||||
|
||||
ps := NewPeerstore()
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
ach := ps.AddrStream(ctx, pid)
|
||||
@@ -195,8 +221,11 @@ func TestAddrStreamDuplicates(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestPeerstoreProtoStore(t *testing.T) {
|
||||
ps := NewPeerstore()
|
||||
func TestAddrStreamDuplicates(t *testing.T) {
|
||||
runTestWithPeerstores(t, testAddrStreamDuplicates)
|
||||
}
|
||||
|
||||
func testPeerstoreProtoStore(t *testing.T, ps Peerstore) {
|
||||
p1 := peer.ID("TESTPEER")
|
||||
|
||||
protos := []string{"a", "b", "c", "d"}
|
||||
@@ -250,20 +279,23 @@ func TestPeerstoreProtoStore(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestBasicPeerstore(t *testing.T) {
|
||||
ps := NewPeerstore()
|
||||
func TestPeerstoreProtoStore(t *testing.T) {
|
||||
runTestWithPeerstores(t, testAddrStreamDuplicates)
|
||||
}
|
||||
|
||||
func testBasicPeerstore(t *testing.T, ps Peerstore) {
|
||||
var pids []peer.ID
|
||||
addrs := getAddrs(t, 10)
|
||||
for i, a := range addrs {
|
||||
p := peer.ID(fmt.Sprint(i))
|
||||
for _, a := range addrs {
|
||||
priv, _, _ := crypto.GenerateKeyPair(crypto.RSA, 512)
|
||||
p, _ := peer.IDFromPrivateKey(priv)
|
||||
pids = append(pids, p)
|
||||
ps.AddAddr(p, a, PermanentAddrTTL)
|
||||
}
|
||||
|
||||
peers := ps.Peers()
|
||||
if len(peers) != 10 {
|
||||
t.Fatal("expected ten peers")
|
||||
t.Fatal("expected ten peers, got", len(peers))
|
||||
}
|
||||
|
||||
pinfo := ps.PeerInfo(pids[0])
|
||||
@@ -272,6 +304,10 @@ func TestBasicPeerstore(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestBasicPeerstore(t *testing.T) {
|
||||
runTestWithPeerstores(t, testBasicPeerstore)
|
||||
}
|
||||
|
||||
func BenchmarkBasicPeerstore(b *testing.B) {
|
||||
ps := NewPeerstore()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user