@cryptotaxi247 / kubo / commits / 41c124a28

worked on gathering data for diagnostic messages and some other misc cleanup

Jeromy committed Aug 6, 2014 at 18:37 UTC 41c124a2825e7b1860464f090089799e4b3863d8
10 files changed +107 -49
peer/peer.go
+2
@@ -2,6 +2,7 @@ package peer
2
3 import (
4 "encoding/hex"
5 + "time"
6
7 u "github.com/jbenet/go-ipfs/util"
8 ma "github.com/jbenet/go-multiaddr"
@@ -30,6 +31,7 @@ type Map map[u.Key]*Peer
31 type Peer struct {
32 ID ID
33 Addresses []*ma.Multiaddr
34 + Distance time.Duration
35 }
36
37 // Key returns the ID as a Key (string) for maps.
routing/dht/bucket.go
+5
@@ -61,3 +61,8 @@ func (b *Bucket) Split(cpl int, target ID) *Bucket {
61 }
62 return (*Bucket)(out)
63 }
64 +
65 +func (b *Bucket) getIter() *list.Element {
66 + bucket_list := (*list.List)(b)
67 + return bucket_list.Front()
68 +}
routing/dht/dht.go
+49 -38
@@ -34,7 +34,7 @@ type IpfsDHT struct {
34
35 // Map keys to peers that can provide their value
36 // TODO: implement a TTL on each of these keys
37 - providers map[u.Key][]*peer.Peer
37 + providers map[u.Key][]*providerInfo
38 providerLock sync.RWMutex
39
40 // map of channels waiting for reply messages
@@ -43,6 +43,9 @@ type IpfsDHT struct {
43
44 // Signal to shutdown dht
45 shutdown chan struct{}
46 +
47 + // When this peer started up
48 + birth time.Time
49 }
50
51 // Create a new DHT object with the given peer as the 'local' host
@@ -61,9 +64,10 @@ func NewDHT(p *peer.Peer) (*IpfsDHT, error) {
64 dht.datastore = ds.NewMapDatastore()
65 dht.self = p
66 dht.listeners = make(map[uint64]chan *swarm.Message)
64 - dht.providers = make(map[u.Key][]*peer.Peer)
67 + dht.providers = make(map[u.Key][]*providerInfo)
68 dht.shutdown = make(chan struct{})
69 dht.routes = NewRoutingTable(20, convertPeerID(p.ID))
70 + dht.birth = time.Now()
71 return dht, nil
72 }
73
@@ -121,6 +125,8 @@ func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) (*peer.Peer, error) {
125 // NOTE: this function is just a quick sketch
126 func (dht *IpfsDHT) handleMessages() {
127 u.DOut("Begin message handling routine")
128 +
129 + checkTimeouts := time.NewTicker(time.Minute * 5)
130 for {
131 select {
132 case mes,ok := <-dht.network.Chan.Incoming:
@@ -157,7 +163,10 @@ func (dht *IpfsDHT) handleMessages() {
163 }
164 //
165
160 - u.DOut("Got message type: '%s' [id = %x]", DHTMessage_MessageType_name[int32(pmes.GetType())], pmes.GetId())
166 + u.DOut("[peer: %s]", dht.self.ID.Pretty())
167 + u.DOut("Got message type: '%s' [id = %x, from = %s]",
168 + DHTMessage_MessageType_name[int32(pmes.GetType())],
169 + pmes.GetId(), mes.Peer.ID.Pretty())
170 switch pmes.GetType() {
171 case DHTMessage_GET_VALUE:
172 dht.handleGetValue(mes.Peer, pmes)
@@ -171,35 +180,57 @@ func (dht *IpfsDHT) handleMessages() {
180 dht.handleGetProviders(mes.Peer, pmes)
181 case DHTMessage_PING:
182 dht.handlePing(mes.Peer, pmes)
183 + case DHTMessage_DIAGNOSTIC:
184 + // TODO: network diagnostic messages
185 }
186
187 case err := <-dht.network.Chan.Errors:
188 u.DErr("dht err: %s", err)
189 case <-dht.shutdown:
190 + checkTimeouts.Stop()
191 return
192 + case <-checkTimeouts.C:
193 + dht.providerLock.Lock()
194 + for k,parr := range dht.providers {
195 + var cleaned []*providerInfo
196 + for _,v := range parr {
197 + if time.Since(v.Creation) < time.Hour {
198 + cleaned = append(cleaned, v)
199 + }
200 + }
201 + dht.providers[k] = cleaned
202 + }
203 + dht.providerLock.Unlock()
204 }
205 }
206 }
207
208 func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *DHTMessage) {
209 dskey := ds.NewKey(pmes.GetKey())
210 + var resp *pDHTMessage
211 i_val, err := dht.datastore.Get(dskey)
212 if err == nil {
188 - resp := &pDHTMessage{
213 + resp = &pDHTMessage{
214 Response: true,
215 Id: *pmes.Id,
216 Key: *pmes.Key,
217 Value: i_val.([]byte),
218 + Success: true,
219 }
194 -
195 - mes := swarm.NewMessage(p, resp.ToProtobuf())
196 - dht.network.Chan.Outgoing <- mes
220 } else if err == ds.ErrNotFound {
221 // Find closest peer(s) to desired key and reply with that info
199 - // TODO: this will need some other metadata in the protobuf message
200 - // to signal to the querying peer that the data its receiving
201 - // is actually a list of other peer
222 + closer := dht.routes.NearestPeer(convertKey(u.Key(pmes.GetKey())))
223 + resp = &pDHTMessage{
224 + Response: true,
225 + Id: *pmes.Id,
226 + Key: *pmes.Key,
227 + Value: closer.ID,
228 + Success: false,
229 + }
230 }
231 +
232 + mes := swarm.NewMessage(p, resp.ToProtobuf())
233 + dht.network.Chan.Outgoing <- mes
234 }
235
236 // Store a value in this peer local storage
@@ -263,14 +294,14 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *DHTMessage) {
294 // This is just a quick hack, formalize method of sending addrs later
295 addrs := make(map[u.Key]string)
296 for _,prov := range providers {
266 - ma := prov.NetAddress("tcp")
297 + ma := prov.Value.NetAddress("tcp")
298 str,err := ma.String()
299 if err != nil {
300 u.PErr("Error: %s", err)
301 continue
302 }
303
273 - addrs[prov.Key()] = str
304 + addrs[prov.Value.Key()] = str
305 }
306
307 data,err := json.Marshal(addrs)
@@ -290,6 +321,11 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *DHTMessage) {
321 dht.network.Chan.Outgoing <-mes
322 }
323
324 +type providerInfo struct {
325 + Creation time.Time
326 + Value *peer.Peer
327 +}
328 +
329 func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *DHTMessage) {
330 //TODO: need to implement TTLs on providers
331 key := u.Key(pmes.GetKey())
@@ -324,35 +360,10 @@ func (dht *IpfsDHT) Halt() {
360 dht.network.Close()
361 }
362
327 -// Ping a peer, log the time it took
328 -func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
329 - // Thoughts: maybe this should accept an ID and do a peer lookup?
330 - u.DOut("Enter Ping.")
331 -
332 - pmes := pDHTMessage{Id: GenerateMessageID(), Type: DHTMessage_PING}
333 - mes := swarm.NewMessage(p, pmes.ToProtobuf())
334 -
335 - before := time.Now()
336 - response_chan := dht.ListenFor(pmes.Id)
337 - dht.network.Chan.Outgoing <- mes
338 -
339 - tout := time.After(timeout)
340 - select {
341 - case <-response_chan:
342 - roundtrip := time.Since(before)
343 - u.POut("Ping took %s.", roundtrip.String())
344 - return nil
345 - case <-tout:
346 - // Timed out, think about removing peer from network
347 - u.DOut("Ping peer timed out.")
348 - return u.ErrTimeout
349 - }
350 -}
351 -
363 func (dht *IpfsDHT) addProviderEntry(key u.Key, p *peer.Peer) {
364 u.DOut("Adding %s as provider for '%s'", p.Key().Pretty(), key)
365 dht.providerLock.Lock()
366 provs := dht.providers[key]
356 - dht.providers[key] = append(provs, p)
367 + dht.providers[key] = append(provs, &providerInfo{time.Now(), p})
368 dht.providerLock.Unlock()
369 }
routing/dht/dht_test.go
-4
@@ -6,13 +6,9 @@ import (
6 ma "github.com/jbenet/go-multiaddr"
7 u "github.com/jbenet/go-ipfs/util"
8
9 - "fmt"
10 -
9 "time"
10 )
11
14 -var _ = fmt.Println
15 -
12 func TestPing(t *testing.T) {
13 u.Debug = false
14 addr_a,err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/1234")
routing/dht/messages.pb.go
+8
@@ -74,6 +74,7 @@ type DHTMessage struct {
74 Value []byte `protobuf:"bytes,3,opt,name=value" json:"value,omitempty"`
75 Id *uint64 `protobuf:"varint,4,req,name=id" json:"id,omitempty"`
76 Response *bool `protobuf:"varint,5,opt,name=response" json:"response,omitempty"`
77 + Success *bool `protobuf:"varint,6,opt,name=success" json:"success,omitempty"`
78 XXX_unrecognized []byte `json:"-"`
79 }
80
@@ -116,6 +117,13 @@ func (m *DHTMessage) GetResponse() bool {
117 return false
118 }
119
120 +func (m *DHTMessage) GetSuccess() bool {
121 + if m != nil && m.Success != nil {
122 + return *m.Success
123 + }
124 + return false
125 +}
126 +
127 func init() {
128 proto.RegisterEnum("dht.DHTMessage_MessageType", DHTMessage_MessageType_name, DHTMessage_MessageType_value)
129 }
routing/dht/messages.proto
+1
@@ -22,4 +22,5 @@ message DHTMessage {
22
23 // Signals whether or not this message is a response to another message
24 optional bool response = 5;
25 + optional bool success = 6;
26 }
routing/dht/pDHTMessage.go
+2
@@ -7,6 +7,7 @@ type pDHTMessage struct {
7 Value []byte
8 Response bool
9 Id uint64
10 + Success bool
11 }
12
13 func (m *pDHTMessage) ToProtobuf() *DHTMessage {
@@ -19,6 +20,7 @@ func (m *pDHTMessage) ToProtobuf() *DHTMessage {
20 pmes.Key = &m.Key
21 pmes.Response = &m.Response
22 pmes.Id = &m.Id
23 + pmes.Success = &m.Success
24
25 return pmes
26 }
routing/dht/routing.go
+27
@@ -208,3 +208,30 @@ func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error
208 return found_peer, nil
209 }
210 }
211 +
212 +// Ping a peer, log the time it took
213 +func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
214 + // Thoughts: maybe this should accept an ID and do a peer lookup?
215 + u.DOut("Enter Ping.")
216 +
217 + pmes := pDHTMessage{Id: GenerateMessageID(), Type: DHTMessage_PING}
218 + mes := swarm.NewMessage(p, pmes.ToProtobuf())
219 +
220 + before := time.Now()
221 + response_chan := dht.ListenFor(pmes.Id)
222 + dht.network.Chan.Outgoing <- mes
223 +
224 + tout := time.After(timeout)
225 + select {
226 + case <-response_chan:
227 + roundtrip := time.Since(before)
228 + p.Distance = roundtrip //TODO: This isnt threadsafe
229 + u.POut("Ping took %s.", roundtrip.String())
230 + return nil
231 + case <-tout:
232 + // Timed out, think about removing peer from network
233 + u.DOut("Ping peer timed out.")
234 + dht.Unlisten(pmes.Id)
235 + return u.ErrTimeout
236 + }
237 +}
routing/dht/table.go
+11 -5
@@ -1,12 +1,10 @@
1 package dht
2
3 import (
4 - "encoding/hex"
4 "container/list"
5 "sort"
6
7 peer "github.com/jbenet/go-ipfs/peer"
9 - u "github.com/jbenet/go-ipfs/util"
8 )
9
10 // RoutingTable defines the routing table.
@@ -114,7 +112,6 @@ func (rt *RoutingTable) NearestPeer(id ID) *peer.Peer {
112
113 // Returns a list of the 'count' closest peers to the given ID
114 func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
117 - u.POut("Searching table, size = %d", rt.Size())
115 cpl := xor(id, rt.local).commonPrefixLen()
116
117 // Get bucket at cpl index or last bucket
@@ -148,8 +145,6 @@ func (rt *RoutingTable) NearestPeers(id ID, count int) []*peer.Peer {
145 var out []*peer.Peer
146 for i := 0; i < count && i < peerArr.Len(); i++ {
147 out = append(out, peerArr[i].p)
151 - u.POut("peer out: %s - %s", peerArr[i].p.ID.Pretty(),
152 - hex.EncodeToString(xor(id, convertPeerID(peerArr[i].p.ID))))
148 }
149
150 return out
@@ -163,3 +158,14 @@ func (rt *RoutingTable) Size() int {
158 }
159 return tot
160 }
161 +
162 +// NOTE: This is potentially unsafe... use at your own risk
163 +func (rt *RoutingTable) listpeers() []*peer.Peer {
164 + var peers []*peer.Peer
165 + for _,buck := range rt.Buckets {
166 + for e := buck.getIter(); e != nil; e = e.Next() {
167 + peers = append(peers, e.Value.(*peer.Peer))
168 + }
169 + }
170 + return peers
171 +}
routing/dht/util.go
+2 -2
@@ -11,8 +11,8 @@ import (
11 // ID for IpfsDHT should be a byte slice, to allow for simpler operations
12 // (xor). DHT ids are based on the peer.IDs.
13 //
14 -// NOTE: peer.IDs are biased because they are multihashes (first bytes
15 -// biased). Thus, may need to re-hash keys (uniform dist). TODO(jbenet)
14 +// The type dht.ID signifies that its contents have been hashed from either a
15 +// peer.ID or a util.Key. This unifies the keyspace
16 type ID []byte
17
18 func (id ID) Equal(other ID) bool {