From 3778829de85b66e36bfe44c62c60ecdc25e924ac Mon Sep 17 00:00:00 2001 From: Cole Brown Date: Thu, 14 Jun 2018 16:17:50 -0400 Subject: [PATCH] Abstract Peerstore tests, fix bugs --- addr_manager.go | 15 ++++++++-- addr_manager_ds.go | 6 ++-- addr_manager_test.go | 29 +++++++++++------- peerstore_test.go | 70 +++++++++++++++++++++++++++++++++----------- 4 files changed, 87 insertions(+), 33 deletions(-) diff --git a/addr_manager.go b/addr_manager.go index c916197..09319af 100644 --- a/addr_manager.go +++ b/addr_manager.go @@ -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)) diff --git a/addr_manager_ds.go b/addr_manager_ds.go index ed39709..125d6f0 100644 --- a/addr_manager_ds.go +++ b/addr_manager_ds.go @@ -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 { diff --git a/addr_manager_test.go b/addr_manager_test.go index c7b9b5f..1f49188 100644 --- a/addr_manager_test.go +++ b/addr_manager_test.go @@ -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 } diff --git a/peerstore_test.go b/peerstore_test.go index 55a8e90..1626132 100644 --- a/peerstore_test.go +++ b/peerstore_test.go @@ -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()