@cryptotaxi247 / kubo / commits / 46aa22e94

some better logging and cleanup

Jeromy committed Dec 28, 2014 at 23:46 UTC 46aa22e949c0884dd67f255ae7fbcfa6781e5cf1
3 files changed +20 -71
routing/dht/dht.go
+5 -25
@@ -102,8 +102,8 @@ func (dht *IpfsDHT) Connect(ctx context.Context, npeer peer.ID) error {
102 return nil
103 }
104
105 -// putValueToNetwork stores the given key/value pair at the peer 'p'
106 -func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p peer.ID,
105 +// putValueToPeer stores the given key/value pair at the peer 'p'
106 +func (dht *IpfsDHT) putValueToPeer(ctx context.Context, p peer.ID,
107 key u.Key, rec *pb.Record) error {
108
109 pmes := pb.NewMessage(pb.Message_PUT_VALUE, string(key), 0)
@@ -237,12 +237,12 @@ func (dht *IpfsDHT) Update(ctx context.Context, p peer.ID) {
237 }
238
239 // FindLocal looks for a peer with a given ID connected to this dht and returns the peer and the table it was found in.
240 -func (dht *IpfsDHT) FindLocal(id peer.ID) (peer.PeerInfo, *kb.RoutingTable) {
240 +func (dht *IpfsDHT) FindLocal(id peer.ID) peer.PeerInfo {
241 p := dht.routingTable.Find(id)
242 if p != "" {
243 - return dht.peerstore.PeerInfo(p), dht.routingTable
243 + return dht.peerstore.PeerInfo(p)
244 }
245 - return peer.PeerInfo{}, nil
245 + return peer.PeerInfo{}
246 }
247
248 // findPeerSingle asks peer 'p' if they know where the peer with id 'id' is
@@ -256,26 +256,6 @@ func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p peer.ID, key u.Ke
256 return dht.sendRequest(ctx, p, pmes)
257 }
258
259 -func (dht *IpfsDHT) addProviders(key u.Key, pbps []*pb.Message_Peer) []peer.ID {
260 - peers := pb.PBPeersToPeerInfos(pbps)
261 -
262 - var provArr []peer.ID
263 - for _, pi := range peers {
264 - p := pi.ID
265 -
266 - // Dont add outselves to the list
267 - if p == dht.self {
268 - continue
269 - }
270 -
271 - log.Debugf("%s adding provider: %s for %s", dht.self, p, key)
272 - // TODO(jbenet) ensure providers is idempotent
273 - dht.providers.AddProvider(key, p)
274 - provArr = append(provArr, p)
275 - }
276 - return provArr
277 -}
278 -
259 // nearestPeersToQuery returns the routing tables closest peers.
260 func (dht *IpfsDHT) nearestPeersToQuery(pmes *pb.Message, count int) []peer.ID {
261 key := u.Key(pmes.GetKey())
routing/dht/handlers.go
+2 -5
@@ -39,7 +39,7 @@ func (dht *IpfsDHT) handlerForMsgType(t pb.Message_MessageType) dhtHandler {
39 }
40
41 func (dht *IpfsDHT) handleGetValue(ctx context.Context, p peer.ID, pmes *pb.Message) (*pb.Message, error) {
42 - log.Debugf("%s handleGetValue for key: %s\n", dht.self, pmes.GetKey())
42 + log.Debugf("%s handleGetValue for key: %s", dht.self, pmes.GetKey())
43
44 // setup response
45 resp := pb.NewMessage(pmes.GetType(), pmes.GetKey(), pmes.GetClusterLevel())
@@ -127,7 +127,7 @@ func (dht *IpfsDHT) handlePutValue(ctx context.Context, p peer.ID, pmes *pb.Mess
127 }
128
129 err = dht.datastore.Put(dskey, data)
130 - log.Debugf("%s handlePutValue %v\n", dht.self, dskey)
130 + log.Debugf("%s handlePutValue %v", dht.self, dskey)
131 return pmes, err
132 }
133
@@ -137,9 +137,6 @@ 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 -
140 resp := pb.NewMessage(pmes.GetType(), "", pmes.GetClusterLevel())
141 var closest []peer.ID
142
routing/dht/routing.go
+13 -41
@@ -50,7 +50,7 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
50 wg.Add(1)
51 go func(p peer.ID) {
52 defer wg.Done()
53 - err := dht.putValueToNetwork(ctx, p, key, rec)
53 + err := dht.putValueToPeer(ctx, p, key, rec)
54 if err != nil {
55 log.Errorf("failed putting value to peer: %s", err)
56 }
@@ -125,12 +125,18 @@ func (dht *IpfsDHT) Provide(ctx context.Context, key u.Key) error {
125 return err
126 }
127
128 + wg := sync.WaitGroup{}
129 for p := range peers {
129 - err := dht.putProvider(ctx, p, string(key))
130 - if err != nil {
131 - log.Error(err)
132 - }
130 + wg.Add(1)
131 + go func(p peer.ID) {
132 + defer wg.Done()
133 + err := dht.putProvider(ctx, p, string(key))
134 + if err != nil {
135 + log.Error(err)
136 + }
137 + }(p)
138 }
139 + wg.Wait()
140 return nil
141 }
142
@@ -144,7 +150,6 @@ func (dht *IpfsDHT) FindProviders(ctx context.Context, key u.Key) ([]peer.PeerIn
150 }
151
152 func (dht *IpfsDHT) getClosestPeers(ctx context.Context, key u.Key, count int) (<-chan peer.ID, error) {
147 - log.Error("Get Closest Peers")
153 tablepeers := dht.routingTable.NearestPeers(kb.ConvertKey(key), AlphaValue)
154 if len(tablepeers) == 0 {
155 return nil, kb.ErrLookupFailure
@@ -170,15 +175,12 @@ func (dht *IpfsDHT) getClosestPeers(ctx context.Context, key u.Key, count int) (
175 go func() {
176 wg.Wait()
177 close(out)
173 - log.Error("Closing closest peer chan")
178 }()
179
180 return out, nil
181 }
182
183 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")
184 closer, err := dht.closerPeersSingle(ctx, key, p)
185 if err != nil {
186 log.Errorf("error getting closer peers: %s", err)
@@ -204,8 +206,6 @@ func (dht *IpfsDHT) getClosestPeersRecurse(ctx context.Context, key u.Key, p pee
206 }
207
208 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
@@ -236,6 +236,7 @@ func (dht *IpfsDHT) FindProvidersAsync(ctx context.Context, key u.Key, count int
236
237 func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, count int, peerOut chan peer.PeerInfo) {
238 defer close(peerOut)
239 + defer log.Event(ctx, "findProviders end", &key)
240 log.Debugf("%s FindProviders %s", dht.self, key)
241
242 ps := pset.NewLimited(count)
@@ -294,40 +295,11 @@ func (dht *IpfsDHT) findProvidersAsyncRoutine(ctx context.Context, key u.Key, co
295 }
296 }
297
297 -func (dht *IpfsDHT) addPeerListAsync(ctx context.Context, k u.Key, peers []*pb.Message_Peer, ps *pset.PeerSet, count int, out chan peer.PeerInfo) {
298 - var wg sync.WaitGroup
299 - peerInfos := pb.PBPeersToPeerInfos(peers)
300 - for _, pi := range peerInfos {
301 - wg.Add(1)
302 - go func(pi peer.PeerInfo) {
303 - defer wg.Done()
304 -
305 - p := pi.ID
306 - if err := dht.ensureConnectedToPeer(ctx, p); err != nil {
307 - log.Errorf("%s", err)
308 - return
309 - }
310 -
311 - dht.providers.AddProvider(k, p)
312 - if ps.TryAdd(p) {
313 - select {
314 - case out <- pi:
315 - case <-ctx.Done():
316 - return
317 - }
318 - } else if ps.Size() >= count {
319 - return
320 - }
321 - }(pi)
322 - }
323 - wg.Wait()
324 -}
325 -
298 // FindPeer searches for a peer with given ID.
299 func (dht *IpfsDHT) FindPeer(ctx context.Context, id peer.ID) (peer.PeerInfo, error) {
300
301 // Check if were already connected to them
330 - if pi, _ := dht.FindLocal(id); pi.ID != "" {
302 + if pi := dht.FindLocal(id); pi.ID != "" {
303 return pi, nil
304 }
305