@cryptotaxi247 / kubo / commits / b1d11ccfc

peerstore constructs peers

Now, all peers should be retrieved from the Peerstore, which will construct the peers accordingly. This ensures there's only one peer object per peer (opposite would be bad: things get out sync) cc @whyrusleeping

Juan Batiz-Benet committed Oct 19, 2014 at 23:40 UTC b1d11ccfcbcd75abdc3601d79e0557206f07c397
7 files changed +141 -93
core/core.go
+17 -22
@@ -12,7 +12,6 @@ import (
12
13 bserv "github.com/jbenet/go-ipfs/blockservice"
14 config "github.com/jbenet/go-ipfs/config"
15 - ci "github.com/jbenet/go-ipfs/crypto"
15 diag "github.com/jbenet/go-ipfs/diagnostics"
16 exchange "github.com/jbenet/go-ipfs/exchange"
17 bitswap "github.com/jbenet/go-ipfs/exchange/bitswap"
@@ -92,8 +91,7 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
91 }
92
93 peerstore := peer.NewPeerstore()
95 -
96 - local, err := initIdentity(cfg, online)
94 + local, err := initIdentity(cfg, peerstore, online)
95 if err != nil {
96 return nil, err
97 }
@@ -179,7 +177,7 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
177 }, nil
178 }
179
182 -func initIdentity(cfg *config.Config, online bool) (*peer.Peer, error) {
180 +func initIdentity(cfg *config.Config, peers peer.Peerstore, online bool) (*peer.Peer, error) {
181 if cfg.Identity.PeerID == "" {
182 return nil, errors.New("Identity was not set in config (was ipfs init run?)")
183 }
@@ -188,22 +186,23 @@ func initIdentity(cfg *config.Config, online bool) (*peer.Peer, error) {
186 return nil, errors.New("No peer ID in config! (was ipfs init run?)")
187 }
188
189 + // get peer from peerstore (so it is constructed there)
190 + id := peer.ID(b58.Decode(cfg.Identity.PeerID))
191 + peer, err := peers.Get(id)
192 + if err != nil {
193 + return nil, err
194 + }
195 +
196 // address is optional
192 - var addresses []ma.Multiaddr
197 if len(cfg.Addresses.Swarm) > 0 {
198 maddr, err := ma.NewMultiaddr(cfg.Addresses.Swarm)
199 if err != nil {
200 return nil, err
201 }
202
199 - addresses = []ma.Multiaddr{maddr}
203 + peer.AddAddress(maddr)
204 }
205
202 - var (
203 - sk ci.PrivKey
204 - pk ci.PubKey
205 - )
206 -
206 // when not online, don't need to parse private keys (yet)
207 if online {
208 skb, err := base64.StdEncoding.DecodeString(cfg.Identity.PrivKey)
@@ -211,20 +210,12 @@ func initIdentity(cfg *config.Config, online bool) (*peer.Peer, error) {
210 return nil, err
211 }
212
214 - sk, err = ci.UnmarshalPrivateKey(skb)
215 - if err != nil {
213 + if err := peer.LoadAndVerifyKeyPair(skb); err != nil {
214 return nil, err
215 }
218 -
219 - pk = sk.GetPublic()
216 }
217
222 - return &peer.Peer{
223 - ID: peer.ID(b58.Decode(cfg.Identity.PeerID)),
224 - Addresses: addresses,
225 - PrivKey: sk,
226 - PubKey: pk,
227 - }, nil
218 + return peer, nil
219 }
220
221 func initConnections(ctx context.Context, cfg *config.Config, pstore peer.Peerstore, route *dht.IpfsDHT) {
@@ -240,7 +231,11 @@ func initConnections(ctx context.Context, cfg *config.Config, pstore peer.Peerst
231 }
232
233 // setup peer
243 - npeer := &peer.Peer{ID: peer.DecodePrettyID(p.PeerID)}
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 {
core/mock.go
+10 -3
@@ -12,11 +12,18 @@ import (
12 mdht "github.com/jbenet/go-ipfs/routing/mock"
13 )
14
15 +// NewMockNode constructs an IpfsNode for use in tests.
16 func NewMockNode() (*IpfsNode, error) {
17 nd := new(IpfsNode)
18
19 //Generate Identity
19 - nd.Identity = &peer.Peer{ID: []byte("TESTING")}
20 + nd.Peerstore = peer.NewPeerstore()
21 + var err error
22 + nd.Identity, err = nd.Peerstore.Get(peer.ID("TESTING"))
23 + if err != nil {
24 + return nil, err
25 + }
26 +
27 pk, sk, err := ci.GenerateKeyPair(ci.RSA, 1024)
28 if err != nil {
29 return nil, err
@@ -40,13 +47,13 @@ func NewMockNode() (*IpfsNode, error) {
47 return nil, err
48 }
49
43 - nd.DAG = &mdag.DAGService{bserv}
50 + nd.DAG = &mdag.DAGService{Blocks: bserv}
51
52 // Namespace resolver
53 nd.Namesys = nsys.NewNameSystem(dht)
54
55 // Path resolver
49 - nd.Resolver = &path.Resolver{nd.DAG}
56 + nd.Resolver = &path.Resolver{DAG: nd.DAG}
57
58 return nd, nil
59 }
crypto/spipe/handshake.go
+5 -33
@@ -348,41 +348,13 @@ func getOrConstructPeer(peers peer.Peerstore, rpk ci.PubKey) (*peer.Peer, error)
348 }
349
350 npeer, err := peers.Get(rid)
351 - if err != nil || npeer == nil {
352 - if err != peer.ErrNotFound {
353 - return nil, err // unexpected error happened.
354 - }
355 -
356 - // dont have peer, so construct it + add it to peerstore.
357 - npeer = &peer.Peer{ID: rid, PubKey: rpk}
358 - if err := peers.Put(npeer); err != nil {
359 - return nil, err
360 - }
361 -
362 - // done, return the newly constructed peer.
363 - return npeer, nil
364 - }
365 -
366 - // did have it locally.
367 -
368 - // let's verify ID
369 - if !npeer.ID.Equal(rid) {
370 - e := "Expected peer.ID does not match sent pubkey's hash: %v - %v"
371 - return nil, fmt.Errorf(e, npeer, rid)
372 - }
373 -
374 - if npeer.PubKey == nil {
375 - // didn't have a pubkey, just set it.
376 - npeer.PubKey = rpk
377 - return npeer, nil
351 + if err != nil {
352 + return nil, err // unexpected error happened.
353 }
354
380 - // did have pubkey, let's verify it's really the same.
381 - // this shouldn't ever happen, given we hashed, etc, but it could mean
382 - // expected code (or protocol) invariants violated.
383 - if !npeer.PubKey.Equals(rpk) {
384 - log.Error("WARNING: PubKey mismatch: %v", npeer)
385 - panic("secure channel pubkey mismatch")
355 + // public key verification happens in Peer.VerifyAndSetPubKey
356 + if err := npeer.VerifyAndSetPubKey(rpk); err != nil {
357 + return nil, err // pubkey mismatch or other problem
358 }
359 return npeer, nil
360 }
peer/peer.go
+64
@@ -1,6 +1,7 @@
1 package peer
2
3 import (
4 + "fmt"
5 "sync"
6 "time"
7
@@ -13,6 +14,8 @@ import (
14 "bytes"
15 )
16
17 +var log = u.Logger("peer")
18 +
19 // ID is a byte slice representing the identity of a peer.
20 type ID mh.Multihash
21
@@ -122,3 +125,64 @@ func (p *Peer) SetLatency(laten time.Duration) {
125 }
126 p.Unlock()
127 }
128 +
129 +// LoadAndVerifyKeyPair unmarshalls, loads a private/public key pair.
130 +// Error if (a) unmarshalling fails, or (b) pubkey does not match id.
131 +func (p *Peer) LoadAndVerifyKeyPair(marshalled []byte) error {
132 +
133 + sk, err := ic.UnmarshalPrivateKey(marshalled)
134 + if err != nil {
135 + return fmt.Errorf("Failed to unmarshal private key: %v", err)
136 + }
137 +
138 + // construct and assign pubkey. ensure it matches this peer
139 + if err := p.VerifyAndSetPubKey(sk.GetPublic()); err != nil {
140 + return err
141 + }
142 +
143 + // if we didn't have the priavte key, assign it
144 + if p.PrivKey == nil {
145 + p.PrivKey = sk
146 + return nil
147 + }
148 +
149 + // if we already had the keys, check they're equal.
150 + if p.PrivKey.Equals(sk) {
151 + return nil // as expected. keep the old objects.
152 + }
153 +
154 + // keys not equal. invariant violated. this warrants a panic.
155 + // these keys should be _the same_ because peer.ID = H(pk)
156 + // this mismatch should never happen.
157 + log.Error("%s had PrivKey: %v -- got %v", p, p.PrivKey, sk)
158 + panic("invariant violated: unexpected key mismatch")
159 +}
160 +
161 +// VerifyAndSetPubKey sets public key, given it matches the peer.ID
162 +func (p *Peer) VerifyAndSetPubKey(pk ic.PubKey) error {
163 + pkid, err := IDFromPubKey(pk)
164 + if err != nil {
165 + return fmt.Errorf("Failed to hash public key: %v", err)
166 + }
167 +
168 + if !p.ID.Equal(pkid) {
169 + return fmt.Errorf("Public key does not match peer.ID.")
170 + }
171 +
172 + // if we didn't have the keys, assign them.
173 + if p.PubKey == nil {
174 + p.PubKey = pk
175 + return nil
176 + }
177 +
178 + // if we already had the pubkey, check they're equal.
179 + if p.PubKey.Equals(pk) {
180 + return nil // as expected. keep the old objects.
181 + }
182 +
183 + // keys not equal. invariant violated. this warrants a panic.
184 + // these keys should be _the same_ because peer.ID = H(pk)
185 + // this mismatch should never happen.
186 + log.Error("%s had PubKey: %v -- got %v", p, p.PubKey, pk)
187 + panic("invariant violated: unexpected key mismatch")
188 +}
peer/peerstore.go
+19 -10
@@ -9,10 +9,6 @@ import (
9 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
10 )
11
12 -// ErrNotFound signals a peer wasn't found. this is here to avoid having to
13 -// leak the ds abstraction to clients of Peerstore, just for the error.
14 -var ErrNotFound = ds.ErrNotFound
15 -
12 // Peerstore provides a threadsafe collection for peers.
13 type Peerstore interface {
14 Get(ID) (*Peer, error)
@@ -39,15 +35,28 @@ func (p *peerstore) Get(i ID) (*Peer, error) {
35
36 k := u.Key(i).DsKey()
37 val, err := p.peers.Get(k)
42 - if err != nil {
38 + switch err {
39 +
40 + // some other datastore error
41 + default:
42 return nil, err
44 - }
43
46 - peer, ok := val.(*Peer)
47 - if !ok {
48 - return nil, errors.New("stored value was not a Peer")
44 + // not found, construct it ourselves, add it to datastore, and return.
45 + case ds.ErrNotFound:
46 + peer := &Peer{ID: i}
47 + if err := p.peers.Put(k, peer); err != nil {
48 + return nil, err
49 + }
50 + return peer, nil
51 +
52 + // no error, got it back fine
53 + case nil:
54 + peer, ok := val.(*Peer)
55 + if !ok {
56 + return nil, errors.New("stored value was not a Peer")
57 + }
58 + return peer, nil
59 }
50 - return peer, nil
60 }
61
62 func (p *peerstore) Put(peer *Peer) error {
peer/peerstore_test.go
+5 -4
@@ -56,8 +56,8 @@ func TestPeerstore(t *testing.T) {
56 }
57
58 _, err = ps.Get(ID("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a33"))
59 - if err == nil {
60 - t.Error(errors.New("should've been an error here"))
59 + if err != nil {
60 + t.Error(errors.New("should not have an error here"))
61 }
62
63 err = ps.Delete(ID("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31"))
@@ -65,9 +65,10 @@ func TestPeerstore(t *testing.T) {
65 t.Error(err)
66 }
67
68 + // reconstruct!
69 _, err = ps.Get(ID("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a31"))
69 - if err == nil {
70 - t.Error(errors.New("should've been an error here"))
70 + if err != nil {
71 + t.Error(errors.New("should not have an error anyway. reconstruct!"))
72 }
73
74 p22, err = ps.Get(ID("11140beec7b5ea3f0fdbc95d0dd47f3c5bc275da8a32"))
routing/dht/dht.go
+21 -21
@@ -300,10 +300,9 @@ func (dht *IpfsDHT) addPeer(pb *Message_Peer) (*peer.Peer, error) {
300 }
301
302 // check if we already have this peer.
303 - pr, _ := dht.peerstore.Get(peer.ID(pb.GetId()))
304 - if pr == nil {
305 - pr = &peer.Peer{ID: peer.ID(pb.GetId())}
306 - dht.peerstore.Put(pr)
303 + pr, err := dht.getPeer(peer.ID(pb.GetId()))
304 + if err != nil {
305 + return nil, err
306 }
307 pr.AddAddress(addr) // idempotent
308
@@ -481,6 +480,16 @@ func (dht *IpfsDHT) betterPeersToQuery(pmes *Message, count int) []*peer.Peer {
480 return filtered
481 }
482
483 +func (dht *IpfsDHT) getPeer(id peer.ID) (*peer.Peer, error) {
484 + p, err := dht.peerstore.Get(id)
485 + if err != nil {
486 + err = fmt.Errorf("Failed to get peer from peerstore: %s", err)
487 + log.Error("%s", err)
488 + return nil, err
489 + }
490 + return p, nil
491 +}
492 +
493 func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (*peer.Peer, error) {
494
495 id := peer.ID(pbp.GetId())
@@ -490,26 +499,16 @@ func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (*peer.Peer, error) {
499 return nil, errors.New("found self")
500 }
501
493 - p, _ := dht.peerstore.Get(id)
494 - if p == nil {
495 - p, _ = dht.FindLocal(id)
496 - if p != nil {
497 - panic("somehow peer not getting into peerstore")
498 - }
502 + p, err := dht.getPeer(id)
503 + if err != nil {
504 + return nil, err
505 }
506
501 - if p == nil {
502 - maddr, err := pbp.Address()
503 - if err != nil {
504 - return nil, err
505 - }
506 -
507 - // create new Peer
508 - p = &peer.Peer{ID: id}
509 - p.AddAddress(maddr)
510 - dht.peerstore.Put(p)
511 - log.Info("dht found new peer: %s %s", p, maddr)
507 + maddr, err := pbp.Address()
508 + if err != nil {
509 + return nil, err
510 }
511 + p.AddAddress(maddr)
512 return p, nil
513 }
514
@@ -541,6 +540,7 @@ func (dht *IpfsDHT) loadProvidableKeys() error {
540 return nil
541 }
542
543 +// PingRoutine periodically pings nearest neighbors.
544 func (dht *IpfsDHT) PingRoutine(t time.Duration) {
545 tick := time.Tick(t)
546 for {