@cryptotaxi247 / kubo / commits / ac62d13e4

peerstore Put -> Add

Changed lots of peer use, and changed the peerstore to ensure there is only ever one peer in use. Fixed #174

Juan Batiz-Benet committed Oct 20, 2014 at 06:37 UTC ac62d13e426b25d611be3a60e474163a00830bd6
11 files changed +107 -29
core/core.go
+1 -6
@@ -233,15 +233,10 @@ func initConnections(ctx context.Context, cfg *config.Config, pstore peer.Peerst
233 // setup peer
234 npeer, err := pstore.Get(peer.DecodePrettyID(p.PeerID))
235 if err != nil {
236 - log.Error("%s", err)
237 - continue
238 - }
239 - npeer.AddAddress(maddr)
240 -
241 - if err = pstore.Put(npeer); err != nil {
236 log.Error("Bootstrapping error: %v", err)
237 continue
238 }
239 + npeer.AddAddress(maddr)
240
241 if _, err = route.Connect(ctx, npeer); err != nil {
242 log.Error("Bootstrapping error: %v", err)
core/mock.go
+6 -2
@@ -22,12 +22,16 @@ func NewMockNode() (*IpfsNode, error) {
22 return nil, err
23 }
24
25 - nd.Identity, err = peer.WithKeyPair(sk, pk)
25 + p, err := peer.WithKeyPair(sk, pk)
26 if err != nil {
27 return nil, err
28 }
29 +
30 nd.Peerstore = peer.NewPeerstore()
30 - nd.Peerstore.Put(nd.Identity)
31 + nd.Identity, err = nd.Peerstore.Add(p)
32 + if err != nil {
33 + return nil, err
34 + }
35
36 // Temp Datastore
37 dstore := ds.NewMapDatastore()
net/conn/dial.go
+5 -4
@@ -23,6 +23,11 @@ func (d *Dialer) Dial(ctx context.Context, network string, remote peer.Peer) (Co
23 return nil, fmt.Errorf("No remote address for network %s", network)
24 }
25
26 + remote, err := d.Peerstore.Add(remote)
27 + if err != nil {
28 + log.Error("Error putting peer into peerstore: %s", remote)
29 + }
30 +
31 // TODO: try to get reusing addr/ports to work.
32 // madialer := manet.Dialer{LocalAddr: laddr}
33 madialer := manet.Dialer{}
@@ -33,10 +38,6 @@ func (d *Dialer) Dial(ctx context.Context, network string, remote peer.Peer) (Co
38 return nil, err
39 }
40
36 - if err := d.Peerstore.Put(remote); err != nil {
37 - log.Error("Error putting peer into peerstore: %s", remote)
38 - }
39 -
41 c, err := newSingleConn(ctx, d.LocalPeer, remote, maconn)
42 if err != nil {
43 return nil, err
net/conn/dial_test.go
+14 -4
@@ -68,13 +68,18 @@ func setupConn(t *testing.T, ctx context.Context, a1, a2 string) (a, b Conn) {
68 t.Fatal("Listen address is nil.")
69 }
70
71 - l1, err := Listen(ctx, laddr, p1, peer.NewPeerstore())
71 + ps1 := peer.NewPeerstore()
72 + ps2 := peer.NewPeerstore()
73 + ps1.Add(p1)
74 + ps2.Add(p2)
75 +
76 + l1, err := Listen(ctx, laddr, p1, ps1)
77 if err != nil {
78 t.Fatal(err)
79 }
80
81 d2 := &Dialer{
77 - Peerstore: peer.NewPeerstore(),
82 + Peerstore: ps2,
83 LocalPeer: p2,
84 }
85
@@ -108,7 +113,12 @@ func TestDialer(t *testing.T) {
113 t.Fatal("Listen address is nil.")
114 }
115
111 - l, err := Listen(ctx, laddr, p1, peer.NewPeerstore())
116 + ps1 := peer.NewPeerstore()
117 + ps2 := peer.NewPeerstore()
118 + ps1.Add(p1)
119 + ps2.Add(p2)
120 +
121 + l, err := Listen(ctx, laddr, p1, ps1)
122 if err != nil {
123 t.Fatal(err)
124 }
@@ -116,7 +126,7 @@ func TestDialer(t *testing.T) {
126 go echoListen(ctx, l)
127
128 d := &Dialer{
119 - Peerstore: peer.NewPeerstore(),
129 + Peerstore: ps2,
130 LocalPeer: p2,
131 }
132
net/conn/multiconn.go
+4 -4
@@ -72,10 +72,10 @@ func (c *MultiConn) Add(conns ...Conn) {
72 log.Error("%s", c2)
73 c.Unlock() // ok to unlock (to log). panicing.
74 log.Error("%s", c)
75 - log.Error("c.LocalPeer: %s %#v", c.LocalPeer(), c.LocalPeer())
76 - log.Error("c2.LocalPeer: %s %#v", c2.LocalPeer(), c2.LocalPeer())
77 - log.Error("c.RemotePeer: %s %#v", c.RemotePeer(), c.RemotePeer())
78 - log.Error("c2.RemotePeer: %s %#v", c2.RemotePeer(), c2.RemotePeer())
75 + log.Error("c.LocalPeer: %s %p", c.LocalPeer(), c.LocalPeer())
76 + log.Error("c2.LocalPeer: %s %p", c2.LocalPeer(), c2.LocalPeer())
77 + log.Error("c.RemotePeer: %s %p", c.RemotePeer(), c.RemotePeer())
78 + log.Error("c2.RemotePeer: %s %p", c2.RemotePeer(), c2.RemotePeer())
79 c.Lock() // gotta relock to avoid lock panic from deferring.
80 panic("connection addresses mismatch")
81 }
net/conn/multiconn_test.go
+2
@@ -95,6 +95,8 @@ func setupMultiConns(t *testing.T, ctx context.Context) (a, b *MultiConn) {
95 // peerstores
96 p1ps := peer.NewPeerstore()
97 p2ps := peer.NewPeerstore()
98 + p1ps.Add(p1)
99 + p2ps.Add(p2)
100
101 // listeners
102 listen := func(addr ma.Multiaddr, p peer.Peer, ps peer.Peerstore) Listener {
net/swarm/swarm.go
+1 -1
@@ -116,7 +116,7 @@ func (s *Swarm) Dial(peer peer.Peer) (conn.Conn, error) {
116 }
117
118 // check if we don't have the peer in Peerstore
119 - err := s.peers.Put(peer)
119 + peer, err := s.peers.Add(peer)
120 if err != nil {
121 return nil, err
122 }
peer/peer.go
+35 -1
@@ -1,6 +1,7 @@
1 package peer
2
3 import (
4 + "errors"
5 "fmt"
6 "sync"
7 "time"
@@ -83,6 +84,9 @@ type Peer interface {
84 // Get/SetLatency manipulate the current latency measurement.
85 GetLatency() (out time.Duration)
86 SetLatency(laten time.Duration)
87 +
88 + // Update with the data of another peer instance
89 + Update(Peer) error
90 }
91
92 type peer struct {
@@ -125,7 +129,9 @@ func (p *peer) PubKey() ic.PubKey {
129 // Addresses returns the peer's multiaddrs
130 func (p *peer) Addresses() []ma.Multiaddr {
131 cp := make([]ma.Multiaddr, len(p.addresses))
132 + p.RLock()
133 copy(cp, p.addresses)
134 + defer p.RUnlock()
135 return cp
136 }
137
@@ -182,7 +188,6 @@ func (p *peer) SetLatency(laten time.Duration) {
188 // LoadAndVerifyKeyPair unmarshalls, loads a private/public key pair.
189 // Error if (a) unmarshalling fails, or (b) pubkey does not match id.
190 func (p *peer) LoadAndVerifyKeyPair(marshalled []byte) error {
185 -
191 sk, err := ic.UnmarshalPrivateKey(marshalled)
192 if err != nil {
193 return fmt.Errorf("Failed to unmarshal private key: %v", err)
@@ -199,6 +204,9 @@ func (p *peer) VerifyAndSetPrivKey(sk ic.PrivKey) error {
204 return err
205 }
206
207 + p.Lock()
208 + defer p.Unlock()
209 +
210 // if we didn't have the priavte key, assign it
211 if p.privKey == nil {
212 p.privKey = sk
@@ -224,6 +232,9 @@ func (p *peer) VerifyAndSetPubKey(pk ic.PubKey) error {
232 return fmt.Errorf("Failed to hash public key: %v", err)
233 }
234
235 + p.Lock()
236 + defer p.Unlock()
237 +
238 if !p.id.Equal(pkid) {
239 return fmt.Errorf("Public key does not match peer.ID.")
240 }
@@ -246,6 +257,29 @@ func (p *peer) VerifyAndSetPubKey(pk ic.PubKey) error {
257 panic("invariant violated: unexpected key mismatch")
258 }
259
260 +func (p *peer) Update(other Peer) error {
261 + if !p.ID().Equal(other.ID()) {
262 + return errors.New("peer ids do not match")
263 + }
264 +
265 + for _, a := range other.Addresses() {
266 + p.AddAddress(a)
267 + }
268 +
269 + sk := other.PrivKey()
270 + pk := other.PubKey()
271 + p.Lock()
272 + if p.privKey == nil {
273 + p.privKey = sk
274 + }
275 +
276 + if p.pubKey == nil {
277 + p.pubKey = pk
278 + }
279 + defer p.Unlock()
280 + return nil
281 +}
282 +
283 // WithKeyPair returns a Peer object with given keys.
284 func WithKeyPair(sk ic.PrivKey, pk ic.PubKey) (Peer, error) {
285 if sk == nil && pk == nil {
peer/peerstore.go
+28 -3
@@ -12,7 +12,7 @@ import (
12 // Peerstore provides a threadsafe collection for peers.
13 type Peerstore interface {
14 Get(ID) (Peer, error)
15 - Put(Peer) error
15 + Add(Peer) (Peer, error)
16 Delete(ID) error
17 All() (*Map, error)
18 }
@@ -63,12 +63,37 @@ func (p *peerstore) Get(i ID) (Peer, error) {
63 }
64 }
65
66 -func (p *peerstore) Put(peer Peer) error {
66 +func (p *peerstore) Add(peer Peer) (Peer, error) {
67 p.Lock()
68 defer p.Unlock()
69
70 k := peer.Key().DsKey()
71 - return p.peers.Put(k, peer)
71 + val, err := p.peers.Get(k)
72 + switch err {
73 + // some other datastore error
74 + default:
75 + return nil, err
76 +
77 + // not found? just add and return.
78 + case ds.ErrNotFound:
79 + err := p.peers.Put(k, peer)
80 + return peer, err
81 +
82 + // no error, already here.
83 + case nil:
84 + peer2, ok := val.(Peer)
85 + if !ok {
86 + return nil, errors.New("stored value was not a Peer")
87 + }
88 +
89 + if peer == peer2 {
90 + return peer, nil
91 + }
92 +
93 + // must do some merging.
94 + peer2.Update(peer)
95 + return peer2, nil
96 + }
97 }
98
99 func (p *peerstore) Delete(i ID) error {
peer/peerstore_test.go
+9 -2
@@ -27,11 +27,15 @@ func TestPeerstore(t *testing.T) {
27 // p31, _ := setupPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a33", "/ip4/127.0.0.1/tcp/3456")
28 // p41, _ := setupPeer("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a34", "/ip4/127.0.0.1/tcp/4567")
29
30 - err := ps.Put(p11)
30 + p13, err := ps.Add(p11)
31 if err != nil {
32 t.Error(err)
33 }
34
35 + if p13 != p11 {
36 + t.Error("these should be the same")
37 + }
38 +
39 p12, err := ps.Get(ID("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31"))
40 if err != nil {
41 t.Error(err)
@@ -41,10 +45,13 @@ func TestPeerstore(t *testing.T) {
45 t.Error(errors.New("peers should be the same"))
46 }
47
44 - err = ps.Put(p21)
48 + p23, err := ps.Add(p21)
49 if err != nil {
50 t.Error(err)
51 }
52 + if p23 != p21 {
53 + t.Error("These should be the same")
54 + }
55
56 p22, err := ps.Get(ID("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a32"))
57 if err != nil {
routing/dht/ext_test.go
+2 -2
@@ -203,7 +203,7 @@ func TestNotFound(t *testing.T) {
203
204 local := peer.WithIDString("test_peer")
205 peerstore := peer.NewPeerstore()
206 - peerstore.Put(local)
206 + peerstore.Add(local)
207
208 d := NewDHT(ctx, local, peerstore, fn, fs, ds.NewMapDatastore())
209
@@ -269,7 +269,7 @@ func TestLessThanKResponses(t *testing.T) {
269 fs := &fauxSender{}
270 local := peer.WithIDString("test_peer")
271 peerstore := peer.NewPeerstore()
272 - peerstore.Put(local)
272 + peerstore.Add(local)
273
274 d := NewDHT(ctx, local, peerstore, fn, fs, ds.NewMapDatastore())
275