newMessage and more impl.
Juan Batiz-Benet committed
Sep 16, 2014 at 07:17 UTC
15a823d05815a6e8d4cb03048ef89c14a5c4e84a
4 files changed
+59
-104
routing/dht/Message.go
+15
-5
@@ -6,6 +6,15 @@ import (
6
u "github.com/jbenet/go-ipfs/util"
7
)
8
9
+func newMessage(typ Message_MessageType, key string, level int) *Message {
10
+ m := &Message{
11
+ Type: &typ,
12
+ Key: &key,
13
+ }
14
+ m.SetClusterLevel(level)
15
+ return m
16
+}
17
+
18
func peerToPBPeer(p *peer.Peer) *Message_Peer {
19
pbp := new(Message_Peer)
20
if len(p.Addresses) == 0 || p.Addresses[0] == nil {
@@ -24,7 +33,7 @@ func peerToPBPeer(p *peer.Peer) *Message_Peer {
33
}
34
35
func peersToPBPeers(peers []*peer.Peer) []*Message_Peer {
27
- pbpeers = make([]*Message_Peer, len(peers))
36
+ pbpeers := make([]*Message_Peer, len(peers))
37
for i, p := range peers {
38
pbpeers[i] = peerToPBPeer(p)
39
}
@@ -34,18 +43,19 @@ func peersToPBPeers(peers []*peer.Peer) []*Message_Peer {
43
// GetClusterLevel gets and adjusts the cluster level on the message.
44
// a +/- 1 adjustment is needed to distinguish a valid first level (1) and
45
// default "no value" protobuf behavior (0)
37
-func (m *Message) GetClusterLevel() int32 {
46
+func (m *Message) GetClusterLevel() int {
47
level := m.GetClusterLevelRaw() - 1
48
if level < 0 {
49
u.PErr("handleGetValue: no routing level specified, assuming 0\n")
50
level = 0
51
}
43
- return level
52
+ return int(level)
53
}
54
55
// SetClusterLevel adjusts and sets the cluster level on the message.
56
// a +/- 1 adjustment is needed to distinguish a valid first level (1) and
57
// default "no value" protobuf behavior (0)
49
-func (m *Message) SetClusterLevel(level int32) {
50
- m.ClusterLevelRaw = &level
58
+func (m *Message) SetClusterLevel(level int) {
59
+ lvl := int32(level)
60
+ m.ClusterLevelRaw = &lvl
61
}
routing/dht/dht.go
+30
-77
@@ -246,11 +246,7 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p *peer.Peer,
246
func (dht *IpfsDHT) getValueSingle(ctx context.Context, p *peer.Peer,
247
key u.Key, level int) (*Message, error) {
248
249
- typ := Message_GET_VALUE
250
- skey := string(key)
251
- pmes := &Message{Type: &typ, Key: &skey}
252
- pmes.SetClusterLevel(int32(level))
253
-
249
+ pmes := newMessage(Message_GET_VALUE, string(key), level)
250
return dht.sendRequest(ctx, p, pmes)
251
}
252
@@ -262,7 +258,7 @@ func (dht *IpfsDHT) getFromPeerList(ctx context.Context, key u.Key,
258
peerlist []*Message_Peer, level int) ([]byte, error) {
259
260
for _, pinfo := range peerlist {
265
- p, err := dht.peerFromInfo(pinfo)
261
+ p, err := dht.ensureConnectedToPeer(pinfo)
262
if err != nil {
263
u.DErr("getFromPeers error: %s\n", err)
264
continue
@@ -334,34 +330,9 @@ func (dht *IpfsDHT) Find(id peer.ID) (*peer.Peer, *kb.RoutingTable) {
330
return nil, nil
331
}
332
337
-func (dht *IpfsDHT) findPeerSingle(p *peer.Peer, id peer.ID, timeout time.Duration, level int) (*Message, error) {
338
- pmes := Message{
339
- Type: Message_FIND_NODE,
340
- Key: string(id),
341
- ID: swarm.GenerateMessageID(),
342
- Value: []byte{byte(level)},
343
- }
344
-
345
- mes := swarm.NewMessage(p, pmes.ToProtobuf())
346
- listenChan := dht.listener.Listen(pmes.ID, 1, time.Minute)
347
- t := time.Now()
348
- dht.netChan.Outgoing <- mes
349
- after := time.After(timeout)
350
- select {
351
- case <-after:
352
- dht.listener.Unlisten(pmes.ID)
353
- return nil, u.ErrTimeout
354
- case resp := <-listenChan:
355
- roundtrip := time.Since(t)
356
- resp.Peer.SetLatency(roundtrip)
357
- pmesOut := new(Message)
358
- err := proto.Unmarshal(resp.Data, pmesOut)
359
- if err != nil {
360
- return nil, err
361
- }
362
-
363
- return pmesOut, nil
364
- }
333
+func (dht *IpfsDHT) findPeerSingle(ctx context.Context, p *peer.Peer, id peer.ID, level int) (*Message, error) {
334
+ pmes := newMessage(Message_FIND_NODE, string(id), level)
335
+ return dht.sendRequest(ctx, p, pmes)
336
}
337
338
func (dht *IpfsDHT) printTables() {
@@ -370,54 +341,27 @@ func (dht *IpfsDHT) printTables() {
341
}
342
}
343
373
-func (dht *IpfsDHT) findProvidersSingle(p *peer.Peer, key u.Key, level int, timeout time.Duration) (*Message, error) {
374
- pmes := Message{
375
- Type: Message_GET_PROVIDERS,
376
- Key: string(key),
377
- ID: swarm.GenerateMessageID(),
378
- Value: []byte{byte(level)},
379
- }
380
-
381
- mes := swarm.NewMessage(p, pmes.ToProtobuf())
382
-
383
- listenChan := dht.listener.Listen(pmes.ID, 1, time.Minute)
384
- dht.netChan.Outgoing <- mes
385
- after := time.After(timeout)
386
- select {
387
- case <-after:
388
- dht.listener.Unlisten(pmes.ID)
389
- return nil, u.ErrTimeout
390
- case resp := <-listenChan:
391
- u.DOut("FindProviders: got response.\n")
392
- pmesOut := new(Message)
393
- err := proto.Unmarshal(resp.Data, pmesOut)
394
- if err != nil {
395
- return nil, err
396
- }
397
-
398
- return pmesOut, nil
399
- }
344
+func (dht *IpfsDHT) findProvidersSingle(ctx context.Context, p *peer.Peer, key u.Key, level int) (*Message, error) {
345
+ pmes := newMessage(Message_GET_PROVIDERS, string(key), level)
346
+ return dht.sendRequest(ctx, p, pmes)
347
}
348
349
// TODO: Could be done async
403
-func (dht *IpfsDHT) addPeerList(key u.Key, peers []*Message_PBPeer) []*peer.Peer {
350
+func (dht *IpfsDHT) addProviders(key u.Key, peers []*Message_Peer) []*peer.Peer {
351
var provArr []*peer.Peer
352
for _, prov := range peers {
406
- // Dont add outselves to the list
407
- if peer.ID(prov.GetId()).Equal(dht.self.ID) {
353
+ p, err := dht.peerFromInfo(prov)
354
+ if err != nil {
355
+ u.PErr("error getting peer from info: %v\n", err)
356
continue
357
}
410
- // Dont add someone who is already on the list
411
- p := dht.network.GetPeer(u.Key(prov.GetId()))
412
- if p == nil {
413
- u.DOut("given provider %s was not in our network already.\n", peer.ID(prov.GetId()).Pretty())
414
- var err error
415
- p, err = dht.peerFromInfo(prov)
416
- if err != nil {
417
- u.PErr("error connecting to new peer: %s\n", err)
418
- continue
419
- }
358
+
359
+ // Dont add outselves to the list
360
+ if p.ID.Equal(dht.self.ID) {
361
+ continue
362
}
363
+
364
+ // TODO(jbenet) ensure providers is idempotent
365
dht.providers.AddProvider(key, p)
366
provArr = append(provArr, p)
367
}
@@ -450,6 +394,7 @@ func (dht *IpfsDHT) betterPeerToQuery(pmes *Message) *peer.Peer {
394
}
395
396
// self is closer? nil
397
+ key := u.Key(pmes.GetKey())
398
if kb.Closer(dht.self.ID, closer.ID, key) {
399
return nil
400
}
@@ -478,11 +423,19 @@ func (dht *IpfsDHT) peerFromInfo(pbp *Message_Peer) (*peer.Peer, error) {
423
// create new Peer
424
p := &peer.Peer{ID: id}
425
p.AddAddress(maddr)
481
- dht.peerstore.Put(pr)
426
+ dht.peerstore.Put(p)
427
+ }
428
+ return p, nil
429
+}
430
+
431
+func (dht *IpfsDHT) ensureConnectedToPeer(pbp *Message_Peer) (*peer.Peer, error) {
432
+ p, err := dht.peerFromInfo(pbp)
433
+ if err != nil {
434
+ return nil, err
435
}
436
437
// dial connection
485
- err = dht.network.Dial(p)
438
+ err = dht.network.DialPeer(p)
439
return p, err
440
}
441
@@ -497,7 +450,7 @@ func (dht *IpfsDHT) loadProvidableKeys() error {
450
return nil
451
}
452
500
-// Builds up list of peers by requesting random peer IDs
453
+// Bootstrap builds up list of peers by requesting random peer IDs
454
func (dht *IpfsDHT) Bootstrap() {
455
id := make([]byte, 16)
456
rand.Read(id)
routing/dht/handlers.go
+13
-21
@@ -40,11 +40,8 @@ func (dht *IpfsDHT) handlerForMsgType(t Message_MessageType) dhtHandler {
40
41
func (dht *IpfsDHT) putValueToNetwork(p *peer.Peer, key string, value []byte) error {
42
typ := Message_PUT_VALUE
43
- pmes := &Message{
44
- Type: &typ,
45
- Key: &key,
46
- Value: value,
47
- }
43
+ pmes := newMessage(Message_PUT_VALUE, string(key), 0)
44
+ pmes.Value = value
45
46
mes, err := msg.FromObject(p, pmes)
47
if err != nil {
@@ -57,10 +54,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
54
u.DOut("handleGetValue for key: %s\n", pmes.GetKey())
55
56
// setup response
60
- resp := &Message{
61
- Type: pmes.Type,
62
- Key: pmes.Key,
63
- }
57
+ resp := newMessage(pmes.GetType(), pmes.GetKey(), pmes.GetClusterLevel())
58
59
// first, is the key even a key?
60
key := pmes.GetKey()
@@ -113,24 +107,22 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
107
}
108
109
// Store a value in this peer local storage
116
-func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *Message) {
110
+func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *Message) (*Message, error) {
111
dht.dslock.Lock()
112
defer dht.dslock.Unlock()
113
dskey := ds.NewKey(pmes.GetKey())
114
err := dht.datastore.Put(dskey, pmes.GetValue())
121
- if err != nil {
122
- // For now, just panic, handle this better later maybe
123
- panic(err)
124
- }
115
+ return nil, err
116
}
117
118
func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *Message) (*Message, error) {
119
u.DOut("[%s] Responding to ping from [%s]!\n", dht.self.ID.Pretty(), p.ID.Pretty())
129
- return &Message{Type: pmes.Type}, nil
120
+
121
+ return newMessage(pmes.GetType(), "", int(pmes.GetClusterLevel())), nil
122
}
123
124
func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) (*Message, error) {
133
- resp := &Message{Type: pmes.Type}
125
+ resp := newMessage(pmes.GetType(), "", pmes.GetClusterLevel())
126
var closest *peer.Peer
127
128
// if looking for self... special case where we send it on CloserPeers.
@@ -156,10 +148,7 @@ func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *Message) (*Message, error
148
}
149
150
func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *Message) (*Message, error) {
159
- resp := &Message{
160
- Type: pmes.Type,
161
- Key: pmes.Key,
162
- }
151
+ resp := newMessage(pmes.GetType(), pmes.GetKey(), pmes.GetClusterLevel())
152
153
// check if we have this value, to add ourselves as provider.
154
has, err := dht.datastore.Has(ds.NewKey(pmes.GetKey()))
@@ -193,11 +182,14 @@ type providerInfo struct {
182
Value *peer.Peer
183
}
184
196
-func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *Message) {
185
+func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *Message) (*Message, error) {
186
key := u.Key(pmes.GetKey())
187
+
188
u.DOut("[%s] Adding [%s] as a provider for '%s'\n",
189
dht.self.ID.Pretty(), p.ID.Pretty(), peer.ID(key).Pretty())
190
+
191
dht.providers.AddProvider(key, p)
192
+ return nil, nil
193
}
194
195
// Halt stops all communications from this peer and shut down
routing/dht/routing.go
+1
-1
@@ -261,7 +261,7 @@ func (dht *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Pee
261
}
262
if pmes.GetSuccess() {
263
u.DOut("Got providers back from findProviders call!\n")
264
- provs := dht.addPeerList(key, pmes.GetPeers())
264
+ provs := dht.addProviders(key, pmes.GetPeers())
265
ll.Success = true
266
return provs, nil
267
}