rewrite of provides to better select peers to send RPCs to
refactor test peer creation to be deterministic and reliable a bit of cleanup trying to figure out TestGetFailure add test to verify deterministic peer creation switch put RPC over to use getClosestPeers rm 0xDEADC0DE fix queries not searching peer if its not actually closer
Jeromy committed
Dec 14, 2014 at 00:50 UTC
07b064010e584f046234681de6f47d79d09b9689
17 files changed
+230
-129
cmd/ipfs/init.go
+2
-1
@@ -2,6 +2,7 @@ package main
2
3
import (
4
"bytes"
5
+ "crypto/rand"
6
"encoding/base64"
7
"fmt"
8
"os"
@@ -252,7 +253,7 @@ func identityConfig(nbits int) (config.Identity, error) {
253
}
254
255
fmt.Printf("generating key pair...")
255
- sk, pk, err := ci.GenerateKeyPair(ci.RSA, nbits)
256
+ sk, pk, err := ci.GenerateKeyPair(ci.RSA, nbits, rand.Reader)
257
if err != nil {
258
return ident, err
259
}
cmd/seccat/seccat.go
+1
-1
@@ -115,7 +115,7 @@ func setupPeer(a args) (peer.ID, peer.Peerstore, error) {
115
}
116
117
out("generating key pair...")
118
- sk, pk, err := ci.GenerateKeyPair(ci.RSA, a.keybits)
118
+ sk, pk, err := ci.GenerateKeyPair(ci.RSA, a.keybits, u.NewTimeSeededRand())
119
if err != nil {
120
return "", nil, err
121
}
core/mock.go
+3
-1
@@ -1,7 +1,9 @@
1
package core
2
3
import (
4
+ "crypto/rand"
5
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
6
+
7
ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
8
syncds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
9
"github.com/jbenet/go-ipfs/blocks/blockstore"
@@ -28,7 +30,7 @@ func NewMockNode() (*IpfsNode, error) {
30
nd := new(IpfsNode)
31
32
// Generate Identity
31
- sk, pk, err := ci.GenerateKeyPair(ci.RSA, 1024)
33
+ sk, pk, err := ci.GenerateKeyPair(ci.RSA, 1024, rand.Reader)
34
if err != nil {
35
return nil, err
36
}
crypto/key.go
+3
-2
@@ -8,6 +8,7 @@ import (
8
"encoding/base64"
9
"errors"
10
"fmt"
11
+ "io"
12
13
"crypto/elliptic"
14
"crypto/hmac"
@@ -75,10 +76,10 @@ type PubKey interface {
76
type GenSharedKey func([]byte) ([]byte, error)
77
78
// Generates a keypair of the given type and bitsize
78
-func GenerateKeyPair(typ, bits int) (PrivKey, PubKey, error) {
79
+func GenerateKeyPair(typ, bits int, src io.Reader) (PrivKey, PubKey, error) {
80
switch typ {
81
case RSA:
81
- priv, err := rsa.GenerateKey(rand.Reader, bits)
82
+ priv, err := rsa.GenerateKey(src, bits)
83
if err != nil {
84
return nil, nil, err
85
}
crypto/key_test.go
+3
-2
@@ -2,11 +2,12 @@ package crypto
2
3
import (
4
"bytes"
5
+ u "github.com/jbenet/go-ipfs/util"
6
"testing"
7
)
8
9
func TestRsaKeys(t *testing.T) {
9
- sk, pk, err := GenerateKeyPair(RSA, 512)
10
+ sk, pk, err := GenerateKeyPair(RSA, 512, u.NewTimeSeededRand())
11
if err != nil {
12
t.Fatal(err)
13
}
@@ -90,7 +91,7 @@ func testKeyEquals(t *testing.T, k Key) {
91
t.Fatal("Key not equal to key with same bytes.")
92
}
93
93
- sk, pk, err := GenerateKeyPair(RSA, 512)
94
+ sk, pk, err := GenerateKeyPair(RSA, 512, u.NewTimeSeededRand())
95
if err != nil {
96
t.Fatal(err)
97
}
namesys/resolve_test.go
+1
-1
@@ -15,7 +15,7 @@ func TestRoutingResolve(t *testing.T) {
15
resolver := NewRoutingResolver(d)
16
publisher := NewRoutingPublisher(d)
17
18
- privk, pubk, err := ci.GenerateKeyPair(ci.RSA, 512)
18
+ privk, pubk, err := ci.GenerateKeyPair(ci.RSA, 512, u.NewTimeSeededRand())
19
if err != nil {
20
t.Fatal(err)
21
}
net/mock/mock_net.go
+1
-1
@@ -41,7 +41,7 @@ func New(ctx context.Context) Mocknet {
41
}
42
43
func (mn *mocknet) GenPeer() (inet.Network, error) {
44
- sk, _, err := testutil.RandKeyPair(512)
44
+ sk, _, err := testutil.SeededKeyPair(512, int64(len(mn.nets)))
45
if err != nil {
46
return nil, err
47
}
routing/dht/dht.go
+10
-23
@@ -103,9 +103,8 @@ func (dht *IpfsDHT) Connect(ctx context.Context, npeer peer.ID) error {
103
}
104
105
// putValueToNetwork stores the given key/value pair at the peer 'p'
106
-// meaning: it sends a PUT_VALUE message to p
106
func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p peer.ID,
108
- key string, rec *pb.Record) error {
107
+ key u.Key, rec *pb.Record) error {
108
109
pmes := pb.NewMessage(pb.Message_PUT_VALUE, string(key), 0)
110
pmes.Record = rec
@@ -285,7 +284,7 @@ func (dht *IpfsDHT) nearestPeersToQuery(pmes *pb.Message, count int) []peer.ID {
284
}
285
286
// betterPeerToQuery returns nearestPeersToQuery, but iff closer than self.
288
-func (dht *IpfsDHT) betterPeersToQuery(pmes *pb.Message, count int) []peer.ID {
287
+func (dht *IpfsDHT) betterPeersToQuery(pmes *pb.Message, p peer.ID, count int) []peer.ID {
288
closer := dht.nearestPeersToQuery(pmes, count)
289
290
// no node? nil
@@ -302,11 +301,16 @@ func (dht *IpfsDHT) betterPeersToQuery(pmes *pb.Message, count int) []peer.ID {
301
}
302
303
var filtered []peer.ID
305
- for _, p := range closer {
304
+ for _, clp := range closer {
305
+ // Dont send a peer back themselves
306
+ if p == clp {
307
+ continue
308
+ }
309
+
310
// must all be closer than self
311
key := u.Key(pmes.GetKey())
308
- if !kb.Closer(dht.self, p, key) {
309
- filtered = append(filtered, p)
312
+ if !kb.Closer(dht.self, clp, key) {
313
+ filtered = append(filtered, clp)
314
}
315
}
316
@@ -323,23 +327,6 @@ func (dht *IpfsDHT) ensureConnectedToPeer(ctx context.Context, p peer.ID) error
327
return dht.network.DialPeer(ctx, p)
328
}
329
326
-//TODO: this should be smarter about which keys it selects.
327
-func (dht *IpfsDHT) loadProvidableKeys() error {
328
- kl, err := dht.datastore.KeyList()
329
- if err != nil {
330
- return err
331
- }
332
- for _, dsk := range kl {
333
- k := u.KeyFromDsKey(dsk)
334
- if len(k) == 0 {
335
- log.Errorf("loadProvidableKeys error: %v", dsk)
336
- }
337
-
338
- dht.providers.AddProvider(k, dht.self)
339
- }
340
- return nil
341
-}
342
-
330
// PingRoutine periodically pings nearest neighbors.
331
func (dht *IpfsDHT) PingRoutine(t time.Duration) {
332
defer dht.Children().Done()
routing/dht/dht_test.go
+9
-10
@@ -14,7 +14,6 @@ import (
14
dssync "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore/sync"
15
ma "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-multiaddr"
16
17
- // ci "github.com/jbenet/go-ipfs/crypto"
17
inet "github.com/jbenet/go-ipfs/net"
18
peer "github.com/jbenet/go-ipfs/peer"
19
routing "github.com/jbenet/go-ipfs/routing"
@@ -33,9 +32,9 @@ func init() {
32
}
33
}
34
36
-func setupDHT(ctx context.Context, t *testing.T, addr ma.Multiaddr) *IpfsDHT {
35
+func setupDHT(ctx context.Context, t *testing.T, addr ma.Multiaddr, seed int64) *IpfsDHT {
36
38
- sk, pk, err := testutil.RandKeyPair(512)
37
+ sk, pk, err := testutil.SeededKeyPair(512, seed)
38
if err != nil {
39
t.Fatal(err)
40
}
@@ -71,7 +70,7 @@ func setupDHTS(ctx context.Context, n int, t *testing.T) ([]ma.Multiaddr, []peer
70
71
for i := 0; i < n; i++ {
72
addrs[i] = testutil.RandLocalTCPAddress()
74
- dhts[i] = setupDHT(ctx, t, addrs[i])
73
+ dhts[i] = setupDHT(ctx, t, addrs[i], int64(i))
74
peers[i] = dhts[i].self
75
}
76
@@ -120,8 +119,8 @@ func TestPing(t *testing.T) {
119
addrA := testutil.RandLocalTCPAddress()
120
addrB := testutil.RandLocalTCPAddress()
121
123
- dhtA := setupDHT(ctx, t, addrA)
124
- dhtB := setupDHT(ctx, t, addrB)
122
+ dhtA := setupDHT(ctx, t, addrA, 1)
123
+ dhtB := setupDHT(ctx, t, addrB, 2)
124
125
peerA := dhtA.self
126
peerB := dhtB.self
@@ -153,8 +152,8 @@ func TestValueGetSet(t *testing.T) {
152
addrA := testutil.RandLocalTCPAddress()
153
addrB := testutil.RandLocalTCPAddress()
154
156
- dhtA := setupDHT(ctx, t, addrA)
157
- dhtB := setupDHT(ctx, t, addrB)
155
+ dhtA := setupDHT(ctx, t, addrA, 1)
156
+ dhtB := setupDHT(ctx, t, addrB, 2)
157
158
defer dhtA.Close()
159
defer dhtB.Close()
@@ -642,8 +641,8 @@ func TestConnectCollision(t *testing.T) {
641
addrA := testutil.RandLocalTCPAddress()
642
addrB := testutil.RandLocalTCPAddress()
643
645
- dhtA := setupDHT(ctx, t, addrA)
646
- dhtB := setupDHT(ctx, t, addrB)
644
+ dhtA := setupDHT(ctx, t, addrA, int64((rtime*2)+1))
645
+ dhtB := setupDHT(ctx, t, addrB, int64((rtime*2)+2))
646
647
peerA := dhtA.self
648
peerB := dhtB.self
routing/dht/ext_test.go
+3
-4
@@ -47,9 +47,8 @@ func TestGetFailures(t *testing.T) {
47
t.Fatal("Did not get expected error!")
48
}
49
50
- msgs := make(chan *pb.Message, 100)
50
+ t.Log("Timeout test passed.")
51
52
- // u.POut("NotFound Test\n")
52
// Reply with failures to every message
53
nets[1].SetHandler(inet.ProtocolDHT, func(s inet.Stream) {
54
defer s.Close()
@@ -68,8 +67,6 @@ func TestGetFailures(t *testing.T) {
67
if err := pbw.WriteMsg(resp); err != nil {
68
panic(err)
69
}
71
-
72
- msgs <- resp
70
})
71
72
// This one should fail with NotFound
@@ -83,6 +80,8 @@ func TestGetFailures(t *testing.T) {
80
t.Fatal("expected error, got none.")
81
}
82
83
+ t.Log("ErrNotFound check passed!")
84
+
85
// Now we test this DHT's handleGetValue failure
86
{
87
typ := pb.Message_GET_VALUE
routing/dht/handlers.go
+7
-4
@@ -93,7 +93,7 @@ func (dht *IpfsDHT) handleGetValue(ctx context.Context, p peer.ID, pmes *pb.Mess
93
}
94
95
// Find closest peer on given cluster to desired key and reply with that info
96
- closer := dht.betterPeersToQuery(pmes, CloserPeerCount)
96
+ closer := dht.betterPeersToQuery(pmes, p, CloserPeerCount)
97
closerinfos := peer.PeerInfos(dht.peerstore, closer)
98
if closer != nil {
99
for _, pi := range closerinfos {
@@ -137,6 +137,9 @@ func (dht *IpfsDHT) handlePing(_ context.Context, p peer.ID, pmes *pb.Message) (
137
}
138
139
func (dht *IpfsDHT) handleFindPeer(ctx context.Context, p peer.ID, pmes *pb.Message) (*pb.Message, error) {
140
+ log.Errorf("handle find peer %s start", p)
141
+ defer log.Errorf("handle find peer %s end", p)
142
+
143
resp := pb.NewMessage(pmes.GetType(), "", pmes.GetClusterLevel())
144
var closest []peer.ID
145
@@ -144,11 +147,11 @@ func (dht *IpfsDHT) handleFindPeer(ctx context.Context, p peer.ID, pmes *pb.Mess
147
if peer.ID(pmes.GetKey()) == dht.self {
148
closest = []peer.ID{dht.self}
149
} else {
147
- closest = dht.betterPeersToQuery(pmes, CloserPeerCount)
150
+ closest = dht.betterPeersToQuery(pmes, p, CloserPeerCount)
151
}
152
153
if closest == nil {
151
- log.Debugf("handleFindPeer: could not find anything.")
154
+ log.Warningf("handleFindPeer: could not find anything.")
155
return resp, nil
156
}
157
@@ -189,7 +192,7 @@ func (dht *IpfsDHT) handleGetProviders(ctx context.Context, p peer.ID, pmes *pb.
192
}
193
194
// Also send closer peers.
192
- closer := dht.betterPeersToQuery(pmes, CloserPeerCount)
195
+ closer := dht.betterPeersToQuery(pmes, p, CloserPeerCount)
196
if closer != nil {
197
infos := peer.PeerInfos(dht.peerstore, providers)
198
resp.CloserPeers = pb.PeerInfosToPBPeers(dht.network, infos)
routing/dht/query.go
+8
-23
@@ -7,8 +7,8 @@ import (
7
peer "github.com/jbenet/go-ipfs/peer"
8
queue "github.com/jbenet/go-ipfs/peer/queue"
9
"github.com/jbenet/go-ipfs/routing"
10
- kb "github.com/jbenet/go-ipfs/routing/kbucket"
10
u "github.com/jbenet/go-ipfs/util"
11
+ pset "github.com/jbenet/go-ipfs/util/peerset"
12
todoctr "github.com/jbenet/go-ipfs/util/todocounter"
13
14
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
@@ -71,7 +71,7 @@ type dhtQueryRunner struct {
71
peersToQuery *queue.ChanQueue
72
73
// peersSeen are all the peers queried. used to prevent querying same peer 2x
74
- peersSeen peer.Set
74
+ peersSeen *pset.PeerSet
75
76
// rateLimit is a channel used to rate limit our processing (semaphore)
77
rateLimit chan struct{}
@@ -97,7 +97,7 @@ func newQueryRunner(ctx context.Context, q *dhtQuery) *dhtQueryRunner {
97
query: q,
98
peersToQuery: queue.NewChanQueue(ctx, queue.NewXORDistancePQ(q.key)),
99
peersRemaining: todoctr.NewSyncCounter(),
100
- peersSeen: peer.Set{},
100
+ peersSeen: pset.New(),
101
rateLimit: make(chan struct{}, q.concurrency),
102
cg: ctxgroup.WithContext(ctx),
103
}
@@ -117,7 +117,7 @@ func (r *dhtQueryRunner) Run(peers []peer.ID) (*dhtQueryResult, error) {
117
118
// add all the peers we got first.
119
for _, p := range peers {
120
- r.addPeerToQuery(r.cg.Context(), p, "") // don't have access to self here...
120
+ r.addPeerToQuery(r.cg.Context(), p)
121
}
122
123
// go do this thing.
@@ -153,32 +153,17 @@ func (r *dhtQueryRunner) Run(peers []peer.ID) (*dhtQueryResult, error) {
153
return nil, err
154
}
155
156
-func (r *dhtQueryRunner) addPeerToQuery(ctx context.Context, next peer.ID, benchmark peer.ID) {
156
+func (r *dhtQueryRunner) addPeerToQuery(ctx context.Context, next peer.ID) {
157
// if new peer is ourselves...
158
if next == r.query.dialer.LocalPeer() {
159
return
160
}
161
162
- // if new peer further away than whom we got it from, don't bother (loops)
163
- // TODO----------- this benchmark should be replaced by a heap:
164
- // we should be doing the s/kademlia "continue to search"
165
- // (i.e. put all of them in a heap sorted by dht distance and then just
166
- // pull from the the top until a) you exhaust all peers you get,
167
- // b) you succeed, c) your context expires.
168
- if benchmark != "" && kb.Closer(benchmark, next, r.query.key) {
162
+ if !r.peersSeen.TryAdd(next) {
163
+ log.Debug("query peer was already seen")
164
return
165
}
166
172
- // if already seen, no need.
173
- r.Lock()
174
- _, found := r.peersSeen[next]
175
- if found {
176
- r.Unlock()
177
- return
178
- }
179
- r.peersSeen[next] = struct{}{}
180
- r.Unlock()
181
-
167
log.Debugf("adding peer to query: %v", next)
168
169
// do this after unlocking to prevent possible deadlocks.
@@ -278,7 +263,7 @@ func (r *dhtQueryRunner) queryPeer(cg ctxgroup.ContextGroup, p peer.ID) {
263
}
264
265
r.query.dialer.Peerstore().AddAddresses(next.ID, next.Addrs)
281
- r.addPeerToQuery(cg.Context(), next.ID, p)
266
+ r.addPeerToQuery(cg.Context(), next.ID)
267
log.Debugf("PEERS CLOSER -- worker for: %v added %v (%v)", p, next.ID, next.Addrs)
268
}
269
} else {
routing/dht/routing.go
+107
-19
@@ -40,19 +40,24 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
40
return err
41
}
42
43
- peers := dht.routingTable.NearestPeers(kb.ConvertKey(key), KValue)
44
-
45
- query := newQuery(key, dht.network, func(ctx context.Context, p peer.ID) (*dhtQueryResult, error) {
46
- log.Debugf("%s PutValue qry part %v", dht.self, p)
47
- err := dht.putValueToNetwork(ctx, p, string(key), rec)
48
- if err != nil {
49
- return nil, err
50
- }
51
- return &dhtQueryResult{success: true}, nil
52
- })
43
+ pchan, err := dht.getClosestPeers(ctx, key, KValue)
44
+ if err != nil {
45
+ return err
46
+ }
47
54
- _, err = query.Run(ctx, peers)
55
- return err
48
+ wg := sync.WaitGroup{}
49
+ for p := range pchan {
50
+ wg.Add(1)
51
+ go func(p peer.ID) {
52
+ defer wg.Done()
53
+ err := dht.putValueToNetwork(ctx, p, key, rec)
54
+ if err != nil {
55
+ log.Errorf("failed putting value to peer: %s", err)
56
+ }
57
+ }(p)
58
+ }
59
+ wg.Wait()
60
+ return nil
61
}
62
63
// GetValue searches for the value corresponding to given Key.
@@ -111,18 +116,19 @@ func (dht *IpfsDHT) GetValue(ctx context.Context, key u.Key) ([]byte, error) {
116
// Provide makes this node announce that it can provide a value for the given key
117
func (dht *IpfsDHT) Provide(ctx context.Context, key u.Key) error {
118
119
+ log.Event(ctx, "Provide Value start", &key)
120
+ defer log.Event(ctx, "Provide Value end", &key)
121
dht.providers.AddProvider(key, dht.self)
115
- peers := dht.routingTable.NearestPeers(kb.ConvertKey(key), PoolSize)
116
- if len(peers) == 0 {
117
- return nil
122
+
123
+ peers, err := dht.getClosestPeers(ctx, key, KValue)
124
+ if err != nil {
125
+ return err
126
}
127
120
- //TODO FIX: this doesn't work! it needs to be sent to the actual nearest peers.
121
- // `peers` are the closest peers we have, not the ones that should get the value.
122
- for _, p := range peers {
128
+ for p := range peers {
129
err := dht.putProvider(ctx, p, string(key))
130
if err != nil {
125
- return err
131
+ log.Error(err)
132
}
133
}
134
return nil
@@ -137,6 +143,87 @@ func (dht *IpfsDHT) FindProviders(ctx context.Context, key u.Key) ([]peer.PeerIn
143
return providers, nil
144
}
145
146
+func (dht *IpfsDHT) getClosestPeers(ctx context.Context, key u.Key, count int) (<-chan peer.ID, error) {
147
+ log.Error("Get Closest Peers")
148
+ tablepeers := dht.routingTable.NearestPeers(kb.ConvertKey(key), AlphaValue)
149
+ if len(tablepeers) == 0 {
150
+ return nil, kb.ErrLookupFailure
151
+ }
152
+
153
+ out := make(chan peer.ID, count)
154
+ peerset := pset.NewLimited(count)
155
+
156
+ for _, p := range tablepeers {
157
+ out <- p
158
+ peerset.Add(p)
159
+ }
160
+
161
+ wg := sync.WaitGroup{}
162
+ for _, p := range tablepeers {
163
+ wg.Add(1)
164
+ go func(p peer.ID) {
165
+ dht.getClosestPeersRecurse(ctx, key, p, peerset, out)
166
+ wg.Done()
167
+ }(p)
168
+ }
169
+
170
+ go func() {
171
+ wg.Wait()
172
+ close(out)
173
+ log.Error("Closing closest peer chan")
174
+ }()
175
+
176
+ return out, nil
177
+}
178
+
179
+func (dht *IpfsDHT) getClosestPeersRecurse(ctx context.Context, key u.Key, p peer.ID, peers *pset.PeerSet, peerOut chan<- peer.ID) {
180
+ log.Error("closest peers recurse")
181
+ defer log.Error("closest peers recurse end")
182
+ closer, err := dht.closerPeersSingle(ctx, key, p)
183
+ if err != nil {
184
+ log.Errorf("error getting closer peers: %s", err)
185
+ return
186
+ }
187
+
188
+ wg := sync.WaitGroup{}
189
+ for _, p := range closer {
190
+ if kb.Closer(p, dht.self, key) && peers.TryAdd(p) {
191
+ select {
192
+ case peerOut <- p:
193
+ case <-ctx.Done():
194
+ return
195
+ }
196
+ wg.Add(1)
197
+ go func(p peer.ID) {
198
+ dht.getClosestPeersRecurse(ctx, key, p, peers, peerOut)
199
+ wg.Done()
200
+ }(p)
201
+ }
202
+ }
203
+ wg.Wait()
204
+}
205
+
206
+func (dht *IpfsDHT) closerPeersSingle(ctx context.Context, key u.Key, p peer.ID) ([]peer.ID, error) {
207
+ log.Errorf("closest peers single %s %s", p, key)
208
+ defer log.Errorf("closest peers single end %s %s", p, key)
209
+ pmes, err := dht.findPeerSingle(ctx, p, peer.ID(key))
210
+ if err != nil {
211
+ return nil, err
212
+ }
213
+
214
+ var out []peer.ID
215
+ for _, pbp := range pmes.GetCloserPeers() {
216
+ pid := peer.ID(pbp.GetId())
217
+ dht.peerstore.AddAddresses(pid, pbp.Addresses())
218
+ err := dht.ensureConnectedToPeer(ctx, pid)
219
+ if err != nil {
220
+ return nil, err
221
+ }
222
+ out = append(out, pid)
223
+ }
224
+ return out, nil
225
+}
226
+
227
// FindProvidersAsync is the same thing as FindProviders, but returns a channel.
228
// Peers will be returned on the channel as soon as they are found, even before
229
// the search query completes.
@@ -182,6 +269,7 @@ func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, co
269
// Add unique providers from request, up to 'count'
270
for _, prov := range provs {
271
if ps.TryAdd(prov.ID) {
272
+ dht.peerstore.AddAddresses(prov.ID, prov.Addrs)
273
select {
274
case peerOut <- prov:
275
case <-ctx.Done():
routing/kbucket/sorting.go
new
+59
@@ -0,0 +1,59 @@
1
+package kbucket
2
+
3
+import (
4
+ "container/list"
5
+ peer "github.com/jbenet/go-ipfs/peer"
6
+ "sort"
7
+)
8
+
9
+// A helper struct to sort peers by their distance to the local node
10
+type peerDistance struct {
11
+ p peer.ID
12
+ distance ID
13
+}
14
+
15
+// peerSorterArr implements sort.Interface to sort peers by xor distance
16
+type peerSorterArr []*peerDistance
17
+
18
+func (p peerSorterArr) Len() int { return len(p) }
19
+func (p peerSorterArr) Swap(a, b int) { p[a], p[b] = p[b], p[a] }
20
+func (p peerSorterArr) Less(a, b int) bool {
21
+ return p[a].distance.less(p[b].distance)
22
+}
23
+
24
+//
25
+
26
+func copyPeersFromList(target ID, peerArr peerSorterArr, peerList *list.List) peerSorterArr {
27
+ for e := peerList.Front(); e != nil; e = e.Next() {
28
+ p := e.Value.(peer.ID)
29
+ pID := ConvertPeerID(p)
30
+ pd := peerDistance{
31
+ p: p,
32
+ distance: xor(target, pID),
33
+ }
34
+ peerArr = append(peerArr, &pd)
35
+ if e == nil {
36
+ log.Debug("list element was nil")
37
+ return peerArr
38
+ }
39
+ }
40
+ return peerArr
41
+}
42
+
43
+func SortClosestPeers(peers []peer.ID, target ID) []peer.ID {
44
+ var psarr peerSorterArr
45
+ for _, p := range peers {
46
+ pID := ConvertPeerID(p)
47
+ pd := &peerDistance{
48
+ p: p,
49
+ distance: xor(target, pID),
50
+ }
51
+ psarr = append(psarr, pd)
52
+ }
53
+ sort.Sort(psarr)
54
+ var out []peer.ID
55
+ for _, p := range psarr {
56
+ out = append(out, p.p)
57
+ }
58
+ return out
59
+}
routing/kbucket/table.go
-35
@@ -2,7 +2,6 @@
2
package kbucket
3
4
import (
5
- "container/list"
5
"fmt"
6
"sort"
7
"sync"
@@ -103,40 +102,6 @@ func (rt *RoutingTable) nextBucket() peer.ID {
102
return ""
103
}
104
106
-// A helper struct to sort peers by their distance to the local node
107
-type peerDistance struct {
108
- p peer.ID
109
- distance ID
110
-}
111
-
112
-// peerSorterArr implements sort.Interface to sort peers by xor distance
113
-type peerSorterArr []*peerDistance
114
-
115
-func (p peerSorterArr) Len() int { return len(p) }
116
-func (p peerSorterArr) Swap(a, b int) { p[a], p[b] = p[b], p[a] }
117
-func (p peerSorterArr) Less(a, b int) bool {
118
- return p[a].distance.less(p[b].distance)
119
-}
120
-
121
-//
122
-
123
-func copyPeersFromList(target ID, peerArr peerSorterArr, peerList *list.List) peerSorterArr {
124
- for e := peerList.Front(); e != nil; e = e.Next() {
125
- p := e.Value.(peer.ID)
126
- pID := ConvertPeerID(p)
127
- pd := peerDistance{
128
- p: p,
129
- distance: xor(target, pID),
130
- }
131
- peerArr = append(peerArr, &pd)
132
- if e == nil {
133
- log.Debug("list element was nil")
134
- return peerArr
135
- }
136
- }
137
- return peerArr
138
-}
139
-
105
// Find a specific peer by ID or return nil
106
func (rt *RoutingTable) Find(id peer.ID) peer.ID {
107
srch := rt.NearestPeers(ConvertPeerID(id), 1)
util/testutil/gen.go
+6
-2
@@ -17,7 +17,11 @@ import (
17
)
18
19
func RandKeyPair(bits int) (ci.PrivKey, ci.PubKey, error) {
20
- return ci.GenerateKeyPair(ci.RSA, bits)
20
+ return ci.GenerateKeyPair(ci.RSA, bits, crand.Reader)
21
+}
22
+
23
+func SeededKeyPair(bits int, seed int64) (ci.PrivKey, ci.PubKey, error) {
24
+ return ci.GenerateKeyPair(ci.RSA, bits, u.NewSeededRand(seed))
25
}
26
27
// RandPeerID generates random "valid" peer IDs. it does not NEED to generate
@@ -120,7 +124,7 @@ func RandPeerNetParams() (*PeerNetParams, error) {
124
var p PeerNetParams
125
var err error
126
p.Addr = RandLocalTCPAddress()
123
- p.PrivKey, p.PubKey, err = ci.GenerateKeyPair(ci.RSA, 512)
127
+ p.PrivKey, p.PubKey, err = ci.GenerateKeyPair(ci.RSA, 512, u.NewTimeSeededRand())
128
if err != nil {
129
return nil, err
130
}
util/util.go
+7
@@ -107,6 +107,13 @@ func NewTimeSeededRand() io.Reader {
107
}
108
}
109
110
+func NewSeededRand(seed int64) io.Reader {
111
+ src := rand.NewSource(seed)
112
+ return &randGen{
113
+ Rand: *rand.New(src),
114
+ }
115
+}
116
+
117
func (r *randGen) Read(p []byte) (n int, err error) {
118
for i := 0; i < len(p); i++ {
119
p[i] = byte(r.Rand.Intn(255))