@cryptotaxi247 / kubo / commits / 69dd26023

refactor peerSet

Jeromy committed Dec 16, 2014 at 18:33 UTC 69dd26023637965f0232a7d095378b33e74d74fe
3 files changed +22 -9
exchange/bitswap/bitswap.go
+1 -1
@@ -222,7 +222,7 @@ func (bs *bitswap) sendWantlistToProviders(ctx context.Context, wantlist *wl.Wan
222 providers := bs.routing.FindProvidersAsync(child, k, maxProvidersPerRequest)
223
224 for prov := range providers {
225 - if ps.AddIfSmallerThan(prov, -1) { //Do once per peer
225 + if ps.TryAdd(prov) { //Do once per peer
226 bs.send(ctx, prov, message)
227 }
228 }
routing/dht/routing.go
+4 -4
@@ -141,11 +141,11 @@ func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int
141 func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, count int, peerOut chan peer.Peer) {
142 defer close(peerOut)
143
144 - ps := pset.NewPeerSet()
144 + ps := pset.NewLimitedPeerSet(count)
145 provs := dht.providers.GetProviders(ctx, key)
146 for _, p := range provs {
147 // NOTE: assuming that this list of peers is unique
148 - if ps.AddIfSmallerThan(p, count) {
148 + if ps.TryAdd(p) {
149 select {
150 case peerOut <- p:
151 case <-ctx.Done():
@@ -176,7 +176,7 @@ func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, co
176
177 // Add unique providers from request, up to 'count'
178 for _, prov := range provs {
179 - if ps.AddIfSmallerThan(prov, count) {
179 + if ps.TryAdd(prov) {
180 select {
181 case peerOut <- prov:
182 case <-ctx.Done():
@@ -226,7 +226,7 @@ func (dht *IpfsDHT) addPeerListAsync(ctx context.Context, k u.Key, peers []*pb.M
226 }
227
228 dht.providers.AddProvider(k, p)
229 - if ps.AddIfSmallerThan(p, count) {
229 + if ps.TryAdd(p) {
230 select {
231 case out <- p:
232 case <-ctx.Done():
util/peerset/peerset.go
+17 -4
@@ -7,13 +7,22 @@ import (
7
8 // PeerSet is a threadsafe set of peers
9 type PeerSet struct {
10 - ps map[string]bool
11 - lk sync.RWMutex
10 + ps map[string]bool
11 + lk sync.RWMutex
12 + size int
13 }
14
15 func NewPeerSet() *PeerSet {
16 ps := new(PeerSet)
17 ps.ps = make(map[string]bool)
18 + ps.size = -1
19 + return ps
20 +}
21 +
22 +func NewLimitedPeerSet(size int) *PeerSet {
23 + ps := new(PeerSet)
24 + ps.ps = make(map[string]bool)
25 + ps.size = -1
26 return ps
27 }
28
@@ -36,10 +45,14 @@ func (ps *PeerSet) Size() int {
45 return len(ps.ps)
46 }
47
39 -func (ps *PeerSet) AddIfSmallerThan(p peer.Peer, maxsize int) bool {
48 +// TryAdd Attempts to add the given peer into the set.
49 +// This operation can fail for one of two reasons:
50 +// 1) The given peer is already in the set
51 +// 2) The number of peers in the set is equal to size
52 +func (ps *PeerSet) TryAdd(p peer.Peer) bool {
53 var success bool
54 ps.lk.Lock()
42 - if _, ok := ps.ps[string(p.ID())]; !ok && len(ps.ps) < maxsize {
55 + if _, ok := ps.ps[string(p.ID())]; !ok && (len(ps.ps) < ps.size || ps.size == -1) {
56 success = true
57 ps.ps[string(p.ID())] = true
58 }