@cryptotaxi247 / kubo / commits / 5af562580

use query for getClosestPeers

Jeromy committed Dec 29, 2014 at 06:32 UTC 5af5625805a968b2e3524090d41e517971b8883c
1 file changed +32 -36
routing/dht/routing.go
+32 -36
@@ -159,52 +159,48 @@ func (dht *IpfsDHT) getClosestPeers(ctx context.Context, key u.Key, count int) (
159 peerset := pset.NewLimited(count)
160
161 for _, p := range tablepeers {
162 - out <- p
162 + select {
163 + case out <- p:
164 + case <-ctx.Done():
165 + return nil, ctx.Err()
166 + }
167 peerset.Add(p)
168 }
169
166 - wg := sync.WaitGroup{}
167 - for _, p := range tablepeers {
168 - wg.Add(1)
169 - go func(p peer.ID) {
170 - dht.getClosestPeersRecurse(ctx, key, p, peerset, out)
171 - wg.Done()
172 - }(p)
173 - }
170 + query := newQuery(key, dht.network, func(ctx context.Context, p peer.ID) (*dhtQueryResult, error) {
171 + closer, err := dht.closerPeersSingle(ctx, key, p)
172 + if err != nil {
173 + log.Errorf("error getting closer peers: %s", err)
174 + return nil, err
175 + }
176 +
177 + var filtered []peer.PeerInfo
178 + for _, p := range closer {
179 + if kb.Closer(p, dht.self, key) && peerset.TryAdd(p) {
180 + select {
181 + case out <- p:
182 + case <-ctx.Done():
183 + return nil, ctx.Err()
184 + }
185 + filtered = append(filtered, dht.peerstore.PeerInfo(p))
186 + }
187 + }
188 +
189 + return &dhtQueryResult{closerPeers: filtered}, nil
190 + })
191
192 go func() {
176 - wg.Wait()
177 - close(out)
193 + defer close(out)
194 + // run it!
195 + _, err := query.Run(ctx, tablepeers)
196 + if err != nil {
197 + log.Errorf("closestPeers query run error: %s", err)
198 + }
199 }()
200
201 return out, nil
202 }
203
183 -func (dht *IpfsDHT) getClosestPeersRecurse(ctx context.Context, key u.Key, p peer.ID, peers *pset.PeerSet, peerOut chan<- peer.ID) {
184 - closer, err := dht.closerPeersSingle(ctx, key, p)
185 - if err != nil {
186 - log.Errorf("error getting closer peers: %s", err)
187 - return
188 - }
189 -
190 - wg := sync.WaitGroup{}
191 - for _, p := range closer {
192 - if kb.Closer(p, dht.self, key) && peers.TryAdd(p) {
193 - select {
194 - case peerOut <- p:
195 - case <-ctx.Done():
196 - return
197 - }
198 - wg.Add(1)
199 - go func(p peer.ID) {
200 - dht.getClosestPeersRecurse(ctx, key, p, peers, peerOut)
201 - wg.Done()
202 - }(p)
203 - }
204 - }
205 - wg.Wait()
206 -}
207 -
204 func (dht *IpfsDHT) closerPeersSingle(ctx context.Context, key u.Key, p peer.ID) ([]peer.ID, error) {
205 pmes, err := dht.findPeerSingle(ctx, p, peer.ID(key))
206 if err != nil {